Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 53 additions & 0 deletions benchmarks/engine/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
# Engine microbenchmarks

Standalone microbenchmarks for engine-internal hot paths. No model or GPU
required — they isolate a single mechanism so the effect of a change is
measurable on its own.

## `queue_handoff_latency.py`

Quantifies the Orchestrator → server handoff over the janus output queue, the
path changed when the server's output reader moved from a 1 ms busy-poll to an
event-driven `await`.

It mirrors production: a producer **thread** (the Orchestrator runs in its own
thread) emits messages to an **async consumer** (the server event loop), and
compares two consumer strategies:

- **OLD** — `sync_q.get_nowait()`, return `None` on empty, `asyncio.sleep(1ms)`,
retry.
- **NEW** — `await async_q.get()`, park until a message arrives.

```bash
python benchmarks/engine/queue_handoff_latency.py
python benchmarks/engine/queue_handoff_latency.py --poll-ms 1.0 --messages 3000 --runs 3
```

### Representative result

```text
== per-message wakeup latency (producer thread -> async consumer) ==
OLD poll : mean= 573.8us p50= 574.7us p99= 1103.6us max= 1280.8us
NEW await: mean= 144.7us p50= 96.3us p99= 398.3us max= 585.1us

== idle cost (no messages flowing) ==
OLD poll : wakeups= 1795 ( 897.5/s) cpu= 111.6 ms over 2s idle
NEW await: wakeups= 1 ( 0.5/s) cpu= 0.4 ms over 2s idle
```

### How to read it

- **Latency** — the old poll adds ~0.5 ms on average and up to ~1 ms per message
(a message arriving at a uniformly random point in a 1 ms poll cycle waits
half a cycle, hence p50 ≈ 0.5 ms, p99 ≈ the full interval). The new path cuts
mean latency ~4× and p50 ~6×. The residual ~100 µs in the new path is the
inherent cost of waking an asyncio task from another thread; polling cannot
beat that.
- **Idle cost** — the old loop wakes ~900×/s and burns ~5% of a core per idle
engine; the new one parks and costs almost nothing.

The latency win only materializes while the consumer is actually waiting — but
that is the normal streaming regime, where token/audio-chunk generation is far
slower than a queue op, so the consumer waits for essentially every output.
Numbers vary with machine and event-loop timer granularity; run it on the
target host for a local baseline.
162 changes: 162 additions & 0 deletions benchmarks/engine/queue_handoff_latency.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,162 @@
#!/usr/bin/env python3
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM-Omni project
"""Microbenchmark for the engine output-queue handoff.

Isolates the Orchestrator -> server handoff over the janus output queue (no
model or GPU required) to quantify the difference between two consumer
strategies:

* ``OLD`` -- ``sync_q.get_nowait()`` returning ``None`` on empty, then
``asyncio.sleep(poll_interval)`` before retrying (the pre-change behaviour,
``_FINAL_OUTPUT_IDLE_SLEEP_S = 1ms``).
* ``NEW`` -- ``await async_q.get()``, parking until a message arrives.

It mirrors production: a producer **thread** (the Orchestrator runs in its own
thread) emits messages to an **async consumer** (the server event loop). Two
things are measured:

1. Per-message wakeup latency, with messages arriving at a random phase of the
poll cycle (the streaming regime, where token/chunk generation is far slower
than queue ops so the consumer is almost always waiting).
2. Idle cost (wakeups/s and CPU) while no messages are flowing.

Usage::

python benchmarks/engine/queue_handoff_latency.py
python benchmarks/engine/queue_handoff_latency.py --poll-ms 1.0 --messages 3000 --runs 3
"""

from __future__ import annotations

import argparse
import asyncio
import queue
import random
import statistics
import threading
import time

import janus


def _producer(q: janus.Queue, n: int, gap_max_s: float, emit_ts: list[float], ready: threading.Event) -> None:
"""Emit ``n`` messages from a background thread, like the Orchestrator."""
ready.wait()
for i in range(n):
# A random inter-message gap models generation being far slower than
# queue ops, so the consumer is genuinely waiting most of the time.
time.sleep(random.uniform(0.0, gap_max_s))
emit_ts.append(time.perf_counter())
q.sync_q.put_nowait(i)


async def _consume_old(q: janus.Queue, n: int, poll_s: float, recv_ts: list[float]) -> None:
got = 0
while got < n:
try:
q.sync_q.get_nowait()
except queue.Empty:
await asyncio.sleep(poll_s)
continue
recv_ts.append(time.perf_counter())
got += 1


async def _consume_new(q: janus.Queue, n: int, poll_s: float, recv_ts: list[float]) -> None:
for _ in range(n):
await q.async_q.get()
recv_ts.append(time.perf_counter())


async def _measure_latency(consumer, *, n: int, poll_s: float, gap_max_s: float) -> list[float]:
q: janus.Queue = janus.Queue()
emit_ts: list[float] = []
recv_ts: list[float] = []
ready = threading.Event()
th = threading.Thread(target=_producer, args=(q, n, gap_max_s, emit_ts, ready), daemon=True)
th.start()
task = asyncio.ensure_future(consumer(q, n, poll_s, recv_ts))
ready.set()
await task
th.join()
q.shutdown()
return sorted((r - e) * 1e6 for e, r in zip(emit_ts, recv_ts)) # microseconds


async def _measure_idle(strategy: str, *, poll_s: float, secs: float) -> tuple[int, float]:
q: janus.Queue = janus.Queue()
wakeups = 0
cpu0 = time.process_time()
if strategy == "old":
end = time.perf_counter() + secs
while time.perf_counter() < end:
try:
q.sync_q.get_nowait()
except queue.Empty:
wakeups += 1
await asyncio.sleep(poll_s)
else:
task = asyncio.ensure_future(q.async_q.get()) # parks once, never wakes
wakeups = 1
await asyncio.sleep(secs)
task.cancel()
try:
await task
except asyncio.CancelledError:
pass
cpu_ms = (time.process_time() - cpu0) * 1e3
q.shutdown()
return wakeups, cpu_ms


def _summarize(label: str, lat: list[float]) -> None:
print(
f" {label}: mean={statistics.mean(lat):7.1f}us p50={lat[len(lat) // 2]:7.1f}us "
f"p99={lat[int(len(lat) * 0.99)]:7.1f}us max={lat[-1]:7.1f}us"
)


async def _main_async(args: argparse.Namespace) -> None:
poll_s = args.poll_ms / 1e3
gap_max_s = args.gap_max_ms / 1e3
print(
f"poll interval = {args.poll_ms:.1f} ms, messages/run = {args.messages}, "
f"runs = {args.runs}, max inter-message gap = {args.gap_max_ms:.1f} ms\n"
)

print("== per-message wakeup latency (producer thread -> async consumer) ==")
for label, consumer in (("OLD poll ", _consume_old), ("NEW await", _consume_new)):
agg: list[float] = []
for _ in range(args.runs):
agg.extend(await _measure_latency(consumer, n=args.messages, poll_s=poll_s, gap_max_s=gap_max_s))
_summarize(label, sorted(agg))

print("\n== idle cost (no messages flowing) ==")
for label, strategy in (("OLD poll ", "old"), ("NEW await", "new")):
wakeups, cpu_ms = await _measure_idle(strategy, poll_s=poll_s, secs=args.idle_secs)
print(
f" {label}: wakeups={wakeups:6d} ({wakeups / args.idle_secs:7.1f}/s) "
f"cpu={cpu_ms:6.1f} ms over {args.idle_secs:.0f}s idle"
)


def main() -> None:
parser = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
parser.add_argument("--poll-ms", type=float, default=1.0, help="Old poll interval in ms (default: 1.0).")
parser.add_argument("--messages", type=int, default=3000, help="Messages per latency run (default: 3000).")
parser.add_argument("--runs", type=int, default=3, help="Latency runs to aggregate (default: 3).")
parser.add_argument(
"--gap-max-ms", type=float, default=2.0, help="Max random inter-message gap in ms (default: 2.0)."
)
parser.add_argument(
"--idle-secs", type=float, default=2.0, help="Idle measurement window in seconds (default: 2.0)."
)
parser.add_argument("--seed", type=int, default=0, help="RNG seed for reproducibility (default: 0).")
args = parser.parse_args()
random.seed(args.seed)
asyncio.run(_main_async(args))


if __name__ == "__main__":
main()
2 changes: 1 addition & 1 deletion docs/configuration/environment_variables.md
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ depends on the installed kernels and model path.
| `SPEAKER_MAX_UPLOADED` | Integer; default `1000` | Speech server; read when speaker storage initializes | Environment-only setting. A non-integer logs a warning and uses `1000`; range is not otherwise validated. | Stable |
| `VLLM_OMNI_ABORT_TIMEOUT` | Float seconds; default `2` | Engine abort wait for Videos DELETE and `generate()` cancel/error cleanup; read when `async_omni` imports | Environment-only setting. A non-float raises `ValueError` during import. | Experimental |
| `VLLM_OMNI_ASYNC_OUTPUT_TIMEOUT` | Float seconds; default `600` | Diffusion engine async-output wait in `step_streaming`; resolved per call on the request path, not at import | Environment-only setting. A non-float or `<=0` value warns once and uses the default. | Experimental |
| `VLLM_OMNI_EVENT_DRIVEN_ORCH` | `1`, `true`, `yes` or `on` enables; defaults to on for Qwen3-TTS and off for other pipelines | Pipeline default computed from `pipeline_config.model_type` at engine initialization; env override resolved at `Orchestrator` construction and separately when the serving-side final-output drain starts | An explicit env value wins; otherwise both consumers use the engine's pipeline default. Values are stripped and case-normalized; any unrecognized value selects the legacy poll loop. Set before server startup. | Experimental |
| `VLLM_OMNI_EVENT_DRIVEN_ORCH` | `1`, `true`, `yes` or `on` enables; defaults to on for Qwen3-TTS and off for other pipelines | Pipeline default computed from `pipeline_config.model_type` at engine initialization; env override resolved at `Orchestrator` construction | An explicit env value wins; otherwise the engine's pipeline default is used. Values are stripped and case-normalized; any unrecognized value selects the legacy poll loop. Only the orchestration loop is affected; the serving-side final-output drain is always event-driven. Set before server startup. | Experimental |
| `VLLM_OMNI_INPUT_WAIT_TIMEOUT_S` | Float seconds; default `600`; `<=0` disables | Full-payload input coordinator, not async-chunk transfer; read when the scheduler module imports in each worker | Environment-only setting. A non-float logs a warning and uses `600`. | Stable operational control |
| `VLLM_OMNI_ORCH_MONITOR_PATH` | Filesystem path; default `<current-working-directory>/vllm_omni_orch_monitor_<timestamp>.json` | Orchestrator monitor enabled by `--enable-orch-monitor`; read when the monitor is created | Environment-only path override. Parent directories are created; write errors are logged. | Diagnostic |
| `VLLM_OMNI_SPEAKER_REGISTRATION_POLICY` | `overwrite` or `immutable`; default `overwrite` | Speech server; read when speaker storage initializes | Environment-only setting. `immutable` rejects re-registering an existing uploaded voice name until it is deleted; any other value raises `ValueError` at startup. | Experimental |
Expand Down
9 changes: 5 additions & 4 deletions docs/serving/speech_api.md
Original file line number Diff line number Diff line change
Expand Up @@ -973,15 +973,16 @@ Use `/v1/audio/voices` to list available voices for the loaded model.
Multi-stage omni deployments route stage outputs through a single orchestrator
loop. The legacy loop polls every stage replica on a 1 ms cadence. The
event-driven mode replaces the poll with one reader task per live stage
replica awaiting its client directly, and switches the serving-side
final-output drain to a condition-variable wakeup at the same time. Qwen3-TTS
uses this mode by default; other pipelines keep the legacy poll unless enabled.
replica awaiting its client directly. Qwen3-TTS uses this mode by default;
other pipelines keep the legacy poll unless enabled. The serving-side
final-output drain (orchestrator to API server) is always event-driven and is
not affected by this flag.

**Configuration (environment variables):**

| Variable | Default | Description |
| --- | --- | --- |
| `VLLM_OMNI_EVENT_DRIVEN_ORCH` | On for Qwen3-TTS; off for other pipelines | Switches the orchestration loop and the final-output drain from the legacy 1 ms poll to event-driven wakeups. An explicit value wins; otherwise the pipeline default computed at engine initialization is used. The override is resolved at orchestrator construction and separately when the final-output drain starts. `1`, `true`, `yes`, or `on` enables it, ignoring case and surrounding whitespace; other values select the legacy poll loop. |
| `VLLM_OMNI_EVENT_DRIVEN_ORCH` | On for Qwen3-TTS; off for other pipelines | Switches the orchestration loop from the legacy 1 ms poll to event-driven wakeups. An explicit value wins; otherwise the pipeline default computed at engine initialization is used. The override is resolved at orchestrator construction. `1`, `true`, `yes`, or `on` enables it, ignoring case and surrounding whitespace; other values select the legacy poll loop. |

Set it on the process that runs the orchestrator (stage 0 of an omni deployment)
before starting the server when overriding the pipeline default. For example,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -398,11 +398,16 @@ async def add_request_async(
)
)

async def try_get_output_async(self) -> Any | None:
try:
return self._fixture.output_sync_q.get_nowait()
except queue.Empty:
return None
async def try_get_output_async(self) -> Any:
# Mirror the real AsyncOmniEngine: park until a message is available
# instead of returning None on empty. The AsyncOmni output loop no
# longer polls, so a synchronous None return would hot-spin the loop
# without yielding and starve the event loop.
while True:
try:
return self._fixture.output_sync_q.get_nowait()
except queue.Empty:
await asyncio.sleep(0.001)

def get_stage_metadata(self, stage_id: int) -> StageRuntimeInfo:
return self.stage_metadata[stage_id]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -569,7 +569,7 @@ async def submit_interaction(self, *args: Any, **kwargs: Any) -> None:
output_queue: asyncio.Queue[ErrorMessage] = asyncio.Queue()
orchestrator = object.__new__(Orchestrator)
orchestrator.stage_pools = [_RejectingStagePool()] # pyright: ignore[reportAttributeAccessIssue]
orchestrator.output_async_queue = output_queue # pyright: ignore[reportAttributeAccessIssue]
orchestrator.output_sync_queue = output_queue # pyright: ignore[reportAttributeAccessIssue]
orchestrator.request_states = {"req-1": OrchestratorRequestState(request_id="req-1")}

await orchestrator._handle_interaction(
Expand Down Expand Up @@ -604,7 +604,7 @@ async def test_runner_prompt_update_failure_surfaces_non_fatal_error(self, pipel
orchestrator.stage_pools = [ # pyright: ignore[reportAttributeAccessIssue]
StagePool(0, cast(StagePoolClient, inline_client))
]
orchestrator.output_async_queue = output_queue # pyright: ignore[reportAttributeAccessIssue]
orchestrator.output_sync_queue = output_queue # pyright: ignore[reportAttributeAccessIssue]
orchestrator.request_states = {"req-1": OrchestratorRequestState(request_id="req-1")}

try:
Expand Down Expand Up @@ -814,11 +814,16 @@ async def submit_interaction_async(
)
)

async def try_get_output_async(self) -> Any | None:
try:
return self._fixture.output_sync_q.get_nowait()
except queue.Empty:
return None
async def try_get_output_async(self) -> Any:
# Mirror the real engine: park until a message is available
# instead of returning None on empty. The AsyncOmni output loop no
# longer polls, so a synchronous None return would hot-spin the
# loop without yielding and starve the event loop.
while True:
try:
return self._fixture.output_sync_q.get_nowait()
except queue.Empty:
await asyncio.sleep(0.001)

def get_stage_metadata(self, stage_id: int) -> StageRuntimeInfo:
return self.stage_metadata[stage_id]
Expand Down
Loading
Loading