Skip to content

perf(kv-index): tune index DashMap shard count - #607

Closed
slin1237 wants to merge 10 commits into
mainfrom
slin/flash-indexer-perf
Closed

slin1237 wants to merge 10 commits into
mainfrom
slin/flash-indexer-perf

Conversation

@slin1237

@slin1237 slin1237 commented Mar 4, 2026 •

Copy link
Copy Markdown
Member

Summary

  • Set explicit INDEX_SHARD_COUNT=256 for the main index DashMap (was default num_cpus*4)
  • Iterative experiment to find optimal shard count for concurrent read+write workloads

What changed

  • kv_index/src/event_tree.rs: Added INDEX_SHARD_COUNT constant, used in PositionalIndexer::new()

Why

DashMap lookup cost dominates the query path (~50-100ns per shard lock acquisition). More shards reduce the probability of reader-writer collision on the same shard. The default shard count is a general-purpose heuristic — we're testing whether a higher explicit value improves our access pattern.

Test plan

  • All 158 kv-index tests pass
  • CI benchmark will provide comparison data against perf_history.md baseline

Summary by CodeRabbit

  • Refactor
    • Restructured event indexing system with improved concurrency handling for enhanced performance under concurrent load
    • Optimized internal worker identification representation for better memory efficiency

slin1237 added 10 commits March 2, 2026 22:51
Gap 1 — Per-position DashMaps:
  Replace DashMap<(usize, ContentHash), SeqEntry> with
  RwLock<Vec<DashMap<ContentHash, SeqEntry, FxBuildHasher>>>.
  One DashMap per block position eliminates cross-position shard
  contention and improves cache locality for hot prefix positions
  (0-3). Vec grows lazily via ensure_levels(); read lock held for
  the duration of queries (only blocks rare Vec growth).

Gap 2 — Eliminate set cloning in query path:
  linear_scan_drain now accesses DashMap entries directly via Ref
  and checks worker membership in-place without cloning the
  FxHashSet<WorkerId>. Previously cloned the set at every position
  during drain (up to jump_size=64 clones per drain call).

Gap 3 — Intern worker IDs:
  Add DashMap<WorkerId, (), FxBuildHasher> intern table. First call
  per worker does Arc::from() (heap alloc); subsequent calls return
  Arc::clone() (atomic increment only). Lookup by &str works via
  Arc<str>: Borrow<str>. Eliminates per-call allocation overhead in
  apply_stored, apply_removed, and remove_or_clear_worker.

Gap 4 — Lazy rolling hash for Single entries:
  Add SeqEntry::workers_if_single() that returns the worker set
  directly for Single entries without requiring rolling hash
  computation. Since content hash collisions at 64-bit XXH3 are
  ~2^-64 (practically impossible), a matching content_hash at the
  same position is unambiguous. Applied in get_workers_lazy,
  count_workers_at, and linear_scan_drain — the three query-path
  helpers that previously always called ensure_seq_hash_computed.

All query helpers (get_workers_lazy, count_workers_at,
linear_scan_drain) converted to associated functions taking
&[DashMap<...>] to share a single read guard acquired once in
jump_search_matches. Explicit drop(levels) before acquiring
worker_blocks.read() maintains lock ordering and prevents deadlock.

Lock ordering: worker_blocks → index → DashMap shards → LevelIndex.

Public API unchanged — all 158 kv-index tests pass.

Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
…er IDs

Three tightly coupled optimizations targeting the 10-115x gap with Dynamo's
Flash Indexer on 128-CPU benchmarks. Dynamo deep dive confirmed their actual
code uses pre-allocated index (no RwLock), per-thread worker_blocks (zero
contention), and Copy WorkerWithDpRank (u64+u32).

Change 1 — Remove RwLock from index:
- Replace RwLock<Vec<DashMap>> with Box<[DashMap]> pre-allocated at
  construction via new max_num_blocks parameter (default 2048)
- Eliminates atomic fetch_add cache-line bouncing on 128 cores (~50-200ns
  per acquisition vs ~5ns uncontended)
- Remove ensure_levels(); blocks beyond capacity are truncated with warning
- jump_search_matches accesses &self.index directly (no read guard)

Change 2 — worker_blocks to DashMap:
- Replace RwLock<FxHashMap<WorkerId, LevelIndex>> with
  DashMap<u32, LevelIndex, FxBuildHasher>
- apply_stored reduced from 3 separate lock acquisitions to clean DashMap
  entry/get pattern
- remove_or_clear_worker uses DashMap::remove() instead of global write lock
  that previously serialized ALL worker operations

