Skip to content

fix(kv_index): remove the 2048-worker cap and surface KV subscription failures - #1706

Merged
slin1237 merged 1 commit into
mainfrom
fix/kv-index-worker-cap
Jun 14, 2026
Merged

slin1237 merged 1 commit into
mainfrom
fix/kv-index-worker-cap

Conversation

@slin1237

@slin1237 slin1237 commented Jun 12, 2026 •

Copy link
Copy Markdown
Member

Summary

Fixes a live correctness bug found while mining dynamo for designs applicable to the #1685 incident cluster: cache-aware routing silently ignores workers above interned id 2048.

PositionalIndexer::intern_worker asserted id < MAX_WORKERS (2048). The assert fired inside per-worker KV-event subscription tasks, and the resulting panic was discarded (let _ = sub.handle.await), so there was no log, no metric, no restart — the index just stopped covering those workers. Two ways to hit it:

  • Fleet size: the [Performance]: High CPU usage in spawn_worker_collector #1685 reporter runs 4000 workers — roughly half their fleet never feeds the cache-aware index.
  • Churn: interned ids are monotonic and never recycled (a URL keeps its id forever; removals free nothing), so any long-lived gateway with enough pod replacement eventually crosses 2048 regardless of fleet size.

Changes

  • Growable TreeSizes: the fixed 2048-slot Vec<AtomicUsize> becomes segmented, doubling, lazily-allocated storage (OnceLock<Box<[AtomicUsize]>> per segment). Reads stay a lock-free array index on the query hot path — the property the fixed Vec existed for — while covering the entire u32 id space. Memory: 8 B per interned id, ≤2× the high-water id, segments allocated on first write.
  • Non-panicking interning: intern_worker returns Result<WorkerId, WorkerIdExhausted>; the only error is exhausting the u32 id space (detected via a u64 counter rather than wrapping). Nothing is inserted on the error path.
  • Failures surfaced, three layers deep: panics are caught inside the subscription task (catch_unwind) and logged with worker context at the moment they happen (not at worker removal); JoinErrors are logged on both stop paths; intern failure logs and returns cleanly. All paths increment a new smg_kv_event_subscription_failures_total{worker, reason} counter (panic / join_error / intern_failed).

Tests

Segment-math coverage across the full u32 space (boundaries + contiguity), interning 5000 workers with stores/matches probed on both sides of the old cap, id-lifecycle semantics (no recycling) made explicit, no-alloc reset for never-written segments, u32-exhaustion edge, and 4 threads racing interning across the 2048 segment boundary.

cargo clippy --workspace --all-targets --all-features -- -D warnings clean; workspace tests 3557 passed / 0 failed.

Refs #1685. First PR of the wave informed by the dynamo design review.

Summary by CodeRabbit

  • New Features

    • Unbounded worker capacity for event indexing (segmented tree sizing).
    • Added Prometheus metric for KV event subscription failures, labeled by worker and reason.
  • Bug Fixes

    • Worker ID interning is now fallible with explicit exhaustion handling (includes a dedicated exhaustion error).
    • KV event subscription tasks now handle panics and record a failure metric instead of stalling.
    • Shutdown and worker-removal flows now reliably await subscription task completion (with join-error reporting).

@slin1237
slin1237 requested a review from CatherineSue as a code owner June 12, 2026 22:56
@coderabbitai

coderabbitai Bot commented Jun 12, 2026 •

Copy link
Copy Markdown

Review Change Stack

📝 Walkthrough

Walkthrough

Worker ID interning in PositionalIndexer transitions from a fixed u32-based MAX_WORKERS cap to unbounded u64 monotonic allocation with explicit exhaustion handling. Storage shifts from fixed-size vector to segmented TreeSizes with lock-free reads. intern_worker becomes fallible, returning Result<WorkerId, WorkerIdExhausted>. KvEventMonitor gains panic isolation via catch_unwind, task join error tracking, and metrics recording for subscription failures.

Changes

Worker ID Unbounding and Failure Handling

