Repository navigation
Conversation
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
There was a problem hiding this comment.
Pull request overview
Hardens MultiprocDiffusionExecutor’s async result-pump so a single bad result message can’t permanently wedge a worker queue, and so pump-thread death becomes externally visible via check_health()—addressing the “zombie engine” failure mode described in #5793/#5821.
Changes:
- Wrap per-message dispatch in the pump loop, extracting
_dispatch_result()so exceptions during dispatch drop only the offending message instead of killing the pump thread. - Teach
check_health()to treat dead result-pump threads asEngineDeadError(while not flagging intentional shutdown). - Stop caching outputs for already-cancelled/abandoned waiters (single-output and batch-split), avoiding long-lived memory pinning in
_completed_outputs. - Add targeted unit tests for guarded dispatch, abandoned-waiter output dropping, and dead-pump detection.
Reviewed changes
Copilot reviewed 2 out of 2 changed files in this pull request and generated 1 comment.
| File | Description |
|---|---|
vllm_omni/diffusion/executor/multiproc_executor.py |
Makes result dispatch exception-safe, drops outputs for cancelled waiters, and surfaces pump-thread death via health checks. |
tests/diffusion/test_result_pump.py |
Adds/extends CPU-only tests to verify guarded dispatch, abandoned-waiter handling, and health-check behavior. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| elif msg.kind == AsyncOutputKind.OUTPUT_READY: | ||
| batch_id = msg.async_output_id | ||
| with self._futures_lock: | ||
| per_req_map = self._batch_split_map.pop(batch_id, None) if batch_id else None | ||
| if per_req_map is not None: |
|
This PR appears to belong to: docs/design/module/diffusion/continuous_batching.md, docs/design/module/diffusion/index.md, docs/design/module/diffusion/diffusion_runtime.md. Module owners: @Isotr0py @princepride @SamitHuang @wtomin @ZJY0516 @RuixiangMa @david6666666 @xuechendi @fhfuih @ivanusto, please review your own changes and leave a short self-review comment describing what you checked. PRs without author self-review may not be assigned a reviewer. Please take a look when you have a chance. If you would like an automated review, mention @vllm-omni-review-bot in a comment. |
98cc8ef to
7f85b71
Compare
async_output_id is what routes a result back to its waiter, so a message without one strands its request until that request's own timeout. The branch dropped it silently; say so at error level instead, and return early so the lookup below no longer needs the 'if batch_id' guard. Raised in review of vllm-project#6253. Signed-off-by: ivanusto <ivanusto@gmail.com>
|
Self-review. What I checked
Where I would welcome a second opinion
Not covered No GPU test. Everything here is CPU-only unit-level; I have not re-run an end-to-end video generation against this branch. The original failure was reproduced on a single GB10 with MiniMax-H3 FL2VA (#5821), and the equivalent patch has been running in that deployment, but that patch predates this refactor. |
|
Status update, since this one has been sitting a while. All checks are green on The one review comment on the diff has been addressed. The bot flagged that the elif msg.kind == AsyncOutputKind.OUTPUT_READY:
batch_id = msg.async_output_id
if not batch_id:
# async_output_id is what routes the result back to its waiter.
# Without it there is nobody to resolve and nothing to cache, so
# the request hangs until its own timeout — say so rather than
# dropping the message silently.
logger.error("Dropping OUTPUT_READY with no async_output_id; its request cannot be resolved")
returnThat is the same principle as the PR itself: the failure was never that the pump did the wrong thing, it was that it died or discarded work without leaving a trace, so the symptom showed up much later as a request that simply never completed. @Isotr0py @princepride @SamitHuang @wtomin @ZJY0516 @RuixiangMa @david6666666 @xuechendi @fhfuih — when one of you has a moment. Happy to rebase if |
| pending = self._output_futures.pop(per_req_id, None) | ||
| if pending is not None and not pending.done(): | ||
| try_set_result(pending, per_req_result) | ||
| elif pending is not None: |
There was a problem hiding this comment.
[P3] One corrupt per-request result drops its whole batch's siblings — vllm_omni/diffusion/executor/multiproc_executor.py:954
get_request_output() raising for one request escapes _deliver_batch_split mid-loop; the split map was already popped, so every later request in the same batch loses its output and hangs until _ASYNC_OUTPUT_TIMEOUT even though only one request was corrupt. Message-level containment now holds, but inside a batch split one bad request still costs the whole batch. Smallest fix: a per-request try/except in that loop resolving the failing request with DiffusionOutput(error=...) so siblings still complete.
There was a problem hiding this comment.
Agreed, and fixed in f81871b1f.
The extraction half of the loop body is now wrapped, and a request whose result cannot be read is resolved with a DiffusionOutput(error=...) of its own, so its siblings still complete:
for per_req_id, req_id in per_req_map.items():
per_req_result: DiffusionOutput
try:
req_output = batch_output.get_request_output(req_id) if batch_output is not None else None
...
except Exception as e:
# The split map has already been popped, so an exception escaping
# here would take every request later in this batch with it: they
# would never be resolved and would hang until their own timeout.
# One corrupt request costs one request.
logger.exception("Failed to extract batch output for request %s", req_id)
per_req_result = DiffusionOutput(error=f"Failed to extract batch output for request {req_id}: {e}")
with self._futures_lock:
... # resolve / drop / cache block unchangedOne deliberate choice worth flagging: the with self._futures_lock block below is left outside the try. try_set_result() already swallows InvalidStateError (#5983), and dict writes plus set_result() on a freshly constructed Future cannot raise, so widening the guard would only add a second place that has to decide what to do with a half-resolved request. The per-message guard in the pump loop stays the backstop for anything unforeseen. Happy to widen it if you would rather have belt and braces.
Tests: test_batch_split_failure_does_not_kill_the_pump previously asserted not doomed.done() with the comment "that one message is lost, as expected"; that assertion is now the opposite, since the corrupt request is failed rather than left hanging. Added test_one_corrupt_request_does_not_drop_its_batch_siblings, a three-request batch whose middle entry raises, asserting all three futures resolve and that r-2, delivered after the corrupt r-1, still gets its real result. Both fail on the previous commit and pass on this one.
| # dies, every process is still alive and every request still runs on the | ||
| # GPU — the results simply never come back. Without this check the | ||
| # engine reports healthy forever and only a restart recovers it. | ||
| if self._pump_running and not self._pump_stop.is_set(): |
There was a problem hiding this comment.
[P3] Sync async_diffusion_output.md's reliability section with the new pump semantics — vllm_omni/diffusion/executor/multiproc_executor.py:1020
Dead-pump → EngineDeadError (health/readiness flip) and the cancelled-waiter drop in _dispatch_result extend the "Reliability, Lifecycle & Timeout Behavior" contract owned by docs/design/feature/async_diffusion_output.md, but the page still describes the pump only as "resolves waiting futures or populates _completed_outputs". A short same-PR paragraph there (dispatch containment, pump-death detection, drop-on-abandoned-waiter) keeps the active design page authoritative.
There was a problem hiding this comment.
Agreed, and done in fa9b58bad.
Added item 5 to "Reliability, Lifecycle & Timeout Behavior" rather than editing items 1 to 4, so the diff stays small and #6255's rewrite of item 2 is untouched:
- Pump Fault Containment & Liveness:
A pump thread is the sole reader of its worker's result queue, so anything escaping it costs every later result rather than one message. Dispatch of each dequeued message is therefore wrapped: a bad message is logged with its traceback and dropped, and the pump keeps draining. Within a batch split the containment is per request, because the split map is popped before delivery starts. Outputs whose waiter has already been cancelled or aborted are dropped instead of cached, since nobody will ever collect them and_completed_outputshas no TTL or abort cleanup (a decoded video would stay pinned until shutdown). Finally,check_health()treats a dead pump thread asEngineDeadError: the processes are all alive and the GPU work still runs, so without that check the engine reports healthy forever while no result ever comes back. Intentional shutdown is excluded, sinceshutdown()sets the stop event before joining.
Also extended the "Batch Split" paragraph with one sentence on the per-request containment from the other comment, since that is where a reader would look for it.
_deliver_batch_split pops the split map before it starts delivering, so an exception raised while extracting one request's result escaped mid-loop and took every request after it in the same batch with it: no future resolved, no cache entry, nothing left to route a later result through. Those requests then hung until their own timeout even though only one entry was corrupt. Wrap the extraction so a request that cannot be read from the batch output is resolved with a DiffusionOutput(error=...) of its own and its siblings still complete. The resolve/cache block below is deliberately left outside the try: try_set_result already swallows InvalidStateError, and dict writes and set_result on a fresh Future cannot raise, so the per-message guard in the pump loop stays the backstop for anything unforeseen. Raised in review of vllm-project#6253. Signed-off-by: ivanusto <ivanusto@gmail.com>
The reliability section still described the pump only as resolving waiting futures or populating _completed_outputs, which no longer covers what it does. Add the dispatch containment, the per-request containment inside a batch split, the drop of outputs whose waiter has been abandoned, and the dead-pump health check, so the design page stays authoritative for this contract. Raised in review of vllm-project#6253. Signed-off-by: ivanusto <ivanusto@gmail.com>
dda442d to
fa9b58b
Compare
async_output_id is what routes a result back to its waiter, so a message without one strands its request until that request's own timeout. The branch dropped it silently; say so at error level instead, and return early so the lookup below no longer needs the 'if batch_id' guard. Raised in review of vllm-project#6253. Signed-off-by: ivanusto <ivanusto@gmail.com>
_deliver_batch_split pops the split map before it starts delivering, so an exception raised while extracting one request's result escaped mid-loop and took every request after it in the same batch with it: no future resolved, no cache entry, nothing left to route a later result through. Those requests then hung until their own timeout even though only one entry was corrupt. Wrap the extraction so a request that cannot be read from the batch output is resolved with a DiffusionOutput(error=...) of its own and its siblings still complete. The resolve/cache block below is deliberately left outside the try: try_set_result already swallows InvalidStateError, and dict writes and set_result on a fresh Future cannot raise, so the per-message guard in the pump loop stays the backstop for anything unforeseen. Raised in review of vllm-project#6253. Signed-off-by: ivanusto <ivanusto@gmail.com>
The reliability section still described the pump only as resolving waiting futures or populating _completed_outputs, which no longer covers what it does. Add the dispatch containment, the per-request containment inside a batch split, the drop of outputs whose waiter has been abandoned, and the dead-pump health check, so the design page stays authoritative for this contract. Raised in review of vllm-project#6253. Signed-off-by: ivanusto <ivanusto@gmail.com>
|
@hsliuustc0106 thanks for the review. Both P3s are addressed, and the branch is now rebased onto current P3-1, one corrupt per-request result drops its whole batch's siblings ( Result extraction inside P3-2, sync the design page ( Added item 5, "Pump Fault Containment & Liveness", to "Reliability, Lifecycle & Timeout Behavior" in Verification
One ask This PR still has no |
Document diffusion_kv_metadata on NewRequestData, correct multiproc shutdown so completed async outputs drop with the executor, and tag abort/consumer-drop async-output reclaim as pending vllm-project#6253/vllm-project#6439/vllm-project#6580. Signed-off-by: Huang, Zeyu <11222265+fhfuih@users.noreply.github.com>
Omni ReviewBot: no human activity for 16 days@ivanusto this pull request has had no human commit, comment or review since 2026-08-22. Per repository policy it may be closed if it stays inactive. To keep it moving, any one of these is enough: push an update, reply to the open blocker, or post the current plan and timeline. |
|
Status update, since the bot flagged inactivity. This one is not blocked on me. The state as of
There are no open review threads left on my side, so what it needs is a maintainer look rather than another push from me. @hsliuustc0106 if you have a moment, this is ready for a re-review, or for the Happy to rebase again if it goes stale in the meantime. |
Omni ReviewBot: no human activity for 15 days@ivanusto this pull request has had no human commit, comment or review since 2026-09-08. Please consider marking this PR as draft until work can resume. The author or a maintainer decides whether to change the PR state. To keep it moving, any one of these is enough: push an update, reply to the open blocker, or post the current plan and timeline. |
… outputs The pump thread is the sole reader of a worker result queue, so an exception escaping the dispatch costs the whole queue rather than the one message that caused it: every later result goes undelivered while the server keeps reporting healthy. Route each message through _dispatch_result() inside a try, so a bad message is logged and dropped on its own. Two deliveries could also go missing without a trace. An OUTPUT_READY with no async_output_id has nobody to resolve and nothing to cache, and was discarded silently; it is now logged. And inside a batch split the map has already been popped, so a failure while extracting one request's result took every request after it in the same batch with it; extraction is now per request, and the one that fails resolves with its own error output. Signed-off-by: ivanusto <ivanusto@gmail.com>
fa9b58b to
6a4b91d
Compare
|
Rebased onto current What What is still missing on
Behaviour on Tests. Three new cases in
Whole file: 39 passed. |
Omni ReviewBot: no human activity for 13 days@ivanusto this pull request has had no human commit, comment or review since 2026-09-24. Please confirm the current plan and next step. The author or a maintainer decides whether to change the PR state. To keep it moving, any one of these is enough: push an update, reply to the open blocker, or post the current plan and timeline. |
Omni ReviewBot routing recordAssigned Strict on zcode (GLM-5.3-Flash) under experiment |
Omni ReviewBot attempt recordReview attempt ended as failed (step 'review' (agent.review_diff): unhandled error: RuntimeError: zcode exited 1 without a result event: statusCode: undefined } Error: Turn execution failed (traceId: 2b1ba3c3-d97a-4a4c-88f5-92a22f32f245) — check |
Purpose
A
DiffusionResultPumpthread is the sole reader of one worker's result queue. #5983 fixed theInvalidStateErrorthat was killing it, and that fix is correct — but it left three gaps around itself, all of which were raised in #5793 / #5821 and none of which are onmaintoday.1. Only the dequeue half of the pump loop is guarded.
_result_pumpwrapsdequeueintry/except; everything after it runs unguarded, and_start_result_pumppassestarget=self._result_pumpwith no wrapper. So anything raised while dispatching still escapes and kills the thread — and because the thread is the only reader of that queue, the cost is not one message but every message after it._deliver_batch_split'sbatch_output.get_request_output()and the_sync_result_buffer.put()path are both live examples. This PR extracts the dispatch body into_dispatch_result()and wraps the call, so one bad message costs one message.2.
check_health()never looked at the pump. This is the part that makes the failure so unpleasant operationally: when a pump thread dies, every worker process is still alive and every request still runs to completion on the GPU — the results simply never come back./healthand/v1/modelsanswer 200 throughout, so an operator's only signal is noticing that requests stopped completing. Now a dead pump is reported asEngineDeadError, the same as a dead worker process. Guarded on_pump_running and not _pump_stop.is_set()so a deliberate shutdown isn't reported as a failure.3. A result whose waiter had already cancelled was cached forever.
_completed_outputsexists for the race where a result arrives before the caller asks for it. A caller that cancelled — an inter-output timeout, or an abort — never asks, so the entry is never popped:_completed_outputsis only drained bywait_output_ready()and the early-arrival adopt path, and there is no TTL and no abort-side cleanup. The pinned value is the unpacked output, which for at2varequest is a whole decoded video. Dropping it is also what theelifmakes explicit; thepending is Noneearly-arrival path is untouched.The broad
except Exceptionadded in (1) is deliberate rather than an oversight: it sits at a top-level thread boundary, and it logs the full traceback vialogger.exceptioninstead of swallowing. Per the code-quality guidance, that is the sanctioned form at such a boundary — the alternative is the status quo, where the thread dies.Most of the diff is the dedent from extracting
_dispatch_result(). The logic change is the three points above; reviewing with whitespace ignored makes that much clearer.Reported in #5793 and #5821. #6255 covers the other still-open item from those reports — the hardcoded 30 s
_ASYNC_OUTPUT_TIMEOUTthat triggers the cancellation in the first place. The two are independent and can land in either order.Test Plan
Added three test classes to
tests/diffusion/test_result_pump.py(CPU-only, no GPU), and generalized the existing_feed_one_msg_to_pumphelper into_feed_msgs_to_pumpso a test can feed a bad message followed by a good one and assert the good one still lands.TestResultPumpDispatchIsGuarded— a message that raises during dispatch (sync-buffer path, and a corrupt batch output) must not stop the next message from being delivered by the same pump thread.TestResultPumpDropsOutputForAbandonedWaiter— a cancelled waiter leaves_completed_outputsempty, in both the single-request and batch-split paths, while a result with no waiter at all is still cached as before.TestCheckHealthDetectsDeadPump— a dead pump thread raisesEngineDeadErrorand sets_is_failed; a live one is healthy; a stopped one (shutdown) is not reported as dead.vLLM Version: 0.26.1rc1.dev608+g99a10304d
vLLM-Omni Commit: baba7d1
Test Result
With the
multiproc_executor.pychange reverted and the tests kept, the five tests that assert the new behaviour fail and the three that assert preserved behaviour still pass:ruff checkandruff format --checkpass on both changed files.