Change 3 — Internal u32 worker IDs:
- SeqEntry uses FxHashSet<u32> internally — FxHash on u32 is 1 instruction
  vs 8+ for Arc<str> (20-byte string hash + pointer deref)
- Eliminates Arc::clone atomic refcount bouncing in D×W inner query loop
- intern_worker() returns u32 with double-checked locking pattern
- Public API unchanged: &str params, OverlapScores returns Arc<str> keys
- Conversion at API boundary is O(W), negligible vs inner loop savings

Files changed:
- kv_index/src/event_tree.rs: all three structural changes
- model_gateway/src/core/kv_event_monitor.rs: pass max_num_blocks=2048
- model_gateway/src/policies/cache_aware.rs: update test constructors
- model_gateway/benches/radix_tree_benchmark.rs: update bench constructors

Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
…apScores

What changed:
- kv_index/src/event_tree.rs: Remove inner RwLock from worker_blocks
  (DashMap shard-level locking only), add atomic tree_sizes tracking
  (O(1) reads vs O(n) locked iteration), add retain guard in
  linear_scan_drain (skip when workers.len() >= active.len()), remove
  id_to_worker Vec and RwLock (OverlapScores uses u32 keys directly),
  add early_exit parameter to find_matches, add WorkerId = u32 type
  alias. Convert all tests from string-key to u32 worker ID API.
- model_gateway/src/policies/cache_aware.rs: Update score_overlap to
  use indexer.worker_id() for u32 lookups into OverlapScores, add
  early_exit=false to find_matches call.
- model_gateway/benches/radix_tree_benchmark.rs: Add early_exit=false
  to find_matches calls in bench_indexer_match and
  bench_indexer_concurrent macros.

Why: Close 5 performance gaps between our PositionalIndexer and
Dynamo's Flash Indexer: (1) eliminate inner RwLock contention on
worker_blocks, (2) atomic O(1) tree_sizes instead of locked iteration,
(3) retain guard to skip O(active) iteration when all workers still
match, (4) remove worker ID interning overhead (no Arc<str> conversion
at return time), (6) early_exit support for fast existence checks.

How: worker_blocks changed from DashMap<u32, RwLock<FxHashMap>> to
DashMap<u32, FxHashMap> — plain HashMaps behind DashMap's shard locks.
tree_sizes tracked via DashMap<u32, AtomicUsize> with fetch_add/sub on
store/remove. OverlapScores now keyed by u32 internally; consumers use
PositionalIndexer::worker_id() to map URLs to IDs. The workers_if_single
optimization (fix 5, already present) is retained as an advantage over
Dynamo.

Signed-off-by: Simon Lin <simon@nvidia.com>
Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Phase F: Replace FxHashSet<u32> active set with Vec<u32> in
jump_search_matches to avoid hash table clone overhead at position 0
and improve cache locality during retain/drain operations.

- get_workers_lazy returns Option<Vec<u32>> instead of Option<FxHashSet<u32>>
- linear_scan_drain takes &mut Vec<u32> with swap_remove retain pattern
  instead of FxHashSet::retain(), avoiding per-element hashing
- Drain paths use active.iter() + active.clear() instead of active.drain()
- jump_search_matches iterates with for &w in &active

Phase G: Change default jump_size from 64 to 32 to match Dynamo's
default. Smaller jumps reduce the linear scan window on failed jumps
at the cost of more jump probes.

Also fixes pre-existing clippy lint: worker_map.iter() → worker_map.values()

Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Update benchmark JUMP_SIZE constant and STORE macro to use 32 instead
of 64, matching the new default in PositionalIndexer::default().
This ensures benchmarks exercise the same jump_size as production.

Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
DashMap defaults to num_cpus * 4 shards. On 128-core machines this
creates 512 shards per map (2048 total across 4 DashMaps), most of
which are empty — wasting memory and polluting CPU caches.

Constrain shard counts to match the approach used in token_tree and
string_tree:
- index: 32 shards (hot path, many entries, benefits from concurrency)
- worker_blocks/tree_sizes/worker_to_id: 8 shards each (keyed by
  worker_id, at most ~500 entries, low contention)

Also fix pre-existing clippy len_zero warnings in tests:
  assert!(indexer.index.len() > 0) → assert!(!indexer.index.is_empty())

Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
32 shards caused ~20-25% concurrent throughput regression with 32
benchmark threads — every shard was contended on every query. 64
shards gives 2:1 shard-to-thread ratio while still being 8x less
than the 512 default on 128-core machines.

Worker maps stay at 8 shards (low contention, few entries).

Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
64 shards on the index DashMap still caused concurrent throughput
regression (100w: 253K vs 387K baseline). The mixed read+write
concurrent benchmark needs high shard counts to avoid contention.

Revert index to DashMap default (num_cpus * 4). Only constrain the
three worker-keyed maps (worker_blocks, tree_sizes, worker_to_id) to
8 shards — they hold at most ~500 entries and have low contention.

Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Set explicit INDEX_SHARD_COUNT=256 for the main index DashMap instead of
relying on the default (num_cpus * 4). Higher shard count reduces per-shard
contention probability under concurrent read+write workloads.

Iteration 1 of shard tuning experiment — comparing against baseline in
perf_history.md to find optimal value.

Signed-off-by: Simon Lin <simon.lin@nvidia.com>
Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
@chatgpt-codex-connector

Copy link
Copy Markdown

Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits.
Repo admins can enable using credits for code reviews in their settings.

@github-actions github-actions Bot added benchmarks Benchmark changes model-gateway Model gateway crate changes labels Mar 4, 2026
@coderabbitai

coderabbitai Bot commented Mar 4, 2026 •

Copy link
Copy Markdown
📝 Walkthrough

Walkthrough

The pull request refactors WorkerId from Arc to u32 across the indexing system, introduces DashMap-based per-worker storage with atomic size tracking, updates PositionalIndexer's find_matches signature to include an early_exit parameter, and adjusts cache_aware.rs and benchmarks to accommodate these internal API changes.

Changes

Cohort / File(s) Summary
Core Indexing Refactoring
kv_index/src/event_tree.rs
WorkerId type changed from Arc to u32; OverlapScores now maps u32 IDs instead of WorkerIds; added worker_id(&str) -> Option lookup method; SeqEntry updated to use FxHashSet and FxHashMap<u32, ...>; replaced per-worker RwLock with DashMap-based structures and per-worker atomic counters; find_matches signature updated to require early_exit boolean parameter; internal helpers reworked for numeric IDs and DashMap access patterns.
Benchmark Updates
model_gateway/benches/radix_tree_benchmark.rs
Updated PositionalIndexer constructor from jump size 64 to 32; added false boolean parameter to find_matches call sites in bench_indexer_match and bench_indexer_concurrent benchmarks.
Policy Layer Adaptation
model_gateway/src/policies/cache_aware.rs
Updated overlap score/tree_size lookups to use indexer.worker_id() for URL-to-ID mapping, then query overlap.scores and overlap.tree_sizes by numeric ID; adjusted tie-breaker calculations and filtering logic to use mapped worker IDs.

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~60 minutes

Possibly related PRs

Suggested labels

model-gateway, benchmarks

Suggested reviewers

  • CatherineSue
  • key4ng
  • XinyueZhang369

Poem

🐰 From strings we spring to numbers bright,
Arc takes flight, u32 shines right,
DashMaps dance where RwLocks once dwelled,
Atomic counters ring the bell—
Refactored swift, concurrent spelled!

🚥 Pre-merge checks | ✅ 3
✅ Passed checks (3 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The PR title 'perf(kv-index): tune index DashMap shard count' directly describes the main change: tuning the DashMap shard count to 256 for performance optimization. However, the title is quite narrow and refers to a specific implementation detail rather than the broader refactoring scope evident in the changeset, which includes WorkerId type changes, OverlapScores restructuring, and API adjustments across multiple files.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
  • 📝 Generate docstrings (stacked PR)
  • 📝 Generate docstrings (commit on current branch)
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Post copyable unit tests in a comment
  • Commit unit tests in branch slin/flash-indexer-perf

Comment @coderabbitai help to get the list of available commands and usage tips.

@gemini-code-assist

Copy link
Copy Markdown
Contributor

Summary of Changes

Hello, 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 refactors the kv-index component to enhance performance and concurrency. By explicitly configuring DashMap shard counts, optimizing worker ID representation from Arc<str> to u32 with an interning mechanism, and streamlining concurrency primitives, the changes aim to reduce overhead and contention. These improvements are designed to make the positional indexer more efficient and scalable under high read and write loads.

Highlights

  • DashMap Shard Count Tuning: The main index DashMap now explicitly sets its shard count to 256, moving away from the default num_cpus*4 heuristic. This is an iterative tuning effort aimed at reducing contention in concurrent read+write workloads.
  • Worker ID Optimization: The WorkerId type has been changed from Arc<str> to u32. Worker URLs are now interned to u32 IDs, eliminating Arc<str> hashing and atomic refcount bouncing in hot query paths, improving performance.
  • Concurrency Model Refinement: The worker_blocks and tree_sizes structures within PositionalIndexer now leverage DashMap's internal sharding directly, removing the parking_lot::RwLock previously used. Atomic tree_sizes provide O(1) size queries.
  • Query Path Enhancements: The find_matches method now accepts an early_exit parameter, allowing queries to return immediately upon finding any match at position 0, which is useful for quick existence checks. Helper functions for querying workers have been updated to skip rolling hash computation for 'Single' entries, further optimizing the query path.
  • Default Jump Size Adjustment: The default jump_size for PositionalIndexer has been adjusted from 64 to 32, which is reflected in both the core logic and benchmark configurations.
Changelog
  • kv_index/src/event_tree.rs
    • Updated documentation comments for DashMap usage and performance characteristics.
    • Removed parking_lot::RwLock import and added AtomicU32, AtomicUsize, and Ordering imports.
    • Introduced INDEX_SHARD_COUNT (256) and WORKER_SHARD_COUNT (8) constants for DashMap initialization.
    • Changed WorkerId type alias from Arc<str> to u32.
    • Updated OverlapScores and SeqEntry to use the new u32 WorkerId type.
    • Added a new workers_if_single method to SeqEntry for optimized worker retrieval.
    • Replaced LevelIndex type alias with WorkerBlockMap.
    • Refactored PositionalIndexer fields to use DashMap<u32, WorkerBlockMap> for worker_blocks, and added tree_sizes, worker_to_id, and next_worker_id for improved concurrency and worker management.
    • Modified PositionalIndexer::new to initialize DashMaps with explicit shard amounts and new atomic fields.
    • Added PositionalIndexer::worker_id method to retrieve interned u32 IDs from worker URLs.
    • Updated apply_stored and apply_removed methods to utilize interned u32 worker IDs and atomically update tree_sizes.
    • Modified current_size to sum atomic tree_sizes for O(1) complexity.
    • Changed find_matches method signature to include an early_exit boolean parameter.
    • Refactored internal query helper functions (get_workers_lazy, count_workers_at, linear_scan_drain) to accept &DashMap directly, use u32 worker IDs, and incorporate the early_exit parameter.
    • Adjusted the default jump_size in PositionalIndexer::default from 64 to 32.
    • Updated PositionalIndexer's fmt::Debug implementation to reflect new internal structure.
    • Modified numerous test cases to pass the new early_exit parameter to find_matches and to use indexer.worker_id() for worker ID assertions.
  • model_gateway/benches/radix_tree_benchmark.rs
    • Updated PositionalIndexer::new calls within benchmarks to use a jump_size of 32.
    • Modified indexer.find_matches calls in benchmarks to pass false for the new early_exit parameter.
    • Changed the JUMP_SIZE constant from 64 to 32.
  • model_gateway/src/policies/cache_aware.rs
    • Updated indexer.find_matches calls to pass false for the new early_exit parameter.
    • Modified the logic for filtering and selecting the best worker to use indexer.worker_id() for looking up scores and tree sizes, adapting to the u32 worker ID change.
Using Gemini Code Assist

The 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 /gemini <command> or @gemini-code-assist <command>. Below is a summary of the supported commands on the current page.

Feature Command Description
Code Review /gemini review Performs a code review for the current pull request in its current state.
Pull Request Summary /gemini summary Provides a summary of the current pull request in its current state.
Comment @gemini-code-assist Responds in comments when explicitly tagged, both in pull request comments and review comments.
Help /gemini help Displays a list of available commands.

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 .gemini/ folder in the base of the repository. Detailed instructions can be found here.

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

  1. Review the Privacy Notices, Generative AI Prohibited Use Policy, Terms of Service, and learn how to configure Gemini Code Assist in GitHub here. Gemini can make mistakes, so double check it and use code with caution. ↩

@mergify

mergify Bot commented Mar 4, 2026

Copy link
Copy Markdown
Contributor

Hi @slin1237, this PR has merge conflicts that must be resolved before it can be merged. Please rebase your branch:

git fetch origin main
git rebase origin/main
# resolve any conflicts, then:
git push --force-with-lease

@mergify mergify Bot added the needs-rebase PR has merge conflicts that need to be resolved label Mar 4, 2026

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request introduces significant performance improvements to the PositionalIndexer by refactoring its core data structures and algorithms. The changes go well beyond tuning the DashMap shard count as suggested by the title. Key improvements include:

  • Replacing Arc<str> worker IDs with u32 integers, reducing overhead in hot paths.
  • Replacing a coarse-grained RwLock with DashMap's fine-grained sharding for better concurrency.
  • Introducing atomic counters for O(1) size queries.
  • Optimizing the find_matches algorithm with an early exit path and more efficient linear scanning.

The code is well-structured and the performance optimizations are solid. The test suite has also been comprehensively updated to reflect these major changes.

I have one suggestion regarding the linear_scan_drain function to improve its maintainability by refactoring its large number of arguments into a context struct.

Comment on lines +587 to 597
#[expect(clippy::too_many_arguments)]
fn linear_scan_drain(
&self,
index: &DashMap<(usize, ContentHash), SeqEntry, FxBuildHasher>,
sequence: &[ContentHash],
seq_hashes: &mut Vec<SequenceHash>,
active: &mut FxHashSet<WorkerId>,
scores: &mut OverlapScores,
active: &mut Vec<u32>,
internal_scores: &mut FxHashMap<u32, u32>,
lo: usize,
hi: usize,
early_exit: bool,
) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

The linear_scan_drain function has a large number of arguments (8), which makes it harder to read and maintain. The clippy::too_many_arguments lint is suppressed, but it would be cleaner to group these arguments into a context struct. This would improve readability and make it easier to pass state around if more helpers are added in the future.

For example, you could introduce a LinearScanContext struct:

struct LinearScanContext<'a> {
    index: &'a DashMap<(usize, ContentHash), SeqEntry, FxBuildHasher>,
    sequence: &'a [ContentHash],
    seq_hashes: &'a mut Vec<SequenceHash>,
    active: &'a mut Vec<u32>,
    internal_scores: &'a mut FxHashMap<u32, u32>,
    early_exit: bool,
}

impl<'a> LinearScanContext<'a> {
    fn run(&mut self, lo: usize, hi: usize) {
        // ... function body here ...
    }
}

The call site in jump_search_matches would then be simplified.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 6

🤖 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 766-767: The Default implementation now constructs the type with
Self::new(32) but the struct and function docs still claim the default jump-size
is 64; update all relevant documentation/comments (the struct-level comment and
the new()/Default() method docs in event_tree.rs) to state the default jump-size
is 32 and adjust any doc examples or references that mention 64 so they reflect
the runtime behavior (search for Default, Self::new, and any comments mentioning
"jump-size" or "64" in that file).
- Around line 354-361: apply_stored currently does a relative increment of
self.tree_sizes for worker_id (adding blocks.len()), which can drift when
wb_ref.insert replaces existing seq_hash entries; instead synchronize atomically
to the absolute number of entries in that worker's writebuffer: read
wb_ref.len() after applying inserts/replacements and set the per-worker
AtomicUsize in self.tree_sizes to that absolute value (use
entry(worker_id).and_modify(|size| size.store(wb_len,
Ordering::Relaxed)).or_insert(AtomicUsize::new(wb_len)) or equivalent atomic
update), and make the same change for the other occurrence noted around lines
404-407 so current_size() and routing tie-breaks reflect wb_ref.len() exactly.
- Around line 377-379: The remove/clear branches call intern_worker which
allocates new IDs and grows state; change these paths to perform a lookup-only
so unknown workers are treated as no-ops. Concretely, stop calling intern_worker
in the remove/clear code and instead try to find an existing id without creating
one (e.g. check your existing worker-id map or index for a mapping and only
proceed if found), then use that id to access self.worker_blocks (replace
get_mut after intern_worker with a conditional get/get_mut guarded by the
lookup). Apply the same lookup-only change to the apply_cleared call sites (and
the other similar spot around the 453-454 region) so removals/clears do not
create empty per-worker state.
- Around line 571-573: The fast-path using entry.value().workers_if_single()
bypasses required prefix/seq_hash validation and can return stale active
workers; change the fast-path so that after obtaining workers (via
workers_if_single() or workers()), you explicitly verify the
entry.value().seq_hash (or equivalent prefix-hash field) matches the expected
prefix history for the current scan/count request before returning
workers.len(), and otherwise fall back to the full filtering path that checks
each worker's prefix. Apply the same fix to the other occurrence in the nearby
scan/count code (the block around the 613-631 range) so both fast paths preserve
prefix disambiguation.

In `@model_gateway/benches/radix_tree_benchmark.rs`:
- Line 752: Benchmark uses inconsistent jump sizes: replace the hardcoded 64
passed into build_populated_indexer(...) in the concurrent indexer setup with
the JUMP_SIZE constant so all modes use the same jump size; locate the call to
build_populated_indexer that currently supplies 64 and change that argument to
JUMP_SIZE to make summary/store/match and concurrent indexer setups comparable.

In `@model_gateway/src/policies/cache_aware.rs`:
- Around line 572-590: The hot path calls indexer.worker_id repeatedly for each
candidate (in the filter closure, max_by_key tie-break, and later debug),
causing redundant DashMap lookups and URL hashing; refactor to iterate healthy
workers once and precompute a small struct/tuple (e.g., (idx, wid, score,
tree_size, load)) by calling indexer.worker_id(workers[idx].url()) a single
time, derive score from overlap.scores and tree_size from overlap.tree_sizes
with unwrap_or(0), then use that precomputed tuple for filtering, max_by_key
comparisons, and any debug/logging to eliminate duplicate lookups and reduce
latency.

