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
2 changes: 0 additions & 2 deletions docs/design/fullduplex.md
Original file line number Diff line number Diff line change
Expand Up @@ -179,7 +179,6 @@ vllm_omni/
│ │ ├── realtime_input.py RealtimeEnvelope (query rules, first message), parse_resume_request
│ │ ├── session_attachment.py DuplexSessionAttachmentRegistry (resume tokens, replay journal)
│ │ ├── audio_encoding.py encode_audio, injected into DuplexOmniEngine for the plugin's data plane
│ │ ├── chat_completions.py DuplexChatCompletionsAdapter (/v1/chat/completions on a session per request)
│ │ └── websocket.py websocket send/close/receive helpers
│ └── openai/api_server.py builds DuplexOmni for duplex models; session-backed app state
├── protocol/ SHARED WIRE CODEC (no engine / no model / no transport)
Expand Down Expand Up @@ -430,7 +429,6 @@ follow-up PRs port them (RFC vllm-omni#7181, PR 2/3).
`tests/entrypoints/openai/test_duplex_session_attachment.py`,
`tests/entrypoints/openai_api/test_duplex_api_server.py` (the duplex
server: pipeline probe, app state, routes, warmup gate),
`tests/entrypoints/duplex/test_chat_completions_adapter.py`,
`tests/clients/**`, `tests/engine/test_duplex_import_boundary.py`,
`tests/model_executor/models/minicpmo_4_5/duplex/**`,
`tests/worker/test_native_duplex_hooks.py`.
Expand Down
12 changes: 6 additions & 6 deletions docs/serving/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -155,12 +155,12 @@ declares a `duplex_plugin` and the deploy configuration sets
`session_mode: duplex`); a stock Realtime client needs no vendor query
parameter, and `duplex=0` is the explicit opt-out. Such a server serves the
websocket route plus `POST /v1/chat/completions`, and no other turn-based
HTTP route. On a duplex server that chat route is not the turn-based path:
each request runs on a short-lived duplex session, so it holds one of
`duplex_session.max_sessions` for its lifetime and answers at the model's
real-time pace -- see [Full Duplex](full_duplex_api.md). Clients should
also verify the duplex capability payload because the query-parameter form
falls back to the ordinary realtime handler when duplex is unavailable.
HTTP route. That chat route is the ordinary chat service running on the
duplex engine, served alongside the live sessions when the model's plugin
declares `supports_chat_completions` -- see [Full Duplex](full_duplex_api.md).
Clients should also verify the duplex capability payload because the
query-parameter form falls back to the ordinary realtime handler when duplex
is unavailable.

## Related Endpoints

Expand Down
31 changes: 12 additions & 19 deletions docs/serving/realtime_duplex_api.md
Original file line number Diff line number Diff line change
Expand Up @@ -482,11 +482,11 @@ Semantic divergences hidden behind shared names:
the Tier 3 surface on top of the same Tier 1 events, so both kinds of
client can share one deployment and one event log.

### Turn-based Chat audio metadata
### Chat audio metadata

For turn-based models served through Chat fallback, non-empty audio choices
produced by the Omni Chat audio encoder include `audio_metadata`, in both
streaming and complete Chat responses. Text-only choices do not require it.
On `/v1/chat/completions`, non-empty audio choices produced by the Omni Chat
audio encoder include `audio_metadata`, in both streaming and complete Chat
responses. Text-only choices do not require it.

| Field | Meaning |
| --- | --- |
Expand All @@ -501,17 +501,10 @@ The encoder uses a model-reported sample rate when available; Qwen3-Omni
currently omits it and uses the existing 24 kHz default. This metadata does
not independently determine a model's sample rate.

Chat fallback requests the session's output format explicitly. Its internal
`response.output_audio.delta` carries `format`, `sample_rate_hz`, `channels`,
and `audio_duration_ms`. The Realtime endpoint exposes output `format` and
`sample_rate_hz` on `response.audio.delta`, with duration in
`metadata.audio_duration_ms`; it does not currently forward `channels`.
The duration is cumulative for the response:
integer frames are summed before converting to milliseconds and rounding down.
Non-empty audio without valid metadata, or a sample-rate change within one
response, fails that fallback response with `response_error`.

In this Chat fallback path, `response.done.response.metadata.playback` reports:
On the duplex route the audio delta itself carries the output `format` and
`sample_rate_hz`, with the response's cumulative duration in
`metadata.audio_duration_ms`; `channels` is not forwarded.
`response.done.response.metadata.playback` reports:

- `generated_ms`: cumulative source-waveform duration.
- `sent_ms`: cumulative audio accepted into the server output path, including
Expand All @@ -520,10 +513,10 @@ In this Chat fallback path, `response.done.response.metadata.playback` reports:
- `played_ms`: cumulative playback reported by the client.
- `committed_ms`: playback position used for history reconciliation.