Layer / File(s) Summary
Core unbounding infrastructure: error type and TreeSizes
crates/kv_index/src/event_tree.rs
New WorkerIdExhausted error type and TreeSizes segmented storage structure backed by OnceLock<Box<[AtomicUsize]>> per segment. Imports and constants updated to support unbounded u64-based worker ID interning.
PositionalIndexer refactor: struct, initialization, interning, and tree-size methods
crates/kv_index/src/event_tree.rs
Struct fields transition from Vec<AtomicUsize> to TreeSizes and AtomicU32 to AtomicU64. intern_worker rewritten to return Result<WorkerId, WorkerIdExhausted> with u64 monotonic counter and exhaustion check. Write-path methods (apply_stored, apply_removed, apply_cleared, remove_worker, current_size) and query paths (jump_search_matches) updated to use segmented storage API. Existing tests adapted and new tests validate unbounded behavior and segment boundary correctness.
Public API surface: error type re-export
crates/kv_index/src/lib.rs
WorkerIdExhausted added to public crate re-exports.
Observability infrastructure: KV event subscription failure metrics
model_gateway/src/observability/metrics.rs
New smg_kv_event_subscription_failures_total Prometheus counter registered with worker and reason labels (panic, join_error, intern_failed). record_kv_event_subscription_failure helper method added to record failures.
KvEventMonitor resilience: panic isolation, task join handling, and failure metrics
model_gateway/src/worker/kv_event_monitor.rs
Subscription loop wrapped with catch_unwind to isolate panics and record metrics. on_worker_removed and stop now await task handles and record join_error metrics. subscription_loop treats intern_worker as fallible, logs failures, records intern_failed metric, and exits early. Imports updated with FutureExt and tracing macros. Unit tests updated to unwrap calls.
Caller updates: benchmarks and test suites
crates/kv_index/benches/throughput_bench.rs, model_gateway/benches/radix_tree_benchmark.rs, model_gateway/src/policies/cache_aware.rs
All benchmark and test call sites updated to handle fallible intern_worker by unwrapping the Result.

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~60 minutes

Possibly related PRs

  • lightseekorg/smg#579: Updates to model_gateway/benches/radix_tree_benchmark.rs to handle now-fallible PositionalIndexer::intern_worker(...).unwrap() are directly tied to benchmark suite addition in the same file.
  • lightseekorg/smg#946: Both PRs modify crates/kv_index/src/event_tree.rs around PositionalIndexer::apply_stored tree-size update logic that is refactored in this PR's TreeSizes transition.
  • lightseekorg/smg#584: Both PRs modify worker-ID interning and assignment in crates/kv_index/src/event_tree.rs (intern_worker, next_worker_id, and ID mapping structures).

Suggested reviewers

  • CatherineSue
  • key4ng

Poem

🐰 A worker blooms where once a cap stood tall,
Now boundless IDs march from zero to all,
TreeSizes whisper in segments so keen,
While monitors catch the panics unseen—
Metrics dance, errors found and embraced,
The unbounded kingdom, forever replaced! ✨

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title accurately summarizes the two main changes: removing the 2048-worker cap and surfacing KV subscription failures.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.

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

✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch fix/kv-index-worker-cap

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

@github-actions github-actions Bot added benchmarks Benchmark changes model-gateway Model gateway crate changes kv-index KV index crate changes labels Jun 12, 2026
@claude

claude Bot commented Jun 12, 2026 •

Copy link
Copy Markdown

👋 The PR description doesn't fully follow
PULL_REQUEST_TEMPLATE.md:

  • Missing header: ## Description
  • Missing header: ### Problem
  • Missing header: ### Solution
  • Missing header: ## Test Plan (section exists as ## Tests but the template requires ## Test Plan)

Please update the PR description so reviewers have the context they need.

@chatgpt-codex-connector chatgpt-codex-connector 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.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: ff20ff860c

ℹ️ 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".

