Skip to content

[Bugfix][Diffusion] Release async outputs abandoned by consumer streams - #6580

Open
rahul-steiger-nv wants to merge 14 commits into
vllm-project:mainfrom
rahul-steiger-nv:mp-async-leak
Open

rahul-steiger-nv wants to merge 14 commits into
vllm-project:mainfrom
rahul-steiger-nv:mp-async-leak

Conversation

@rahul-steiger-nv

@rahul-steiger-nv rahul-steiger-nv commented Aug 24, 2026 •

Copy link
Copy Markdown
Contributor

Purpose

Related to #6413 and a follow-up to the merged #6439.

#6439 releases async diffusion outputs abandoned by aborted requests. A consumer can also disappear without an abort, leaving materialized image or video tensors cached indefinitely. This PR tracks unclaimed async outputs by stream queue and releases them when no consumer remains.

  • Key ownership by queue so cleanup from an old stream cannot drop a new stream's outputs when a request ID is reused.
  • Materialize async outputs in the output stream and handle cancellation, timeout, and stream closure without recreating retired executor waiters.
  • Preserve eager request admission: registration and admission errors still occur when async_add_req_and_stream_response() is called. An inner generator uses aclosing to close its output stream when it exits.
  • Keep engine-internal materialization wait time in a separate field feeding output_ready_wait_time_ms, rather than exposing it as a model stage in stage_durations.
  • During engine shutdown, retain queued output futures for registered consumers before executor teardown clears cached results. Completed results remain deliverable, unfinished outputs receive a shutdown error, and new waits after executor shutdown fail immediately instead of timing out.

Test Plan

Run the focused cleanup, RPC-routing, result-pump, and shutdown-race suites:

python -m pytest -q \
  tests/diffusion/test_diffusion_engine_cleanup.py \
  tests/diffusion/test_diffusion_engine_rpc_routing.py \
  tests/diffusion/test_result_pump.py \
  tests/diffusion/test_multiproc_shutdown_race.py

Coverage includes reused request IDs, missing and closed streams, cancellation during and after materialization, eager admission and call-site errors, queued-result delivery through real executor shutdown, and immediate errors for waits after shutdown.

Timing-specific validation also covered:

python -m pytest -q \
  tests/diffusion/test_diffusion_engine_metrics.py \
  tests/diffusion/test_diffusion_engine.py::test_step_streaming_excludes_output_wait_from_execution_time

Test Result

  • 108 tests passed across the four focused suites above (68 cleanup/RPC-routing tests and 40 result-pump/shutdown-race tests).
  • Cleanup, metrics, and execution-timing tests passed together earlier in this revision series: 53 passed.
  • The tightened exact shutdown-error assertion was rerun separately: 1 passed.
  • Evaluation used vllm-omni-v0.30.0-arm64.sqsh from ../enroot-images, with the local checkout mounted into Enroot. The environment reported a vLLM-Omni 0.25.0 / vLLM 0.30.0 version mismatch warning.
  • git diff --check passed. No full diffusion-suite or model-inference validation is claimed.

The PR is based on main at b63e35ab4, after #6439 merged. Earlier rebase checks passed all applicable pre-commit checks except mypy, which reported four errors also reproduced on unmodified main; those checks were not rerun for the latest revisions.

@chatgpt-codex-connector

Copy link
Copy Markdown

Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits.
Credits must be used to enable repository wide code reviews.

@vllm-omni-review-bot

Copy link
Copy Markdown

This PR appears to belong to: docs/design/module/diffusion/index.md.

Module owners: @Isotr0py @princepride @SamitHuang @wtomin @ZJY0516 @RuixiangMa @david6666666 @xuechendi

@rahul-steiger-nv, 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.

@fhfuih fhfuih left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Thanks for identifying this problem and submit a PR.

So as I understand, you also move the waiting-for-output from step_streaming to get_output_stream. This also seems cleaner in terms of ownership (at least for me) because the stream yields materialized output, not a handle.

Comment thread vllm_omni/diffusion/diffusion_engine.py Outdated
Comment thread vllm_omni/diffusion/diffusion_engine.py Outdated
Comment on lines +800 to +801
try:
output = await asyncio.wait_for(asyncio.wrap_future(fut), timeout=_ASYNC_OUTPUT_TIMEOUT)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

On rebase onto current main, keep _async_output_timeout()

The wait you moved into get_output_stream still uses _ASYNC_OUTPUT_TIMEOUT. main already replaced that with _async_output_timeout(). Don’t restore the old constant when you unstack from #6439.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Sounds good! I will wait until #6439 merges and then use _async_output_timeout()

@hsliuustc0106 hsliuustc0106 added bug Something isn't working diffusion codes related to diffusion models labels Aug 25, 2026
@hsliuustc0106

Copy link
Copy Markdown
Collaborator

resolve conflicts

Comment thread vllm_omni/diffusion/diffusion_engine.py
Comment thread vllm_omni/diffusion/diffusion_engine.py Outdated
@rahul-steiger-nv

Copy link
Copy Markdown
Contributor Author

resolve conflicts

I will resolve with a rebase once #6439 merges.

@hsliuustc0106 hsliuustc0106 added the high priority high priority issue, needs to be done asap label Aug 28, 2026
fhfuih added a commit to fhfuih/vllm-omni that referenced this pull request Sep 3, 2026
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>
@vllm-omni-review-bot

Copy link
Copy Markdown

Omni ReviewBot: no human activity for 10 days

@rahul-steiger-nv this pull request has had no human commit, comment or review since 2026-08-28. 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.

@rahul-steiger-nv

Copy link
Copy Markdown
Contributor Author

I will wait until #6439 merges.

Comment thread vllm_omni/diffusion/diffusion_engine.py Outdated
Comment thread vllm_omni/diffusion/diffusion_engine.py Outdated
@vllm-omni-review-bot

vllm-omni-review-bot commented Sep 17, 2026 •

Copy link
Copy Markdown

Omni ReviewBot triage note

Resolved as of 0213b6909b48: the high-priority or low-quality signal noted on an earlier commit no longer applies.

@vllm-omni-review-bot

Copy link
Copy Markdown

Omni ReviewBot: no human activity for 7 days

@rahul-steiger-nv this pull request has had no human commit, comment or review since 2026-09-17. 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.

Rahul Steiger and others added 7 commits September 27, 2026 01:11
Signed-off-by: Rahul Steiger <rsteiger@aws-cmh-slurm-1-vscode-04.cm.cluster>
Signed-off-by: Rahul Steiger <rsteiger@aws-cmh-slurm-1-vscode-04.cm.cluster>
Signed-off-by: Rahul Steiger <rsteiger@aws-cmh-slurm-1-vscode-04.cm.cluster>
Signed-off-by: Rahul Steiger <rsteiger@nvl72d228-T10.cm.cluster>
Signed-off-by: Rahul Steiger <rsteiger@nvl72d228-T10.cm.cluster>
Signed-off-by: Rahul Steiger <rsteiger@nvl72d228-T10.cm.cluster>
Signed-off-by: Rahul Steiger <rsteiger@nvl72d228-T10.cm.cluster>
Signed-off-by: Rahul Steiger <rsteiger@nvl72d228-T10.cm.cluster>
@rahul-steiger-nv

rahul-steiger-nv commented Sep 27, 2026 •

Copy link
Copy Markdown
Contributor Author

@hsliuustc0106 as #6439 has been merged, could we try to merge this as well? Thanks in advance.

Copy link
Copy Markdown
Collaborator

Reviewed head 0213b6909b487d2e87a8865f64cf6063b668a060 against base b63e35ab4ffae2f78b556150943ec0e08a444030.

No additional actionable findings. The earlier issues are addressed: admission happens on first generator iteration; completed/exceptional waiters retire their IDs; cancellation after delivery avoids recreating executor waiters. Queued and late outputs are reclaimed when the consumer disappears, and the current environment-controlled timeout is preserved.

All five changed Python files pass static syntax parsing and the diff passes whitespace checks; reported CI is green. This was a static review of the fork snapshot. I did not run the lifecycle tests, and the PR correctly discloses that post-rebase collection is blocked by the installed vLLM version.

@princepride princepride left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Thanks for tightening up cleanup of abandoned async outputs. I'm requesting changes on four points (details inline):

  1. Request-ID reuse can drop a newer request's outputs. get_output_stream's finally pops _unclaimed_async_outputs[request_id] unconditionally, unlike the adjacent _out_streams.get(request_id) is queue guard.
  2. close() turns already-delivered results into "aborted" errors. It drops every unclaimed async output, including ones whose results are already cached and queued ahead of the closed sentinel for a live consumer.
  3. async_add_req_and_stream_response is now an async generator, so admission (add_request) is deferred until the first __anext__. That's a silent API behavior change for callers that abort/inspect before iterating or that expect admission errors at the call site.
  4. output_ready_wait is written into output.stage_durations, which flows into OmniRequestOutput.stage_durations, API responses and benchmark stage aggregation, so an engine-internal wait shows up as a model stage.

Comment thread vllm_omni/diffusion/diffusion_engine.py Outdated
Comment thread vllm_omni/diffusion/diffusion_engine.py
Comment thread vllm_omni/diffusion/diffusion_engine.py Outdated
Comment thread vllm_omni/diffusion/diffusion_engine.py Outdated
@vllm-omni-review-bot

Copy link
Copy Markdown
Omni ReviewBot routing record

Assigned Strict on cursor (cursor-grok-4.6-high) under experiment fleet-strict-cursor-grok46-zcode-glm53flash-5050-c5-z10-20261002.

Signed-off-by: Rahul Steiger <rsteiger@nvidia.com>
Signed-off-by: Rahul Steiger <rsteiger@nvidia.com>
Signed-off-by: Rahul Steiger <rsteiger@nvidia.com>
Signed-off-by: Rahul Steiger <rsteiger@nvidia.com>
Signed-off-by: Rahul Steiger <rsteiger@nvidia.com>
Signed-off-by: Rahul Steiger <rsteiger@nvidia.com>
@rahul-steiger-nv

Copy link
Copy Markdown
Contributor Author

@princepride Thanks for the detailed review. I updated the PR accordingly. Let me know if there is anything else I can do!

@vllm-omni-review-bot vllm-omni-review-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.

Omni ReviewBot review

4 actionable finding(s).

CI at 496d546c3b98 (2026-10-10T03:03:25.997245+00:00): required check(s) blocking: buildkite/vllm-omni (missing).

Full review analysis

Scan:

Category Result
Tests / verification 2 finding(s) below
Security no finding reported
Docs / comments 2 finding(s) below
Behavior / compatibility no finding reported
Correctness no finding reported

Validated:

  • [resolved] princepride stage_durations: wait lives on DiffusionOutput.output_ready_wait_time (data.py:1816) and metrics output_ready_wait_time_ms (diffusion_engine.py:511), not stage_durations.
  • [resolved] output_ready_wait not written into stage_durations; residual: test_step_streaming_excludes_output_wait_from_execution_time patches get_output_stream so it never proves the producer
  • [resolved] princepride request-id reuse / close() aborting live cached results / deferred admission / stage_durations: head keys unclaimed by queue (1206), transfers wait_output_ready for live streams (1621-1623), admits in add_request before aclosing (1212-1215), and writes output_ready_wait_time not a stage key. Residuals already kept (step_streaming aclosing, docs, field order).
  • [claim-verified] Eager admission, separate wait field, shutdown transfer, and post-shutdown fail-fast are in the head diff as described.
  • [claim-verified] Test-plan files exist and new cases pin reused IDs, missing streams, cancel-during/after materialization, eager admission, and queued-result delivery; cannot re-confirm the reported 108/53/1 pass counts (read-only, no shell).
  • [claim-refuted] diffusion_runtime.md:387 still says pending #6580 and :400 says cached outputs are dropped with the executor — contradicts close() salvage at this head

4 actionable finding(s).

Verdict: COMMENT

Findings

  • **[P2] This PR lands the async-output lifecycle that docs/design/module/diffusion/dif…** — docs/design/module/diffusion/diffusion_runtime.md`
Evidence for This PR lands the async-output lifecycle that `docs/design/module/diffusion/dif…

This PR lands the async-output lifecycle that docs/design/module/diffusion/diffusion_runtime.md still marks *(pending #6253/#6439/#6580)*. get_output_stream() now calls wait_output_ready(), drop_outputs abandoned ids on consumer close, and close() copies live-stream futures into _shutdown_output_futures before executor teardown — so cached OUTPUT_READY is no longer dropped with the executor for a still-registered consumer. docs/design/feature/async_diffusion_output.md:183 still lists async waiting on step_streaming(), which now only reads output.output_ready_wait_time. Update those sections to this contract.

Evidence: Trigger: this head is PR #6580, which implements the cancellation/shutdown protocol the design docs still label pending. Unmet requirement: those docs still describe the old cache-until-teardown contract, so a reader of the module/feature design will follow a protocol this code no longer has. Unchanged by this diff, present in the PR-time tree: docs/design/module/diffusion/diffusion_runtime.md:387 *(pending #6253/#6439/#6580)*. Today, an unconsumed OUTPUT_READY is cached; docs/design/module/diffusion/diffusion_runtime.md:399-401 Completed async outputs / cached for later wait_output_ready() calls are dropped with the executor / object.; docs/design/module/diffusion/diffusion_runtime.md:449 consumer drop before consumption is also required *(pending #6253/#6439/#6580)*.; docs/design/feature/async_diffusion_output.md:183 - vllm_omni/diffusion/diffusion_engine.py: step()/step_streaming()/add_req_and_wait_for_response() async output waiting. Decisive new contract in this diff: vllm_omni/diffusion/diffusion_engine.py:1165 fut = self.executor.wait_output_ready(async_output_id); :1206-1209 abandoned_ids = self._unclaimed_async_outputs.pop(queue, set()) then self.executor.drop_output(async_output_id); :1621-1623 self._shutdown_output_futures[stream] = { output_id: self.executor.wait_output_ready(output_id) ...}; :481-485 generator = self.get_output_stream(request_id) / output_ready_wait_time = output.output_ready_wait_time / The stream now materializes async outputs before yielding — step_streaming no longer waits.

Suggestion: get_output_stream() materializes via wait_output_ready() and releases unclaimed ids on consumer close (#6580). close() transfers queued futures for live streams before executor teardown; completed results stay deliverable and new waits after shutdown fail immediately.

- **[P2] `MultiprocDiffusionExecutor.drop_output` still documents abort-only (`An aborte…** — `vllm_omni/diffusion/executor/multiproc_executor.py:1168`
Evidence for `MultiprocDiffusionExecutor.drop_output` still documents abort-only (`An aborte…

MultiprocDiffusionExecutor.drop_output still documents abort-only (An aborted request never calls wait_output_ready; bounded under abort traffic) even though this PR updated abstract.drop_output to consumer-unavailable and now also calls drop_output on cancel, timeout, consumer-close, and missing-stream. Update the override docstring to those release paths.

Evidence: This PR changed abstract.drop_output from abort-only to consumer-unavailable: abstract.py:141 """Reclaim an async output whose consumer is no longer available. The override still teaches abort-only: multiproc_executor.py:1168 An aborted request never calls :meth:wait_output_ready, so a late and 1171 bounded under abort traffic. Trigger: the same PR's engine now also drops on cancel/timeout/aclose/missing-stream — diffusion_engine.py:1182 self.executor.drop_output(async_output_id), 1186 self.executor.drop_output(async_output_id), 1209 self.executor.drop_output(async_output_id), 1593 self.executor.drop_output(async_output_id). Adverse effect: the implementation docstring maintainers read still claims abort-only, so it does not match the updated abstract contract or the new required release paths.

- **[P2] `step_streaming` still does a bare `async for` over `get_output_stream` (`diffu…** — `vllm_omni/diffusion/diffusion_engine.py`
Evidence for `step_streaming` still does a bare `async for` over `get_output_stream` (`diffu…

step_streaming still does a bare async for over get_output_stream (diffusion_engine.py:480-481) and parks at yield formatted_outputs (:527). This PR moved materialization/drop_output into get_output_stream.finally (:1202-1209) and wrapped only async_add_req_and_stream_response with async with aclosing(self.get_output_stream(request_id)) (:1215). Streaming serving (unchanged) iterates step_streaming and yields without aclosing (stage_diffusion_proc.py:223-227; also inline_stage_diffusion_client.py:158-162). While parked after a non-finished output, a later queued async_output_id stays in _unclaimed_async_outputs; abandoning or aclose()ing step_streaming does not enter that finally, so those ids are not drop_output'd. Wrap step_streaming the same way: async with aclosing(self.get_output_stream(request_id)) as generator:.

Evidence: This PR’s new drop contract lives in get_output_stream.finally: vllm_omni/diffusion/diffusion_engine.py:1206 abandoned_ids = self._unclaimed_async_outputs.pop(queue, set()) and :1208-1209 for async_output_id in abandoned_ids: self.executor.drop_output(async_output_id). Only async_add_req_and_stream_response forces that finally on aclose: vllm_omni/diffusion/diffusion_engine.py:1215 async with aclosing(self.get_output_stream(request_id)) as stream:. step_streaming does not: :480 generator = self.get_output_stream(request_id) then :481 async for output in generator: then parks at :527 yield formatted_outputs. Trigger (unchanged by this diff, present in the PR-time tree): vllm_omni/diffusion/stage_diffusion_proc.py:223 async for results in self._engine.step_streaming(request): and :227 yield result — after a non-finished yield the inner gen stays at :1196 yield output, so a later _put_output(..., async_output_id=...) is never dropped if the serving/step_streaming consumer is abandoned without aclose. Adverse effect: unclaimed async outputs leak (the gap test_closing_response_stream_closes_inner_stream covers only the aclosing wrapper, not step_streaming).

Suggestion: async with aclosing(self.get_output_stream(request_id)) as generator:
async for output in generator:


🤖 This review was generated by InferMatrix Copilot, an open-source repo-maintenance agent for PR review, CI debugging and issue triage. Try it on your own repo, and ⭐ star it if it helped!

_check_and_start_background_loop=mocker.AsyncMock(),
_prepare_request_for_admission=mocker.Mock(return_value=request),
_add_prepared_request=mocker.Mock(return_value="timed-request"),
get_output_stream=output_stream,

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] test_step_streaming_excludes_output_wait_from_execution_time patches get_output…

Evidence and suggested fix

test_step_streaming_excludes_output_wait_from_execution_time patches get_output_stream and yields a DiffusionOutput already populated with output_ready_wait_time (tests/diffusion/test_diffusion_engine.py:880,895), so it only proves step_streaming subtraction. Production writes the field in get_output_stream after wait_output_ready (vllm_omni/diffusion/diffusion_engine.py:1177). The only real-stream check is test_async_output_is_claimed_after_materialization's >= 0.0 (tests/diffusion/test_diffusion_engine_cleanup.py:336), which also holds for the dataclass default 0.0 (vllm_omni/diffusion/data.py:1816). Deleting the assignment would not fail tests, and diffusion_engine_exec_time_ms would include the wait. Drive a non-zero wait through real get_output_stream (completed executor future + patched perf_counter) and assert output.output_ready_wait_time / output_ready_wait_time_ms and diffusion_engine_exec_time_ms.

Evidence: Trigger: the new timing test never exercises the production writer. tests/diffusion/test_diffusion_engine.py:880 output = DiffusionOutput(stage_durations={"denoise": 10.0}, output_ready_wait_time=output_wait) and :895 get_output_stream=output_stream, inject the field and mock the stream. Production assignment is vllm_omni/diffusion/diffusion_engine.py:1177 output.output_ready_wait_time = time.perf_counter() - output_ready_wait_start_time. Unmet requirement: tests/diffusion/test_diffusion_engine_cleanup.py:336 assert materialized_output.output_ready_wait_time >= 0.0 still passes if that assignment is deleted, because vllm_omni/diffusion/data.py:1816 output_ready_wait_time: float = 0.0 is the default — then step_streaming exec_total_time = time.perf_counter() - exec_start_time - output_ready_wait_time would count materialization wait inside diffusion_engine_exec_time_ms with no test failure.

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

Labels

bug Something isn't working diffusion codes related to diffusion models high priority high priority issue, needs to be done asap

Projects

None yet

Development

Successfully merging this pull request may close these issues.

6 participants