The native model path still records audio before sending; it does not yet
share Chat fallback's output-acceptance timing. Only client playback reports
establish what was played. The existing duration-proportional text truncation
remains an approximation, not word-level audio alignment.
The engine records audio as sent when it emits the delta, before the socket
delivers it; only client playback reports establish what was played. The
existing duration-proportional text truncation remains an approximation, not
word-level audio alignment.

### Event catalogue

Expand Down
17 changes: 16 additions & 1 deletion tests/e2e/online_serving/test_minicpmo_4_5_duplex.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
import base64
import json
from pathlib import Path
from typing import TypedDict

import pytest
import websockets
Expand Down Expand Up @@ -50,6 +51,13 @@ def _assert_request_metrics(metrics: object, *, expected_count: int) -> None:
assert isinstance(request["response_id"], str)
assert request["ttft_ms"] is not None and request["ttft_ms"] >= 0
assert request["ttfp_ms"] >= 0
# TTFT/TTFP are anchored on the server's model-turn request start,
# which the engine announces on response.created and every delta;
# the client only falls back to its own receive clock without it.
assert request["source"] == "server_request_start_and_client_receive", request
origin = request["measurement_origin"]
assert origin["ttft"].startswith("native model-turn request execution start"), origin
assert origin["ttfp"].startswith("native model-turn request execution start"), origin
assert request["rtf"] is not None and request["rtf"] >= 0
assert request["audio_generation_ms"] >= 0
assert request["audio_duration_ms"] > 0
Expand Down Expand Up @@ -174,6 +182,13 @@ async def receive_outcome() -> dict[str, object]:
return await asyncio.wait_for(receive_outcome(), timeout=timeout_s)


class _SeededTurnResult(TypedDict):
audio_bytes: int
transcript: str
output_text: str
event_types: list[str]