Comment on lines +302 to +304
fn segment_len(segment: usize) -> usize {
FIRST_SEGMENT_LEN << segment
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P2 Badge Cap the final tree-size segment length

When the monotonic worker id reaches the final segment (id >= 4,294,965,248, which intern_worker can still return successfully before exhaustion), this calculation returns 2048 << 21 == 4,294,967,296, so the first stored block for any of the last 2048 valid IDs tries to allocate billions of AtomicUsize slots even though only 2048 u32 IDs remain. This turns the intended graceful u32 exhaustion path into an OOM/panic at the end of the id space; the last segment needs to be clamped to the remaining u32-id range.

Useful? React with 👍 / 👎.

@claude claude 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.

Clean fix. The segmented TreeSizes design correctly preserves lock-free array-index reads on the query hot path while removing the 2048-worker cap. The intern_worker error handling is sound, and the catch_unwind + metrics additions give operators visibility into failures that were previously silent. Tests are thorough — covers segment-math boundaries, growth past 2048, id lifecycle, u32 exhaustion, and concurrent interning across segment boundaries. No issues found.

@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 replaces the fixed-size worker limit with a growable, lock-free segmented array (TreeSizes) to support unbounded worker counts, and updates intern_worker to return a Result to handle ID exhaustion safely. The review feedback identifies a critical OOM vulnerability where allocating segments up to u32::MAX can trigger a 32 GiB allocation; capping the segment count to 14 is recommended. Additionally, the feedback points out an O(N) traversal on the query hot path in TreeSizes::total and suggests optimizing it to O(1) by maintaining a global atomic block counter across additions, removals, and resets.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment on lines +54 to +55
/// Number of doubling segments needed to cover the entire u32 worker-id space.
const SEGMENT_COUNT: usize = (u32::BITS - FIRST_SEGMENT_BITS + 1) as usize;

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.

critical

Critical: Exponential Memory Growth and OOM Vulnerability

The current implementation of SEGMENT_COUNT dynamically calculates the number of segments to cover the entire u32 space (which results in 22 segments). However, because segment sizes double with each index, the last segment (segment = 21) has a length of 2^32 entries.

When a worker ID near u32::MAX is interned and used, TreeSizes::slot will attempt to allocate 2^32 elements of AtomicUsize. On a 64-bit target, this requires 32 GiB of contiguous virtual memory for a single segment, which will immediately trigger an Out-Of-Memory (OOM) panic on almost any system. Furthermore, on 32-bit targets, FIRST_SEGMENT_LEN << 21 overflows usize and will panic or wrap.

Since no real-world deployment will ever have more than a few million workers over its lifetime, we should cap SEGMENT_COUNT at a reasonable limit (e.g., 14, which supports up to 33,552,383 workers while capping the maximum segment allocation at a perfectly safe 128 MiB).

Suggested change
/// Number of doubling segments needed to cover the entire u32 worker-id space.
const SEGMENT_COUNT: usize = (u32::BITS - FIRST_SEGMENT_BITS + 1) as usize;
/// Number of doubling segments needed to cover a reasonable worker-id space.
/// Capped at 14 to prevent exponential memory growth and OOM on large IDs.
const SEGMENT_COUNT: usize = 14;
/// Maximum worker ID supported by the capped SEGMENT_COUNT.
/// Equal to (1 << (FIRST_SEGMENT_BITS + SEGMENT_COUNT)) - FIRST_SEGMENT_LEN - 1.
const MAX_WORKER_ID: u32 = 33_552_383;

Comment on lines +277 to +286
struct TreeSizes {
segments: [OnceLock<Box<[AtomicUsize]>>; SEGMENT_COUNT],
}

impl TreeSizes {
fn new() -> Self {
Self {
segments: std::array::from_fn(|_| OnceLock::new()),
}
}

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.

critical

Critical: O(N) Traversal on the Query Hot Path

Currently, TreeSizes::total (which is called by PositionalIndexer::current_size) iterates over all elements of all allocated segments to sum the counters.

Since current_size() is called via has_event_indexer on every single request on the query hot path, this O(N) traversal will completely destroy performance as the number of workers grows (e.g., iterating over millions of elements).

We can optimize this to O(1) by maintaining a single global total_blocks: AtomicUsize counter in TreeSizes. Every time we add or subtract blocks, we update this global counter. When resetting a worker's count, we can use swap(0, ...) to get the previous count and subtract it from the global total.

Suggested change
struct TreeSizes {
segments: [OnceLock<Box<[AtomicUsize]>>; SEGMENT_COUNT],
}
impl TreeSizes {
fn new() -> Self {
Self {
segments: std::array::from_fn(|_| OnceLock::new()),
}
}
struct TreeSizes {
segments: [OnceLock<Box<[AtomicUsize]>>; SEGMENT_COUNT],
total_blocks: AtomicUsize,
}
impl TreeSizes {
fn new() -> Self {
Self {
segments: std::array::from_fn(|_| OnceLock::new()),
total_blocks: AtomicUsize::new(0),
}
}

Comment on lines +293 to +299
fn locate(id: u32) -> (usize, usize) {
let virtual_idx = id as u64 + FIRST_SEGMENT_LEN as u64;
let msb = 63 - virtual_idx.leading_zeros();
let segment = (msb - FIRST_SEGMENT_BITS) as usize;
let offset = (virtual_idx - (1u64 << msb)) as usize;
(segment, offset)
}

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.

high

Add a defensive bounds check to locate to cap the ID at MAX_WORKER_ID. This ensures that TreeSizes is completely panic-safe and will never index out of bounds even if called with arbitrary or un-interned IDs.

Suggested change
fn locate(id: u32) -> (usize, usize) {
let virtual_idx = id as u64 + FIRST_SEGMENT_LEN as u64;
let msb = 63 - virtual_idx.leading_zeros();
let segment = (msb - FIRST_SEGMENT_BITS) as usize;
let offset = (virtual_idx - (1u64 << msb)) as usize;
(segment, offset)
}
#[inline]
fn locate(id: u32) -> (usize, usize) {
let id = id.min(MAX_WORKER_ID);
let virtual_idx = id as u64 + FIRST_SEGMENT_LEN as u64;
let msb = 63 - virtual_idx.leading_zeros();
let segment = (msb - FIRST_SEGMENT_BITS) as usize;
let offset = (virtual_idx - (1u64 << msb)) as usize;
(segment, offset)
}

Comment on lines +327 to +343
/// Reset a worker's count to 0 without allocating its segment if absent.
fn reset(&self, id: u32) {
let (segment, offset) = Self::locate(id);
if let Some(entries) = self.segments[segment].get() {
entries[offset].store(0, Ordering::Relaxed);
}
}

/// Sum of all counters across allocated segments.
fn total(&self) -> usize {
self.segments
.iter()
.filter_map(OnceLock::get)
.flat_map(|entries| entries.iter())
.map(|size| size.load(Ordering::Relaxed))
.sum()
}

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.

high

Implement the O(1) reset and total methods using the new total_blocks counter.

    /// Reset a worker's count to 0 without allocating its segment if absent.
    /// Returns the previous count.
    fn reset(&self, id: u32) -> usize {
        let (segment, offset) = Self::locate(id);
        if let Some(entries) = self.segments[segment].get() {
            entries[offset].swap(0, Ordering::Relaxed)
        } else {
            0
        }
    }

    /// Sum of all counters across allocated segments (O(1) atomic load).
    #[inline]
    fn total(&self) -> usize {
        self.total_blocks.load(Ordering::Relaxed)
    }

Comment on lines 471 to 475
if num_new_blocks > 0 {
self.tree_sizes[worker_id as usize].fetch_add(num_new_blocks, Ordering::Relaxed);
self.tree_sizes
.slot(worker_id)
.fetch_add(num_new_blocks, Ordering::Relaxed);
}

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.

high

Update apply_stored to atomically increment the global total_blocks counter.

        if num_new_blocks > 0 {
            self.tree_sizes
                .slot(worker_id)
                .fetch_add(num_new_blocks, Ordering::Relaxed);
            self.tree_sizes
                .total_blocks
                .fetch_add(num_new_blocks, Ordering::Relaxed);
        }

Comment on lines 513 to 517
if num_removed > 0 {
self.tree_sizes[worker_id as usize].fetch_sub(num_removed, Ordering::Relaxed);
self.tree_sizes
.slot(worker_id)
.fetch_sub(num_removed, Ordering::Relaxed);
}

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.

high

Update apply_removed to atomically decrement the global total_blocks counter.

        if num_removed > 0 {
            self.tree_sizes
                .slot(worker_id)
                .fetch_sub(num_removed, Ordering::Relaxed);
            self.tree_sizes
                .total_blocks
                .fetch_sub(num_removed, Ordering::Relaxed);
        }

}
}
self.tree_sizes[worker_id as usize].store(0, Ordering::Relaxed);
self.tree_sizes.reset(worker_id);

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.

high

Update apply_cleared to subtract the previous worker count from the global total_blocks counter.

        let prev_size = self.tree_sizes.reset(worker_id);
        self.tree_sizes.total_blocks.fetch_sub(prev_size, Ordering::Relaxed);

}
}
self.tree_sizes[worker_id as usize].store(0, Ordering::Relaxed);
self.tree_sizes.reset(worker_id);

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.

high

Update remove_worker to subtract the previous worker count from the global total_blocks counter.

        let prev_size = self.tree_sizes.reset(worker_id);
        self.tree_sizes.total_blocks.fetch_sub(prev_size, Ordering::Relaxed);

Comment on lines +634 to +639
Entry::Vacant(entry) => {
let id = self.next_worker_id.fetch_add(1, Ordering::Relaxed);
let id = u32::try_from(id).map_err(|_| WorkerIdExhausted)?;
entry.insert(id);
Ok(id)
}

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.

high

Add a check in intern_worker to return Err(WorkerIdExhausted) if the assigned ID exceeds MAX_WORKER_ID. Avoid defensive overflow checks like try_from since the ID is guaranteed to be far below u32::MAX.

            Entry::Vacant(entry) => {
                let id = self.next_worker_id.fetch_add(1, Ordering::Relaxed) as u32;
                if id > MAX_WORKER_ID {
                    return Err(WorkerIdExhausted);
                }
                entry.insert(id);
                Ok(id)
            }
References
  1. Do not add defensive overflow checks (e.g., checked_add) when the domain logic guarantees that the values are physically constrained to be far below the type's maximum limit (e.g., vocabulary size vs u32::MAX).

Comment on lines +2251 to +2261
indexer
.next_worker_id
.store(u32::MAX as u64, Ordering::Relaxed);
assert_eq!(indexer.intern_worker("http://last:8000").unwrap(), u32::MAX);
// The id space is now exhausted: new URLs error, known URLs still resolve.
assert_eq!(
indexer.intern_worker("http://one-too-many:8000"),
Err(WorkerIdExhausted)
);
assert_eq!(indexer.intern_worker("http://last:8000").unwrap(), u32::MAX);
assert_eq!(indexer.worker_id("http://last:8000"), Some(u32::MAX));

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

Update the exhaustion test to use MAX_WORKER_ID instead of u32::MAX to align with the new capped limit.

Suggested change
indexer
.next_worker_id
.store(u32::MAX as u64, Ordering::Relaxed);
assert_eq!(indexer.intern_worker("http://last:8000").unwrap(), u32::MAX);
// The id space is now exhausted: new URLs error, known URLs still resolve.
assert_eq!(
indexer.intern_worker("http://one-too-many:8000"),
Err(WorkerIdExhausted)
);
assert_eq!(indexer.intern_worker("http://last:8000").unwrap(), u32::MAX);
assert_eq!(indexer.worker_id("http://last:8000"), Some(u32::MAX));
indexer
.next_worker_id
.store(MAX_WORKER_ID as u64, Ordering::Relaxed);
assert_eq!(indexer.intern_worker("http://last:8000").unwrap(), MAX_WORKER_ID);
// The id space is now exhausted: new URLs error, known URLs still resolve.
assert_eq!(
indexer.intern_worker("http://one-too-many:8000"),
Err(WorkerIdExhausted)
);
assert_eq!(indexer.intern_worker("http://last:8000").unwrap(), MAX_WORKER_ID);
assert_eq!(indexer.worker_id("http://last:8000"), Some(MAX_WORKER_ID));

@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: 2

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@model_gateway/src/observability/metrics.rs`:
- Around line 954-960: The function record_kv_event_subscription_failure
currently accepts free-form reason: &'static str which lets label values drift;
change it to accept a closed enum (e.g. KvSubFailureReason) instead of &'static
str, define explicit variants (e.g. Panic, JoinError, InternFailed) and
implement a stable label mapping (via Display or a method) that returns
canonical strings like "panic", "join_error", "intern_failed", then use that
mapped &str when constructing the counter label inside
record_kv_event_subscription_failure; update callers to pass the enum variants
(or provide conversion helpers) so dashboards/alerts remain stable.

In `@model_gateway/src/worker/kv_event_monitor.rs`:
- Around line 153-180: The panic catch currently wraps the call to
Self::subscription_loop so cleanup code in that task (specifically
indexer.remove_worker(worker_id, worker_blocks)) can be skipped on panic; move
the catch_unwind inside the task scope where worker_id and worker_blocks are in
scope (i.e. inside subscription_loop or a small wrapper run-by-task) so you can
always run cleanup in a finally-like path regardless of panic or normal return;
ensure you still AssertUnwindSafe the inner future, preserve the existing error
logging and Metrics::record_kv_event_subscription_failure on panic, and call
indexer.remove_worker(worker_id, worker_blocks) (or equivalent cleanup) in both
the Err and Ok exit paths before returning.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: c99e6173-a5b5-42e4-b9e6-50ea03b1b6f8

📥 Commits

Reviewing files that changed from the base of the PR and between 09cee1d and ff20ff8.

📒 Files selected for processing (7)
  • crates/kv_index/benches/throughput_bench.rs
  • crates/kv_index/src/event_tree.rs
  • crates/kv_index/src/lib.rs
  • model_gateway/benches/radix_tree_benchmark.rs
  • model_gateway/src/observability/metrics.rs
  • model_gateway/src/policies/cache_aware.rs
  • model_gateway/src/worker/kv_event_monitor.rs

Comment on lines +954 to +960
pub fn record_kv_event_subscription_failure(worker_url: &str, reason: &'static str) {
let worker_interned = intern_string(worker_url);
counter!(
"smg_kv_event_subscription_failures_total",
"worker" => worker_interned,
"reason" => reason
)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Constrain reason to a closed set instead of free-form strings.

reason: &'static str makes label taxonomy easy to drift over time. Use an enum (or shared constants) and map to stable label values (panic|join_error|intern_failed) to keep dashboards/alerts consistent.

♻️ Suggested refactor
+pub enum KvSubscriptionFailureReason {
+    Panic,
+    JoinError,
+    InternFailed,
+}
+
+impl KvSubscriptionFailureReason {
+    const fn as_label(&self) -> &'static str {
+        match self {
+            Self::Panic => "panic",
+            Self::JoinError => "join_error",
+            Self::InternFailed => "intern_failed",
+        }
+    }
+}
+
- pub fn record_kv_event_subscription_failure(worker_url: &str, reason: &'static str) {
+ pub fn record_kv_event_subscription_failure(
+     worker_url: &str,
+     reason: KvSubscriptionFailureReason,
+ ) {
     let worker_interned = intern_string(worker_url);
     counter!(
         "smg_kv_event_subscription_failures_total",
         "worker" => worker_interned,
-        "reason" => reason
+        "reason" => reason.as_label()
     )
     .increment(1);
 }
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
pub fn record_kv_event_subscription_failure(worker_url: &str, reason: &'static str) {
let worker_interned = intern_string(worker_url);
counter!(
"smg_kv_event_subscription_failures_total",
"worker" => worker_interned,
"reason" => reason
)
pub enum KvSubscriptionFailureReason {
Panic,
JoinError,
InternFailed,
}
impl KvSubscriptionFailureReason {
const fn as_label(&self) -> &'static str {
match self {
Self::Panic => "panic",
Self::JoinError => "join_error",
Self::InternFailed => "intern_failed",
}
}
}
pub fn record_kv_event_subscription_failure(
worker_url: &str,
reason: KvSubscriptionFailureReason,
) {
let worker_interned = intern_string(worker_url);
counter!(
"smg_kv_event_subscription_failures_total",
"worker" => worker_interned,
"reason" => reason.as_label()
)
.increment(1);
}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/observability/metrics.rs` around lines 954 - 960, The
function record_kv_event_subscription_failure currently accepts free-form
reason: &'static str which lets label values drift; change it to accept a closed
enum (e.g. KvSubFailureReason) instead of &'static str, define explicit variants
(e.g. Panic, JoinError, InternFailed) and implement a stable label mapping (via
Display or a method) that returns canonical strings like "panic", "join_error",
"intern_failed", then use that mapped &str when constructing the counter label
inside record_kv_event_subscription_failure; update callers to pass the enum
variants (or provide conversion helpers) so dashboards/alerts remain stable.

