Repository navigation
[Perf][Engine] Event-driven orchestration loop (opt-in) — S1 of #4855 - #5221
Conversation
…DRIVEN_ORCH S1 of the scheduling redesign (RFC vllm-project#4855 problem 3): the orchestrator's coordination layer is poll-driven — one coroutine serially scans every stage replica via asyncio.wait_for(client.get_output_async(), 0.001) and sleeps 1 ms when idle, and the serving-side final-output drain does get_nowait() + 1 ms sleep. Every audio chunk crosses two-to-three 1 ms poll boundaries of pure jitter, and the wait_for emulation creates and cancels ~2 tasks per replica per millisecond. With VLLM_OMNI_EVENT_DRIVEN_ORCH=1 (default off, legacy loop unchanged): - Orchestrator: one reader task per live LLM stage replica awaits client.get_output_async() directly — the same pattern vLLM's own AsyncLLM output handler uses — and feeds a single dispatch queue that this loop consumes serially, so routing/handling semantics are identical to the legacy loop. Readers are reconciled against live_replica_ids() (and client swaps) on a 0.5 s idle tick. Diffusion stages keep their nowait-poll contract via a per-pool poller task. - Serving drain: AsyncOmniEngine.get_output_blocking_async() blocks on the janus queue's condition variable in a dedicated drain thread (immediate wakeup per message); the 1 s timeout only bounds the existing orchestrator-liveness check. The shared per-output handling is extracted into _process_llm_stage_outputs()/_fatal_engine_dead() so both loops run the exact same code; the legacy loop's control flow is otherwise untouched. Tests: tests/engine/test_orchestrator_event_driven.py re-runs the legacy scenario matrix (two-stage LLM, diffusion, async-chunk, abort, shutdown, multi-replica, engine-dead fatal) through the event-driven loop, plus reader-reconcile-on-client-swap and blocking-drain unit tests. Signed-off-by: Yueqian Lin <gao@yueqian.us>
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
Resolves conflicts in the orchestration loop against three upstream changes that landed in the same code the S1 commit refactored. `_process_llm_stage_outputs()` adopts main's body from vllm-project#4974 (per-stage metrics fix: `iteration_stats` now gates on `raw_outputs.outputs` rather than `scheduler_stats`, and the record gate widened to either) and from the streaming-segment work (`segment_token_ids`, `segment_output_metadata`). `_fatal_engine_dead()` is replaced by `_handle_engine_dead()`, which carries main's per-replica fault isolation: evict the dead replica, close its duplex sessions, and only escalate to `_fatal_error` when the stage has no available replicas left. The helper returns whether the failure was fatal so both loops keep main's control flow -- continue on a survivable death, re-raise otherwise. The event-driven loop's reader reconcile and diffusion poller now walk `available_replica_ids()` instead of `live_replica_ids()`. Eviction adds to `_unavailable_replicas` without clearing `clients[replica_id]`, so a `live_replica_ids()` walk would respawn a reader for an evicted replica on every reconcile tick and re-raise the same EngineDeadError indefinitely. The legacy loop already iterated `available_replica_ids()`; this restores parity. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
… loops Neither suite exercised the non-fatal branch of the EngineDeadError handler: a stage replica dying while others remain must report fatal=False, evict the replica, and keep serving. Added to the shared error-handling suite and to the event-driven parity matrix so both loops are held to it. The test also pins that the evicted replica stops being polled, which is what distinguishes an available_replica_ids() walk from a live_replica_ids() one. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
Rebuilds the S1 event-driven loop on top of vllm-project#4285 (per-replica fault isolation), which landed the same shared-handler extraction this branch was carrying and changed what a dead replica means. Upstream now evicts a dead replica via pool.evict_replica() -- which clears the client slot -- and keeps the orchestrator alive, failing only the requests stranded on that replica. The branch's own EngineDeadError handling assumed the older fail-stop contract and is dropped in favor of upstream's _handle_dead_replica(); both loops call it and neither tears the server down. What the branch still adds on top of upstream: - _process_llm_stage_outputs() extracts upstream's LLM poll-handling block so the event-driven dispatcher runs byte-identical per-output handling. The EngineDeadError catch stays with the callers, since eviction needs the replica id. - _orchestration_loop_event_driven() replaces the 1 ms poll cadence with one reader task per available LLM replica awaiting get_output_async() directly, feeding a single serial dispatch queue. - The diffusion poller task carries its own EngineDeadError catch. Upstream's legacy loop covers the diffusion poll under the same try as the LLM poll; once the poll moves into a separate task, that catch has to move with it or a diffusion replica death escapes as a raw reader failure. - Reader reconcile and the diffusion poller walk available_replica_ids() for parity with the legacy loop, so an evicted or membership-disabled replica stops being read. Eviction is reconciled immediately rather than at the next 0.5 s idle tick. Tests: the parity matrix now re-runs upstream's eviction suite through the event-driven loop, replacing the branch's own dead-replica test -- upstream's coverage is broader and encodes the current contract. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
The diffusion poller task reports a whole-task failure with replica_id=-1, since it cannot attribute one. Routing that into _handle_dead_replica() calls pool.evict_replica(-1), which raises ValueError for an out-of-range id -- so a survivable stage death would escape the dispatcher as a ValueError and tear the orchestrator down, the opposite of vllm-project#4285's contract. Only evict when the failure carries a real replica id; an unattributed one falls through to the existing reader-failure path. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
…eview Five fixes to the event-driven orchestration loop, all parity gaps against the legacy poll loop rather than issues with the merge itself. Stats-only batches were dropped. The reader skipped any batch with no request outputs, but StagePool._poll_stage_raw deliberately keeps scheduler-only and finished-only batches -- omni schedulers emit SchedulerStats on throttled ticks with no output, and dropping those loses the KV/queue gauges for that interval. The reader now mirrors _poll_stage_raw's exact condition. The four stats tests from the legacy suite are added to the parity matrix, which is what catches it. Reconcile could starve. It ran only when the dispatch queue went idle, so under continuous traffic asyncio.wait always returned on an output and reconcile never fired -- a replica registering at runtime would never get a reader and its outputs would never drain. Reconcile now runs on a wall-clock schedule. Diffusion pollers were spawned once at startup. A pool with no live replica yet reports stage_type None, so a stage whose first replica registers at runtime got no poller at all. Poller lifecycle moved into reconcile alongside readers. A crashed reader could be respawned before its error was dispatched, since Task.done() is true for a failed task; the respawn re-raised against the same dead client and queued a duplicate eviction. Reconcile now leaves a done reader alone until eviction removes the slot, and both eviction sites skip a replica that is no longer live. Readers and pollers retired mid-run (client swap, eviction) were cancelled but never awaited. They are collected and gathered with the rest on teardown. Also moves note_loop(idle=False) after routing so an evicted replica does not count as an active tick, matching where the legacy loop sets idle=False. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
| # replica awaits `client.get_output_async()` directly — the same pattern vLLM's | ||
| # own AsyncLLM output handler uses — and feeds a single serial dispatch queue. | ||
| # Default off; the legacy poll loop remains the fallback. | ||
| _EVENT_DRIVEN_ORCH_ENV = "VLLM_OMNI_EVENT_DRIVEN_ORCH" |
There was a problem hiding this comment.
[P2] This flag is user-facing but appears nowhere in docs/ or recipes/. It selects a different serving runtime mode (orchestrator loop + final-output drain flip together), so operators need at least: identifier, accepted values (1/true/yes/on — the PR body only says =1), default off, scope, and the PR body's caveats as known limitations (A/B predates the rebuild on the fault-isolation work; diffusion-poller branch unit-tested only).
Repo precedent documents env vars next to their feature (VLLM_OMNI_ORCH_MONITOR_PATH in docs/contributing/profiling.md, SPEAKER_* table in docs/serving/speech_api.md). docs/serving/speech_api.md — streaming TTS serving is the measured path — looks like the natural home; a one-liner on the draft engine-orchestration design page would also fit.
There was a problem hiding this comment.
Documented in 5f737f5.
docs/serving/speech_api.md gains an "Orchestration Loop (experimental)" section following the SPEAKER_* precedent you pointed at: identifier, default 0 (off), the full accepted set, where to set it (the stage-0 orchestrator process), and the PR body's caveats as known limitations. Worth noting the accepted set is wider than the =1 the PR body claimed: 1, true, yes and on, matched case-insensitively after surrounding whitespace is stripped.
Also took the second suggestion: docs/design/module/engine_orchestration.md now carries a one-liner under Contract status recording that the mode changes poll cadence only, so the routing, ordering and terminal-state contracts on that page hold identically on both loops.
One note on the diff: docs/serving/speech_api.md already fails markdownlint-cli2 with 34 pre-existing errors, and that hook is in the CI SKIP list as historical debt. I kept this change purely additive rather than let --fix reformat 30 unrelated table lines into the PR.
| next_reconcile = _time.monotonic() + _ORCH_READER_RECONCILE_INTERVAL_S | ||
| _reconcile_readers() | ||
| if not done: | ||
| self._orch_monitor.note_loop(idle=True) |
There was a problem hiding this comment.
[P3] Under this loop the orch-monitor counts change meaning: docs/contributing/profiling.md documents loop_idle/loop_active as poll-loop iterations per 1 s window, but idle ticks drop from ~1000/s to one per 0.5 s reconcile timeout (busy still records one per routed output), so window counts and loop_active_pct are not comparable across modes. Worth one sentence in profiling.md alongside the flag docs above.
There was a problem hiding this comment.
Added in 5f737f5.
docs/contributing/profiling.md now states that loop_idle and loop_active are iteration counts whose absolute scale is tied to which loop is running: roughly one iteration per millisecond when the default poll loop is idle, versus one per 0.5 s reconcile timeout under the event-driven loop, with a busy orchestrator still recording one per routed output on either. It closes with the point you made, that window counts and the loop_active_pct summary are not comparable across the two modes, and links to the flag section in the Speech API page.
…nitor counts The flag selects a different serving runtime mode but appeared nowhere in docs/ or recipes/, so operators had no way to learn its name, accepted values, default, or scope. docs/serving/speech_api.md gains an "Orchestration Loop" section following the SPEAKER_* precedent: the variable, its default, the full set of accepted truthy values, where to set it, and the PR's caveats as known limitations. docs/contributing/profiling.md notes that loop_idle/loop_active are iteration counts whose scale is tied to the loop mode, so window counts and loop_active_pct are not comparable between the poll loop and the event-driven loop. docs/design/module/engine_orchestration.md gets a one-line pointer recording that the opt-in mode changes poll cadence only, leaving the routing, ordering and terminal-state contracts intact. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
…ventory vllm-project#6217 landed an environment-variable inventory with drift tests after this branch documented its flag in the Speech API page. Three of those tests fail for an unregistered name: the static scanner finds the `os.environ.get` in `orchestrator.py` and requires a classification, every public `VLLM_OMNI_*` must appear on the reference page, and the reviewed snapshot pins the `PUBLIC_OMNI` count. Classified `PUBLIC_OMNI` — it is an operator-facing serving switch, same bracket as `VLLM_OMNI_ORCH_MONITOR_PATH`. Reference-page row records the full accepted set, that unrecognized values leave the legacy poll loop selected, and an Experimental lifecycle. Snapshot count 22 -> 23. The narrative section in `docs/serving/speech_api.md` stays: the reference page is the registry, that page is where an operator meets the feature. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
Main registered VLLM_OMNI_ASYNC_OUTPUT_TIMEOUT and dropped the stale ROSVOT_SOURCE_DIR entry after this branch last merged. Both sides inserted one name at the same spot in _PUBLIC_OMNI and one row at the same spot in the reference table, so the merge kept both in sorted order. The snapshot count collapsed to a single 22 -> 23 edit; with both names present it is 24. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
…ion loops vllm-project#6150 added per-output metrics to the LLM poll block this branch extracts into `_process_llm_stage_outputs()`, so the conflict resolves by folding them into the helper: the `kv_wait_s` observation (kept ahead of the `req_state is None` continue, as its comment requires, so it lands in all modes) and the waiting- count update derived from `scheduler_stats.num_waiting_reqs`. The diffusion side needed more than a merge. vllm-project#6150 also began piggybacking scheduler metrics on diffusion outputs and emitting a metrics-only sentinel under `DIFFUSION_METRICS_ONLY_REQUEST_ID`. The legacy loop drains and skips it inline; the event-driven poller had no such branch and would have routed the sentinel to `_handle_processed_outputs` as a real request output. Extracted the drain into `_absorb_diffusion_metrics()` and called it from both loops, marking the tick active before skipping so the monitor accounting matches. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
The test read Orchestrator._orchestration_loop and asserted the source positions of three literals. Diffusion metric absorption has since moved into _absorb_diffusion_metrics so both orchestration loops share one implementation, so _update_stage_replica_waiting( is no longer inside that function and .index() raised ValueError: substring not found. Assert the invariant where it now lives. The snapshot-before-verdict ordering is checked inside the helper, and the absorb-before-route guard is checked in every loop that polls diffusion output, which covers the event-driven loop the old single-function assertion could not see. Verified by mutation: dropping the guard from the event-driven loop, routing before absorbing in the legacy loop, and snapshotting after the verdict in the helper are each caught. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
… into both loops vllm-project#6564 changed `_apply_raw_terminal_stage_finish` to report whether it finished a raw terminal request, and the legacy loop now collects those ids per poll and drains them through `_finish_raw_terminal_requests` after routing. That accumulator lives across the extracted helper, so `_process_llm_stage_outputs` takes it as a parameter and fills it. The event-driven dispatcher creates and drains one per dispatched item, matching the legacy loop's per-poll scope -- without it, realtime audio completions would simply never fire under the flag, silently and only for that mode. Signed-off-by: Yueqian Lin <linyueqian@outlook.com>
hsliuustc0106
left a comment
There was a problem hiding this comment.
lgtm, I think we need to evaluate more model perf gains later
|
我重新对了最终 head
有两个边界需要说清楚:
我当前环境没有 |
…project#4855 (vllm-project#5221) Signed-off-by: Yueqian Lin <gao@yueqian.us> Signed-off-by: Yueqian Lin <linyueqian@outlook.com> Co-authored-by: Yueqian Lin <gao@yueqian.us>
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>
…project#4855 (vllm-project#5221) Signed-off-by: Yueqian Lin <gao@yueqian.us> Signed-off-by: Yueqian Lin <linyueqian@outlook.com> Co-authored-by: Yueqian Lin <gao@yueqian.us>
Purpose
S1 of the scheduling redesign in RFC #4855 (problem 3: the orchestrator coordination layer is poll-driven and keeps the GPU idle under load). This makes the orchestration loop event-driven behind an opt-in flag (
VLLM_OMNI_EVENT_DRIVEN_ORCH=1, default off). The legacy poll loop is untouched when the flag is unset.What is slow today
The
Orchestratorruns one coroutine that serially scans every stage replica withasyncio.wait_for(client.get_output_async(), timeout=0.001)andasyncio.sleep(0.001)when idle (orchestrator.py), and the serving-side final-output drain busy-pollsget_nowait()+sleep(0.001)(entrypoints/async_omni.py). Two consequences:wait_for(..., 0.001)emulation creates and cancels ~2 tasks per replica per millisecond, so the process spins at ~1 kHz even when completely idle.The underlying stage client is vLLM's own
AsyncMPClient, whoseget_output_async()is natively awaitable and surfacesEngineDeadErrorthrough the same queue — upstreamAsyncLLMconsumes it with a plainawait. This change restores that pattern for the omni orchestrator.Changes
vllm_omni/engine/orchestrator.py— new_orchestration_loop_event_driven()(selected by the flag): one reader task per available LLM stage replicaawaitsclient.get_output_async()directly and feeds a single dispatch queue that the loop consumes serially, so routing/handling semantics are identical to the legacy loop; only the poll cadence is removed. Diffusion stages keep their nowait-poll contract via a per-pool poller task. Readers and pollers are reconciled againstavailable_replica_ids()and client swaps on a 0.5 s wall-clock schedule, which covers elastic membership and replica eviction. The per-output handling is extracted into_process_llm_stage_outputs()and reused by both loops, so there is one code path for behavior. A one-line enabled-log announces the loop mode + reader/poller counts.vllm_omni/engine/async_omni_engine.py—get_output_blocking_async(): blocks on the janus queue's condition variable in a dedicated drain thread (immediate wakeup per message), returningNoneon a 1 s timeout so the caller keeps its liveness check. The executor is lazily created and torn down inshutdown().vllm_omni/entrypoints/async_omni.py— the final-output drain usesget_output_blocking_async(1 s)under the flag (condition-variable wakeup) instead ofget_nowait()+ 1 ms sleep; legacy path unchanged. Falls back to the legacy drain if the engine variant lacks the new method.tests/engine/test_orchestrator_event_driven.py— re-runs the legacy scenario matrix through the event-driven loop (two-stage LLM, diffusion, async-chunk, abort, shutdown, multi-replica, the stats-recording cases, and the per-replica eviction suite), plus reader-reconcile-on-client-swap, blocking-drain units, and flag parsing. A fixture guard asserts the flag actually took effect so parity can't silently exercise the old loop.Default off; no behavior change unless
VLLM_OMNI_EVENT_DRIVEN_ORCH=1.Rebuilt on top of #4285
This branch predates #4285 (per-replica fault isolation), which rewrote the same handler the branch had refactored and landed an equivalent shared-handler extraction. The branch has been rebuilt on top of it rather than merged around it: its own
EngineDeadErrorhandling is gone, and both loops now call upstream's_handle_dead_replica(), so a dead replica is evicted and the server stays up on either loop. The event-driven parity matrix runs upstream's eviction suite, including the diffusion-poller death case.Two things the event-driven loop has to do that the legacy loop gets for free. Upstream wraps the diffusion poll and the LLM poll in one
try, so moving the diffusion poll into its own task means carrying thatEngineDeadErrorcatch with it. And a whole-task poller failure carries no replica id, so it must not reachevict_replica(), which rejects an out-of-range id.Test Plan
vLLM Version: 0.27.0
vLLM-Omni Commit: merged with
3ad64c6a1(current main at the time of writing)pytest tests/engine/.Test Result
Unit: 506 passed, 1 skipped across
tests/engine/, reproduced on two consecutive runs. The parity suite re-runs 22 legacy scenarios through the event-driven loop.e2e (H20, env-toggle A/B, ON = event-driven, OFF = legacy):
Idle CPU of the API-server process (
pidstat, warmed, no requests):Serving A/B:
Correctness: seeded generation is byte-for-byte identical across arms (c=1 total audio exactly equal; per-request wav byte-lengths equal), whisper transcripts match targets on both arms, and no reader crash-loops / stuck requests / missing chunks across 3 server boots (2x Qwen3-TTS, 1x GLM-TTS).
Read: idle-CPU spin removed (~35x); TTFP tails improve at moderate concurrency (c=8) where poll jitter is a visible fraction of TTFP; c=1 and c=32 are parity (at c=32 TTFP is dominated by admission queueing, which is S3 in the RFC, not this change); throughput unchanged.
Honest caveats
stage_type: diffusion), so the diffusion-poller branch is exercised by the unit tests, not by a TTS smoke on the box.cc @Sy0307 for review — this is the S1 piece of the #4855 scheduling redesign; happy to share the profiling harness and raw e2e artifacts.