feat(policies): add least_load load-balancing policy - #1629
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 (16)
📝 WalkthroughWalkthroughAdds a new LeastLoad routing policy and integrates it across config, CLI/Python bindings, the policy factory, policy registry, worker monitor load polling, and worker-removal cache cleanup; unit tests verify selection and cache behavior. ChangesLeast-Load Balancing Policy
Sequence DiagramsequenceDiagram
participant WorkerMonitor
participant LeastLoadPolicy
participant Worker
WorkerMonitor->>LeastLoadPolicy: update_loads(group reports)
LeastLoadPolicy->>LeastLoadPolicy: compute score(worker) across healthy workers
LeastLoadPolicy->>Worker: increment_processed() on selected worker
Estimated Code Review Effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly Related PRs
Suggested Labels
Suggested Reviewers
Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 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 least_load load-balancing policy that routes requests based on a combination of active in-flight requests and KV-cache pressure. The changes integrate this policy into the configuration, validation, CLI, factory, registry, and worker monitor. The review feedback identifies several important improvements: robustly handling potential NaN values in token usage calculation to avoid routing lockups, using parking_lot::RwLock instead of std::sync::RwLock for codebase consistency and to avoid lock poisoning boilerplate, ensuring that processed request metrics are incremented in the single-worker fast path, and optimizing lock acquisition patterns.
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.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 037ad40a15
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
There was a problem hiding this comment.
Actionable comments posted: 3
🤖 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 326-345: The branch added validation for PolicyConfig::LeastLoad
but lacks direct regression tests for the edge cases; add unit tests that
exercise validation.rs's PolicyConfig::LeastLoad path to assert
ConfigError::InvalidValue is returned when load_check_interval_secs == 0 and
when lambda is non-finite (NaN/INFINITY) or negative. Create tests that
construct a PolicyConfig::LeastLoad with load_check_interval_secs = 0 and with
lambda = f64::NAN, f64::INFINITY, and lambda = -0.1, call the same validation
function used by the module (the validator that emits
ConfigError::InvalidValue), and assert the returned error contains field
"load_check_interval_secs" or "lambda" and the expected reason strings. Ensure
tests live alongside other config validation tests and use explicit assertions
(not panics) so these regression cases are covered in CI.
In `@model_gateway/src/main.rs`:
- Around line 931-934: The CLI branch that maps "least_load" to a PolicyConfig
currently hardcodes load_check_interval_secs: 5 which conflicts with
PolicyConfig::LeastLoad's default of 10 in config/types.rs; update parse_policy
(the code that returns PolicyConfig::LeastLoad for "least_load") to use the same
default (10) or to call Default::default() for the LeastLoad variant so CLI and
serde-config share the same defaults.
In `@model_gateway/src/policies/factory.rs`:
- Around line 22-24: Add unit tests that assert the factory returns a
LeastLoadPolicy when given a PolicyConfig::LeastLoad and when requesting the
policy by name; specifically, in the factory tests call create_from_config with
a PolicyConfig::LeastLoad (including a sample lambda) and assert the returned
Arc inner type corresponds to LeastLoadPolicy created via
LeastLoadPolicy::with_lambda, and similarly call create_by_name for the
least_load variant and assert the same; ensure the tests check the lambda value
is preserved (or that the concrete type is LeastLoadPolicy) to lock down the new
dispatch paths in create_from_config and create_by_name.
🪄 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: 2388af81-069a-425a-801a-f1e58177d91f
📒 Files selected for processing (8)
model_gateway/src/config/types.rsmodel_gateway/src/config/validation.rsmodel_gateway/src/main.rsmodel_gateway/src/policies/factory.rsmodel_gateway/src/policies/least_load.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/policies/registry.rsmodel_gateway/src/worker/monitor.rs
037ad40 to
c98a059
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: c98a059d9d
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
There was a problem hiding this comment.
Actionable comments posted: 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 `@model_gateway/src/policies/least_load.rs`:
- Around line 129-133: The update_loads method in LeastLoadPolicy currently
extends cached_loads and never removes stale entries; modify the policy to prune
removed workers by either adding a remove_worker method to the
LoadBalancingPolicy trait and implement it in LeastLoadPolicy to delete entries
from cached_loads, or change update_loads (LeastLoadPolicy::update_loads) to
accept the full current worker set and replace the cache (write-locked
cached_loads = new_map) instead of extend so stale keys are dropped; ensure the
chosen approach integrates with PolicyRegistry::remove_worker_from_cache_aware
and WorkerMonitor eviction so LeastLoadPolicy entries are cleaned when workers
are removed.
🪄 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: fa8c6695-a434-48ce-8567-cedd17a548c1
📒 Files selected for processing (14)
bindings/python/src/lib.rsbindings/python/src/smg/router.pybindings/python/src/smg/router_args.pybindings/python/tests/test_router_config.pybindings/python/tests/test_startup_sequence.pybindings/python/tests/test_validation.pymodel_gateway/src/config/types.rsmodel_gateway/src/config/validation.rsmodel_gateway/src/main.rsmodel_gateway/src/policies/factory.rsmodel_gateway/src/policies/least_load.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/policies/registry.rsmodel_gateway/src/worker/monitor.rs
c98a059 to
fd663d0
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: fd663d09ce
ℹ️ 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".
| // ==================== Routing Policy ==================== | ||
| /// Load balancing policy to use | ||
| #[arg(long, default_value = "cache_aware", value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "prefix_hash", "consistent_hashing", "manual", "bucket"], help_heading = "Routing Policy")] | ||
| #[arg(long, default_value = "cache_aware", value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "least_load", "prefix_hash", "consistent_hashing", "manual", "bucket"], help_heading = "Routing Policy")] |
There was a problem hiding this comment.
Add least_load to Helm values schema
When this policy is selected via the Helm chart, chart values validation still rejects it because deploy/helm/smg/values.schema.json line 10 only enumerates cache_aware, round_robin, power_of_two, manual, random, and prefix_hash (I checked the Helm chart policy schema and values comments). This means Kubernetes/Helm users cannot enable the new least_load policy even though the gateway CLI and Python args now accept it; update the chart schema (and related values comment) alongside this new accepted policy.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 3
♻️ Duplicate comments (1)
model_gateway/src/policies/least_load.rs (1)
129-133:⚠️ Potential issue | 🟠 Major | ⚡ Quick win
update_loadskeeps scoring from stale load reports.Lines 129-133 merge new reports into
cached_loads, but they never clear keys that disappear from a poll. After a monitor miss, that worker keeps its old KV utilization indefinitely, so routing no longer matches the documented “missing report =>k = 0” behavior.Suggested fix
fn update_loads(&self, loads: &HashMap<String, WorkerLoadResponse>) { if let Ok(mut cached) = self.cached_loads.write() { - cached.extend(loads.iter().map(|(k, v)| (k.clone(), v.clone()))); + cached.clone_from(loads); } }This follows the PR contract that workers without a current load report should be scored as if
k = 0.🤖 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/least_load.rs` around lines 129 - 133, The update_loads function currently merges new reports into cached_loads but never removes keys for workers that disappeared, causing stale scores; modify update_loads (and the cached_loads write) to replace the cache contents with the new loads set (or explicitly remove keys not present in the incoming loads) so that workers missing from a poll are treated as absent (k = 0); locate the update_loads method and change the write logic to clear or overwrite cached_loads based on the keys of the provided loads HashMap instead of using extend.
🤖 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/types.rs`:
- Around line 376-387: Add unit tests that assert the PolicyConfig::LeastLoad
variant reports the correct name via its name() method, round-trips through
serde (serialize -> deserialize equals original), and that omitted fields
default to default_least_load_interval and default_least_load_lambda. Create
tests that construct LeastLoad with defaults omitted and with explicit values,
verify deserialized structs equal expected values, and assert the
load_check_interval_secs and lambda match the default_* functions when not
provided; reference symbols: PolicyConfig::LeastLoad, name(),
default_least_load_interval, default_least_load_lambda, and serde
(serde_json::to_string/from_str).
In `@model_gateway/src/main.rs`:
- Line 152: The CLI accepts "bucket" in the value_parser but parse_policy lacks
a "bucket" match arm so inputs silently fall back to "round_robin"; update the
parse_policy function to explicitly handle "bucket" (either map it to the
correct RoutingPolicy variant or return a parse error to fail closed), ensuring
the match in parse_policy covers the "bucket" string and returns a Result/Err
instead of defaulting to the fallback; reference the value_parser and
parse_policy identifiers and add the "bucket" branch (or explicit Err) to keep
CLI parsing consistent and visible to users.
In `@model_gateway/src/policies/least_load.rs`:
- Around line 98-100: The early return for the single-worker fast path skips the
shared post-selection bookkeeping, so change the block in least_load.rs that
currently does `if healthy.len() == 1 { return Some(healthy[0]); }` to instead
perform the same post-selection steps as the multi-worker path: call
`increment_processed()` for the selected worker (healthy[0]) and run any
remaining shared bookkeeping before returning the selection; ensure you
reference and invoke the same helper(s) used by the multi-worker path so the
single-worker path updates the processed counter identically.
---
Duplicate comments:
In `@model_gateway/src/policies/least_load.rs`:
- Around line 129-133: The update_loads function currently merges new reports
into cached_loads but never removes keys for workers that disappeared, causing
stale scores; modify update_loads (and the cached_loads write) to replace the
cache contents with the new loads set (or explicitly remove keys not present in
the incoming loads) so that workers missing from a poll are treated as absent (k
= 0); locate the update_loads method and change the write logic to clear or
overwrite cached_loads based on the keys of the provided loads HashMap instead
of using extend.
🪄 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: e7ea00da-ab68-4a1a-917c-7187b0bc4ef7
📒 Files selected for processing (16)
bindings/python/src/lib.rsbindings/python/src/smg/router.pybindings/python/src/smg/router_args.pybindings/python/tests/test_router_config.pybindings/python/tests/test_startup_sequence.pybindings/python/tests/test_validation.pymodel_gateway/src/config/types.rsmodel_gateway/src/config/validation.rsmodel_gateway/src/main.rsmodel_gateway/src/policies/factory.rsmodel_gateway/src/policies/least_load.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/policies/power_of_two.rsmodel_gateway/src/policies/registry.rsmodel_gateway/src/worker/monitor.rsmodel_gateway/src/workflow/steps/local/remove_from_policy_registry.rs
| // ==================== Routing Policy ==================== | ||
| /// Load balancing policy to use | ||
| #[arg(long, default_value = "cache_aware", value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "prefix_hash", "consistent_hashing", "manual", "bucket"], help_heading = "Routing Policy")] | ||
| #[arg(long, default_value = "cache_aware", value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "least_load", "prefix_hash", "consistent_hashing", "manual", "bucket"], help_heading = "Routing Policy")] |
There was a problem hiding this comment.
bucket is accepted by CLI but silently parsed as round_robin.
value_parser allows "bucket" (Lines 152/217/221), but parse_policy has no "bucket" match arm, so Line 949 fallback applies. This causes incorrect routing policy selection without user-visible error.
Suggested fix (fail closed until bucket parsing is explicit)
- #[arg(long, default_value = "cache_aware", value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "least_load", "prefix_hash", "consistent_hashing", "manual", "bucket"], help_heading = "Routing Policy")]
+ #[arg(long, default_value = "cache_aware", value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "least_load", "prefix_hash", "consistent_hashing", "manual"], help_heading = "Routing Policy")]
@@
- #[arg(long, value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "least_load", "prefix_hash", "consistent_hashing", "manual", "bucket"], help_heading = "PD Disaggregation")]
+ #[arg(long, value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "least_load", "prefix_hash", "consistent_hashing", "manual"], help_heading = "PD Disaggregation")]
@@
- #[arg(long, value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "least_load", "prefix_hash", "consistent_hashing", "manual", "bucket"], help_heading = "PD Disaggregation")]
+ #[arg(long, value_parser = ["random", "round_robin", "cache_aware", "power_of_two", "least_load", "prefix_hash", "consistent_hashing", "manual"], help_heading = "PD Disaggregation")]Also applies to: 217-222, 916-950
🤖 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` at line 152, The CLI accepts "bucket" in the
value_parser but parse_policy lacks a "bucket" match arm so inputs silently fall
back to "round_robin"; update the parse_policy function to explicitly handle
"bucket" (either map it to the correct RoutingPolicy variant or return a parse
error to fail closed), ensuring the match in parse_policy covers the "bucket"
string and returns a Result/Err instead of defaulting to the fallback; reference
the value_parser and parse_policy identifiers and add the "bucket" branch (or
explicit Err) to keep CLI parsing consistent and visible to users.
fd663d0 to
859a9c5
Compare
Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
859a9c5 to
4e61087
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 4e6108715e
ℹ️ 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".
| let k = loads | ||
| .and_then(|m| m.get(worker.url())) | ||
| .map(|l| l.effective_token_usage().clamp(0.0, 0.999)) | ||
| .unwrap_or(0.0); |
There was a problem hiding this comment.
Ignore non-finite load samples before scoring
In gRPC mode the backend's token_usage is a proto double, so a worker can report NaN; f64::clamp leaves NaN unchanged. If the first healthy worker's cached sample is NaN, best_score becomes NaN and every later s < best_score comparison is false, causing least_load to keep selecting that first worker until a valid sample arrives. Treat non-finite usage as missing or clamp it only after checking is_finite().
Useful? React with 👍 / 👎.
|
👋 The PR description doesn't fully follow
Please update the PR description so reviewers have the context they need. |
Description
Problem
SMG has no load-aware routing policy that accounts for KV-cache pressure.
power_of_twocompares two random workers on a single metric, and the count-based policies ignore how full each worker's KV cache is — so traffic can pile onto a worker that is about to hit the KV preemption/recompute cliff, spiking TTFT.Solution
Add
least_load: a policy that scores every healthy workerChecklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit