feat(cache-aware): KV-aware imbalance triggers (spread + overload) via load monitor - #1621
Conversation
|
Caution Review failedThe pull request is closed. ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (13)
📝 WalkthroughWalkthroughAdds two token-usage thresholds and wires an optional worker load-snapshot stream into cache-aware routing so CacheAwarePolicy can use KV token-usage (spread and overload) alongside request-count for imbalance detection; exposes thresholds in CLI and Python bindings and updates configs, registry, wiring, and tests. ChangesCache-aware load balancing with token usage thresholds
Sequence Diagram (high-level)sequenceDiagram
participant AppContext
participant WorkerMonitor
participant PolicyRegistry
participant CacheAwarePolicy
participant KvEventMonitor
AppContext->>WorkerMonitor: subscribe() -> load_receiver
AppContext->>PolicyRegistry: set_load_receiver(load_receiver)
AppContext->>PolicyRegistry: register_kv_event_monitor(kv_event_monitor)
PolicyRegistry->>CacheAwarePolicy: set_load_receiver(load_receiver)
PolicyRegistry->>CacheAwarePolicy: inject_kv_event_monitor(kv_event_monitor)
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Code Review
This pull request introduces a new KV-cache utilization threshold (balance_token_usage_threshold) for cache-aware load balancing, allowing the routing policy to trigger rebalancing based on backend token usage (KV-cache saturation) rather than just request counts. The feedback suggests two performance and concurrency improvements: first, cloning the LoadReceiver in max_backend_token_usage to release the RwLock read lock immediately and minimize lock contention on the hot path; second, collecting policies from the DashMap into a temporary Vec in set_load_receiver to avoid holding shard locks while acquiring the write lock on load_rx.
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.
| let guard = self.load_rx.read(); | ||
| let rx = guard.as_ref()?; | ||
| let loads = rx.borrow(); |
There was a problem hiding this comment.
By cloning the LoadReceiver (which is a cheap watch::Receiver clone), we can release the RwLock read lock on self.load_rx immediately. This avoids holding the read lock across the entire loop and the watch::Ref borrow, reducing lock holding time on the hot path.
| let guard = self.load_rx.read(); | |
| let rx = guard.as_ref()?; | |
| let loads = rx.borrow(); | |
| let rx = self.load_rx.read().as_ref()?.clone(); | |
| let loads = rx.borrow(); |
| for entry in self.model_policies.iter() { | ||
| Self::maybe_inject_load_rx(entry.value(), rx.as_ref()); | ||
| } |
There was a problem hiding this comment.
To avoid holding DashMap shard locks during potentially blocking operations, collect the policies into a temporary Vec before iterating. While maybe_inject_load_rx is fast, it acquires a write lock on load_rx which can block if a request thread is currently holding a read lock on the hot path. Collecting the policies first completely avoids holding the DashMap shard locks while waiting for the write lock.
let policies: Vec<_> = self.model_policies.iter().map(|entry| Arc::clone(entry.value())).collect();
for policy in policies {
Self::maybe_inject_load_rx(&policy, rx.as_ref());
}References
- To avoid holding
DashMapshard locks during slow operations, collect the keys and/or values into a temporary collection (e.g.,Vec) before iterating and performing the operations.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: d327d79021
ℹ️ 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".
| // the KV-usage imbalance trigger. `with_worker_monitor` ran | ||
| // earlier in the build chain, so this is already set. | ||
| if let Some(ref worker_monitor) = self.worker_monitor { | ||
| registry.set_load_receiver(Some(worker_monitor.subscribe())); |
There was a problem hiding this comment.
Inject load receiver into PD policies
In prefilling/decoding mode this call runs while building AppContext, but the PD prefill/decode policies are created later in RouterFactory::create_pd_router/create_grpc_pd_router and installed via set_prefill_policy/set_decode_policy, whose setters do not apply the already stored load_rx. As a result, when a PD config uses cache_aware with balance_token_usage_threshold < 1.0, those policies keep load_rx = None and silently fall back to request-count balancing, so the new KV-usage trigger never works in PD routing. The setter path should inject the saved receiver into newly installed PD policies as well.
Useful? React with 👍 / 👎.
| for &idx in healthy_indices { | ||
| if let Some(load) = loads.get(workers[idx].url()) { | ||
| let usage = load.effective_token_usage(); | ||
| max_usage = Some(max_usage.map_or(usage, |m| m.max(usage))); |
There was a problem hiding this comment.
Fall back when any healthy load entry is missing
When load polling fails for one healthy worker, WorkerMonitor prunes that worker's entry but leaves successful entries for the same group; this loop still returns Some(max_usage) as soon as at least one healthy worker has data. With balance_token_usage_threshold < 1.0, is_imbalanced then returns on KV usage and skips the existing request-count check, so a pool with one missing load entry and a large request-count imbalance can keep using cache affinity instead of shortest-queue. Require load data for all healthy workers before using the KV-only branch, otherwise return None so the count fallback runs.
Useful? React with 👍 / 👎.
| /// hottest engine exceeding the ceiling. Otherwise (threshold disabled, or | ||
| /// no snapshot yet) it falls back to the request-count spread, so the policy | ||
| /// is never left without an imbalance signal. | ||
| fn is_imbalanced(&self, workers: &[Arc<dyn Worker>], healthy_indices: &[usize]) -> bool { |
There was a problem hiding this comment.
🟡 Nit: The KV-usage trigger path (balance_token_usage_threshold < 1.0 with a wired LoadReceiver) has zero unit-test coverage. Every test in the diff sets balance_token_usage_threshold: 1.0 (disabled), and no test calls set_load_receiver(Some(...)) with mock load data.
The PR description states "the KV-trigger branch is exercised via the policy's monitor-injection tests," but those tests only exercise the count-based fallback path.
A test like the following would close the gap:
- Create a policy with
balance_token_usage_threshold: 0.8 - Wire a
watch::channelwith aHashMapcontaining a worker at 0.9 token usage - Assert
select_workertriggers the imbalance path (returns the min-load worker) even though request counts are balanced
This is the main new behavior — worth covering to catch regressions.
| }); | ||
| } | ||
|
|
||
| if *balance_token_usage_threshold <= 0.0 { |
There was a problem hiding this comment.
🟡 Nit: f32::NAN <= 0.0 is false in IEEE 754, so NaN passes this validation gate. While the runtime behavior is safe (NaN fails (0.0..1.0).contains() in is_imbalanced, falling through to count mode), it's better to reject invalid input at config time.
Changing the condition to !(*balance_token_usage_threshold > 0.0) would catch NaN, negative zero, and negative values in one check.
| eviction_interval_secs, | ||
| max_tree_size, | ||
| block_size: 16, | ||
| balance_token_usage_threshold: 1.0, |
There was a problem hiding this comment.
🟡 Nit: cache_aware_policy() hardcodes this to 1.0 without exposing it as a parameter. The CLI, YAML, and Python paths all let users configure it, but programmatic RouterConfigBuilder callers (including tests/common/test_config.rs) cannot. Consider adding it as an optional parameter or a chained builder method for parity.
There was a problem hiding this comment.
Clean, well-structured PR that follows existing patterns (mirrors the kv_event_monitor injection). The core logic — KV-usage imbalance trigger with graceful fallback to count-based — is correct. Three 🟡 nits posted inline (test coverage for the KV path, NaN validation, builder parity). None are blocking.
Summary: 0 🔴 Important · 3 🟡 Nit · 0 🟣 Pre-existing
d327d79 to
941d5b4
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
bindings/python/src/lib.rs (1)
471-473:⚠️ Potential issue | 🟠 Major | ⚡ Quick winKeep
_Router(...)positional API stable by appending this new argument at the end.Line 471 explicitly documents that new parameters must be appended to avoid breaking positional callers. Inserting
balance_token_usage_thresholdat Line 775 / Line 886 shifts all following positional arguments and creates a backward-incompatible constructor change.Suggested diff
@@ - block_size = 16, - balance_token_usage_threshold = 1.0, + block_size = 16, max_idle_secs = 14400, @@ - mesh_advertise_host = None, - drain_settle_secs = 5, + mesh_advertise_host = None, + drain_settle_secs = 5, + balance_token_usage_threshold = 1.0, @@ - block_size: usize, - balance_token_usage_threshold: f32, + block_size: usize, max_idle_secs: u64, @@ - mesh_peer_urls: Vec<String>, - mesh_advertise_host: Option<String>, - drain_settle_secs: u64, + mesh_peer_urls: Vec<String>, + mesh_advertise_host: Option<String>, + drain_settle_secs: u64, + balance_token_usage_threshold: f32,Also applies to: 775-776, 886-887
🤖 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 `@bindings/python/src/lib.rs` around lines 471 - 473, The new parameter balance_token_usage_threshold was inserted mid-argument list and breaks the positional API for _Router; move balance_token_usage_threshold so it is appended after drain_settle_secs (i.e., add it to the end of the _Router parameter list and any associated struct/constructor signatures in bindings/python/src/lib.rs), update the trailing-doc comment that enforces append-only parameter additions, and adjust any callers or tests that rely on positional construction to use the new trailing parameter (or switch them to keyword args) so existing positional calls remain compatible.model_gateway/src/config/builder.rs (1)
118-135: 🧹 Nitpick | 🔵 Trivial | ⚡ Quick winExpose
balance_token_usage_thresholdin the builder helper instead of hard-coding1.0.Line 133 hard-disables the KV-usage trigger for all callers that use
RouterConfigBuilder::cache_aware_policy(...). Consider adding a dedicated setter (or helper overload) so programmatic builder users can enable this feature without dropping down to manualPolicyConfigconstruction.🤖 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/config/builder.rs` around lines 118 - 135, The builder currently hard-codes PolicyConfig::CacheAware.balance_token_usage_threshold = 1.0 inside RouterConfigBuilder::cache_aware_policy which prevents callers from enabling the KV-usage trigger; update the API to accept and propagate this value instead of the literal (either by adding a new parameter balance_token_usage_threshold: f32 to cache_aware_policy and assigning it into PolicyConfig::CacheAware, or add a chainable RouterConfigBuilder::balance_token_usage_threshold(f32) setter that stores the value and uses it when cache_aware_policy constructs PolicyConfig::CacheAware); ensure the symbol names mentioned (cache_aware_policy, RouterConfigBuilder, PolicyConfig::CacheAware, balance_token_usage_threshold) are used so callers can set the threshold programmatically.
🤖 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/policies/cache_aware.rs`:
- Around line 299-306: The current logic in the block that computes max_usage
(using healthy_indices, loads, workers, and load.effective_token_usage()) uses
KV token data if any healthy worker has load info, which can mislead when the
snapshot is partial; change it to first verify that loads contains entries for
every healthy worker in healthy_indices (e.g., by checking
loads.get(workers[idx].url()) is Some for all idx), and only then compute and
return the maximum effective_token_usage across all healthy workers; if any
healthy worker is missing from loads, return None (so the caller falls back to
count-based logic) instead of computing with incomplete data.
---
Outside diff comments:
In `@bindings/python/src/lib.rs`:
- Around line 471-473: The new parameter balance_token_usage_threshold was
inserted mid-argument list and breaks the positional API for _Router; move
balance_token_usage_threshold so it is appended after drain_settle_secs (i.e.,
add it to the end of the _Router parameter list and any associated
struct/constructor signatures in bindings/python/src/lib.rs), update the
trailing-doc comment that enforces append-only parameter additions, and adjust
any callers or tests that rely on positional construction to use the new
trailing parameter (or switch them to keyword args) so existing positional calls
remain compatible.
In `@model_gateway/src/config/builder.rs`:
- Around line 118-135: The builder currently hard-codes
PolicyConfig::CacheAware.balance_token_usage_threshold = 1.0 inside
RouterConfigBuilder::cache_aware_policy which prevents callers from enabling the
KV-usage trigger; update the API to accept and propagate this value instead of
the literal (either by adding a new parameter balance_token_usage_threshold: f32
to cache_aware_policy and assigning it into PolicyConfig::CacheAware, or add a
chainable RouterConfigBuilder::balance_token_usage_threshold(f32) setter that
stores the value and uses it when cache_aware_policy constructs
PolicyConfig::CacheAware); ensure the symbol names mentioned
(cache_aware_policy, RouterConfigBuilder, PolicyConfig::CacheAware,
balance_token_usage_threshold) are used so callers can set the threshold
programmatically.
🪄 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: 783fee78-b3c4-4ef0-a929-c4aef16dfe95
📒 Files selected for processing (13)
bindings/python/src/lib.rsbindings/python/src/smg/router_args.pymodel_gateway/src/app_context.rsmodel_gateway/src/config/builder.rsmodel_gateway/src/config/types.rsmodel_gateway/src/config/validation.rsmodel_gateway/src/main.rsmodel_gateway/src/policies/cache_aware.rsmodel_gateway/src/policies/factory.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/policies/registry.rsmodel_gateway/tests/routing/cache_aware_backward_compat_test.rsmodel_gateway/tests/routing/test_pd_routing.rs
| let mut max_usage: Option<f64> = None; | ||
| for &idx in healthy_indices { | ||
| if let Some(load) = loads.get(workers[idx].url()) { | ||
| let usage = load.effective_token_usage(); | ||
| max_usage = Some(max_usage.map_or(usage, |m| m.max(usage))); | ||
| } | ||
| } | ||
| max_usage |
There was a problem hiding this comment.
Require full healthy-worker snapshot coverage before using KV-token imbalance mode.
Lines 299-306 currently return a KV max if any healthy worker has load data. That can mask imbalance when one or more healthy workers are missing from the snapshot. In that partial-snapshot case, fall back to count-based logic instead of using incomplete KV data.
Suggested fix
fn max_backend_token_usage(
&self,
workers: &[Arc<dyn Worker>],
healthy_indices: &[usize],
) -> Option<f64> {
let guard = self.load_rx.read();
let rx = guard.as_ref()?;
let loads = rx.borrow();
let mut max_usage: Option<f64> = None;
for &idx in healthy_indices {
- if let Some(load) = loads.get(workers[idx].url()) {
- let usage = load.effective_token_usage();
- max_usage = Some(max_usage.map_or(usage, |m| m.max(usage)));
- }
+ let load = loads.get(workers[idx].url())?;
+ let usage = load.effective_token_usage();
+ max_usage = Some(max_usage.map_or(usage, |m| m.max(usage)));
}
max_usage
}🤖 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/policies/cache_aware.rs` around lines 299 - 306, The
current logic in the block that computes max_usage (using healthy_indices,
loads, workers, and load.effective_token_usage()) uses KV token data if any
healthy worker has load info, which can mislead when the snapshot is partial;
change it to first verify that loads contains entries for every healthy worker
in healthy_indices (e.g., by checking loads.get(workers[idx].url()) is Some for
all idx), and only then compute and return the maximum effective_token_usage
across all healthy workers; if any healthy worker is missing from loads, return
None (so the caller falls back to count-based logic) instead of computing with
incomplete data.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 941d5b4abb
ℹ️ 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".
| eviction_interval_secs = 120, | ||
| max_tree_size = 2usize.pow(26), | ||
| block_size = 16, | ||
| balance_token_usage_threshold = 1.0, |
There was a problem hiding this comment.
Append the new _Router parameter instead
In the scenario where external Python code calls the native _Router(...) positionally, adding this argument here shifts every later parameter (max_idle_secs, assignment_mode, etc.), so existing callers pass values into the wrong Rust fields or hit type errors. The struct below documents that new parameters must be appended to avoid breaking positional callers, so this should stay as a trailing argument and be populated by keyword/default for backward compatibility.
Useful? React with 👍 / 👎.
| // the KV-usage imbalance trigger. `with_worker_monitor` ran | ||
| // earlier in the build chain, so this is already set. | ||
| if let Some(ref worker_monitor) = self.worker_monitor { | ||
| registry.set_load_receiver(Some(worker_monitor.subscribe())); |
There was a problem hiding this comment.
Wire the receiver for cache-aware PD-only configs
In PD configs where the top-level policy is not cache_aware but prefill_policy or decode_policy is cache_aware with balance_token_usage_threshold < 1.0, the enclosing is_cache_aware check skips this new call, so the registry never stores a load receiver before RouterFactory::create_pd_router/create_grpc_pd_router builds the PD policies. Those cache-aware PD policies then always have load_rx = None and silently fall back to request-count imbalance; the load receiver should be wired whenever any configured main/PD policy is cache-aware, not only when the main policy is.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 1
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (3)
bindings/python/src/lib.rs (1)
471-473:⚠️ Potential issue | 🟠 Major | ⚡ Quick winPreserve positional Python constructor compatibility.
Line 471 explicitly documents append-only constructor growth, but the new argument is inserted mid-signature (Line 775 and Line 886). This shifts subsequent positional arguments and can silently misconfigure existing Python callers using positional
_Router(...).Suggested fix
- block_size = 16, - balance_token_usage_threshold = 1.0, + block_size = 16, max_idle_secs = 14400, ... - mesh_advertise_host = None, + mesh_advertise_host = None, + balance_token_usage_threshold = 1.0, drain_settle_secs = 5,- block_size: usize, - balance_token_usage_threshold: f32, + block_size: usize, max_idle_secs: u64, ... - mesh_advertise_host: Option<String>, + mesh_advertise_host: Option<String>, + balance_token_usage_threshold: f32, drain_settle_secs: u64,Also applies to: 775-887
🤖 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 `@bindings/python/src/lib.rs` around lines 471 - 473, The new parameter was inserted mid-list which breaks positional Python callers of the _Router constructor; move the new field/parameter (drain_settle_secs) to the end of the Router/_Router struct and its constructor/signature (and any factory/new functions that construct _Router) so all existing positional arguments keep their original ordering, and adjust any callers that rely on the new parameter to pass it positionally at the end or via keyword; ensure defaulting/serialization logic for drain_settle_secs is handled where _Router is instantiated (e.g., in the builder/new function and in places that call Router::new/_Router::new) so behavior remains unchanged for existing callers.model_gateway/src/main.rs (1)
922-953:⚠️ Potential issue | 🟠 Major | ⚡ Quick winMap
"bucket"explicitly inparse_policyinstead of silently downgrading.
parse_policyhas no"bucket"arm, so it falls through toRoundRobin(Line 952). This silently ignores user-selected bucket policy and violates the CLI-to-policy contract.Suggested fix
"prefix_hash" => PolicyConfig::PrefixHash { prefix_token_count: self.prefix_token_count, load_factor: self.prefix_hash_load_factor, }, + "bucket" => PolicyConfig::Bucket { + balance_abs_threshold: self.balance_abs_threshold, + balance_rel_threshold: self.balance_rel_threshold, + bucket_adjust_interval_secs: 5, + }, "manual" => PolicyConfig::Manual {🤖 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/main.rs` around lines 922 - 953, parse_policy currently lacks a "bucket" match arm and silently falls back to PolicyConfig::RoundRobin; add a "bucket" => PolicyConfig::Bucket { ... } arm in the parse_policy function so the user's "bucket" selection is honored, and populate the Bucket struct fields from this parser's config fields (use the struct's bucket-related members such as self.bucket_count, self.bucket_size, self.eviction_interval, self.max_idle_secs or other bucket_* fields you have) similar to how "cache_aware" and "manual" arms use self.*; ensure the new arm returns PolicyConfig::Bucket and lives alongside the other match arms in parse_policy.model_gateway/src/app_context.rs (1)
662-677:⚠️ Potential issue | 🟠 Major | ⚡ Quick winWire the load snapshot independently of the default policy type.
registry.set_load_receiver(...)only runs whenconfig.policyis the defaultcache_awarepolicy. If the router default is something else but a model later requestscache_awarevia policy hint,PolicyRegistry::create_policy_from_type()never sees a stored receiver and that policy silently stays count-only even withbalance_token_usage_threshold < 1.0.Suggested fix
fn with_kv_event_monitor(mut self, config: &RouterConfig) -> Self { use crate::config::types::PolicyConfig; + if let (Some(registry), Some(worker_monitor)) = (&self.policy_registry, &self.worker_monitor) { + registry.set_load_receiver(Some(worker_monitor.subscribe())); + } + let is_cache_aware = matches!(config.policy, PolicyConfig::CacheAware { .. }); if is_cache_aware { let monitor = Arc::new(KvEventMonitor::new(None)); debug!("Created KV event monitor for event-driven cache-aware routing"); // Inject monitor into PolicyRegistry — propagates to default_policy // and any other existing cache-aware policies. if let Some(ref registry) = self.policy_registry { registry.set_kv_event_monitor(Some(Arc::clone(&monitor))); - if let Some(ref worker_monitor) = self.worker_monitor { - registry.set_load_receiver(Some(worker_monitor.subscribe())); - } } self.kv_event_monitor = Some(monitor); } self }🤖 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/app_context.rs` around lines 662 - 677, The code only calls registry.set_load_receiver(...) inside the is_cache_aware branch, so when the router default policy is non-cache-aware but per-model hints request cache-aware, PolicyRegistry::create_policy_from_type() never sees a load receiver and new cache-aware policies remain count-only; change the initialization to always wire the load receiver when self.policy_registry and self.worker_monitor exist (i.e., call registry.set_load_receiver(Some(worker_monitor.subscribe())) outside/independent of the is_cache_aware conditional), while keeping registry.set_kv_event_monitor(Some(Arc::clone(&monitor))) limited to the is_cache_aware branch so monitors are only created when needed.
♻️ Duplicate comments (1)
model_gateway/src/policies/cache_aware.rs (1)
291-306:⚠️ Potential issue | 🟠 Major | ⚡ Quick winRequire full healthy-worker coverage before using KV-token imbalance mode.
max_backend_token_usage()still returns a value when only a subset ofhealthy_indicesexists in the watch snapshot. That makesis_imbalanced()switch to token-usage mode on partial data instead of falling back to request-count mode, so a missing healthy worker can hide or spuriously trigger rebalancing.Suggested fix
fn max_backend_token_usage( &self, workers: &[Arc<dyn Worker>], healthy_indices: &[usize], ) -> Option<f64> { let guard = self.load_rx.read(); let rx = guard.as_ref()?; let loads = rx.borrow(); let mut max_usage: Option<f64> = None; for &idx in healthy_indices { - if let Some(load) = loads.get(workers[idx].url()) { - let usage = load.effective_token_usage(); - max_usage = Some(max_usage.map_or(usage, |m| m.max(usage))); - } + let load = loads.get(workers[idx].url())?; + let usage = load.effective_token_usage(); + max_usage = Some(max_usage.map_or(usage, |m| m.max(usage))); } max_usage }🤖 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/policies/cache_aware.rs` around lines 291 - 306, The current max_backend_token_usage() can return a partial value when some healthy workers are missing from the watch snapshot; change it so it requires full coverage: before computing max, verify that every workers[idx].url() for each idx in healthy_indices exists in loads and if any are missing return None so callers (e.g., is_imbalanced()) fall back to request-count mode; implement this by iterating healthy_indices, checking loads.get(workers[idx].url()) and immediately returning None on missing entries, otherwise compute and return the max effective_token_usage() as before.
🤖 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/config/validation.rs`:
- Around line 284-290: The validation for balance_token_usage_threshold in the
function that returns ConfigError::InvalidValue currently only checks <= 0.0 and
thus accepts NaN/inf; update the condition to also require the value is finite
(e.g., use !balance_token_usage_threshold.is_finite() ||
*balance_token_usage_threshold <= 0.0) so NaN/inf are rejected and return the
same ConfigError with field "balance_token_usage_threshold", value set from
balance_token_usage_threshold.to_string(), and the existing reason message.
---
Outside diff comments:
In `@bindings/python/src/lib.rs`:
- Around line 471-473: The new parameter was inserted mid-list which breaks
positional Python callers of the _Router constructor; move the new
field/parameter (drain_settle_secs) to the end of the Router/_Router struct and
its constructor/signature (and any factory/new functions that construct _Router)
so all existing positional arguments keep their original ordering, and adjust
any callers that rely on the new parameter to pass it positionally at the end or
via keyword; ensure defaulting/serialization logic for drain_settle_secs is
handled where _Router is instantiated (e.g., in the builder/new function and in
places that call Router::new/_Router::new) so behavior remains unchanged for
existing callers.
In `@model_gateway/src/app_context.rs`:
- Around line 662-677: The code only calls registry.set_load_receiver(...)
inside the is_cache_aware branch, so when the router default policy is
non-cache-aware but per-model hints request cache-aware,
PolicyRegistry::create_policy_from_type() never sees a load receiver and new
cache-aware policies remain count-only; change the initialization to always wire
the load receiver when self.policy_registry and self.worker_monitor exist (i.e.,
call registry.set_load_receiver(Some(worker_monitor.subscribe()))
outside/independent of the is_cache_aware conditional), while keeping
registry.set_kv_event_monitor(Some(Arc::clone(&monitor))) limited to the
is_cache_aware branch so monitors are only created when needed.
In `@model_gateway/src/main.rs`:
- Around line 922-953: parse_policy currently lacks a "bucket" match arm and
silently falls back to PolicyConfig::RoundRobin; add a "bucket" =>
PolicyConfig::Bucket { ... } arm in the parse_policy function so the user's
"bucket" selection is honored, and populate the Bucket struct fields from this
parser's config fields (use the struct's bucket-related members such as
self.bucket_count, self.bucket_size, self.eviction_interval, self.max_idle_secs
or other bucket_* fields you have) similar to how "cache_aware" and "manual"
arms use self.*; ensure the new arm returns PolicyConfig::Bucket and lives
alongside the other match arms in parse_policy.
---
Duplicate comments:
In `@model_gateway/src/policies/cache_aware.rs`:
- Around line 291-306: The current max_backend_token_usage() can return a
partial value when some healthy workers are missing from the watch snapshot;
change it so it requires full coverage: before computing max, verify that every
workers[idx].url() for each idx in healthy_indices exists in loads and if any
are missing return None so callers (e.g., is_imbalanced()) fall back to
request-count mode; implement this by iterating healthy_indices, checking
loads.get(workers[idx].url()) and immediately returning None on missing entries,
otherwise compute and return the max effective_token_usage() as before.
🪄 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: 14ea9365-786c-4c85-ab97-86091b613fd5
📒 Files selected for processing (13)
bindings/python/src/lib.rsbindings/python/src/smg/router_args.pymodel_gateway/src/app_context.rsmodel_gateway/src/config/builder.rsmodel_gateway/src/config/types.rsmodel_gateway/src/config/validation.rsmodel_gateway/src/main.rsmodel_gateway/src/policies/cache_aware.rsmodel_gateway/src/policies/factory.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/policies/registry.rsmodel_gateway/tests/routing/cache_aware_backward_compat_test.rsmodel_gateway/tests/routing/test_pd_routing.rs
| if *balance_token_usage_threshold <= 0.0 { | ||
| return Err(ConfigError::InvalidValue { | ||
| field: "balance_token_usage_threshold".to_string(), | ||
| value: balance_token_usage_threshold.to_string(), | ||
| reason: "Must be > 0.0 (use >= 1.0 to disable)".to_string(), | ||
| }); | ||
| } |
There was a problem hiding this comment.
Reject non-finite threshold values to avoid silent fallback behavior.
<= 0.0 catches negative/zero values, but NaN/inf still pass and can bypass token-usage mode unexpectedly. Add an is_finite() check here.
Proposed fix
- if *balance_token_usage_threshold <= 0.0 {
+ if !balance_token_usage_threshold.is_finite()
+ || *balance_token_usage_threshold <= 0.0
+ {
return Err(ConfigError::InvalidValue {
field: "balance_token_usage_threshold".to_string(),
value: balance_token_usage_threshold.to_string(),
- reason: "Must be > 0.0 (use >= 1.0 to disable)".to_string(),
+ reason:
+ "Must be finite and > 0.0 (use >= 1.0 to disable)".to_string(),
});
}📝 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.
| if *balance_token_usage_threshold <= 0.0 { | |
| return Err(ConfigError::InvalidValue { | |
| field: "balance_token_usage_threshold".to_string(), | |
| value: balance_token_usage_threshold.to_string(), | |
| reason: "Must be > 0.0 (use >= 1.0 to disable)".to_string(), | |
| }); | |
| } | |
| if !balance_token_usage_threshold.is_finite() | |
| || *balance_token_usage_threshold <= 0.0 | |
| { | |
| return Err(ConfigError::InvalidValue { | |
| field: "balance_token_usage_threshold".to_string(), | |
| value: balance_token_usage_threshold.to_string(), | |
| reason: | |
| "Must be finite and > 0.0 (use >= 1.0 to disable)".to_string(), | |
| }); | |
| } |
🤖 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/config/validation.rs` around lines 284 - 290, The
validation for balance_token_usage_threshold in the function that returns
ConfigError::InvalidValue currently only checks <= 0.0 and thus accepts NaN/inf;
update the condition to also require the value is finite (e.g., use
!balance_token_usage_threshold.is_finite() || *balance_token_usage_threshold <=
0.0) so NaN/inf are rejected and return the same ConfigError with field
"balance_token_usage_threshold", value set from
balance_token_usage_threshold.to_string(), and the existing reason message.
ekzhang
left a comment
There was a problem hiding this comment.
Thank you so much!! Very excited for this, just a minor suggestion since I believe some of our workloads also benefit from the batch-size / concurrent requests balancing :D
| // Token-usage mode: configured (< 1.0) AND a snapshot exists → use the | ||
| // KV signal alone, ignoring request counts. A missing snapshot falls | ||
| // through to the count check below (never left blind). | ||
| let kv_threshold = self.config.balance_token_usage_threshold; | ||
| if (0.0..1.0).contains(&kv_threshold) { | ||
| if let Some(max_usage) = self.max_backend_token_usage(workers, healthy_indices) { | ||
| return max_usage > f64::from(kv_threshold); | ||
| } | ||
| } |
There was a problem hiding this comment.
Could it be more flexible to check based on either the token usage threshold, or the existing --balance-abs-threshold? Reason being that if it's implemented this way, then people can also set the balance abs threshold independently, and you can recover the "sole" token usage behavior by setting --balance-abs-threshold 10000000
| // Token-usage mode: configured (< 1.0) AND a snapshot exists → use the | |
| // KV signal alone, ignoring request counts. A missing snapshot falls | |
| // through to the count check below (never left blind). | |
| let kv_threshold = self.config.balance_token_usage_threshold; | |
| if (0.0..1.0).contains(&kv_threshold) { | |
| if let Some(max_usage) = self.max_backend_token_usage(workers, healthy_indices) { | |
| return max_usage > f64::from(kv_threshold); | |
| } | |
| } | |
| // Token-usage mode: configured (< 1.0) AND a snapshot exists → use the | |
| // KV signal additionally to check for imbalance. | |
| let kv_threshold = self.config.balance_token_usage_threshold; | |
| if (0.0..1.0).contains(&kv_threshold) { | |
| if let Some(max_usage) = self.max_backend_token_usage(workers, healthy_indices) { | |
| if max_usage > f64::from(kv_threshold) { | |
| return true; | |
| } | |
| } | |
| } |
941d5b4 to
83e4e40
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 83e4e40ed8
ℹ️ 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 max_usage > f64::from(self.config.overload_token_usage_threshold) { | ||
| return true; |
There was a problem hiding this comment.
Route KV-triggered rebalances by token usage
When this KV overload/spread branch returns true, the caller falls through to select_worker_min_load, which only compares local in-flight request counts. In the common case where counts are tied or the overloaded backend has the same small count as others, min_by_key can pick the first healthy worker even if that worker is the one above overload_token_usage_threshold, so the new "shed load off that engine" path can continue sending requests to the hot KV backend. The KV-triggered path needs to carry the coldest/non-overloaded candidate through to selection instead of reducing the signal to a boolean.
Useful? React with 👍 / 👎.
| balance_token_usage_threshold: f32, | ||
| /// Backend KV-utilization ceiling (0.0–1.0): a single engine above it | ||
| /// triggers shedding regardless of spread. `>= 1.0` disables (default). | ||
| #[serde(default = "default_balance_token_usage_threshold")] |
There was a problem hiding this comment.
🟡 Nit: overload_token_usage_threshold reuses default_balance_token_usage_threshold as its serde default. Both happen to be 1.0 today, but if someone later changes the balance default they'll silently change the overload default too. A one-line default_overload_token_usage_threshold function eliminates the coupling.
| #[serde(default = "default_balance_token_usage_threshold")] | |
| #[serde(default = "default_overload_token_usage_threshold")] |
(and add fn default_overload_token_usage_threshold() -> f32 { 1.0 } next to the existing default function)
|
|
||
| // Count spread (abs AND rel) over healthy workers. | ||
| let (min_load, max_load) = | ||
| healthy_indices |
There was a problem hiding this comment.
🟡 Nit: Subtle behavioral change — the old code computed (min_load, max_load) over all workers (including unhealthy), while this now scopes to healthy_indices only. The new behavior is arguably more correct (unhealthy workers at load 0 were artificially inflating the spread), but it makes the count-based trigger harder to trip in pools with downed workers. Worth a callout in the PR description so reviewers can sign off on the changed semantics intentionally.
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
model_gateway/src/policies/factory.rs (1)
157-175: 🧹 Nitpick | 🔵 Trivial | ⚡ Quick winAdd coverage for the new
least_loadname aliases intest_create_by_name.The new alias path is untested, so regressions there won’t be caught.
✅ Suggested test additions
assert!(PolicyFactory::create_by_name("power_of_two").is_some()); assert!(PolicyFactory::create_by_name("PowerOfTwo").is_some()); + assert!(PolicyFactory::create_by_name("least_load").is_some()); + assert!(PolicyFactory::create_by_name("LeastLoad").is_some()); assert!(PolicyFactory::create_by_name("cache_aware").is_some()); assert!(PolicyFactory::create_by_name("CacheAware").is_some());🤖 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/policies/factory.rs` around lines 157 - 175, The test test_create_by_name is missing assertions for the new least_load alias; update the test to cover PolicyFactory::create_by_name("least_load") and the capitalized form PolicyFactory::create_by_name("LeastLoad") (both should be is_some()) so the new alias is exercised and future regressions are caught.model_gateway/src/app_context.rs (1)
664-677:⚠️ Potential issue | 🟠 Major | ⚡ Quick winLoad receiver wiring is incorrectly gated by the main policy type.
set_load_receiver()only runs whenconfig.policyiscache_aware, so cache-aware policies created later via policy hints/PD overrides won't receive snapshots and cannot trigger token-usage balancing.💡 Proposed fix
fn with_kv_event_monitor(mut self, config: &RouterConfig) -> Self { use crate::config::types::PolicyConfig; + // Always wire load snapshots when both components exist; only cache-aware + // policies consume it. + if let (Some(registry), Some(worker_monitor)) = (&self.policy_registry, &self.worker_monitor) { + registry.set_load_receiver(Some(worker_monitor.subscribe())); + } + let is_cache_aware = matches!(config.policy, PolicyConfig::CacheAware { .. }); if is_cache_aware { let monitor = Arc::new(KvEventMonitor::new(None)); debug!("Created KV event monitor for event-driven cache-aware routing"); @@ if let Some(ref registry) = self.policy_registry { registry.set_kv_event_monitor(Some(Arc::clone(&monitor))); - // Wire the backend load snapshot so cache-aware policies can use - // the KV-usage imbalance trigger. `with_worker_monitor` ran - // earlier in the build chain, so this is already set. - if let Some(ref worker_monitor) = self.worker_monitor { - registry.set_load_receiver(Some(worker_monitor.subscribe())); - } } self.kv_event_monitor = Some(monitor); }🤖 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/app_context.rs` around lines 664 - 677, The current code only calls registry.set_load_receiver(...) when is_cache_aware is true, so policies created later (via hints/PD overrides) won't receive backend load snapshots; change the wiring so that if self.policy_registry.is_some() and self.worker_monitor.is_some() you always call registry.set_load_receiver(Some(worker_monitor.subscribe())), regardless of is_cache_aware, while keeping the creation/injection of KvEventMonitor (Arc::new(KvEventMonitor::new(None))) and registry.set_kv_event_monitor(...) gated by is_cache_aware; update references to policy_registry, worker_monitor, set_kv_event_monitor, and set_load_receiver accordingly so the load receiver is set whenever a worker_monitor exists.
🤖 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.
Outside diff comments:
In `@model_gateway/src/app_context.rs`:
- Around line 664-677: The current code only calls
registry.set_load_receiver(...) when is_cache_aware is true, so policies created
later (via hints/PD overrides) won't receive backend load snapshots; change the
wiring so that if self.policy_registry.is_some() and
self.worker_monitor.is_some() you always call
registry.set_load_receiver(Some(worker_monitor.subscribe())), regardless of
is_cache_aware, while keeping the creation/injection of KvEventMonitor
(Arc::new(KvEventMonitor::new(None))) and registry.set_kv_event_monitor(...)
gated by is_cache_aware; update references to policy_registry, worker_monitor,
set_kv_event_monitor, and set_load_receiver accordingly so the load receiver is
set whenever a worker_monitor exists.
In `@model_gateway/src/policies/factory.rs`:
- Around line 157-175: The test test_create_by_name is missing assertions for
the new least_load alias; update the test to cover
PolicyFactory::create_by_name("least_load") and the capitalized form
PolicyFactory::create_by_name("LeastLoad") (both should be is_some()) so the new
alias is exercised and future regressions are caught.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: c97312f6-51dd-4f59-89b4-1c88d0093743
📒 Files selected for processing (13)
bindings/python/src/lib.rsbindings/python/src/smg/router_args.pymodel_gateway/src/app_context.rsmodel_gateway/src/config/builder.rsmodel_gateway/src/config/types.rsmodel_gateway/src/config/validation.rsmodel_gateway/src/main.rsmodel_gateway/src/policies/cache_aware.rsmodel_gateway/src/policies/factory.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/policies/registry.rsmodel_gateway/tests/routing/cache_aware_backward_compat_test.rsmodel_gateway/tests/routing/test_pd_routing.rs
👮 Files not reviewed due to content moderation or server errors (6)
- bindings/python/src/lib.rs
- bindings/python/src/smg/router_args.py
- model_gateway/src/main.rs
- model_gateway/src/config/builder.rs
- model_gateway/src/policies/mod.rs
- model_gateway/src/config/validation.rs
|
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 |
83e4e40 to
d6b89c4
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
♻️ Duplicate comments (1)
model_gateway/src/policies/cache_aware.rs (1)
310-329:⚠️ Potential issue | 🟠 Major | ⚡ Quick winRequire full healthy-worker snapshot coverage before applying KV imbalance triggers.
Line 320 currently skips missing healthy-worker entries and still computes KV bounds from a partial snapshot. This can incorrectly trigger or suppress imbalance. For any missing healthy worker, return
Nonesois_imbalancedfalls back to count-based logic.Suggested fix
fn backend_token_usage_bounds( &self, workers: &[Arc<dyn Worker>], healthy_indices: &[usize], ) -> Option<(f64, f64)> { let guard = self.load_rx.read(); let rx = guard.as_ref()?; let loads = rx.borrow(); let mut bounds: Option<(f64, f64)> = None; for &idx in healthy_indices { - if let Some(load) = loads.get(workers[idx].url()) { - let usage = load.effective_token_usage(); - bounds = Some(match bounds { - Some((min, max)) => (min.min(usage), max.max(usage)), - None => (usage, usage), - }); - } + let load = loads.get(workers[idx].url())?; + let usage = load.effective_token_usage(); + bounds = Some(match bounds { + Some((min, max)) => (min.min(usage), max.max(usage)), + None => (usage, usage), + }); } bounds }🤖 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/policies/cache_aware.rs` around lines 310 - 329, In backend_token_usage_bounds, if any healthy worker index is missing from the KV snapshot we must abort and return None instead of computing bounds from a partial view; change the loop in backend_token_usage_bounds (which reads load_rx, gets rx.borrow(), iterates healthy_indices and calls loads.get(workers[idx].url())) so that when loads.get(...) is None for any healthy index the function immediately returns None; otherwise continue computing the min/max usage as before and return the bounds.
🤖 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 `@bindings/python/src/lib.rs`:
- Around line 787-791: The Router Python constructor parameters were inserted in
the middle (balance_token_usage_threshold, overload_token_usage_threshold,
least_load_kv_pressure_weight, least_load_default_throughput,
least_load_mean_prefill_tokens), which shifts existing positional arguments; fix
by moving these new parameters to the end of the Router Python constructor
signature(s) (preserve their default values) so existing positional callers keep
binding the same arguments—update both occurrences of the Router constructor
implementation to append these params rather than inserting them in the middle.
---
Duplicate comments:
In `@model_gateway/src/policies/cache_aware.rs`:
- Around line 310-329: In backend_token_usage_bounds, if any healthy worker
index is missing from the KV snapshot we must abort and return None instead of
computing bounds from a partial view; change the loop in
backend_token_usage_bounds (which reads load_rx, gets rx.borrow(), iterates
healthy_indices and calls loads.get(workers[idx].url())) so that when
loads.get(...) is None for any healthy index the function immediately returns
None; otherwise continue computing the min/max usage as before and return the
bounds.
🪄 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: 3f410fc7-7c39-4c2e-84bc-e9a5707814e3
📒 Files selected for processing (13)
bindings/python/src/lib.rsbindings/python/src/smg/router_args.pymodel_gateway/src/app_context.rsmodel_gateway/src/config/builder.rsmodel_gateway/src/config/types.rsmodel_gateway/src/config/validation.rsmodel_gateway/src/main.rsmodel_gateway/src/policies/cache_aware.rsmodel_gateway/src/policies/factory.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/policies/registry.rsmodel_gateway/tests/routing/cache_aware_backward_compat_test.rsmodel_gateway/tests/routing/test_pd_routing.rs
| balance_token_usage_threshold = 1.0, | ||
| overload_token_usage_threshold = 1.0, | ||
| least_load_kv_pressure_weight = 0.15, | ||
| least_load_default_throughput = 2000.0, | ||
| least_load_mean_prefill_tokens = 1024, |
There was a problem hiding this comment.
Preserve Python positional constructor compatibility.
Lines 787-791 and Lines 902-906 insert new parameters in the middle of Router’s Python constructor argument list, which shifts all following positional arguments (for example, previous positional callers passing max_idle_secs now bind that value to balance_token_usage_threshold). This is a breaking cross-language API contract change.
💡 Suggested fix (append new args to the end of the constructor signature)
- block_size = 16,
- balance_token_usage_threshold = 1.0,
- overload_token_usage_threshold = 1.0,
- least_load_kv_pressure_weight = 0.15,
- least_load_default_throughput = 2000.0,
- least_load_mean_prefill_tokens = 1024,
+ block_size = 16,
max_idle_secs = 14400,
...
- drain_settle_secs = 5,
+ drain_settle_secs = 5,
+ balance_token_usage_threshold = 1.0,
+ overload_token_usage_threshold = 1.0,
+ least_load_kv_pressure_weight = 0.15,
+ least_load_default_throughput = 2000.0,
+ least_load_mean_prefill_tokens = 1024,- block_size: usize,
- balance_token_usage_threshold: f32,
- overload_token_usage_threshold: f32,
- least_load_kv_pressure_weight: f64,
- least_load_default_throughput: f64,
- least_load_mean_prefill_tokens: u32,
+ block_size: usize,
max_idle_secs: u64,
...
- drain_settle_secs: u64,
+ drain_settle_secs: u64,
+ balance_token_usage_threshold: f32,
+ overload_token_usage_threshold: f32,
+ least_load_kv_pressure_weight: f64,
+ least_load_default_throughput: f64,
+ least_load_mean_prefill_tokens: u32,Also applies to: 902-906
🤖 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 `@bindings/python/src/lib.rs` around lines 787 - 791, The Router Python
constructor parameters were inserted in the middle
(balance_token_usage_threshold, overload_token_usage_threshold,
least_load_kv_pressure_weight, least_load_default_throughput,
least_load_mean_prefill_tokens), which shifts existing positional arguments; fix
by moving these new parameters to the end of the Router Python constructor
signature(s) (preserve their default values) so existing positional callers keep
binding the same arguments—update both occurrences of the Router constructor
implementation to append these params rather than inserting them in the middle.
…a load monitor Cache-aware routing decided imbalance — whether to abandon cache affinity for shortest-queue — solely from in-flight request-count spread (w.load()). For long-context workloads (e.g. Kimi K2.5) a few requests saturate one engine's KV without tripping the count thresholds, and the count signal is gateway-local so its effective threshold drifts with the number of router replicas. Consume the backend token_usage signal (0.0-1.0 KV utilization, already polled by WorkerMonitor and used by power_of_two) via the shared subscribe() watch channel, wired through PolicyRegistry like the kv_event_monitor injection. is_imbalanced now ORs three independent triggers: - overload (overload_token_usage_threshold): the hottest engine's KV exceeds a ceiling -> shed a critically-saturated engine regardless of balance. Safety valve, best set high (~0.9). - KV spread (balance_token_usage_threshold): hottest minus coldest KV exceeds the threshold -> a genuine balance signal (not a bare max), and being backend-reported it is invariant to gateway replica count. - count spread: request-count dispersion (abs AND rel) over healthy workers, always evaluated, so high-count / low-KV imbalance is still caught. Both KV knobs default to 1.0 (disabled: utilization and spread are both <= 1.0), so default behavior is unchanged. Exposed via CLI/YAML, the PyO3 binding, and the Python router_args. The rebalance target (select_worker_min_load) still picks minimum request-count; making it KV-aware (spill toward the coldest engine) is a separate follow-up. Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com> Co-authored-by: Sam Shleifer <6045025+sshleifer@users.noreply.github.com>
d6b89c4 to
66d7903
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 66d79035fc
ℹ️ 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".
| // the KV-usage imbalance trigger. `with_worker_monitor` ran | ||
| // earlier in the build chain, so this is already set. | ||
| if let Some(ref worker_monitor) = self.worker_monitor { | ||
| registry.set_load_receiver(Some(worker_monitor.subscribe())); |
There was a problem hiding this comment.
Include cache-aware KV triggers in load polling
When the router runs only cache_aware with balance_token_usage_threshold or overload_token_usage_threshold enabled, this receiver stays attached to an empty watch map because WorkerMonitor::group_monitor_loop skips polling unless get_all_load_aware_policies() is non-empty, and that helper still classifies only power_of_two/least_load as load-aware. In that common cache-aware-only configuration, no /loads/GetLoads snapshots are ever fetched, so backend_token_usage_bounds() returns None and the new KV-usage triggers silently never take effect.
Useful? React with 👍 / 👎.
Description
Problem
Cache-aware routing decides imbalance — whether to abandon cache affinity and fall back to shortest-queue — solely from in-flight request-count spread (
w.load()). Two failure modes:balance_abs_threshold(64) /balance_rel_threshold(1.5) gate. Queue time piles onto the hot engine while others sit idle.w.load()is per-gateway in-flight count. With N router replicas each sees ~1/N of the true load, so the effective imbalance threshold silently scales with replica count.Solution
Consume the backend
token_usagesignal (0.0–1.0 KV utilization, already polled byWorkerMonitorand used bypower_of_two) via the sharedsubscribe()watch channel, wired intoCacheAwarePolicythroughPolicyRegistrylike the existingkv_event_monitorinjection. Imbalance detection becomes a 3-term OR that adds KV-awareness without replacing the existing request-count logic.How cache-aware routing works — and when it falls back
For every request,
select_workerruns in this order:Normal path (balanced pool). The policy scores how much of the request's prefix each worker already holds in its KV cache — via one of three indexers (event-driven KV-events, approximate token tree, or approximate string tree) — and routes to the best-matching worker when the overlap clears
cache_threshold. This is what maximizes cache hits.Fallback path. The policy abandons cache affinity and routes to the shortest queue (
select_worker_min_load→ fewest in-flight requests) in two cases:match_rate ≤ cache_threshold).The imbalance guard
is_imbalancedis a 3-term ORThis PR's core change. The pool is imbalanced if any of these fire:
>ceilingoverload_token_usage_threshold(1.0= off)>thresholdbalance_token_usage_threshold(1.0= off)> absand max> rel·minbalance_abs_threshold(64) /balance_rel_threshold(1.5)1.0, which disables them (utilization and spread are both ≤ 1.0). With defaults, only (c) runs → behavior is byte-for-byte unchanged.--balance-abs-thresholdon its own, and recover "KV-only" behavior by setting it very high.)max_kv > tis a saturation signal, not a balance one — it can't distinguish "one hot engine, idle neighbors" (rebalance) from "all engines uniformly hot" (keep affinity; spilling just destroys cache hits when the cluster is busiest). The overload ceiling (a) is kept separately as an explicit safety valve.Fallback target
When imbalanced,
select_worker_min_loadpicks the shortest queue by request count. Making that spill toward the coldest-KV engine (so a KV-triggered imbalance is relieved on the KV axis) is a separate, coherent follow-up.Changes
policies/cache_aware.rs:is_imbalanced→ 3-term OR;backend_token_usage_boundsreturns(min, max)KV over healthy workers; the count fold runs overhealthy_indices. 5 unit tests pin the distinctions (uniform-high vs one-hot-rest-idle vs overload vs high-count/low-KV vs default-disabled).policies/registry.rs: inject the monitor'ssubscribe()receiver (mirrorskv_event_monitor).config/*,main.rs,policies/{mod,factory}.rs,app_context.rs: threadbalance_token_usage_threshold(spread) +overload_token_usage_threshold(ceiling), both default1.0.bindings/python(lib.rs+router_args.py): expose--balance-token-usage-thresholdand--overload-token-usage-threshold.Test Plan
Local gate (
maturinunavailable, so the binding is verified withcargo build):cargo +nightly fmt --allcargo clippy --all-targets --all-features -- -D warningsFinished— zerocargo test -p smg --lib cache_aware32 passedcargo test -p smg --test routing_tests92 passedcargo build -p smg-pythonFinishedRebased onto
main@63b532a7(incl. #1629least_load) — no conflicts.Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit