Repository navigation
[Bugfix] Prevent DiffusionResultPump crash on cancelled futures (#5793) - #5983
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. |
0c8c2ed to
028d28e
Compare
|
This PR appears to belong to: docs/design/module/diffusion/continuous_batching.md. Module owners: @Isotr0py @princepride Please take a look when you have a chance. If you would like an automated review, mention @vllm-omni-review-bot in a comment. |
|
@Isotr0py @princepride PTAL |
|
hi @Isotr0py @princepride can you take a look at this |
…-project#5793) A late result for an aborted/timed-out request could race the result-pump thread's fut.done() check and set_result()/set_exception() call, raising InvalidStateError with nothing to catch it. Since the pump is a singleton, this permanently kills it: the engine stays "healthy" but never completes any job again. Wrap each set_result/set_exception call on a shared future in a narrow try/except InvalidStateError that drops the late resolution instead of crashing the thread. Rebased onto main: vllm-project#6023 rewrote _result_pump/_deliver_batch_split and reverted to plain set_result/set_exception calls, so this race was reproducible again. Reapplied the same guard on the new call sites and added the regression test to tests/diffusion/test_result_pump.py (the test file vllm-project#6023 introduced) instead of a separate file. Signed-off-by: ANURAG KANADE <anuragkanade6@gmail.com>
d69bba4 to
ae7b3bc
Compare
|
rebased onto main. #6023 rewrote reapplied the same fix on top of #6023's new code, and added the regression test to verified on a real gpu: with the fix reverted, all 3 new tests fail with the exact |
…ect#5793) Move try_set_result/try_set_exception into vllm_omni/diffusion/utils/future_utils.py per review. Signed-off-by: ANURAG KANADE <anuragkanade6@gmail.com>
|
thankyouu!! @Gaohan123 |
A pump thread is the sole reader of one worker's result queue. vllm-project#5983 fixed the InvalidStateError that killed it, but three gaps around that fix remain: 1. Only the dequeue half of the loop was guarded. Anything raised while dispatching still escapes the thread, and the cost is not one message but the whole queue -- every later result goes undelivered. Extract the dispatch body into _dispatch_result() and wrap the call, so one bad message costs one message. 2. check_health() never looked at the pump. When a pump thread dies every worker process is still alive and every request still runs on the GPU, so the engine reported healthy forever while no request could return; only a restart recovered it. Report a dead pump as EngineDeadError. 3. A result whose waiter had already cancelled was cached in _completed_outputs. That cache exists for results that arrive before the caller asks for them; a cancelled caller never asks, so the entry is never popped and the unpacked output -- a whole decoded video, for a t2va request -- stays pinned until the process exits. Drop it instead. The broad 'except Exception' added in (1) is deliberate: this is a top-level thread boundary, and it logs the full traceback rather than swallowing. Most of the diff is the dedent from extracting _dispatch_result(); the logic change is the three points above. Reported in vllm-project#5793 and vllm-project#5821. Signed-off-by: ivanusto <ivanusto@gmail.com>
A pump thread is the sole reader of one worker's result queue. vllm-project#5983 fixed the InvalidStateError that killed it, but three gaps around that fix remain: 1. Only the dequeue half of the loop was guarded. Anything raised while dispatching still escapes the thread, and the cost is not one message but the whole queue -- every later result goes undelivered. Extract the dispatch body into _dispatch_result() and wrap the call, so one bad message costs one message. 2. check_health() never looked at the pump. When a pump thread dies every worker process is still alive and every request still runs on the GPU, so the engine reported healthy forever while no request could return; only a restart recovered it. Report a dead pump as EngineDeadError. 3. A result whose waiter had already cancelled was cached in _completed_outputs. That cache exists for results that arrive before the caller asks for them; a cancelled caller never asks, so the entry is never popped and the unpacked output -- a whole decoded video, for a t2va request -- stays pinned until the process exits. Drop it instead. The broad 'except Exception' added in (1) is deliberate: this is a top-level thread boundary, and it logs the full traceback rather than swallowing. Most of the diff is the dedent from extracting _dispatch_result(); the logic change is the three points above. Reported in vllm-project#5793 and vllm-project#5821. Signed-off-by: ivanusto <ivanusto@gmail.com>
…-project#5793) (vllm-project#5983) Signed-off-by: ANURAG KANADE <anuragkanade6@gmail.com>
What broke
MultiprocDiffusionExecutor._result_pumpis the sole reader of the diffusion worker result queue, running as a singleton background thread. Every dispatch branch resolves a shared future with:The
not fut.done()check and theset_result()/set_exception()call are not atomic. A consumer awaiting the future viaasyncio.wrap_future()underasyncio.wait_for()(an inter-output timeout, or a request abort) can cancel it in the gap between the two. When that happens,set_result()/set_exception()raisesconcurrent.futures.InvalidStateError, which nothing catches — killing the pump thread permanently. Since the pump is a singleton, the engine then accepts and "runs" jobs that silently never complete (/healthstays 200, GPU sits idle).Repro
Reported originally in #5793 via two real triggers: a worker-side generation error aborting a request, and a client-side inter-output timeout (
_ASYNC_OUTPUT_TIMEOUT) firing while the worker was still running. Both land on the same race.Root cause
Check-then-act race on
concurrent.futures.Futurestate across threads: the pump'sdone()check and its resolve call aren't atomic, so a concurrentcancel()from the consumer side can land in between.Fix
Added
_try_set_result/_try_set_exceptionhelpers that wrapset_result/set_exceptionin a narrowexcept concurrent.futures.InvalidStateError, dropping the late resolution instead of crashing the thread. Applied at every call site in_result_pumpthat resolves a future pulled from the shared_rpc_futures/_output_futuresdicts, plus the same pattern inshutdown()'s cleanup loop for consistency. Left untouched the twoset_resultcalls on freshly-created localFuture()objects that aren't shared with any other thread yet — no race possible there.Testing
Added
tests/diffusion/test_multiproc_result_pump.py(CPU-only, no GPU needed): a_RacyFuturesubclass whosedone()deterministically lies (False) while genuinely cancelled underneath, reproducing the exact race window without relying on real thread timing.Verified on an RTX 3090 (vLLM 0.26.0, vllm-omni installed per docs):
tests/diffusion/test_multiproc_result_pump.py— 3 passed.InvalidStateErrortraceback from the original report, at the samefut.set_result(msg.output)call site.tests/diffusion/test_multiproc_engine_concurrency.py— 60 passed, no regressions.ruff check,ruff format --check, and the fullpre-commithook suite pass on both changed files.Fixes #5793