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. |
e573a92 to
de7e57b
Compare
|
@tzhouam would you mind taking a look when you get a chance? This PR replaces the 1ms busy-poll on the output queue with an event-driven |
|
The queue-direction fix is correct and the tests look good. One minor convention note: the PR title uses |
| if self.rpc_output_queue is not None: | ||
| self.rpc_output_queue.sync_q.put_nowait(error_msg) | ||
| self.rpc_output_queue.shutdown() | ||
| except Exception: |
There was a problem hiding this comment.
The single except Exception: pass covers both queue blocks. If output_queue.shutdown() raises (e.g. queue already shut down), the rpc_output_queue block is skipped entirely, so its put_nowait+shutdown() never run. Consider splitting into two independent try/except blocks so a failure on one queue does not skip the other.
There was a problem hiding this comment.
thanks for pointing it out! I just updated this part
got it! updated the title, thanks! |
|
Hi @hsliuustc0106 , both review comments have been addressed, would you mind taking another look when you get a chance? Thanks! |
…queue shutdown Replace the 1ms busy-poll on the engine output queue with an event-driven await. The async consumer (`_final_output_loop`) read the queue via `sync_q.get_nowait()`, returned `None` on empty, and slept 1ms before retrying. That added up to one sleep-interval of latency per message and burned idle CPU. It now parks on `async_q.get()` and wakes when a message arrives. To make this correct and hang-free: - `_bootstrap_orchestrator` hands the Orchestrator the queue's *sync* side, and it pushes outputs with `put_nowait()` (it never awaits output), so `async_q` binds to the consumer loop that awaits it. Renamed the param/attr `output_async_queue` -> `output_sync_queue` and narrowed the type to `janus.SyncQueue`, which contrasts with the still-async `request_async_queue` / `rpc_async_queue`. - On a crash, `_bootstrap_orchestrator` enqueues the fatal `ErrorMessage` and then calls `output_queue.shutdown()`. With `immediate=False`, shutdown drains the fatal message first (so the caller still sees `fatal=True`), then wakes any still-parked getter with `QueueShutDown`. So `await get()` can never hang on a dead engine, whether or not the put landed. `try_get_output*` and `collective_rpc` convert `QueueShutDown` to the existing RuntimeError. Tests: rewrite the async output tests over a real janus queue (event-driven wakeup, fatal-drain-then-raise, shutdown wakes a parked getter), update the orchestrator harness to pass the sync side, and add `put_nowait` to the kv-sender mock queue. Merge integration: the new MembershipController pushes its replica-disappeared error through the same sync side (put_nowait), and the crash handler shuts down each queue independently so a failure on one does not skip the other. Signed-off-by: zeningc <zening.chen@yahoo.com>
Standalone, GPU-free microbenchmark isolating the Orchestrator -> server handoff over the janus output queue. Mirrors production (producer thread -> async consumer) and compares the old 1ms busy-poll against the event-driven await, reporting per-message wakeup latency and idle CPU/wakeups. Representative result: mean latency 574us -> 145us (p50 575us -> 96us), and idle wakeups ~900/s (~5% of a core) -> ~0. Signed-off-by: zeningc <zening.chen@yahoo.com>
283bf7c to
26d97a7
Compare
|
Hi @hsliuustc0106! Just wanted to follow up on this when you have a moment. The review comments are addressed, it's rebased onto the latest |
|
we have updated with vllm 0.24.0, can you help update the use case results? |
hsliuustc0106
left a comment
There was a problem hiding this comment.
The output queue rename is incomplete and leaves live failure paths/tests broken.
| ) -> None: | ||
| self.request_async_queue = request_async_queue | ||
| self.output_async_queue = output_async_queue | ||
| self.output_sync_queue = output_sync_queue |
There was a problem hiding this comment.
Please update the remaining output_async_queue references too. git grep still finds live calls at lines 1251 and 1419; since this constructor no longer sets that attr, those paths raise AttributeError instead of sending the terminal output.
There was a problem hiding this comment.
Thanks! I just fixed it, sorry for the oversight, this was pulled from new code when I updated the current branch.
| orchestrator = Orchestrator( | ||
| request_async_queue=request_queue.async_q, | ||
| output_async_queue=output_queue.async_q, | ||
| output_sync_queue=output_queue.sync_q, |
There was a problem hiding this comment.
Please also update tests/engine/test_orchestrator_stage_input_bridge.py; it still passes output_async_queue=... and is marked core_model/cpu, so it fails with the new constructor signature.
There was a problem hiding this comment.
Thanks, fixed as well!
Merging main pulled in new call sites that used the pre-rename output queue API, which this branch had already renamed output_async_queue -> output_sync_queue (put_nowait on the sync side). These were semantic conflicts: git auto-merged them without a textual conflict, so they landed referencing an attribute the constructor no longer sets. - orchestrator.py:1251 (from vllm-project#4079, diffusion request-level batching) and orchestrator.py:1419 (from vllm-project#4257, Aura non-async-chunk path) still called `await self.output_async_queue.put(...)`, which would raise AttributeError on those terminal-output error/edge paths. Convert to `self.output_sync_queue.put_nowait(...)`. - tests/engine/test_orchestrator_stage_input_bridge.py (new file from vllm-project#4257, marked core_model/cpu) constructed Orchestrator with the old `output_async_queue=` kwarg and failed against the new signature. Update to `output_sync_queue=output_q.sync_q`. Tested on vLLM 0.24.0: engine + orchestrator unit tests (41 passed) and a Qwen3-TTS-12Hz-0.6B streaming /v1/audio/speech smoke test. Signed-off-by: zeningc <zening.chen@yahoo.com>
sure! I've ran all the tests with 0.24.0 and updated the details in PR description |
The event-driven output loop change made `AsyncOmniEngine.try_get_output_async` park on `await async_q.get()` and never return None; the loop's per-empty `await asyncio.sleep()` yield point was removed accordingly. Two test doubles still modeled the old contract, returning None synchronously on an empty queue: `FakeAsyncOmniEngine` (tests/entrypoints/test_omni_entrypoints.py) and `_OrchestratorBridgeEngine` (tests/diffusion/test_diffusion_streaming_output.py). Against those doubles `while True: msg = await try_get_output_async()` spins with no real suspension point once the queue drains, starving the event loop so `generate()` never resumes — the loop hangs and CI kills the step on timeout (exit 124: test_async_omni_yields_only_final_stage_outputs and test_streaming_output_reaches_async_omni_from_mock_pipeline). Update both doubles to honor the new contract: park (poll the backing queue and `await asyncio.sleep(0.001)` on empty) until a real message arrives, instead of returning None. Death is still signaled by enqueuing a fatal ErrorMessage, and shutdown() cancels the loop task, so no None sentinel is needed. Production stays event-driven and unchanged. Signed-off-by: zeningc <zening.chen@yahoo.com>
Two gaps in how a dying orchestrator wakes a parked output-queue reader, both now that the reader parks on the queue instead of polling: 1. Death signaling (put fatal ErrorMessage + queue shutdown) lived only in ``_bootstrap_orchestrator``'s ``except Exception``. A BaseException (SystemExit / KeyboardInterrupt / asyncio.CancelledError) or a clean return skipped it, leaving a parked ``await async_q.get()`` hung forever. Move the shutdown() into ``finally`` so any thread exit wakes the reader; the except only records the error text, which finally enqueues as the fatal message first (drained before QueueShutDown, so the caller still sees fatal=True). This also de-duplicates the per-queue loop. 2. janus raises different, unrelated classes on the two sides of a shut-down queue: ``ShutDown`` from the sync side and ``QueueShutDown`` from the async side. The sync readers (``try_get_output``, ``collective_rpc``) only caught ``QueueShutDown``, so a real shutdown leaked a raw ``janus.ShutDown`` instead of mapping to RuntimeError. Catch both via a shared ``_QUEUE_SHUTDOWN_ERRORS`` tuple at all sites. (Surfaced by the janus 2.0 bump; the old test masked it by mocking side_effect=QueueShutDown.) Tests: drive ``_bootstrap_orchestrator`` over real janus queues and assert a BaseException exit still shuts the queues (reader raises RuntimeError, not a hang), an Exception exit drains the fatal message then shuts down, and rewrite the queue-shutdown test to use a real queue instead of a mocked exception. Signed-off-by: zeningc <zening.chen@yahoo.com>
The output queue is now handed to the MembershipController as the queue's sync side (Orchestrator pushes via put_nowait), but the annotations still read asyncio.Queue[EngineQueueMessage], which is neither the real runtime type nor even a janus type. Narrow the attribute and both handler params to janus.SyncQueue[EngineQueueMessage] to match reality and the Orchestrator's own output_sync_queue typing. Signed-off-by: zeningc <zening.chen@yahoo.com>
Omni ReviewBot triage noteAutomated triage of commit
These are automated triage suggestions only — the final decision belongs to the maintainers. |
Omni ReviewBot: no human activity for 15 days@zeningc this pull request has had no human commit, comment or review since 2026-08-31. 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. |
Omni ReviewBot: no human activity for 23 days@zeningc this pull request has had no human commit, comment or review since 2026-08-31. 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. |
Re-port the event-driven output drain onto the refactored engine layout on main (OmniEngineBase / AsyncOmniBase from vllm-project#7413) and on top of the flag-gated event-driven orchestration loop from vllm-project#5221. - omni_engine_base.py: try_get_output_async parks on output_queue.async_q.get(); a shut-down queue maps to RuntimeError on both the sync and async readers; _bootstrap_orchestrator enqueues the fatal ErrorMessage and shuts the output queue down in `finally` so a parked reader wakes on any thread exit. get_output_blocking_async and its drain-thread executor are removed. - orchestrator.py / duplex_orchestrator.py / duplex session manager / membership_controller.py: push outputs with sync_q.put_nowait; output_async_queue -> output_sync_queue (janus.SyncQueue). - async_omni_base.py: the final-output drain is a single `await engine.try_get_output_async()`; the poll branch, the VLLM_OMNI_EVENT_DRIVEN_ORCH drain branch, and both interval constants are gone. The flag now governs only the orchestration loop. - docs: VLLM_OMNI_EVENT_DRIVEN_ORCH describes only the orchestration loop; the drain is always event-driven. - tests: engine output tests over a real janus queue; orchestrator fixtures pass the sync side; vllm-project#5221's drain-thread unit tests removed with the method; fake engines park instead of returning None. Signed-off-by: zeningc <zening.chen@yahoo.com>
|
Rebased onto current What changed in the re-port:
Validation on 0.30.0: 1450 unit tests pass across |
Omni ReviewBot: no human activity for 7 days@zeningc 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 cursor (cursor-grok-4.6-high) under experiment |
Summary
Replace the busy-poll / drain-thread read of the engine output queue with a
plain event-driven
awaiton the queue's async side, and route orchestratordeath through the queue via
shutdown(). The motivation is async correctness,a cleaner death/hang-safety contract, and lower idle CPU.
Background
The Orchestrator drives the multi-stage pipeline in its own background thread
with its own asyncio loop. The server (
OmniEngineBaseplus the API layer)runs on the main asyncio loop.
They communicate over
janusqueues, which expose one underlying queue as async_q(for thread code) and anasync_q(for event-loop code). The rule:whichever side waits for the next message should
await async_q, and the siderunning in a thread should push with
sync_q.put_nowait.The request queue (server → orchestrator) already follows this: the
orchestrator is the consumer and awaits
async_qon its loop, while the serverpushes with
sync_q.put_nowait. The output queue (orchestrator → server:per-request tokens / audio chunks, the one this PR touches) is the mirror image,
with the server as the consumer. It should await
async_qtoo, but thepre-change code had the orchestrator awaiting the async side, which left the
server unable to await the queue at all.
Why
The output reader is a hot path that every generated token/audio chunk flows
through. Two reasons to change how it works:
The server can await instead of poll (or thread-hop). A
janusqueue'sasync_qattaches to a single event loop (the first one to use it) and canonly be touched from that loop. The old code had the orchestrator do
await async_q.put(), which attachedasync_qto the orchestrator's loop.The server runs on a different loop, so it could not
awaitthe queue atall. Legacy path: poll
sync_q.get_nowait()with a 1 ms sleep. [Perf][Engine] Event-driven orchestration loop (opt-in) — S1 of #4855 #5221 path:block on
sync_q.get(timeout=1)inside a dedicated executor thread. Bothonly exist because the queue was bound to the wrong loop. Now the
orchestrator pushes with
sync_q.put_nowait(), which attaches nothing, soasync_qattaches to the server's loop the first time the server awaits it,and the server reads with a plain
await async_q.get().Idle cost at fleet scale. The poll loop wakes ~900×/s and burns ~5% of a
core per engine while idle. Parking on
get()drops it to ~0, with no extrathread per engine.
Design / what changed
A parked
await get()cannot notice the engine dying on its own, so theorchestrator thread signals through the queue on any exit (Exception,
BaseException, or clean return):
The RPC output queue is unchanged: a crash still enqueues a fatal
ErrorMessagethere, whichCorrelatedRpcClient's router already broadcaststo waiters.
Files:
engine/omni_engine_base.py:try_get_output_asyncparks onasync_q.get()and maps a shut-down queue to the existing
RuntimeError;try_get_output(sync) maps it too;
_bootstrap_orchestratorenqueues the fatal message andshutdown()s the output queue infinally;get_output_blocking_asyncandits lazily created drain executor are removed.
engine/orchestrator.py: pushes outputs withput_nowaiton the sync side;param/attribute renamed
output_async_queue→output_sync_queue, typedjanus.SyncQueue, contrasting with the still-asyncrequest_async_queue/rpc_async_queue.duplex_orchestrator.py/duplex/session/manager.pyfollow the rename (the session manager already preferred
put_nowait).engine/membership_controller.py: pushes withput_nowait; type narrowed.entrypoints/async_omni_base.py: the final-output drain is a singleawait engine.try_get_output_async(); theNone/sleep branch, theVLLM_OMNI_EVENT_DRIVEN_ORCHdrain branch, and both interval constants aregone.
VLLM_OMNI_EVENT_DRIVEN_ORCHnow documents only the orchestrationloop; the drain is always event-driven.
janusqueue (event-driven wakeup,fatal-drain-then-raise, shutdown wakes a parked getter, bootstrap death
signaling on Exception and BaseException); orchestrator fixtures pass the
sync side; the [Perf][Engine] Event-driven orchestration loop (opt-in) — S1 of #4855 #5221 drain-thread unit tests are removed with the method;
fake engines in entrypoint/streaming tests park instead of returning
None.benchmarks/engine/: standalone queue-handoff microbenchmark.Relationship to #5221
#5221 fixed the orchestrator-side poll (one reader task per stage replica) and
is kept as is, flag and all. For the serving-side drain it could not await the
queue because of the loop-binding problem above, so it added a drain thread
gated by the same flag. This PR removes the need for that thread by fixing the
queue direction, so the drain is event-driven for every pipeline with no flag
and no extra thread.
_orchestration_loop_event_drivenand its parity suiteare untouched.
Testing
Re-run on vLLM 0.30.0 after merging
mainat6daf5b30f.Unit tests (all of
tests/engine/, including #5221's event-driven paritysuite, plus the entrypoint, streaming-interaction, and stat-logger suites that
touch the output queue):
raw
pre-commiton every changed file: ruff, format, typos, markdownlint, SPDX,CI-marks and forbidden-imports hooks pass. The
mypy-3.10hook reports thesame 14 pre-existing errors on
mainfor these files (verified on a cleancheckout; none are on changed lines).
Microbenchmark: isolates the queue handoff (producer thread → async consumer),
no model in the loop:
raw
Serving smoke on this branch (Qwen3-TTS 0.6B, 1× RTX 4090, vLLM 0.30.0,
/v1/audio/speech,stream:true, pcm):raw
End-to-end A/B (Qwen3-TTS 0.6B, 1× RTX 4090, 3 × 200 = 600 streaming requests/arm
per cell) from the earlier vLLM 0.24.0 run, confirming no regression: every pooled
TTFP/E2E mean sits within ±1.3%. Not re-run on 0.30.0; the mechanism under test
(queue read on the serving side) is unchanged by the rebase.
raw (baseline = 1 ms poll, candidate = event-driven; vLLM 0.24.0)
Client: N concurrent streaming POSTs to
/v1/audio/speech(stream:true, pcm);TTFP = wall time to first PCM byte, E2E = to last byte. Fixed 8-sentence prompt
set, warmup 6, 3 runs × 200 requests (600 pooled) at concurrency 1 / 8 / 16.
Means below are pooled over the full 600 samples/cell.
(0 errors on either arm. Both arms identical except the queue read mechanism:
baseline polls
sync_q.get_nowait()+ 1 ms sleep, candidate parks onasync_q.get(). p99 tails stay within a few ms of each other and are omitted;at this sub-second, compute-bound scale they're dominated by model time.)
Note
The serving numbers do not justify this on their own; inference is compute-bound
and the gain hides under model time. Merge it for the async correctness (the poll
loop and the drain thread were both workarounds for binding the queue to the
wrong loop), the cleaner death/hang-safety contract, and the idle-CPU saving
that grows with fleet size.