Repository navigation
feat(policies): route least_load by token-work expected-wait - #1647
Conversation
Score each healthy worker by its estimated time-to-drain plus a convex
KV-pressure barrier, and route to the lowest (argmin):
score = (queued_tokens + inflight_tokens) / throughput
+ kv_pressure_weight * k/(1-k)
- queued_tokens: new num_waiting_uncached_tokens load signal, wired through the
sglang scheduler proto + servicer and the shared SchedulerLoadSnapshot;
defaults to 0 for backends that do not report it.
- inflight_tokens: per-worker token-work this router dispatched since the last
poll (stale-snapshot correction), reset on update_loads.
- throughput: the worker's gen_throughput when reported, else the configurable
default_throughput (backends without a live generation rate).
- k/(1-k): M/M/1 KV-pressure barrier. Both terms are in seconds, so they add
directly. Missing signals degrade to in-flight + barrier, then to JSQ.
Token-work, not request count, reflects the load a worker actually carries, so
size-skewed traffic is spread by work rather than by count.
Tuning knobs are exposed with matching defaults via PolicyConfig, the smg CLI,
and the Python router CLI/binding: kv_pressure_weight (0.15s),
default_throughput (2000 tok/s), mean_prefill_tokens (1024).
Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
|
Warning Review limit reached
More reviews will be available in 19 seconds. Learn how PR review limits work. Your organization has run out of usage credits. Purchase more in the billing tab. ⌛ How to resolve this issue?After more reviews become available, a review can be triggered using the We recommend that you space out your commits to avoid hitting the rate limit. 🚦 How do rate limits work?CodeRabbit enforces hourly rate limits for each developer per organization. Our paid plans include higher PR review limits than trial, open-source, and free plans. In all cases, reviews become available again over time. During sustained high-volume PR review activity, CodeRabbit may temporarily slow when the next review becomes available. Please see our Fair Usage Limits Policy for further information. ℹ️ Review info⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (13)
✨ Finishing Touches🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 71c29bc596
ℹ️ 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".
| utilization=result.utilization, | ||
| max_running_requests=result.max_running_requests, | ||
| # Queued token-work: waiting-queue tokens not served from cache. | ||
| num_waiting_uncached_tokens=result.num_waiting_uncached_tokens, |
There was a problem hiding this comment.
Guard the new SGLang load field
With the released sglang versions still allowed by grpc_servicer/pyproject.toml (sglang>=0.5.10), GetLoadsReqOutput does not expose num_waiting_uncached_tokens; when the load monitor calls GetLoads, this constructor raises AttributeError and the RPC returns no load snapshot, so the new least_load path never receives the token-work data it depends on. Please either use a safe fallback such as getattr(..., 0) or raise the dependency to a release that includes this field.
Useful? React with 👍 / 👎.
| 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.
Append the new Python constructor args
These fields are inserted before max_idle_secs even though this struct explicitly says new parameters must be appended to avoid breaking positional _Router(...) callers. Any caller built against the previous signature that passes positional arguments after block_size now has max_idle_secs parsed as least_load_kv_pressure_weight and subsequent arguments shifted, which can produce type errors or silently wrong routing configuration.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Code Review
This pull request transitions the least-load routing policy to a least-token-work policy by incorporating queued token-work, in-flight token corrections, and throughput normalization. It updates protobuf schemas, Python bindings, CLI arguments, and the Rust model gateway to support these new parameters. The review feedback highlights several opportunities to improve robustness, specifically by defensively handling non-finite float values (like NaN or Infinity) in throughput and token usage calculations, and using getattr when accessing new fields on external SGLang objects to prevent crashes on older versions.
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 live_throughput = load.total_gen_throughput(); | ||
| let throughput = if live_throughput > 0.0 { | ||
| live_throughput | ||
| } else { | ||
| self.default_throughput | ||
| }; | ||
| let k = load.effective_token_usage().clamp(0.0, 0.999); | ||
| (queued_tokens + inflight_tokens) / throughput | ||
| + self.kv_pressure_weight * k / (1.0 - k) |
There was a problem hiding this comment.
Defensively check that live_throughput and effective_token_usage are finite numbers. If effective_token_usage() returns NaN (which can happen if the backend reports invalid metrics or under division-by-zero scenarios), calling .clamp(0.0, 0.999) can propagate NaN to the score. This breaks the s < best_score comparison in the selection loop, potentially leading to silent routing failures or extreme load imbalance.
let live_throughput = load.total_gen_throughput();
let throughput = if live_throughput.is_finite() && live_throughput > 0.0 {
live_throughput
} else {
self.default_throughput
};
let usage = load.effective_token_usage();
let k = if usage.is_finite() {
usage.clamp(0.0, 0.999)
} else {
0.0
};
(queued_tokens + inflight_tokens) / throughput
+ self.kv_pressure_weight * k / (1.0 - k)| let (tp_sum, tp_count) = healthy | ||
| .iter() | ||
| .filter_map(|&i| loads.and_then(|m| m.get(workers[i].url()))) | ||
| .map(|l| l.total_gen_throughput()) | ||
| .filter(|t| *t > 0.0) | ||
| .fold((0.0, 0u32), |(s, n), t| (s + t, n + 1)); |
There was a problem hiding this comment.
Defensively filter out non-finite throughput values (such as NaN or Infinity) when calculating the fleet-nominal throughput. If a backend reports an invalid non-finite throughput, it could skew the average or result in an infinite nominal throughput, causing incorrect routing decisions for workers missing a fresh snapshot.
| let (tp_sum, tp_count) = healthy | |
| .iter() | |
| .filter_map(|&i| loads.and_then(|m| m.get(workers[i].url()))) | |
| .map(|l| l.total_gen_throughput()) | |
| .filter(|t| *t > 0.0) | |
| .fold((0.0, 0u32), |(s, n), t| (s + t, n + 1)); | |
| let (tp_sum, tp_count) = healthy | |
| .iter() | |
| .filter_map(|&i| loads.and_then(|m| m.get(workers[i].url()))) | |
| .map(|l| l.total_gen_throughput()) | |
| .filter(|t| t.is_finite() && *t > 0.0) | |
| .fold((0.0, 0u32), |(s, n), t| (s + t, n + 1)); |
| utilization=result.utilization, | ||
| max_running_requests=result.max_running_requests, | ||
| # Queued token-work: waiting-queue tokens not served from cache. | ||
| num_waiting_uncached_tokens=result.num_waiting_uncached_tokens, |
There was a problem hiding this comment.
Use getattr defensively when accessing num_waiting_uncached_tokens on the result object. Since SGLang is an external dependency, older or different versions of SGLang might not expose this field on GetLoadsReqOutput, which would cause an AttributeError and crash the load reporting path.
| num_waiting_uncached_tokens=result.num_waiting_uncached_tokens, | |
| num_waiting_uncached_tokens=getattr(result, "num_waiting_uncached_tokens", 0), |
| utilization=result.utilization, | ||
| max_running_requests=result.max_running_requests, | ||
| # Queued token-work: waiting-queue tokens not served from cache. | ||
| num_waiting_uncached_tokens=result.num_waiting_uncached_tokens, |
There was a problem hiding this comment.
🟡 Nit: num_waiting_uncached_tokens is a new field on GetLoadsReqOutput (from sglang). If the deployed sglang version doesn't expose it yet, this AttributeError crashes the entire _convert_loads_to_protobuf call, which loses ALL load data for the worker — not just the new signal. The Rust scoring side was carefully designed for graceful degradation (queued_tokens = 0 when absent), but a crash here prevents that design from working.
A defensive getattr would keep existing load monitoring intact while the new signal degrades to 0:
| num_waiting_uncached_tokens=result.num_waiting_uncached_tokens, | |
| num_waiting_uncached_tokens=getattr(result, "num_waiting_uncached_tokens", 0), |
| @@ -402,7 +412,15 @@ fn default_least_load_interval() -> u64 { | |||
| } | |||
|
|
|||
| fn default_least_load_kv_pressure_weight() -> f64 { | |||
There was a problem hiding this comment.
🟡 Nit: The default changed from 1.5 to 0.15 and the units changed from "request-equivalents" to "seconds." Since kv_pressure_weight was just introduced (PR #1632), any configs or docs that hardcoded the old default of 1.5 would now produce a ~10× stronger barrier in the new formula than intended. Might be worth calling this out in the PR description or a changelog entry so operators know to re-evaluate explicit kv_pressure_weight values.
There was a problem hiding this comment.
Well-designed change. The token-work scoring with in-flight correction is a solid improvement over raw request-count routing for size-skewed traffic. Lock ordering is consistent (no deadlock risk), validation and constructor defaults align across all 5 config surfaces, and the degradation paths (missing snapshot → time-estimate, dark fleet → JSQ, zero throughput → default) are clean.
Two minor nits posted inline — a defensive getattr on the servicer to match the Rust-side graceful degradation, and a migration note for the kv_pressure_weight unit change. Neither is blocking.
0 🔴 Important · 2 🟡 Nit · 0 🟣 Pre-existing
Description
Problem
least_loadscored workers by in-flight request count plus a KV-pressurebarrier. Under size-skewed traffic that spreads poorly: a 6k-token request and a
200-token request count the same, so a worker already holding a large request
still looks "light" and keeps drawing more large requests until it cliffs.
Solution
Score by token-work expected-wait — estimated time-to-drain (token-work over
throughput) plus the convex KV-pressure barrier — and route to the lowest:
Both terms are in seconds. Token-work reflects the load a worker actually
carries; the in-flight term corrects for stale load polls; the M/M/1 barrier
keeps routing off near-full KV.
Changes
policies/least_load.rs— expected-wait score; per-worker in-flighttoken-work tracking, reset on each poll; graceful degradation (missing
queued_tokens→ 0; missing throughput →default_throughput; no freshsnapshot → drain-time estimate at the fleet-nominal rate; dark fleet → JSQ);
10 unit tests; standalone algorithm doc + tuning-knob reference.
num_waiting_uncached_tokensload signal —SchedulerLoadSnapshot(
protocols),sglang_scheduler.protofield 17, the sglang servicerpopulate, and the gRPC
Fromconversions (tokenspeed → 0).kv_pressure_weight,default_throughput,mean_prefill_tokensadded to
PolicyConfig::LeastLoadwith serde defaults + validation.--least-load-*flags) and the Python router (router_args.pyflags +PyO3 binding).
Test Plan
Pre-PR gate, run locally:
The 10
policies::least_load::testscover: routing by queued token-work,throughput normalization, the KV barrier steering off a near-full worker,
in-flight water-filling within a poll interval,
update_loadsresetting thein-flight estimate, cold-start JSQ, the missing-snapshot drain-time estimate,
and the zero-throughput fallback.
grpc_servicer+router_args.pypasspython3 -m py_compile.Notes
GetLoadsendpoint reportingtoken_usage. Refs feat(grpc): add GetLoads endpoint to the vLLM engine service #1630 — when that merges, its vLLMFrom<SchedulerLoad>needs
num_waiting_uncached_tokens: 0added (additive proto field).consumes and its config/CLI/binding surface). Happy to split if reviewers
prefer.
Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspasses