Repository navigation
Conversation
…te-path contention Move `worker_blocks` (per-worker reverse-lookup FxHashMap) from a shared DashMap on PositionalIndexer to caller-owned task-local storage, matching Dynamo's Flash Indexer architecture for zero contention on the write path. What changed: - kv_index/src/event_tree.rs: Remove `worker_blocks` DashMap from struct. Make `WorkerBlockMap` type and `intern_worker()` public. Change write method signatures (`apply_stored`, `apply_removed`, `apply_cleared`) to accept `&mut WorkerBlockMap` and `worker_id: u32` instead of `&str`. Rewrite `remove_worker` to use O(index_size) scan via `index.retain()` with new `SeqEntry::remove_all_for_worker()` (acceptable since worker removal is rare). Update all ~30 tests with local WorkerBlockMap pattern. - kv_index/src/lib.rs: Export `WorkerBlockMap` type. - model_gateway/src/core/kv_event_monitor.rs: Create task-local `wb` and call `intern_worker()` at subscription_loop start. Thread `wb`/`worker_id` through process_stream → apply_event → apply_stored/apply_removed/ apply_cleared. Update all 20 tests. - model_gateway/src/policies/cache_aware.rs: Update test helpers and 4 inline test sites to use new API (production code unchanged — only calls `find_matches` which is unaffected). - model_gateway/benches/radix_tree_benchmark.rs: Each benchmark thread creates its own WorkerBlockMap — zero DashMap contention on write path. Why: Each subscription task previously acquired a DashMap shard lock on every apply_stored/apply_removed/apply_cleared call to access its worker_blocks entry. With task-local storage, the write path has zero cross-task contention on worker_blocks. This matches Dynamo's thread-local FxHashMap pattern while preserving SMG's tokio-task-per-worker model. How: Surgical API change: callers create `WorkerBlockMap::default()` and call `intern_worker()` once, then pass `(&mut wb, worker_id)` to all write methods. The indexer's shared DashMap for the content index is unchanged. Worker removal (rare path) switches from precise cleanup via worker_blocks to an index scan via `DashMap::retain()`. Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
📝 WalkthroughWalkthroughReplaces internal per-worker DashMap with caller-owned Changes
Sequence Diagram(s)sequenceDiagram
participant Sub as KvEventMonitor
participant WB as WorkerBlockMap
participant Idx as PositionalIndexer
participant Store as PersistentStore
Sub->>Idx: intern_worker(worker_url) -> worker_id
Sub->>WB: create WorkerBlockMap (task-local)
Note right of Sub: For each incoming event
Sub->>Sub: parse event
Sub->>Idx: apply_stored(&mut WB, worker_id, blocks, parent)
Idx->>WB: read/update per-worker block entries
Idx->>Store: lookup/store index structures
Store-->>Idx: ack/result
Idx-->>Sub: Result / ApplyError
Note right of Sub: On removal/clear
Sub->>Idx: apply_removed(&mut WB, worker_id, seq_hashes)
Idx->>WB: remove entries for worker_id
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~28 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches
🧪 Generate unit tests (beta)
Comment |
Summary of ChangesHello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request significantly enhances the performance of the KV indexer by refactoring how per-worker block metadata is managed. By transitioning Highlights
Changelog
Activity
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for Github and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
There was a problem hiding this comment.
Code Review
This pull request effectively refactors the PositionalIndexer to use task-local storage for worker_blocks, which successfully eliminates lock contention on the write path as intended. The changes are consistently applied across the kv_index and model_gateway crates, including updates to tests and benchmarks. The new API is clear and the logic for handling worker state, including connection resets and worker removal, appears robust. I have reviewed the changes and found no issues.
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
kv_index/src/event_tree.rs (1)
283-293:⚠️ Potential issue | 🟡 MinorConflicting default
jump_sizedocumentation.One comment says default is 32 while constructor docs still state 64. Please align docs to the actual default (
Defaultuses 32).Also applies to: 769-772
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@kv_index/src/event_tree.rs` around lines 283 - 293, Documentation for the PositionalIndexer.jump_size is inconsistent: the field docstring says default 32 while the constructor docstring mentions 64; align both to the real default used by the Default implementation. Update the doc comment on the PositionalIndexer struct field `jump_size` and the `pub fn new(jump_size: usize)` docstring (and any other occurrences around the PositionalIndexer docs, e.g., lines noted at 769-772) to state the correct default value (32) to match the `Default` implementation.model_gateway/benches/radix_tree_benchmark.rs (1)
606-623:⚠️ Potential issue | 🟠 MajorConcurrent indexer benchmark writes are mostly no-ops with current parent handling.
Each thread starts with an empty
wbbut writes useSome(parent), so parent lookup fails inapply_storedand the result is ignored. This skews the write-path benchmark.💡 Minimal correction
- let parent = SequenceHash(rng.random_range(1u64..65)); - let _ = - indexer.apply_stored(&mut wb, worker_id, &blocks, Some(parent)); + let _ = + indexer.apply_stored(&mut wb, worker_id, &blocks, None);🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/benches/radix_tree_benchmark.rs` around lines 606 - 623, The benchmark uses an empty WorkerBlockMap `wb` but passes a non-existent `Some(parent)` into `indexer.apply_stored`, causing parent lookup to fail and writes to be no-ops; fix by passing None for the parent (i.e., call `indexer.apply_stored(&mut wb, worker_id, &blocks, None)`) or else create and insert a real parent block before using `Some(parent)` so `apply_stored` can succeed; update the call site around `wb`, `blocks`, `parent`, and `apply_stored` in the loop to ensure writes actually take effect.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@kv_index/src/event_tree.rs`:
- Around line 442-452: remove_worker currently deletes entries from index and
tree_sizes but leaves the URL→ID mapping in worker_to_id, causing worker_id() to
return IDs for removed workers; update remove_worker to also remove the mapping
for the given worker string from the worker_to_id map (use the same worker key
used to compute worker_id), ensuring you reference the existing methods/fields:
remove_worker, worker_id, and the worker_to_id map so the mapping is cleared
when a worker is removed.
---
Outside diff comments:
In `@kv_index/src/event_tree.rs`:
- Around line 283-293: Documentation for the PositionalIndexer.jump_size is
inconsistent: the field docstring says default 32 while the constructor
docstring mentions 64; align both to the real default used by the Default
implementation. Update the doc comment on the PositionalIndexer struct field
`jump_size` and the `pub fn new(jump_size: usize)` docstring (and any other
occurrences around the PositionalIndexer docs, e.g., lines noted at 769-772) to
state the correct default value (32) to match the `Default` implementation.
In `@model_gateway/benches/radix_tree_benchmark.rs`:
- Around line 606-623: The benchmark uses an empty WorkerBlockMap `wb` but
passes a non-existent `Some(parent)` into `indexer.apply_stored`, causing parent
lookup to fail and writes to be no-ops; fix by passing None for the parent
(i.e., call `indexer.apply_stored(&mut wb, worker_id, &blocks, None)`) or else
create and insert a real parent block before using `Some(parent)` so
`apply_stored` can succeed; update the call site around `wb`, `blocks`,
`parent`, and `apply_stored` in the loop to ensure writes actually take effect.
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: a1139f8f-c855-4122-9c38-7a15cb17e810
📒 Files selected for processing (5)
kv_index/src/event_tree.rskv_index/src/lib.rsmodel_gateway/benches/radix_tree_benchmark.rsmodel_gateway/src/core/kv_event_monitor.rsmodel_gateway/src/policies/cache_aware.rs
| pub fn remove_worker(&self, worker: &str) { | ||
| self.remove_or_clear_worker(worker, false); | ||
| let Some(worker_id) = self.worker_id(worker) else { | ||
| return; | ||
| }; | ||
|
|
||
| self.index.retain(|_key, entry| { | ||
| !entry.remove_all_for_worker(worker_id) | ||
| }); | ||
|
|
||
| self.tree_sizes.remove(&worker_id); | ||
| } |
There was a problem hiding this comment.
Remove stale worker_to_id entries during worker cleanup.
remove_worker currently leaves the URL→ID mapping behind. Over time, this can grow stale mappings and make worker_id() return IDs for removed workers.
💡 Suggested fix
pub fn remove_worker(&self, worker: &str) {
- let Some(worker_id) = self.worker_id(worker) else {
+ let Some((_worker_key, worker_id)) = self.worker_to_id.remove(worker) else {
return;
};
self.index.retain(|_key, entry| {
!entry.remove_all_for_worker(worker_id)
});
self.tree_sizes.remove(&worker_id);
}🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@kv_index/src/event_tree.rs` around lines 442 - 452, remove_worker currently
deletes entries from index and tree_sizes but leaves the URL→ID mapping in
worker_to_id, causing worker_id() to return IDs for removed workers; update
remove_worker to also remove the mapping for the given worker string from the
worker_to_id map (use the same worker key used to compute worker_id), ensuring
you reference the existing methods/fields: remove_worker, worker_id, and the
worker_to_id map so the mapping is cleared when a worker is removed.
Run cargo fmt on all modified files to fix line-length and method-chain formatting issues caught by CI. Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
model_gateway/benches/radix_tree_benchmark.rs (1)
606-627:⚠️ Potential issue | 🟠 MajorConcurrent indexer benchmark is dropping writes due to invalid parent chaining.
At Line 607 a fresh
wbis created per thread, but Line 622 always passesSome(parent). Without a seeded parent in that samewb, writes frequently fail, so this benchmark under-measures real write work.💡 Proposed fix
thread::spawn(move || { let mut rng = thread_rng(); let worker_id = indexer.intern_worker(&worker); let mut wb = WorkerBlockMap::default(); + let mut parent: Option<SequenceHash> = None; for i in 0..$ops_per_thread { if i % 3 == 0 { // Read: find_matches let query_tokens = flatten_tokens(&chunks); let content_hashes = compute_request_content_hashes( &query_tokens, block_size, ); black_box(indexer.find_matches(&content_hashes, false)); } else { // Write: apply_stored let new_chunks = generate_token_chunks(4, block_size); let blocks = chunks_to_stored_blocks(&new_chunks); - let parent = SequenceHash(rng.random_range(1u64..65)); let _ = indexer.apply_stored( &mut wb, worker_id, &blocks, - Some(parent), + parent, ); + parent = blocks.last().map(|b| b.seq_hash); } } })🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/benches/radix_tree_benchmark.rs` around lines 606 - 627, The benchmark is creating a fresh WorkerBlockMap (wb) per thread but always passing Some(parent) into indexer.apply_stored, causing writes to fail due to an invalid parent; to fix, seed each thread's wb with a valid parent before the loop (call indexer.apply_stored once with parent None and capture its returned SequenceHash) and then use that returned SequenceHash as the parent for subsequent apply_stored calls (optionally update the parent after each successful apply_stored), referencing worker_id, wb, indexer.apply_stored, and SequenceHash to locate the code to change.
♻️ Duplicate comments (1)
kv_index/src/event_tree.rs (1)
442-451:⚠️ Potential issue | 🟠 Major
remove_workerleaves stale URL→ID mappings inworker_to_id.Line 443 resolves an ID but never removes the mapping itself. This leaks removed worker keys over time and makes
worker_id()return IDs for already-removed workers.💡 Proposed fix
pub fn remove_worker(&self, worker: &str) { - let Some(worker_id) = self.worker_id(worker) else { + let Some((_worker_key, worker_id)) = self.worker_to_id.remove(worker) else { return; }; self.index .retain(|_key, entry| !entry.remove_all_for_worker(worker_id)); self.tree_sizes.remove(&worker_id); }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@kv_index/src/event_tree.rs` around lines 442 - 451, remove_worker currently looks up a worker's ID via worker_id() but never deletes the URL→ID mapping, leaving stale entries in worker_to_id and causing worker_id() to return IDs for removed workers; after you obtain the worker_id in remove_worker (the value produced by worker_id()), explicitly remove the mapping from worker_to_id (e.g. call the map's remove for the worker key) before or as part of the cleanup, then continue clearing entries from index and tree_sizes (keep using entry.remove_all_for_worker and tree_sizes.remove(worker_id)); ensure you reference the same worker string key used to obtain worker_id so the mapping is actually removed.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Outside diff comments:
In `@model_gateway/benches/radix_tree_benchmark.rs`:
- Around line 606-627: The benchmark is creating a fresh WorkerBlockMap (wb) per
thread but always passing Some(parent) into indexer.apply_stored, causing writes
to fail due to an invalid parent; to fix, seed each thread's wb with a valid
parent before the loop (call indexer.apply_stored once with parent None and
capture its returned SequenceHash) and then use that returned SequenceHash as
the parent for subsequent apply_stored calls (optionally update the parent after
each successful apply_stored), referencing worker_id, wb, indexer.apply_stored,
and SequenceHash to locate the code to change.
---
Duplicate comments:
In `@kv_index/src/event_tree.rs`:
- Around line 442-451: remove_worker currently looks up a worker's ID via
worker_id() but never deletes the URL→ID mapping, leaving stale entries in
worker_to_id and causing worker_id() to return IDs for removed workers; after
you obtain the worker_id in remove_worker (the value produced by worker_id()),
explicitly remove the mapping from worker_to_id (e.g. call the map's remove for
the worker key) before or as part of the cleanup, then continue clearing entries
from index and tree_sizes (keep using entry.remove_all_for_worker and
tree_sizes.remove(worker_id)); ensure you reference the same worker string key
used to obtain worker_id so the mapping is actually removed.
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 78ff237d-f5db-4b11-9075-9e66d8e799d2
📒 Files selected for processing (4)
kv_index/src/event_tree.rsmodel_gateway/benches/radix_tree_benchmark.rsmodel_gateway/src/core/kv_event_monitor.rsmodel_gateway/src/policies/cache_aware.rs
… threads The concurrent benchmark threads were creating empty WorkerBlockMap instances, causing all apply_stored calls with parent hashes to fail silently (ApplyError::WorkerNotTracked). The parent blocks were stored during build_populated_indexer into local wb values that were dropped before the threads started. Fix: return per-worker WorkerBlockMap from build_populated_indexer and clone it into each thread, so parent chain lookups succeed. This makes the concurrent benchmark actually exercise the write path correctly. Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/benches/radix_tree_benchmark.rs`:
- Line 608: The benchmark currently calls indexer.intern_worker(&workers[t %
workers.len()]) inside the timed setup; precompute those worker IDs once and
avoid the per-iteration intern lookup. Modify build_populated_indexer to return
the Vec of worker IDs alongside the populated indexer (or add a helper that maps
workers -> worker_id before the timed loop) and replace the inline intern_worker
call with a lookup into the precomputed worker_ids (referencing worker_id,
intern_worker, and build_populated_indexer).
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: b8c60d1e-f786-4e3a-b0cd-fa7970d34d69
📒 Files selected for processing (1)
model_gateway/benches/radix_tree_benchmark.rs
| let worker = workers[t % workers.len()].clone(); | ||
| let chunks = worker_chunks[t % workers.len()].clone(); | ||
| let mut wb = worker_blocks[t % workers.len()].clone(); | ||
| let worker_id = indexer.intern_worker(&workers[t % workers.len()]); |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Precompute worker_id to avoid extra lookup overhead in timed concurrent setup.
Line 608 calls intern_worker during each timed iteration setup. Consider returning worker IDs from build_populated_indexer (or precomputing once) so concurrent benchmark timings focus on read/write paths.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/benches/radix_tree_benchmark.rs` at line 608, The benchmark
currently calls indexer.intern_worker(&workers[t % workers.len()]) inside the
timed setup; precompute those worker IDs once and avoid the per-iteration intern
lookup. Modify build_populated_indexer to return the Vec of worker IDs alongside
the populated indexer (or add a helper that maps workers -> worker_id before the
timed loop) and replace the inline intern_worker call with a lookup into the
precomputed worker_ids (referencing worker_id, intern_worker, and
build_populated_indexer).
Summary
worker_blocks(per-worker reverse-lookupFxHashMap) from a sharedDashMaponPositionalIndexerto caller-owned task-local storage, matching Dynamo's Flash Indexer architectureapply_stored/apply_removed/apply_cleared)WorkerBlockMapat startup — zero cross-task contentionWhat changed
kv_index/src/event_tree.rsworker_blocksDashMap from struct. MakeWorkerBlockMaptype andintern_worker()public. Change write method signatures to accept&mut WorkerBlockMap+worker_id: u32. Rewriteremove_workerto use index scan viaDashMap::retain(). AddSeqEntry::remove_all_for_worker(). Update all ~30 tests.kv_index/src/lib.rsWorkerBlockMaptypemodel_gateway/src/core/kv_event_monitor.rswband callintern_worker()at subscription loop start. Threadwb/worker_idthroughprocess_stream→apply_event→apply_stored/apply_removed/apply_cleared. Update all 20 tests.model_gateway/src/policies/cache_aware.rsfind_matches)model_gateway/benches/radix_tree_benchmark.rsWorkerBlockMap— zero DashMap contention on write pathWhy
Each subscription task previously acquired a DashMap shard lock on every
apply_stored/apply_removed/apply_clearedcall. With task-local storage, the write path has zero cross-task contention onworker_blocks. This mirrors Dynamo's thread-localFxHashMappattern while preserving SMG's tokio-task-per-worker model.Design decisions
remove_workeruses index scan: Sinceworker_blocksis task-local and dropped when the subscription task is aborted,remove_worker(called fromon_worker_removed) scans the index viaDashMap::retain(). This is O(index_size) but acceptable since worker removal is a rare event (backend going down).apply_cleareduses local wb: Called from within the subscription task, so it has access to the localWorkerBlockMapfor precise O(wb_size) cleanup.DashMapfor the content index. Onlyworker_blocksmoves off the struct.Test plan
cargo test -p kv-index— 158 tests passcargo test -p smg --lib -- kv_event_monitor— 20 tests passcargo test -p smg --lib -- cache_aware— 27 tests passcargo clippy -p kv-index -- -D warnings— cleancargo clippy -p smg -- -D warnings— cleancargo build --bench radix_tree_benchmark— compilesSummary by CodeRabbit
New Features
Refactor