async def _run_seeded_text_to_audio(
*,
url: str,
Expand All @@ -183,7 +198,7 @@ async def _run_seeded_text_to_audio(
modalities: tuple[str, ...] = ("audio", "text"),
silence_seconds: float = 12.0,
timeout_s: float = 180.0,
) -> dict[str, object]:
) -> _SeededTurnResult:
"""Speak a seeded text: the duplex route's text-to-speech shape.

A model-native session takes its text once, in the session context
Expand Down
75 changes: 75 additions & 0 deletions tests/engine/duplex/test_engine_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -688,3 +688,78 @@ def test_minicpmo_native_capabilities_do_not_overclaim_single_session_deployment

assert caps["supports_multi_session"] is False
assert caps["supports_multi_session_same_replica"] is False


# ---- per-response request timing (server-side TTFT / TTFP) ----


def test_response_timing_binds_latest_request_start_for_model_turn():
session = _session(config=DuplexSessionConfig(model="test-model", extra_body={"auto_response": True}))
session.mark_model_turn_request_started(0, 10.0)
session.mark_model_turn_request_started(0, 11.0)
session.begin_response(turn_id=0)

first = session.mark_response_first_outputs(observed_at_s=11.2, has_text=True, has_audio=False)
assert first["ttft_ms"] == pytest.approx(200.0)
assert "ttfp_ms" not in first
assert first["source"] == "server_monotonic_request_start"
assert set(first["measurement_origin"]) == {"ttft", "ttfp"}

# Once bound, the active response keeps its anchor; a later start for the
# same turn does not move it.
session.mark_model_turn_request_started(0, 12.0)
second = session.mark_response_first_outputs(observed_at_s=12.3, has_text=False, has_audio=True)
assert second["ttfp_ms"] == pytest.approx(1300.0)
assert second["ttft_ms"] == pytest.approx(200.0)


def test_response_timing_is_empty_without_a_request_start_and_is_scoped_to_one_response():
session = _session()
session.begin_response(turn_id=0)
assert session.mark_response_first_outputs(observed_at_s=1.0, has_text=True, has_audio=True) == {}

# A start recorded while its response is already active binds to it.
session.mark_model_turn_request_started(0, 2.0)
metrics = session.mark_response_first_outputs(observed_at_s=2.5, has_text=True, has_audio=False)
assert metrics["ttft_ms"] == pytest.approx(500.0)
session.end_response()

# The next response measures from its own turn's start; a completed turn's
# start is dropped rather than inherited.
session.mark_model_turn_request_started(0, 3.0)
session.mark_model_turn_request_started(1, 4.0)
session.complete_model_turn(0)
session.begin_response(turn_id=1)
metrics = session.mark_response_first_outputs(observed_at_s=4.25, has_text=True, has_audio=True)
assert metrics["ttft_ms"] == pytest.approx(250.0)
assert metrics["ttfp_ms"] == pytest.approx(250.0)
session.end_response()

# A barge-in aborts the pending starts along with the work they belong to.
session.mark_model_turn_request_started(2, 5.0)
session.barge_in()
session.begin_response(turn_id=2)
assert session.mark_response_first_outputs(observed_at_s=5.5, has_text=True, has_audio=True) == {}


def test_response_tpot_fallback_ignores_single_token_segment():
session = _session(config=DuplexSessionConfig(model="test-model", extra_body={"auto_response": True}))
session.begin_response(turn_id=0)

first_metrics = session.accumulate_response_stage_metrics({"0": {"num_tokens_out": 1, "vllm_tpot_ms": 900.0}})
assert "vllm_tpot_ms" not in first_metrics["0"]

second_metrics = session.accumulate_response_stage_metrics({"0": {"num_tokens_out": 2, "vllm_tpot_ms": 15.0}})
assert second_metrics["0"]["num_tokens_out"] == 3
assert second_metrics["0"]["vllm_tpot_ms"] == 15.0


def test_response_tpot_keeps_token_weighted_value_when_itls_exist():
session = _session(config=DuplexSessionConfig(model="test-model", extra_body={"auto_response": True}))
session.begin_response(turn_id=0)
metrics = session.accumulate_response_stage_metrics(
{"0": {"num_tokens_out": 4, "vllm_tpot_ms": 10.0, "vllm_itls_ms": [30.0]}}
)
assert metrics["0"]["vllm_itls_ms"] == [30.0]
assert metrics["0"]["vllm_itl_ms"] == 30.0
assert metrics["0"]["vllm_tpot_ms"] == 10.0
54 changes: 51 additions & 3 deletions tests/engine/duplex/test_session_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,7 @@ async def run(self, command: DuplexCommand) -> list[DuplexEvent]:

def deliver(
self,
output: object,
output: SimpleNamespace,
*,
stage_id: int = 1,
segment_finished: bool = False,
Expand All @@ -182,7 +182,7 @@ def deliver(
)
return self.runner.on_stage_output(stage_id, output, metrics, request_id=output.request_id, context=context)

async def deliver_and_settle(self, output: object, **kwargs: Any) -> list[DuplexEvent]:
async def deliver_and_settle(self, output: SimpleNamespace, **kwargs: Any) -> list[DuplexEvent]:
self.deliver(output, **kwargs)
return await self.settle()

Expand Down Expand Up @@ -1033,7 +1033,7 @@ async def test_conversation_items_can_be_injected_and_deleted() -> None:
await close_harness(h)


def _stage_metrics_of(event: object) -> dict[str, dict[str, object]]:
def _stage_metrics_of(event: DuplexEvent) -> dict[str, dict[str, object]]:
"""Per-stage engine metrics as the client reads them off one wire event."""
payload = event.to_realtime()
metadata = payload.get("metadata")
Expand Down Expand Up @@ -1193,3 +1193,51 @@ async def test_server_vad_speech_stopped_still_commits_a_turn_mode_session() ->
assert len(_final_submissions(h)) == 1, "the detector's stop commits the turn and starts the response"
finally:
await close_harness(h)


# --------------------------------------------------------------------------- #
# Per-response request timing #
# --------------------------------------------------------------------------- #


@pytest.mark.asyncio
async def test_response_request_metrics_are_measured_from_the_model_turn_request_start(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Per-response TTFT/TTFP start when the runner submits the native request, not at response.created.

The client library and the Omni-DuplexEval benchmark prefer these
server-side numbers to their own receive clock, so the response announces
them on ``response.created`` and every delta repeats them.
"""
from vllm_omni.engine.duplex.session import model_channel as model_channel_module

now = {"monotonic": 10.0}
monkeypatch.setattr(model_channel_module, "time", SimpleNamespace(monotonic=lambda: now["monotonic"]))
h = await open_harness()
try:
await h.run(append_audio()) # the Stage0 request starts executing at t=10.0
request_id = h.stage0_request_id()
now["monotonic"] = 10.25
events = await h.deliver_and_settle(tts_output(request_id, samples=24000, text="hi"))
expected = {
"source": "server_monotonic_request_start",
"measurement_origin": {
"ttft": "native model-turn request execution start to first non-empty text output",
"ttfp": "native model-turn request execution start to first audio output",
},
"ttft_ms": pytest.approx(250.0),
"ttfp_ms": pytest.approx(250.0),
}
created = find(events, "response.created").to_realtime()
assert created["response"]["metadata"]["duplex_event"]["response_request_metrics"] == expected
delta = find(events, "response.output_audio.delta").to_realtime()
assert delta["metadata"]["vllm_omni"]["response_request_metrics"] == expected

# The first outputs fix the numbers; later units of the same response repeat them.
now["monotonic"] = 10.9
events = await h.deliver_and_settle(tts_output(request_id, samples=48000, text="hi there"))
delta = find(events, "response.output_audio.delta").to_realtime()
assert delta["metadata"]["vllm_omni"]["response_request_metrics"] == expected
finally:
await close_harness(h)
Loading
Loading