feat(kv-index): add PositionalIndexer for event-driven cache-aware routing - #560
Conversation
…uting Add a positional indexer that uses DashMap<(position, ContentHash), SeqEntry> for O(1) random access to any depth position, replacing tree pointer-chasing with direct positional lookup. Jump search skips positions in strides of jump_size (default 64), yielding amortized O(D/J + W) matching complexity. What changed: - kv_index/Cargo.toml: add rustc-hash (FxHash) and xxhash-rust (XXH3) deps - kv_index/src/event_tree.rs: new 1250-line module with PositionalIndexer - kv_index/src/lib.rs: export new types (PositionalIndexer, ContentHash, SequenceHash, StoredBlock, OverlapScores, WorkerId, compute_content_hash) Key design decisions: - Dual-hash scheme: ContentHash (XXH3-64, position-independent, from token IDs) for indexing; SequenceHash (position-aware, from backend proto block_hash) for disambiguation - SeqEntry enum with Single/Multi optimization avoids HashMap allocation in the common single-sequence-hash case - Per-worker LevelIndex reverse lookup enables O(1) block removal - FxHashMap/FxHashSet throughout for 3-5x faster hashing on non-adversarial data - DashMap entries cleaned up when last worker is removed (no memory leak) - TOCTOU-safe worker lookup uses graceful fallback instead of unwrap - tracing::warn on unresolvable parent hash, tracing::debug on untracked worker removal - All methods take &self with internal DashMap + parking_lot::RwLock for thread safety Public API: - apply_stored(worker, blocks, parent_seq_hash): process store events - apply_removed(worker, seq_hashes): process remove events - apply_cleared(worker): process cache-clear events - remove_worker(worker): full worker removal - find_matches(content_hashes) -> OverlapScores: jump-search matching - compute_content_hash(token_ids) -> ContentHash: XXH3 hashing - current_size() -> usize: total blocks across all workers Tests: 43 tests covering store/match/remove/clear operations, jump search with various configurations (jump_size=1, 3, 4, 64), concurrent read+write, DashMap cleanup verification, hash computation, and edge cases. Refs: #557, #558 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
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 introduces a new Highlights
Changelog
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
|
📝 WalkthroughWalkthroughA new positional indexer module (event_tree) was added, providing content and sequence hashing plus a DashMap-backed, cache-aware positional indexer with a jump-search overlap scoring algorithm and comprehensive unit tests. Two hashing dependencies were added to Cargo.toml. Changes
Sequence Diagram(s)sequenceDiagram
autonumber
participant Client
participant Indexer as PositionalIndexer
participant Store as DashMap/Index
participant Workers as WorkerBlocks(RwLock)
Client->>Indexer: apply_stored(worker, blocks, parent_seq_hash)
Indexer->>Workers: write worker_blocks (insert/update)
Indexer->>Store: upsert (position, ContentHash) -> SeqEntry
Store-->>Indexer: ack
Client->>Indexer: find_matches(content_hashes)
Indexer->>Store: jump_search over positions
Store-->>Indexer: matching SeqEntry(s) + workers
Indexer->>Workers: read worker metadata for scores
Workers-->>Indexer: worker sizes
Indexer-->>Client: OverlapScores
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes 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 |
There was a problem hiding this comment.
Code Review
This pull request introduces a PositionalIndexer for event-driven cache-aware routing, a significant and well-thought-out feature. The implementation is comprehensive, with a strong focus on performance using DashMap, FxHash, and a jump search algorithm. The test suite is extensive, covering many edge cases and concurrency scenarios. My feedback focuses on further performance optimizations by reducing allocations and improving locking strategies in a few key areas. Overall, this is a high-quality contribution.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: bcae47ceba
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
There was a problem hiding this comment.
Actionable comments posted: 5
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
kv_index/src/lib.rs (1)
1-13: 🧹 Nitpick | 🔵 TrivialModule documentation should mention PositionalIndexer.
The module doc comment (Lines 1-12) describes only
StringTreeandTokenTree. Consider adding a brief description ofPositionalIndexerfor completeness:📝 Suggested documentation update
//! Radix tree implementations for prefix matching and cache-aware routing. //! //! This module provides radix tree data structures that mirror SGLang's -//! scheduler memory management patterns. Two implementations are available: +//! scheduler memory management patterns. Three implementations are available: //! //! - [`StringTree`]: Character-based tree for HTTP router (text input) //! - [`TokenTree`]: Token-based tree for gRPC router (pre-tokenized input) +//! - [`PositionalIndexer`]: Event-driven positional index for cache-aware routing (gRPC) //!🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@kv_index/src/lib.rs` around lines 1 - 13, Update the module-level doc comment to mention PositionalIndexer alongside StringTree and TokenTree; add a one- or two-sentence description of PositionalIndexer (e.g., it provides positional/token indexing utilities used by the token-based tree for efficient positional lookups and prefix matching, and integrates with the LRU/tenant tracking semantics). Locate the top doc block where StringTree and TokenTree are documented and insert the PositionalIndexer description so the module summary lists all three public components.
🤖 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/Cargo.toml`:
- Around line 20-22: The kv_index crate must revert the rustc-hash bump: change
the rustc-hash dependency in kv_index/Cargo.toml from "2.1.1" back to "1.1.0"
because kv_index exposes FxHashMap types (see event_tree.rs -> EventTreeNode)
which are not ABI-compatible with rustc-hash 2.x and conflict with
model_gateway/llm-tokenizer pinning; keep xxhash-rust = "0.8" (with xxh3
feature) as-is. Ensure Cargo.toml for the kv_index crate uses rustc-hash =
"1.1.0" so consumers of EventTreeNode compile without type-mismatch errors.
In `@kv_index/src/event_tree.rs`:
- Around line 230-241: The read-then-write TOCTOU pattern around
self.worker_blocks (checking contains_key then calling
entry(...).or_insert_with(...)) is safe because entry().or_insert_with() handles
concurrent inserts; no change is required, but optionally you can simplify by
taking a single write lock and doing
self.worker_blocks.write().entry(worker_id.clone()).or_insert_with(||
RwLock::new(FxHashMap::default())) before reading—refer to the worker_blocks
field, the worker_id variable, the entry().or_insert_with() call, and the
level_index lookup (wb.get(&worker_id)) when making this change.
- Around line 178-179: The constructor pub fn new(jump_size: usize) should avoid
panicking on 0; replace the assert! with defensive handling by clamping the
value to at least 1 and emitting a warning via tracing::warn! when an
out-of-range value is supplied (e.g., let jump_size = jump_size.max(1); if
original == 0 { tracing::warn!(...); }). Keep the function signature and
invariant enforcement inside new (referencing the new() constructor and
jump_size parameter) so callers are not able to trigger a panic while still
preserving expected behavior.
- Around line 70-73: compute_content_hash currently allocates a Vec<u8> every
call; replace the allocation with an incremental hash update using the xxh3
streaming API to avoid the heap allocation: create a local 4-byte stack buffer
(e.g., let mut b = [0u8;4]) or call to_le_bytes() and feed those bytes directly
into xxhash_rust::xxh3::Xxh3::with_seed(XXH3_SEED) via hasher.update(&b) for
each token_id, then finish the hasher and wrap the u64 in ContentHash (keep
function name compute_content_hash and return type ContentHash unchanged).
- Around line 499-513: The cardinality check using num_workers_at_next ==
active.len() is unsafe because equal counts can mask different worker
identities; instead compare the actual worker sets before deciding to skip the
range: compute the set of workers at next_pos and compare it to the current
active set (the variables active / num_workers_at_next / next_pos / current_pos
are the relevant symbols), and only skip when the sets are equal; otherwise call
linear_scan_drain(content_hashes, &mut seq_hashes, &mut active, &mut scores,
current_pos + 1, next_pos + 1) as before to record exact drain points.
---
Outside diff comments:
In `@kv_index/src/lib.rs`:
- Around line 1-13: Update the module-level doc comment to mention
PositionalIndexer alongside StringTree and TokenTree; add a one- or two-sentence
description of PositionalIndexer (e.g., it provides positional/token indexing
utilities used by the token-based tree for efficient positional lookups and
prefix matching, and integrates with the LRU/tenant tracking semantics). Locate
the top doc block where StringTree and TokenTree are documented and insert the
PositionalIndexer description so the module summary lists all three public
components.
…treaming hasher What changed: - kv_index/src/event_tree.rs: WorkerId type alias changed from String to Arc<str> for cheap cloning on hot paths (SeqEntry insert/remove, linear_scan_drain retain, OverlapScores keys). Matches the existing TenantId = Arc<str> pattern in common.rs. - kv_index/src/event_tree.rs: compute_content_hash now uses streaming Xxh3::with_seed() hasher instead of collecting into Vec<u8> first, eliminating an intermediate allocation per call. - Removed redundant `use std::sync::Arc` imports from two test functions since Arc is now imported at module level. Why: - PR #560 received valid feedback that WorkerId = String causes unnecessary heap allocations on every clone in hot paths. Arc<str> gives O(1) reference-counted cloning. - The Vec<u8> allocation in compute_content_hash was flagged as unnecessary — the streaming XXH3 hasher accumulates bytes internally and produces the same hash output. How: - SeqEntry::new() and SeqEntry::insert() now take &str and intern to Arc<str> internally. SeqEntry::remove() takes &str for lookup (FxHashSet<Arc<str>> supports Borrow<str>). - apply_stored/apply_removed/remove_or_clear_worker convert &str to Arc<str> once at entry point, pass &str to SeqEntry methods. - OverlapScores public fields are now FxHashMap<Arc<str>, _>, queryable with &str via Borrow trait. - All 43 existing tests pass unchanged (string literal lookups work via Arc<str>: Borrow<str>). Refs: PR #560 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 4a75abbcc2
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| if keep_worker { | ||
| wb.insert(worker_id, RwLock::new(FxHashMap::default())); | ||
| } |
There was a problem hiding this comment.
Skip re-adding untracked workers on cache-clear events
apply_cleared routes through remove_or_clear_worker(..., true), and this helper unconditionally executes wb.insert(worker_id, ...) even when wb.remove(&worker_id) returned None. In the presence of late or out-of-order clear events for removed/unknown workers, this turns a no-op into persistent state growth, accumulating empty worker entries in worker_blocks and causing stale worker metadata to linger indefinitely.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
♻️ Duplicate comments (2)
kv_index/src/event_tree.rs (2)
186-188:⚠️ Potential issue | 🟠 MajorAvoid panicking in
new()for invalidjump_size.
assert!here can crash production on config/input mistakes. Prefer clamping (orResult) and keep strictness asdebug_assert!in debug/test builds. (If you clamp, the panic-based test should be updated accordingly.)♻️ Proposed fix
pub fn new(jump_size: usize) -> Self { - assert!(jump_size > 0, "jump_size must be greater than 0"); + if jump_size == 0 { + debug_assert!(jump_size > 0, "jump_size must be greater than 0"); + tracing::warn!("PositionalIndexer::new received jump_size=0; clamping to 1"); + } + let jump_size = jump_size.max(1); Self { index: DashMap::with_hasher(FxBuildHasher), worker_blocks: RwLock::new(FxHashMap::default()), jump_size, } }Based on learnings: production code should avoid panics on invariant breaches; prefer a defensive non-panicking path and use
debug_assert!for debug/test builds.🤖 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 186 - 188, The constructor EventTree::new currently uses assert!(jump_size > 0) which can panic in production; change it to a defensive non-panicking path by replacing assert! with debug_assert! and clamping the input (e.g., let jump_size = if jump_size == 0 { 1 } else { jump_size };) before constructing Self (or alternatively change new to return Result<Self, Error> if you prefer explicit failure); update any tests that expected a panic to reflect the clamped behavior if you choose clamping.
498-507:⚠️ Potential issue | 🔴 CriticalJump-skip check is incorrect when worker identities differ but counts match.
Using only cardinality can skip required scans and over-score/drain the wrong workers. The skip condition should verify the same active worker set at
next_pos, not just equal length.🐛 Proposed fix
let num_workers_at_next = self.count_workers_at( next_pos, content_hashes[next_pos], &mut seq_hashes, content_hashes, ); - if num_workers_at_next == active.len() { + let can_skip = if num_workers_at_next == active.len() { + self.get_workers_lazy( + next_pos, + content_hashes[next_pos], + &mut seq_hashes, + content_hashes, + ) + .is_some_and(|workers_at_next| { + active.iter().all(|w| workers_at_next.contains(w)) + }) + } else { + false + }; + + if can_skip { // All active workers still present at jump destination — skip ahead current_pos = next_pos; } else {🤖 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 498 - 507, The jump-skip currently compares only counts (num_workers_at_next == active.len()) which can falsely allow skipping when the worker identities differ; change the logic to compare the actual worker sets at next_pos versus the current active set: modify or replace count_workers_at to provide the worker identities (e.g., workers_at(next_pos, content_hashes[next_pos], &mut seq_hashes, content_hashes) returning a Vec or HashSet of worker IDs) and then only set current_pos = next_pos when the returned set is equal to the current active set (not just equal length); update references to num_workers_at_next and the skip condition accordingly so the function (count/workers_at), next_pos, active, content_hashes, and seq_hashes are used to compare identity equality instead of cardinality.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Duplicate comments:
In `@kv_index/src/event_tree.rs`:
- Around line 186-188: The constructor EventTree::new currently uses
assert!(jump_size > 0) which can panic in production; change it to a defensive
non-panicking path by replacing assert! with debug_assert! and clamping the
input (e.g., let jump_size = if jump_size == 0 { 1 } else { jump_size };) before
constructing Self (or alternatively change new to return Result<Self, Error> if
you prefer explicit failure); update any tests that expected a panic to reflect
the clamped behavior if you choose clamping.
- Around line 498-507: The jump-skip currently compares only counts
(num_workers_at_next == active.len()) which can falsely allow skipping when the
worker identities differ; change the logic to compare the actual worker sets at
next_pos versus the current active set: modify or replace count_workers_at to
provide the worker identities (e.g., workers_at(next_pos,
content_hashes[next_pos], &mut seq_hashes, content_hashes) returning a Vec or
HashSet of worker IDs) and then only set current_pos = next_pos when the
returned set is equal to the current active set (not just equal length); update
references to num_workers_at_next and the skip condition accordingly so the
function (count/workers_at), next_pos, active, content_hashes, and seq_hashes
are used to compare identity equality instead of cardinality.
Summary
PositionalIndexertokv-indexcrate — a positional hash-map indexer for cache-aware routing using KV cache events from backendsDashMap<(position, ContentHash), SeqEntry>for O(1) positional access with jump search for amortized O(D/J + W) matching complexityWhat changed
kv_index/Cargo.tomlrustc-hash(FxHash) andxxhash-rust(XXH3) dependencieskv_index/src/event_tree.rsPositionalIndexerwith jump search,SeqEntrySingle/Multi optimization, hash computation, 43 testskv_index/src/lib.rsPositionalIndexer,ContentHash,SequenceHash,StoredBlock,OverlapScores,WorkerId,compute_content_hashWhy
The gRPC routing path needs ground-truth KV cache overlap scoring from backend events (not approximated from request text). Backends stream
SubscribeKvEventswith block hashes; this indexer processes those events and answers "which worker has the longest cached prefix for this request?" viafind_matches().A positional indexer with jump search achieves dramatically better throughput than a naive radix tree: O(1) random access to any depth vs O(D) pointer-chasing, and amortized O(D/J + W) matching vs O(D) sequential scan.
How
Dual-hash scheme:
ContentHash(XXH3-64): position-independent, computed from token IDs — used as the DashMap keySequenceHash: position-aware rolling hash from backend protoblock_hash— used for disambiguation when multiple sequences share the same content at a positionData structures:
DashMap<(usize, ContentHash), SeqEntry, FxBuildHasher>— primary indexRwLock<FxHashMap<WorkerId, LevelIndex>>— per-worker reverse lookup for O(1) removalSeqEntry::Single/SeqEntry::Multi— avoids HashMap allocation in common caseJump search algorithm:
jump_size(default 64) positionslinear_scan_drainthe range to find exact drain pointstree_sizesThread safety: All methods are
&self— concurrency via DashMap sharding + parking_lot::RwLock. DashMap entries properly cleaned up when last worker removed (no memory leak). TOCTOU-safe worker lookup with graceful fallback.Test plan
cargo clippy -p kv-index --all-targets -- -D warnings— cleancargo clippy --all-targets --all-features -- -D warnings— clean (full workspace)cargo test -p kv-index -- event_tree— 43/43 tests passmake fmt— cleanTest coverage (43 tests):
Refs: #557, #558
Summary by CodeRabbit
New Features
Tests