ℹ️ Review info
Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: e18db09a-ccc3-4fa7-b38d-9b96da9b7bf1

📥 Commits

Reviewing files that changed from the base of the PR and between dbcd323 and 989adfb.

📒 Files selected for processing (3)
  • kv_index/src/event_tree.rs
  • model_gateway/benches/radix_tree_benchmark.rs
  • model_gateway/src/policies/cache_aware.rs

Comment on lines +354 to +361
// Atomically update tree_sizes.
let num_blocks = blocks.len();
self.tree_sizes
.entry(worker_id)
.and_modify(|size| {
size.fetch_add(num_blocks, Ordering::Relaxed);
})
.or_insert(AtomicUsize::new(num_blocks));

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

tree_sizes accounting can drift under replay/idempotent store patterns.

apply_stored always adds blocks.len(), even when wb_ref.insert replaces existing seq_hash entries. This inflates current_size() and routing tie-break tree_sizes. Use absolute synchronization from wb_ref.len() instead of relative increments/decrements.

🛠️ Proposed fix (store absolute per-worker size)
-        drop(wb_ref);
-
-        // Atomically update tree_sizes.
-        let num_blocks = blocks.len();
-        self.tree_sizes
-            .entry(worker_id)
-            .and_modify(|size| {
-                size.fetch_add(num_blocks, Ordering::Relaxed);
-            })
-            .or_insert(AtomicUsize::new(num_blocks));
+        let worker_len = wb_ref.len();
+        drop(wb_ref);
+
+        self.tree_sizes
+            .entry(worker_id)
+            .and_modify(|size| size.store(worker_len, Ordering::Relaxed))
+            .or_insert(AtomicUsize::new(worker_len));
@@
-        let mut num_removed = 0usize;
         for &seq_hash in seq_hashes {
             let Some((position, content_hash, prefix_hash)) = wb_ref.remove(&seq_hash) else {
                 continue;
             };
@@
-            num_removed += 1;
         }
 
-        drop(wb_ref);
-
-        if num_removed > 0 {
-            if let Some(size) = self.tree_sizes.get(&worker_id) {
-                size.fetch_sub(num_removed, Ordering::Relaxed);
-            }
-        }
+        let worker_len = wb_ref.len();
+        drop(wb_ref);
+
+        self.tree_sizes
+            .entry(worker_id)
+            .and_modify(|size| size.store(worker_len, Ordering::Relaxed))
+            .or_insert(AtomicUsize::new(worker_len));

Also applies to: 404-407

🤖 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 354 - 361, apply_stored currently
does a relative increment of self.tree_sizes for worker_id (adding
blocks.len()), which can drift when wb_ref.insert replaces existing seq_hash
entries; instead synchronize atomically to the absolute number of entries in
that worker's writebuffer: read wb_ref.len() after applying inserts/replacements
and set the per-worker AtomicUsize in self.tree_sizes to that absolute value
(use entry(worker_id).and_modify(|size| size.store(wb_len,
Ordering::Relaxed)).or_insert(AtomicUsize::new(wb_len)) or equivalent atomic
update), and make the same change for the other occurrence noted around lines
404-407 so current_size() and routing tie-breaks reflect wb_ref.len() exactly.

Comment on lines +377 to +379
let worker_id = self.intern_worker(worker);

let wb = self.worker_blocks.read();
let Some(level_index) = wb.get(&worker_id) else {
let Some(mut wb_ref) = self.worker_blocks.get_mut(&worker_id) else {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

Unknown-worker remove/clear paths are no longer true no-ops.

Both paths call intern_worker, which allocates new IDs for unknown workers (and apply_cleared may also create empty per-worker state). This introduces silent state growth from invalid/stale events.

🧩 Proposed fix (lookup-only for remove/clear paths)
     pub fn apply_removed(&self, worker: &str, seq_hashes: &[SequenceHash]) {
-        let worker_id = self.intern_worker(worker);
+        let Some(worker_id) = self.worker_id(worker) else {
+            tracing::debug!(
+                worker = %worker,
+                num_hashes = seq_hashes.len(),
+                "apply_removed: worker not tracked, ignoring"
+            );
+            return;
+        };
@@
     fn remove_or_clear_worker(&self, worker: &str, keep_worker: bool) {
-        let worker_id = self.intern_worker(worker);
+        let Some(worker_id) = self.worker_id(worker) else {
+            return;
+        };

Also applies to: 453-454

🤖 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 377 - 379, The remove/clear branches
call intern_worker which allocates new IDs and grows state; change these paths
to perform a lookup-only so unknown workers are treated as no-ops. Concretely,
stop calling intern_worker in the remove/clear code and instead try to find an
existing id without creating one (e.g. check your existing worker-id map or
index for a mapping and only proceed if found), then use that id to access
self.worker_blocks (replace get_mut after intern_worker with a conditional
get/get_mut guarded by the lookup). Apply the same lookup-only change to the
apply_cleared call sites (and the other similar spot around the 453-454 region)
so removals/clears do not create empty per-worker state.

Comment on lines +571 to +573
if let Some(workers) = entry.value().workers_if_single() {
return workers.len();
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🔴 Critical

Single-entry fast path can produce false matches by bypassing prefix-hash checks.

Using workers_if_single() in scan/count paths skips seq_hash validation. A request can hit the same (position, content_hash) with a different prefix history, causing stale active workers to survive and get inflated scores.

🐛 Proposed fix (preserve prefix disambiguation in scan/count)
     fn count_workers_at(
         index: &DashMap<(usize, ContentHash), SeqEntry, FxBuildHasher>,
         position: usize,
         content_hash: ContentHash,
         seq_hashes: &mut Vec<SequenceHash>,
         sequence: &[ContentHash],
     ) -> usize {
         let Some(entry) = index.get(&(position, content_hash)) else {
             return 0;
         };
-        if let Some(workers) = entry.value().workers_if_single() {
-            return workers.len();
-        }
-        // Multi: need rolling hash to disambiguate
         Self::ensure_seq_hash_computed(seq_hashes, position, sequence);
         entry
             .get(seq_hashes[position])
             .map(|workers| workers.len())
             .unwrap_or(0)
     }
@@
-            // Fast path: Single entry — skip rolling hash, use workers directly.
-            if let Some(workers) = entry.value().workers_if_single() {
-                if workers.len() < active.len() {
-                    let mut i = 0;
-                    while i < active.len() {
-                        if workers.contains(&active[i]) {
-                            i += 1;
-                        } else {
-                            internal_scores.insert(active[i], pos as u32);
-                            active.swap_remove(i);
-                        }
-                    }
-                }
-                if early_exit && !active.is_empty() {
-                    break;
-                }
-                continue;
-            }
-
-            // Multi: need rolling hash to disambiguate.
             Self::ensure_seq_hash_computed(seq_hashes, pos, sequence);
             let seq_hash = seq_hashes[pos];

Also applies to: 613-631

🤖 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 571 - 573, The fast-path using
entry.value().workers_if_single() bypasses required prefix/seq_hash validation
and can return stale active workers; change the fast-path so that after
obtaining workers (via workers_if_single() or workers()), you explicitly verify
the entry.value().seq_hash (or equivalent prefix-hash field) matches the
expected prefix history for the current scan/count request before returning
workers.len(), and otherwise fall back to the full filtering path that checks
each worker's prefix. Apply the same fix to the other occurrence in the nearby
scan/count code (the block around the 613-631 range) so both fast paths preserve
prefix disambiguation.

Comment on lines +766 to 767
Self::new(32)
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

Default jump-size docs are stale.

Default now uses Self::new(32), but the struct/function comments still state default is 64. Please align docs with runtime behavior.

🤖 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 766 - 767, The Default
implementation now constructs the type with Self::new(32) but the struct and
function docs still claim the default jump-size is 64; update all relevant
documentation/comments (the struct-level comment and the new()/Default() method
docs in event_tree.rs) to state the default jump-size is 32 and adjust any doc
examples or references that mention 64 so they reflect the runtime behavior
(search for Default, Self::new, and any comments mentioning "jump-size" or "64"
in that file).

const QUERY_BLOCK_COUNTS: [usize; 3] = [32, 128, 512];
const SHARED_PREFIX_BLOCKS: usize = 8;
const JUMP_SIZE: usize = 64;
const JUMP_SIZE: usize = 32;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

Benchmark uses mixed jump sizes across modes.

JUMP_SIZE is now 32 for summary/store/match setup, but concurrent indexer setup still hardcodes 64 (build_populated_indexer(..., 64)). This makes benchmark modes not directly comparable.

♻️ Proposed fix
-macro_rules! bench_indexer_concurrent {
-    ($group:expr, $num_workers:expr, $block_size:expr, $num_threads:expr, $ops_per_thread:expr) => {{
+macro_rules! bench_indexer_concurrent {
+    ($group:expr, $num_workers:expr, $block_size:expr, $jump_size:expr, $num_threads:expr, $ops_per_thread:expr) => {{
@@
-                    let (indexer, worker_chunks) =
-                        build_populated_indexer(&workers, 64, $block_size, 8, 64);
+                    let (indexer, worker_chunks) =
+                        build_populated_indexer(&workers, 64, $block_size, 8, $jump_size);
@@
-        bench_indexer_concurrent!(group, num_workers, 16, NUM_THREADS, OPS_PER_THREAD);
+        bench_indexer_concurrent!(group, num_workers, 16, JUMP_SIZE, NUM_THREADS, OPS_PER_THREAD);
🤖 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 752, Benchmark uses
inconsistent jump sizes: replace the hardcoded 64 passed into
build_populated_indexer(...) in the concurrent indexer setup with the JUMP_SIZE
constant so all modes use the same jump size; locate the call to
build_populated_indexer that currently supplies 64 and change that argument to
JUMP_SIZE to make summary/store/match and concurrent indexer setups comparable.

Comment on lines +572 to +590
.filter(|&idx| {
indexer
.worker_id(workers[idx].url())
.and_then(|id| overlap.scores.get(&id))
.copied()
.unwrap_or(0)
> 0
})
.max_by_key(|&idx| {
let url = workers[idx].url();
let score = overlap.scores.get(url).copied().unwrap_or(0);
let wid = indexer.worker_id(workers[idx].url());
let score = wid
.and_then(|id| overlap.scores.get(&id))
.copied()
.unwrap_or(0);
let load = workers[idx].load();
let tree_size = overlap.tree_sizes.get(url).copied().unwrap_or(0);
let tree_size = wid
.and_then(|id| overlap.tree_sizes.get(&id))
.copied()
.unwrap_or(0);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick | 🔵 Trivial

Avoid repeated worker_id lookups in the event-driven hot path.

This code resolves worker_id multiple times per candidate worker (filter, tie-break, debug). Precomputing (idx, wid, score, tree_size) once per healthy worker will reduce DashMap reads and URL hashing in a latency-sensitive path.

♻️ Proposed refactor
-        let best_idx = healthy_indices
-            .iter()
-            .copied()
-            .filter(|&idx| {
-                indexer
-                    .worker_id(workers[idx].url())
-                    .and_then(|id| overlap.scores.get(&id))
-                    .copied()
-                    .unwrap_or(0)
-                    > 0
-            })
-            .max_by_key(|&idx| {
-                let wid = indexer.worker_id(workers[idx].url());
-                let score = wid
-                    .and_then(|id| overlap.scores.get(&id))
-                    .copied()
-                    .unwrap_or(0);
-                let load = workers[idx].load();
-                let tree_size = wid
-                    .and_then(|id| overlap.tree_sizes.get(&id))
-                    .copied()
-                    .unwrap_or(0);
-                (score, std::cmp::Reverse(load), std::cmp::Reverse(tree_size))
-            })?;
+        let candidates: Vec<(usize, u32, u32, usize)> = healthy_indices
+            .iter()
+            .copied()
+            .filter_map(|idx| {
+                let wid = indexer.worker_id(workers[idx].url())?;
+                let score = overlap.scores.get(&wid).copied().unwrap_or(0);
+                if score == 0 {
+                    return None;
+                }
+                let tree_size = overlap.tree_sizes.get(&wid).copied().unwrap_or(0);
+                Some((idx, wid, score, tree_size))
+            })
+            .collect();
+
+        let (best_idx, _wid, best_score, _tree_size) = candidates
+            .into_iter()
+            .max_by_key(|(idx, _wid, score, tree_size)| {
+                (
+                    *score,
+                    std::cmp::Reverse(workers[*idx].load()),
+                    std::cmp::Reverse(*tree_size),
+                )
+            })?;

Also applies to: 596-599

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@model_gateway/src/policies/cache_aware.rs` around lines 572 - 590, The hot
path calls indexer.worker_id repeatedly for each candidate (in the filter
closure, max_by_key tie-break, and later debug), causing redundant DashMap
lookups and URL hashing; refactor to iterate healthy workers once and precompute
a small struct/tuple (e.g., (idx, wid, score, tree_size, load)) by calling
indexer.worker_id(workers[idx].url()) a single time, derive score from
overlap.scores and tree_size from overlap.tree_sizes with unwrap_or(0), then use
that precomputed tuple for filtering, max_by_key comparisons, and any
debug/logging to eliminate duplicate lookups and reduce latency.

@slin1237 slin1237 closed this Mar 4, 2026
@lightseek-bot
lightseek-bot deleted the slin/flash-indexer-perf branch March 4, 2026 08:09
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

benchmarks Benchmark changes model-gateway Model gateway crate changes needs-rebase PR has merge conflicts that need to be resolved

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant