Conversation
Load tracking for cache_aware routing paired increment_load() with manually-placed decrement_load() calls on every exit path. When a client disconnects mid-request, axum drops the handler future at its current await point, so the decrement code after the await never runs and the worker leaks one load count per cancelled request. Leaked counts accumulate without bound, and once they diverge across workers (e.g. after one worker restarts and rejoins with an honest counter of 0) the min-load balance fallback funnels all new requests to a single worker. Replace the manual pairing with a WorkerLoadGuard that increments on creation and decrements exactly once on Drop. Rust runs Drop when a future is dropped, so cancellation, early returns, retries and panics all release the count. For streaming responses the guard moves into the stream-forwarding task and is released on [DONE] or when the stream ends. The guard also holds the worker Arc directly, so the decrement no longer silently no-ops if the worker was removed from the registry mid-request. Also guard the periodic load reset in the registry health checker to only fire when all workers are near-idle (max load <= 2), mirroring the existing guarded variant in worker.rs: the unconditional reset wiped real in-flight load for long-running requests every 10 cycles, which permanently disabled cache_aware's balance fallback under agent-style workloads. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Signed-off-by: Chenguoz <chenxingye@hust.edu.cn>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: c323e0caad
ℹ️ 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 mut stream = stream; | ||
| let mut decremented = false; | ||
| let mut guard = Some(guard); | ||
| while let Some(chunk) = stream.next().await { |
There was a problem hiding this comment.
Observe client closure while awaiting the upstream stream
When a streaming client disconnects while stream.next() is pending, dropping the response closes rx but neither cancels this spawned task nor wakes the upstream stream. Because the loop only discovers the closure on a later tx.send, a backend that stalls before its next chunk leaves both the request and WorkerLoadGuard alive indefinitely, recreating the load-counter leak this change is intended to fix. Poll tx.closed() alongside stream.next() so receiver closure promptly drops the upstream stream and guard.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Good catch — fixed in 0db7e14. The forwarding task now selects on tx.closed() alongside stream.next(), so receiver closure promptly drops the upstream stream (cancelling the worker request) and releases the guard instead of waiting for the worker's next chunk. Applied the same pattern to the untracked streaming branch so it also cancels its worker request promptly.
Address review: the stream-forwarding task only noticed client disconnect on the next tx.send, so a worker that stalls before its next chunk would pin the request and the WorkerLoadGuard until the request timeout. Select on tx.closed() alongside stream.next() so receiver closure promptly drops the upstream stream (cancelling the worker request) and releases the guard. Applied to both streaming branches so untracked streams also cancel their worker request promptly. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Signed-off-by: Chenguoz <chenxingye@hust.edu.cn>
|
Could the #216 removes health-driven resets entirely; combining that with this PR's |
Purpose
Fix a load-counter leak in worker load tracking that eventually breaks
cache_awarerouting.Load tracking pairs
increment_load()with manually-placeddecrement_load()calls on every exit path ofroute_typed_request/send_typed_request. This is notcancellation-safe: when a client disconnects mid-request, axum drops the
handler future at whatever await point it is parked on (typically while
awaiting response headers from the worker), so the decrement code after
that await never runs. Every aborted request leaks exactly one load
count, permanently.
On agent-style workloads where clients routinely abort requests
(timeouts, tool-loop cancellations) we measured ~400 leaked counts per
worker per day on a production 10-node vLLM 0.27.1 cluster. The leak
breaks routing, not just metrics, through
cache_aware's min-loadbalance fallback:
the fallback never sees real imbalance;
rejoins with an honest counter of 0 while its peers carry ~395 phantom
counts — the fallback decides the fleet is massively imbalanced and
funnels every new-prefix request to that single worker. We hit
exactly this in production: one node pinned at
--max-num-seqswhilenine idled, ~10x throughput loss until counters were manually reset.
This PR replaces the manual pairing with a
WorkerLoadGuardthatincrements on creation and decrements exactly once on
Drop. Rust runsDropwhen a future is dropped, so cancellation, early returns, retriesand panics all release the count. For streaming responses the guard
moves into the stream-forwarding task and is released on
data: [DONE](same timing as before), with drop-at-task-end as backstop for client
disconnect / upstream error / streams that end without
[DONE].The RAII guard also fixes two adjacent bugs the manual pairing caused:
(the forwarding task decremented when the short error stream ended
and the retry path in
route_typed_requestdecremented again);worker_registry.get_by_url(worker_url)and silently no-op'd if theworker had been removed (e.g. via
/remove_worker) while a requestwas in flight. The guard holds the worker
Arcdirectly.Additionally, the periodic load reset in the registry health checker
(
worker_registry.rs) now only fires when all workers are near-idle(
max_load <= 2), mirroring the already-guarded variant inworker.rs::start_health_checker. The unconditional reset wiped realin-flight load for long-running requests every 10 check cycles, which on
its own disables the balance fallback for agent workloads (requests
running for minutes never survive to be counted). With the RAII guard
the reset is a drift backstop only, so the idle guard is the correct
semantics.
No public API changes;
send_typed_requestis private and itsload_incremented: boolparameter becomesload_guard: Option<WorkerLoadGuard>.Test Plan
cargo checkon the branch, plus a live reproduction against a realvLLM backend. Start the router with
--policy cache_awarein front ofone worker, then:
Test Result
Before (current main / 0.1.15), each abort leaks exactly +1 and the
counter never recovers:
After (this PR), same sequence:
Also running in production on a 10-node vLLM 0.27.1 cluster under an
agent eval workload with frequent client aborts.
Essential Elements of an Effective PR Description Checklist