Skip to content

fix: make worker load tracking cancellation-safe with an RAII guard - #215

Open
Chenguoz wants to merge 2 commits into
vllm-project:mainfrom
Chenguoz:fix/load-counter-leak-raii-guard
Open

Chenguoz wants to merge 2 commits into
vllm-project:mainfrom
Chenguoz:fix/load-counter-leak-raii-guard

Conversation

@Chenguoz

Copy link
Copy Markdown

Purpose

Fix a load-counter leak in worker load tracking that eventually breaks
cache_aware routing.

Load tracking pairs increment_load() with manually-placed
decrement_load() calls on every exit path of
route_typed_request/send_typed_request. This is not
cancellation-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-load
balance fallback:

  • while the leak is uniform across workers, phantom counts dominate and
    the fallback never sees real imbalance;
  • as soon as one worker's counter diverges — e.g. a worker restarts and
    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-seqs while
    nine idled, ~10x throughput loss until counters were manually reset.

This PR replaces 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 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:

  • double decrement for streaming responses with a retryable status
    (the forwarding task decremented when the short error stream ended
    and the retry path in route_typed_request decremented again);
  • lost decrement after worker removal: decrements went through
    worker_registry.get_by_url(worker_url) and silently no-op'd if the
    worker had been removed (e.g. via /remove_worker) while a request
    was in flight. The guard holds the worker Arc directly.

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 in
worker.rs::start_health_checker. The unconditional reset wiped real
in-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_request is private and its
load_incremented: bool parameter becomes
load_guard: Option<WorkerLoadGuard>.

Test Plan

cargo check on the branch, plus a live reproduction against a real
vLLM backend. Start the router with --policy cache_aware in front of
one worker, then:

R=http://127.0.0.1:9155
load() { curl -s $R/workers | python3 -c "import json,sys; print(json.load(sys.stdin)['workers'][0]['load'])"; }

# 1) a completed request must return load to 0
curl -s $R/v1/chat/completions -H 'Content-Type: application/json' \
  -d '{"model":"m","messages":[{"role":"user","content":"say OK"}],"max_tokens":5}' >/dev/null
load

# 2) three aborted requests: client disconnects ~1s in while generation
#    is still running (the cancellation path)
for i in 1 2 3; do
  curl -s --max-time 1 $R/v1/chat/completions -H 'Content-Type: application/json' \
    -d '{"model":"m","messages":[{"role":"user","content":"Count slowly from 1 to 1000, one number per line."}],"max_tokens":800}' >/dev/null
  load
done

# 3) settle 30s (backend finishes/aborts orphaned generations), re-check
sleep 30; load

Test Result

Before (current main / 0.1.15), each abort leaks exactly +1 and the
counter never recovers:

after 1 completed request: 0   (correct)
after abort #1: 1
after abort #2: 2
after abort #3: 3
final load after 30s settle: 3

After (this PR), same sequence:

after 1 completed request: 0
after abort #1: 0
after abort #2: 0
after abort #3: 0
final load after 30s settle: 0

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
  • The purpose of the PR, such as "Fix some issue (link existing issues this PR will resolve)".
  • The test plan, such as providing test command.
  • The test results, such as pasting the results comparison before and after, or e2e results

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>

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 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".

Comment thread src/routers/http/router.rs Outdated
let mut stream = stream;
let mut decremented = false;
let mut guard = Some(guard);
while let Some(chunk) = stream.next().await {

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

P1 Badge 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 👍 / 👎.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

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

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>
@bvolpato

bvolpato commented Sep 3, 2026

Copy link
Copy Markdown

Could the max_load <= 2 reset be removed instead of treated as proof of idleness? One or two long-running requests are valid in-flight load. Resetting them to zero lets an older guard later decrement the count belonging to a newly admitted request.

#216 removes health-driven resets entirely; combining that with this PR's tx.closed() handling would preserve both accounting and stalled-stream cleanup.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants