Skip to content
Merged
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
1 change: 1 addition & 0 deletions docs/configuration/environment_variables.md
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,7 @@ depends on the installed kernels and model path.
| `SPEAKER_SAMPLES_DIR` | Filesystem path; default `~/.cache/vllm-omni/speakers` | Speech server; read when speaker storage initializes | Environment-only setting. The directory is created; filesystem errors propagate. | Stable |
| `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_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; default `0` (off) | Orchestration loop and the serving-side final-output drain; read once when the `Orchestrator` is constructed | Environment-only setting. Values are stripped and case-normalized; any unrecognized value leaves the legacy poll loop selected. | 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_VIDEO_SYNC_TIMEOUT` | Float seconds; default `600` | Synchronous Videos API; read when the API server module imports | Environment-only setting. A non-float raises `ValueError` during import. | Experimental |
Expand Down
9 changes: 9 additions & 0 deletions docs/contributing/profiling.md
Original file line number Diff line number Diff line change
Expand Up @@ -322,6 +322,15 @@ Each 1-second window records:

On shutdown the server also logs a short summary (`loop_active_pct`, per-replica queue averages/maxima).

`loop_idle` and `loop_active` count loop iterations, so their absolute scale is
tied to which orchestration loop is running. Under the default poll loop an idle
orchestrator records roughly one iteration per millisecond. Under the
event-driven loop (`VLLM_OMNI_EVENT_DRIVEN_ORCH=1`, see
[Speech API](../serving/speech_api.md#orchestration-loop-experimental)) an idle
orchestrator wakes only once per 0.5 s reconcile timeout, while a busy one still
records one iteration per routed output. Window counts and the `loop_active_pct`
summary are therefore not comparable across the two modes.

### Relationship to other diagnostics

This monitor is intentionally separate from the existing profiling tools:
Expand Down
6 changes: 6 additions & 0 deletions docs/design/module/engine_orchestration.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,12 @@ the boundary affected by the in-flight stage client/process refactor in
[#5441](https://github.com/vllm-project/vllm-omni/pull/5441). Names and
responsibilities proposed only by that PR are not current contracts.

The orchestration loop also has an opt-in event-driven mode
(`VLLM_OMNI_EVENT_DRIVEN_ORCH=1`, default off) proposed in
[#5221](https://github.com/vllm-project/vllm-omni/pull/5221). It changes poll
cadence only: the routing, ordering, and terminal-state contracts below hold
identically on both loops.

## Ownership boundary

This document owns `AsyncOmniEngine`, `Orchestrator`, request-state creation,
Expand Down
45 changes: 45 additions & 0 deletions docs/serving/speech_api.md
Original file line number Diff line number Diff line change
Expand Up @@ -811,6 +811,51 @@ If you encounter OOM errors:

Use `/v1/audio/voices` to list available voices for the loaded model.

## Orchestration Loop (experimental)

Multi-stage omni deployments route stage outputs through a single orchestrator
loop. By default that loop polls every stage replica on a 1 ms cadence. An
opt-in 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.

**Configuration (environment variables):**

| Variable | Default | Description |
|----------|---------|-------------|
| `VLLM_OMNI_EVENT_DRIVEN_ORCH` | `0` (off) | Switches the orchestration loop and the final-output drain from the legacy 1 ms poll to event-driven wakeups. Enabled by `1`, `true`, `yes`, or `on`, matched case-insensitively after surrounding whitespace is stripped; any other value leaves it off. |

Set it on the process that runs the orchestrator (stage 0 of an omni
deployment) before starting the server:

```bash
export VLLM_OMNI_EVENT_DRIVEN_ORCH=1
vllm serve Qwen/Qwen3-TTS-12Hz-1.7B-Base \
--omni \
--port 8091
```

The server logs the selected loop mode and its reader/poller counts once at
startup, so you can confirm which loop is live.

Routing, output ordering, and terminal-state behavior are identical on both
loops; only the poll cadence changes. Leaving the variable unset keeps the
legacy poll loop, which is the supported default.

**Known limitations:**

- The measured serving A/B (idle CPU 2.43% to 0.07%; TTFP p99 -32% at
concurrency 8) predates the rebuild on the per-replica fault-isolation work
in [#4285](https://github.com/vllm-project/vllm-omni/pull/4285). That work
changed dead-replica handling and reader/poller lifecycle rather than the
steady-state output path, and the parity suite covers it, but the serving
A/B has not been re-run on the current head.
- The diffusion-poller branch is covered by unit tests only. Deployments whose
stages all run as standard engine cores never exercise it, including GLM-TTS,
which deploys its DiT without `stage_type: diffusion`.
- Concurrency 1 and 32 measured at parity with the legacy loop. At 32 the
latency is admission-bound, which this mode does not address.

## Development

Enable debug logging:
Expand Down
2 changes: 1 addition & 1 deletion tests/config/test_environment_variables.py
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,7 @@ def test_inventory_matches_reviewed_snapshot_counts():
"""Make an inventory expansion an explicit review decision."""
category_counts = Counter(item.category for item in ENVIRONMENT_VARIABLE_INVENTORY.values())
assert category_counts == {
EnvironmentVariableCategory.PUBLIC_OMNI: 23,
EnvironmentVariableCategory.PUBLIC_OMNI: 24,
EnvironmentVariableCategory.INHERITED_VLLM: 20,
EnvironmentVariableCategory.PLATFORM_EXTERNAL: 27,
EnvironmentVariableCategory.MODEL_SPECIFIC: 56,
Expand Down
27 changes: 22 additions & 5 deletions tests/diffusion/test_diffusion_engine_metrics.py
Original file line number Diff line number Diff line change
Expand Up @@ -172,12 +172,29 @@ def test_abort_keeps_output_consumers_alive_for_terminal_snapshot(self) -> None:
assert ".cancel(" not in abort_branch

def test_orchestrator_consumes_metrics_only_output_without_routing(self) -> None:
"""A metrics-only sentinel contributes its queue depth and is not routed.

Absorption lives in ``_absorb_diffusion_metrics`` rather than inline in a
loop, so both orchestration loops share one implementation. The ordering
is asserted where it now lives: the snapshot inside the helper, and the
absorb-before-route guard inside every loop that polls diffusion output.
"""
source = _read_source(_ORCHESTRATOR_PATH)
loop_source = _get_function_source(source, "Orchestrator", "_orchestration_loop")
snapshot_pos = loop_source.index("_update_stage_replica_waiting(")
sentinel_pos = loop_source.index("diffusion_output.request_id == DIFFUSION_METRICS_ONLY_REQUEST_ID")
route_pos = loop_source.index("pool.record_output_timestamps([diffusion_output])")
assert snapshot_pos < sentinel_pos < route_pos

# Snapshot before the verdict, so a sentinel still reports its waiting
# depth on the way out instead of being dropped unaccounted.
absorb_source = _get_function_source(source, "Orchestrator", "_absorb_diffusion_metrics")
snapshot_pos = absorb_source.index("_update_stage_replica_waiting(")
sentinel_pos = absorb_source.index("diffusion_output.request_id == DIFFUSION_METRICS_ONLY_REQUEST_ID")
assert snapshot_pos < sentinel_pos

# Absorb before routing, or a sentinel reaches the downstream consumer
# as if it were a real request output.
for loop_name in ("_orchestration_loop", "_orchestration_loop_event_driven"):
loop_source = _get_function_source(source, "Orchestrator", loop_name)
absorb_pos = loop_source.index("self._absorb_diffusion_metrics(")
route_pos = loop_source.index("record_output_timestamps(")
assert absorb_pos < route_pos, loop_name


class TestVaeDecodeEmit:
Expand Down
247 changes: 247 additions & 0 deletions tests/engine/test_orchestrator_event_driven.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,247 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM project
"""Event-driven orchestration loop (``VLLM_OMNI_EVENT_DRIVEN_ORCH=1``) tests.

Parity suite: re-runs the legacy orchestration scenarios from
``test_orchestrator.py`` / ``test_orchestrator_error_handling.py`` with the
event-driven loop selected, so both loops are held to the same behavior. Plus
event-driven-specific coverage: reader reconcile on client swap, the blocking
final-output drain, and flag parsing.
"""

from __future__ import annotations

import asyncio
from types import SimpleNamespace

import janus
import pytest

from vllm_omni.engine.async_omni_engine import AsyncOmniEngine
from vllm_omni.engine.messages import ShutdownRequestMessage
from vllm_omni.engine.orchestrator import _event_driven_orch_enabled

from . import test_orchestrator as legacy
from . import test_orchestrator_error_handling as legacy_errors
from .test_orchestrator import (
FakeOutputProcessor,
FakeStageClient,
OrchestratorFixture,
_build_harness,
_build_request_output,
_engine_core_outputs,
_enqueue_add_request,
_get_output_message,
_sampling_params,
_shutdown_orchestrator,
_wait_for,
)

pytestmark = [pytest.mark.core_model, pytest.mark.cpu]


@pytest.fixture
def orchestrator_factory(monkeypatch):
"""Flag-setting clone of the legacy harness fixture.

Sets ``VLLM_OMNI_EVENT_DRIVEN_ORCH=1`` before any Orchestrator is
constructed and asserts the flag actually took effect, so the parity tests
cannot silently exercise the legacy poll loop.
"""
monkeypatch.setenv("VLLM_OMNI_EVENT_DRIVEN_ORCH", "1")
fixtures: list[OrchestratorFixture] = []

def _factory(*args, **kwargs) -> OrchestratorFixture:
fixture = _build_harness(*args, **kwargs)
assert fixture.orchestrator._event_driven_orch is True
fixtures.append(fixture)
return fixture

yield _factory

for fixture in fixtures:
if fixture.thread.is_alive():
fixture.request_sync_q.put_nowait(ShutdownRequestMessage())
fixture.thread.join(timeout=5)
for q in fixture.queues:
q.close()


# ---------------------------------------------------------------------------
# Parity: the legacy scenario matrix, re-run through the event-driven loop
# ---------------------------------------------------------------------------

_PARITY_TESTS = [
legacy.test_run_two_stage_llm,
legacy.test_run_single_stage_diffusion,
legacy.test_run_single_stage_diffusion_streaming_forwards_intermediate_chunks,
legacy.test_run_llm_to_diffusion,
legacy.test_run_async_chunk,
legacy.test_run_shutdown,
legacy.test_run_abort,
legacy.test_multi_replica_round_robin_distribution,
legacy.test_multi_replica_abort_broadcasts_to_all_replicas,
legacy.test_multi_replica_shutdown_all_replicas,
legacy.test_multi_replica_cfg_companion_inherits_parent_affinity,
# Stats plumbing. The scheduler-stats cases matter most here: a batch with
# no request outputs still carries SchedulerStats on throttled ticks, and a
# reader that drops every output-less batch silently stops reporting
# KV/queue gauges under the event-driven loop.
legacy.test_orchestrator_records_iteration_stats_without_scheduler_stats,
legacy.test_orchestrator_records_scheduler_stats_without_outputs,
legacy.test_orchestrator_does_not_build_iteration_stats_for_finished_only_batch,
legacy.test_orchestrator_does_not_build_iteration_stats_without_stat_logger,
# Per-replica fault isolation (#4285): a dead replica must be evicted and
# the server kept up, on both the LLM reader path and the diffusion poller.
legacy_errors.test_engine_dead_error_evicts_replica_and_keeps_running,
legacy_errors.test_engine_dead_error_fails_only_dead_replica_requests,
legacy_errors.test_forward_to_dead_downstream_stage_fails_request_not_server,
legacy_errors.test_add_request_to_dead_stage_fails_request_not_server,
legacy_errors.test_diffusion_replica_death_on_poll_keeps_server,
legacy_errors.test_diffusion_error_output_routed_as_finished,
legacy_errors.test_diffusion_client_error_output_propagates_status_code,
]


@pytest.mark.asyncio
@pytest.mark.parametrize("legacy_test", _PARITY_TESTS, ids=lambda f: f.__name__)
async def test_event_driven_parity(legacy_test, orchestrator_factory) -> None:
await legacy_test(orchestrator_factory)


# ---------------------------------------------------------------------------
# Event-driven-specific behavior
# ---------------------------------------------------------------------------


def test_flag_parsing(monkeypatch) -> None:
monkeypatch.delenv("VLLM_OMNI_EVENT_DRIVEN_ORCH", raising=False)
assert _event_driven_orch_enabled() is False
for value in ("1", "true", "True", "YES", "on"):
monkeypatch.setenv("VLLM_OMNI_EVENT_DRIVEN_ORCH", value)
assert _event_driven_orch_enabled() is True
for value in ("0", "false", "off", ""):
monkeypatch.setenv("VLLM_OMNI_EVENT_DRIVEN_ORCH", value)
assert _event_driven_orch_enabled() is False


def test_default_is_legacy_loop(monkeypatch) -> None:
"""Without the env flag, the harness runs the legacy poll loop."""
monkeypatch.delenv("VLLM_OMNI_EVENT_DRIVEN_ORCH", raising=False)
fixture = _build_harness([FakeStageClient(stage_type="llm", final_output=True)])
try:
assert fixture.orchestrator._event_driven_orch is False
finally:
fixture.request_sync_q.put_nowait(ShutdownRequestMessage())
fixture.thread.join(timeout=5)
for q in fixture.queues:
q.close()


@pytest.mark.asyncio
async def test_reader_reconcile_picks_up_swapped_client(orchestrator_factory) -> None:
"""Outputs from a replica whose client object was replaced still flow.

The event-driven loop binds one reader task per client object; the
periodic reconcile must respawn the reader when ``pool.clients[replica]``
is swapped (replica replacement), otherwise the new client's outputs
would never be drained.
"""
stage0 = FakeStageClient(stage_type="llm", final_output=True)
processor = FakeOutputProcessor(request_outputs=[_build_request_output("req-swap", token_ids=[3], finished=True)])
orchestrator_fixture = orchestrator_factory([stage0], output_processors=[processor])
request = SimpleNamespace(request_id="req-swap", prompt_token_ids=[1, 2])

try:
await _enqueue_add_request(
orchestrator_fixture,
request_id="req-swap",
prompt=request,
original_prompt={"prompt": "swap"},
sampling_params_list=[_sampling_params()],
final_stage_id=0,
)
await _wait_for(lambda: len(stage0.add_request_calls) == 1)

# Swap in a fresh client for the same replica slot; keep the pool's
# other wiring intact. The reconcile tick (0.5 s) must respawn the
# reader bound to the new client object.
pool = orchestrator_fixture.orchestrator.stage_pools[0]
replacement = FakeStageClient(stage_type="llm", final_output=True)
replacement.stage_id = stage0.stage_id
replacement.replica_id = stage0.replica_id
pool.clients[0] = replacement

replacement.push_engine_core_outputs(_engine_core_outputs("swapped-raw", 1.0))

output_msg = await _get_output_message(orchestrator_fixture, timeout=5.0)
assert output_msg.request_id == "req-swap"
assert output_msg.finished is True
finally:
await _shutdown_orchestrator(orchestrator_fixture)


# ---------------------------------------------------------------------------
# Blocking final-output drain (AsyncOmniEngine.get_output_blocking_async)
# ---------------------------------------------------------------------------


def _drain_engine(alive: bool = True) -> AsyncOmniEngine:
engine = object.__new__(AsyncOmniEngine)
engine.output_queue = janus.Queue()
engine.orchestrator_thread = SimpleNamespace(is_alive=lambda: alive)
return engine


def _drain_cleanup(engine: AsyncOmniEngine) -> None:
if engine._output_drain_executor is not None:
engine._output_drain_executor.shutdown(wait=False)
engine._output_drain_executor = None
engine.output_queue.close()


@pytest.mark.asyncio
async def test_blocking_drain_returns_queued_message() -> None:
engine = _drain_engine()
try:
engine.output_queue.sync_q.put_nowait("msg-1")
assert await engine.get_output_blocking_async(timeout=1.0) == "msg-1"
finally:
_drain_cleanup(engine)


@pytest.mark.asyncio
async def test_blocking_drain_wakes_on_late_message() -> None:
"""A message put after the wait starts wakes the drain, no polling."""
engine = _drain_engine()
try:

async def _delayed_put() -> None:
await asyncio.sleep(0.05)
engine.output_queue.sync_q.put_nowait("late-msg")

put_task = asyncio.create_task(_delayed_put())
msg = await engine.get_output_blocking_async(timeout=5.0)
await put_task
assert msg == "late-msg"
finally:
_drain_cleanup(engine)


@pytest.mark.asyncio
async def test_blocking_drain_timeout_returns_none_when_alive() -> None:
engine = _drain_engine(alive=True)
try:
assert await engine.get_output_blocking_async(timeout=0.05) is None
finally:
_drain_cleanup(engine)


@pytest.mark.asyncio
async def test_blocking_drain_raises_when_orchestrator_dead() -> None:
engine = _drain_engine(alive=False)
try:
with pytest.raises(RuntimeError, match="Orchestrator died"):
await engine.get_output_blocking_async(timeout=0.05)
finally:
_drain_cleanup(engine)
1 change: 1 addition & 0 deletions vllm_omni/config/environment_variable_inventory.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ def is_public_omni(self) -> bool:
"SPEAKER_MAX_UPLOADED",
"SPEAKER_SAMPLES_DIR",
"VLLM_OMNI_ASYNC_OUTPUT_TIMEOUT",
"VLLM_OMNI_EVENT_DRIVEN_ORCH",
"VLLM_OMNI_INPUT_WAIT_TIMEOUT_S",
"VLLM_OMNI_ORCH_MONITOR_PATH",
"VLLM_OMNI_SERVER_STORAGE__FILE_CONCURRENCY",
Expand Down
Loading
Loading