Comment on lines +153 to +180
// Catch panics here so they surface when they happen — a bare
// JoinError would only be observed at worker removal, leaving the
// index silently frozen for this worker until then.
let result = std::panic::AssertUnwindSafe(Self::subscription_loop(
worker,
worker_url,
indexer,
block_sizes,
loop_model_id,
shutdown_rx,
)
))
.catch_unwind()
.await;
if let Err(payload) = result {
let msg = payload
.downcast_ref::<&str>()
.copied()
.map(String::from)
.or_else(|| payload.downcast_ref::<String>().cloned())
.unwrap_or_else(|| "(non-string panic)".into());
error!(
worker_url = %task_url,
panic.message = %msg,
"KV event subscription task panicked; KV events from this \
worker no longer feed cache-aware routing"
);
Metrics::record_kv_event_subscription_failure(&task_url, "panic");
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

Panic isolation currently skips worker-state cleanup and can leave stale index entries.

Catching panic outside subscription_loop logs the failure, but it bypasses the task’s indexer.remove_worker(worker_id, worker_blocks) cleanup path. Since panic is swallowed, later await returns Ok(()), so lifecycle handlers can’t detect and compensate. This can retain stale blocks for that worker and skew cache-aware routing for the model.

Move panic capture into a scope where worker_id + worker_blocks are available and always run cleanup before exit (panic or normal).

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/worker/kv_event_monitor.rs` around lines 153 - 180, The
panic catch currently wraps the call to Self::subscription_loop so cleanup code
in that task (specifically indexer.remove_worker(worker_id, worker_blocks)) can
be skipped on panic; move the catch_unwind inside the task scope where worker_id
and worker_blocks are in scope (i.e. inside subscription_loop or a small wrapper
run-by-task) so you can always run cleanup in a finally-like path regardless of
panic or normal return; ensure you still AssertUnwindSafe the inner future,
preserve the existing error logging and
Metrics::record_kv_event_subscription_failure on panic, and call
indexer.remove_worker(worker_id, worker_blocks) (or equivalent cleanup) in both
the Err and Ok exit paths before returning.

… failures

PositionalIndexer::intern_worker asserted id < MAX_WORKERS(2048). At
fleets past 2048 workers — or any long-lived gateway with enough worker
churn, since interned ids are monotonic and never recycled — the assert
panicked inside the per-worker KV event subscription task, and the
panic was discarded at await time, so cache-aware routing silently
stopped seeing those workers' KV events.

- Replace the fixed 2048-slot tree_sizes Vec with segmented, doubling,
  lazily-allocated storage (OnceLock<Box<[AtomicUsize]>> per segment):
  reads stay a lock-free array index on the query hot path, worker
  count is unbounded (full u32 id space), memory is 8 bytes per
  interned id and at most 2x the high-water id.
- intern_worker returns Result instead of asserting; the only error is
  u32 id-space exhaustion, detected via a u64 counter instead of
  wrapping. Interning can no longer kill a subscription task.
- Surface subscription failures instead of swallowing them: panics are
  caught inside the task (catch_unwind) and logged with worker context
  when they happen, JoinErrors are logged on the stop paths, and all
  failure modes increment the new
  smg_kv_event_subscription_failures_total counter (reasons: panic,
  join_error, intern_failed).
- Tests: segment math across the full u32 space, interning past the
  old cap with matches probed on both sides of the 2048 boundary,
  id-lifecycle (no recycling) semantics, no-alloc reset for never
  written segments, exhaustion edge, and concurrent interning racing
  across a segment boundary.

Refs #1685.

Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
@slin1237
slin1237 force-pushed the fix/kv-index-worker-cap branch from ff20ff8 to b901a2e Compare June 14, 2026 04:05

@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: 1

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@crates/kv_index/src/event_tree.rs`:
- Around line 302-304: The segment_len function performs a bit-shift operation
that could theoretically overflow on 32-bit platforms, even though the project
targets only 64-bit architectures (x86_64 and aarch64). Add a compile-time guard
using a conditional compilation check (such as a compile_error attribute or cfg
assertion) to prevent compilation on 32-bit target platforms. This defensive
measure should be applied at the module or crate level to ensure the code cannot
be compiled for 32-bit targets, providing extra safety even though the practical
risk is minimal given the documented overflow mitigation in the locate()
function.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: 0349125d-c907-4ac9-bf90-40ce6d350bce

📥 Commits

Reviewing files that changed from the base of the PR and between ff20ff8 and b901a2e.

📒 Files selected for processing (7)
  • crates/kv_index/benches/throughput_bench.rs
  • crates/kv_index/src/event_tree.rs
  • crates/kv_index/src/lib.rs
  • model_gateway/benches/radix_tree_benchmark.rs
  • model_gateway/src/observability/metrics.rs
  • model_gateway/src/policies/cache_aware.rs
  • model_gateway/src/worker/kv_event_monitor.rs

Comment on lines +302 to +304
fn segment_len(segment: usize) -> usize {
FIRST_SEGMENT_LEN << segment
}

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

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
set -euo pipefail

echo "== Manifest target hints =="
fd -a Cargo.toml | while read -r manifest; do
  echo "--- $manifest"
  rg -n "target|i686|armv7|thumbv7|x86_64|aarch64" "$manifest" || true
done

echo "== CI target hints =="
if [ -d .github/workflows ]; then
  rg -n "target|i686|armv7|thumbv7|x86_64|aarch64" .github/workflows || true
fi

Repository: lightseekorg/smg

Length of output: 7598


🏁 Script executed:

# First, view the TreeSizes implementation and segment_len function
head -n 330 crates/kv_index/src/event_tree.rs | tail -n 60

Repository: lightseekorg/smg

Length of output: 2253


🏁 Script executed:

# Find FIRST_SEGMENT_LEN constant definition
rg -n "FIRST_SEGMENT_LEN" crates/kv_index/src/event_tree.rs -B 2 -A 2

Repository: lightseekorg/smg

Length of output: 1501


🏁 Script executed:

# Check if there are any explicit target-width requirements already in the crate
rg -n "target_pointer_width|cfg.*target" crates/kv_index/src/

Repository: lightseekorg/smg

Length of output: 42


🏁 Script executed:

# Check the Cargo.toml for the kv_index crate for any platform-specific settings
cat crates/kv_index/Cargo.toml

Repository: lightseekorg/smg

Length of output: 1162


🏁 Script executed:

# Check segment_len usage to confirm it's internal only
rg -n "segment_len" crates/kv_index/src/event_tree.rs

Repository: lightseekorg/smg

Length of output: 288


🏁 Script executed:

# Check if there are any other platform-specific guards in the project
rg -r "target_pointer_width|cfg.*32|compile_error" crates/ --type rust | head -20

Repository: lightseekorg/smg

Length of output: 3280


Segment length bit-shift on 32-bit platforms is technically overflowed but not a practical concern for this project.

Line 303 computes FIRST_SEGMENT_LEN << segment, which would overflow 32-bit usize on the last segment (e.g., segment 22 → 2^33). However:

  • The project targets only x86_64 and aarch64 (both 64-bit) per CI configuration
  • The function is internal-only, not exposed in public API
  • The overflow risk is already documented in code comments (lines 290–291)
  • The locate() function already mitigates ID overflow via u64 casting at the API boundary

A compile_error guard for 32-bit platforms is good defensive practice, but not critical given the platform constraints. If added, it should be treated as an optional safety improvement rather than a required fix.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@crates/kv_index/src/event_tree.rs` around lines 302 - 304, The segment_len
function performs a bit-shift operation that could theoretically overflow on
32-bit platforms, even though the project targets only 64-bit architectures
(x86_64 and aarch64). Add a compile-time guard using a conditional compilation
check (such as a compile_error attribute or cfg assertion) to prevent
compilation on 32-bit target platforms. This defensive measure should be applied
at the module or crate level to ensure the code cannot be compiled for 32-bit
targets, providing extra safety even though the practical risk is minimal given
the documented overflow mitigation in the locate() function.

@slin1237
slin1237 merged commit 54dc4e4 into main Jun 14, 2026
82 of 84 checks passed
@slin1237
slin1237 deleted the fix/kv-index-worker-cap branch June 14, 2026 06:37
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

benchmarks Benchmark changes kv-index KV index crate changes model-gateway Model gateway crate changes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant