From 10ffb1afee7262ab3a48a9035c661801603850c6 Mon Sep 17 00:00:00 2001 From: chickeyton Date: Sun, 20 Sep 2026 12:23:34 +0800 Subject: [PATCH] [Bugfix][MiniCPM-o] Restore the duplex server-side response timing lost in #7413 #7413 rebuilt the duplex serving stack around engine-resident sessions and dropped the per-response request timing that #7242 had added for the MiniCPM-o Omni-DuplexEval / OmniInteract benchmarks: the engine no longer emitted ``response_request_metrics``, so ``DuplexClient`` and the benchmark summaries silently fell back to the client's receive clock, and the TPOT weighting regressed to counting a one-token segment as one interval. Port the timing into the engine session: the model-turn request start is recorded when the model channel submits the append, the first text / audio of the response fix TTFT / TTFP, and the numbers ride response.created (under response.metadata.duplex_event) and every speak / audio delta (under metadata.vllm_omni), exactly where the client already reads them. Restore the interval-weighted TPOT mean. Re-home the MiniCPM-o plugin policy unit tests deleted with tests/engine/duplex/test_duplex_runtime.py, make the duplex e2e gate assert the server-side anchor, fix the pre-existing mypy errors in the two touched test files, and drop three doc passages that still describe the session-per-request chat adapter removed before #7413 merged. Co-Authored-By: Claude Fable 5.1 Signed-off-by: chickeyton --- docs/design/fullduplex.md | 2 - docs/serving/README.md | 12 +- docs/serving/realtime_duplex_api.md | 31 ++--- .../test_minicpmo_4_5_duplex.py | 17 ++- tests/engine/duplex/test_engine_session.py | 75 ++++++++++ tests/engine/duplex/test_session_runner.py | 54 +++++++- .../minicpmo_4_5/duplex/test_plugin_policy.py | 129 ++++++++++++++++++ .../engine/duplex/session/engine_session.py | 107 +++++++++++++-- .../engine/duplex/session/model_channel.py | 31 ++++- 9 files changed, 414 insertions(+), 44 deletions(-) create mode 100644 tests/model_executor/models/minicpmo_4_5/duplex/test_plugin_policy.py diff --git a/docs/design/fullduplex.md b/docs/design/fullduplex.md index 324cd5bdcad..eb370c819b6 100644 --- a/docs/design/fullduplex.md +++ b/docs/design/fullduplex.md @@ -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) @@ -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`. diff --git a/docs/serving/README.md b/docs/serving/README.md index 472ab2a492b..71bbeb14d99 100644 --- a/docs/serving/README.md +++ b/docs/serving/README.md @@ -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 diff --git a/docs/serving/realtime_duplex_api.md b/docs/serving/realtime_duplex_api.md index 029f35f5836..ec780b148a3 100644 --- a/docs/serving/realtime_duplex_api.md +++ b/docs/serving/realtime_duplex_api.md @@ -453,11 +453,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 | | --- | --- | @@ -472,17 +472,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 @@ -491,10 +484,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 diff --git a/tests/e2e/online_serving/test_minicpmo_4_5_duplex.py b/tests/e2e/online_serving/test_minicpmo_4_5_duplex.py index 59b40f5e502..c75268fb2e7 100644 --- a/tests/e2e/online_serving/test_minicpmo_4_5_duplex.py +++ b/tests/e2e/online_serving/test_minicpmo_4_5_duplex.py @@ -9,6 +9,7 @@ import base64 import json from pathlib import Path +from typing import TypedDict import pytest import websockets @@ -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 @@ -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, @@ -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 diff --git a/tests/engine/duplex/test_engine_session.py b/tests/engine/duplex/test_engine_session.py index c90291a8130..5ee3399080c 100644 --- a/tests/engine/duplex/test_engine_session.py +++ b/tests/engine/duplex/test_engine_session.py @@ -610,3 +610,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 diff --git a/tests/engine/duplex/test_session_runner.py b/tests/engine/duplex/test_session_runner.py index 51711f52a91..e3618a8e024 100644 --- a/tests/engine/duplex/test_session_runner.py +++ b/tests/engine/duplex/test_session_runner.py @@ -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, @@ -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() @@ -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") @@ -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) diff --git a/tests/model_executor/models/minicpmo_4_5/duplex/test_plugin_policy.py b/tests/model_executor/models/minicpmo_4_5/duplex/test_plugin_policy.py new file mode 100644 index 00000000000..b2a63545717 --- /dev/null +++ b/tests/model_executor/models/minicpmo_4_5/duplex/test_plugin_policy.py @@ -0,0 +1,129 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the vLLM-Omni project + +"""Engine-policy half of the MiniCPM-o 4.5 duplex plugin, as the session runner calls it. + +Re-homed from ``tests/engine/duplex/test_duplex_runtime.py`` of the +pre-framework runtime extension: the per-stage sampling overrides a session's +runtime config applies, the listen decision on a finished Stage0 segment, and +the scheduler slots one PCM append reserves. +""" + +from __future__ import annotations + +from types import SimpleNamespace + +import pytest +from vllm.sampling_params import SamplingParams + +from vllm_omni.engine.duplex.contracts import DuplexOutputAction, DuplexOutputDecision +from vllm_omni.model_executor.models.minicpmo_4_5.duplex.plugin import ( + MiniCPMO45DuplexPlugin, + duplex_scheduler_token_budget, +) + +pytestmark = [pytest.mark.core_model, pytest.mark.cpu] + +LISTEN_TOKEN_ID = 151705 + + +def _plugin() -> MiniCPMO45DuplexPlugin: + return MiniCPMO45DuplexPlugin(lambda audio, sample_rate_hz, response_format, speed: None) + + +def _decide( + output: object, + *, + segment_token_ids: tuple[int, ...] = (), + segment_output_metadata: dict[str, object] | None = None, +) -> DuplexOutputDecision | None: + return _plugin().decide_output( + stage_id=0, + final_stage_id=1, + segment_finished=True, + segment_token_ids=segment_token_ids, + segment_output_metadata=segment_output_metadata or {}, + output=output, + ) + + +def test_plugin_owns_stage_sampling_overrides_without_mutating_defaults() -> None: + defaults = (SamplingParams(max_tokens=4), SamplingParams(max_tokens=8)) + + configured = _plugin().configure_sampling_params( + runtime_config={ + "duplex_stage_max_tokens": {"0": 20}, + "duplex_stage_sampling_params": {"1": {"stop_token_ids": [151645]}}, + }, + defaults=defaults, + ) + + assert configured[0].max_tokens == 20 + assert configured[1].stop_token_ids == [151645] + assert defaults[0].max_tokens == 4 + assert 151645 not in (defaults[1].stop_token_ids or []) + + +def test_listen_decision_uses_the_raw_streaming_token_snapshot() -> None: + decision = _decide( + SimpleNamespace(outputs=[SimpleNamespace()]), + segment_token_ids=(LISTEN_TOKEN_ID,), + segment_output_metadata={"special_token_ids": {"listen_token_id": LISTEN_TOKEN_ID}}, + ) + + assert decision is not None + assert decision.action is DuplexOutputAction.DIRECT_RESPONSE + assert decision.metadata["duplex_native_decision"] == "listen" + assert decision.metadata["model_listen"] is True + + +@pytest.mark.parametrize("attr", ["token_ids", "cumulative_token_ids"]) +def test_listen_decision_ignores_output_level_token_history(attr: str) -> None: + output = SimpleNamespace( + multimodal_output={"special_token_ids": {"listen_token_id": LISTEN_TOKEN_ID}}, + outputs=[SimpleNamespace()], + **{attr: [42, LISTEN_TOKEN_ID]}, + ) + + assert _decide(output) is None + + +@pytest.mark.parametrize("attr", ["token_ids", "cumulative_token_ids"]) +def test_listen_decision_uses_completion_token_ids(attr: str) -> None: + output = SimpleNamespace( + multimodal_output={"special_token_ids": {"listen_token_id": LISTEN_TOKEN_ID}}, + outputs=[SimpleNamespace(**{attr: [42, LISTEN_TOKEN_ID]})], + ) + + assert _decide(output) is not None + + +def test_listen_decision_uses_completion_stop_reason() -> None: + output = SimpleNamespace( + multimodal_output={"special_token_ids": {"listen_token_id": LISTEN_TOKEN_ID}}, + outputs=[SimpleNamespace(stop_reason=LISTEN_TOKEN_ID)], + ) + + assert _decide(output) is not None + + +def test_a_speak_segment_is_not_a_listen_decision() -> None: + output = SimpleNamespace( + multimodal_output={"special_token_ids": {"listen_token_id": LISTEN_TOKEN_ID}}, + outputs=[SimpleNamespace(token_ids=[42, 43], stop_reason=None)], + ) + + assert _decide(output) is None + + +def test_scheduler_token_budget_estimates_pcm_slots() -> None: + assert duplex_scheduler_token_budget({"audio": "AAAAAA==", "format": "pcm_f32le", "sample_rate_hz": 16000}) == 16 + + +def test_scheduler_token_budget_ignores_client_budget_fields() -> None: + assert ( + duplex_scheduler_token_budget( + {"audio": "AAAAAA==", "format": "pcm_f32le", "duplex_num_input_tokens": 999, "num_input_tokens": 999} + ) + == 16 + ) diff --git a/vllm_omni/engine/duplex/session/engine_session.py b/vllm_omni/engine/duplex/session/engine_session.py index f64ad1f937a..16e585f64c4 100644 --- a/vllm_omni/engine/duplex/session/engine_session.py +++ b/vllm_omni/engine/duplex/session/engine_session.py @@ -102,6 +102,14 @@ class ResponseState: active_request_id: str | None = None active_response_id: str | None = None active_response_turn_id: int | None = None + #: ``time.monotonic()`` at which the native model-turn request that may + #: own a response started executing, keyed by model turn. A response binds + #: its own turn's entry at ``begin_response``; TTFT/TTFP are measured from + #: that anchor rather than from the client's receipt of ``response.created``. + request_started_at_s_by_turn: dict[int, float] = field(default_factory=dict) + active_response_request_started_at_s: float | None = None + active_response_ttft_ms: float | None = None + active_response_ttfp_ms: float | None = None active_response_input_commit_seq: int | None = None active_response_awaits_input_commit: bool = False last_response_id: str | None = None @@ -582,6 +590,56 @@ def clear_request(self, expected_request_id: str | None = None) -> bool: def bind_response_turn(self, turn_id: int | None) -> None: self._response.active_response_turn_id = turn_id + def mark_model_turn_request_started(self, turn_id: int, started_at_s: float) -> None: + """Record the latest native request start that can own one model turn. + + Per-response TTFT/TTFP start here, when the server begins executing + the model-turn request, not when the client receives + ``response.created``. The active response keeps the first start of its + own turn; a start for any other turn waits until that turn's response + begins, the latest one winning. + """ + turn_id = int(turn_id) + started_at_s = float(started_at_s) + if self.active_response_turn_id == turn_id: + if self._response.active_response_request_started_at_s is None: + self._response.active_response_request_started_at_s = started_at_s + return + self._response.request_started_at_s_by_turn[turn_id] = started_at_s + + def mark_response_first_outputs( + self, + *, + observed_at_s: float, + has_text: bool, + has_audio: bool, + ) -> dict[str, object]: + """Return the server-monotonic TTFT/TTFP observed for the active response. + + Empty when no request start was recorded for the response's turn; the + first non-empty text and the first audio each fix their number once. + """ + started_at_s = self._response.active_response_request_started_at_s + if started_at_s is None: + return {} + elapsed_ms = max(0.0, (float(observed_at_s) - started_at_s) * 1000.0) + if has_text and self._response.active_response_ttft_ms is None: + self._response.active_response_ttft_ms = elapsed_ms + if has_audio and self._response.active_response_ttfp_ms is None: + self._response.active_response_ttfp_ms = elapsed_ms + metrics: dict[str, object] = { + "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", + }, + } + if self._response.active_response_ttft_ms is not None: + metrics["ttft_ms"] = self._response.active_response_ttft_ms + if self._response.active_response_ttfp_ms is not None: + metrics["ttfp_ms"] = self._response.active_response_ttfp_ms + return metrics + def active_response_accepts_model_turn(self, turn_id: int | None) -> bool: if self._response.active_response_id is None: return False @@ -675,6 +733,11 @@ def commit_audio_input( def complete_model_turn(self, turn_id: int) -> None: """Advance the model-owned output identity after its terminal signal.""" completed_turn_id = int(turn_id) + self._response.request_started_at_s_by_turn = { + pending_turn_id: started_at_s + for pending_turn_id, started_at_s in self._response.request_started_at_s_by_turn.items() + if pending_turn_id > completed_turn_id + } if completed_turn_id >= self.turn_id: self.turn_id = completed_turn_id + 1 self.sync_fence() @@ -712,6 +775,12 @@ def begin_response(self, *, turn_id: int | None = None) -> str: response_id = f"resp-{self.session_id}-{self.epoch}-{uuid4().hex[:8]}" self._response.active_response_id = response_id self._response.active_response_turn_id = self.turn_id if turn_id is None else int(turn_id) + self._response.active_response_request_started_at_s = self._response.request_started_at_s_by_turn.pop( + self._response.active_response_turn_id, + None, + ) + self._response.active_response_ttft_ms = None + self._response.active_response_ttfp_ms = None self._response.active_response_input_commit_seq = self.input_commit_seq self._response.active_response_awaits_input_commit = self.turn_state == DuplexTurnState.USER_SPEAKING self._response.last_response_id = response_id @@ -730,6 +799,13 @@ def _clear_response_metrics(self) -> None: self._response.stage_metric_tpot_weighted_ms.clear() self._response.stage_metric_tpot_weight.clear() + def _clear_response_timing(self, *, drop_pending: bool) -> None: + if drop_pending: + self._response.request_started_at_s_by_turn.clear() + self._response.active_response_request_started_at_s = None + self._response.active_response_ttft_ms = None + self._response.active_response_ttfp_ms = None + def stash_stage_metrics(self, stage_metrics: Mapping[Any, Any] | None) -> None: """Hold a stage snapshot until a response exists to attribute it to. @@ -825,17 +901,25 @@ def accumulate_response_stage_metrics( tpot_ms = raw_values.get("vllm_tpot_ms") token_count = raw_values.get("num_tokens_out") - if isinstance(tpot_ms, int | float) and tpot_ms > 0: - weight = max(int(token_count) - 1, 1) if isinstance(token_count, int | float) else 1 - self._response.stage_metric_tpot_weighted_ms[stage_id] = ( - self._response.stage_metric_tpot_weighted_ms.get(stage_id, 0.0) + float(tpot_ms) * weight - ) - self._response.stage_metric_tpot_weight[stage_id] = ( - self._response.stage_metric_tpot_weight.get(stage_id, 0) + weight - ) - current["vllm_tpot_ms"] = self._response.stage_metric_tpot_weighted_ms[stage_id] / float( - self._response.stage_metric_tpot_weight[stage_id] + if isinstance(tpot_ms, int | float) and not isinstance(tpot_ms, bool) and tpot_ms > 0: + # Weight by inter-token intervals, not tokens: a one-token + # segment (a bare unit boundary) has none, and its per-token + # time is the whole unit, which must not skew the mean. + weight = ( + max(int(token_count) - 1, 0) + if isinstance(token_count, int | float) and not isinstance(token_count, bool) + else 1 ) + if weight > 0: + self._response.stage_metric_tpot_weighted_ms[stage_id] = ( + self._response.stage_metric_tpot_weighted_ms.get(stage_id, 0.0) + float(tpot_ms) * weight + ) + self._response.stage_metric_tpot_weight[stage_id] = ( + self._response.stage_metric_tpot_weight.get(stage_id, 0) + weight + ) + current["vllm_tpot_ms"] = self._response.stage_metric_tpot_weighted_ms[stage_id] / float( + self._response.stage_metric_tpot_weight[stage_id] + ) for name, value in raw_values.items(): if name not in handled_fields: @@ -1011,6 +1095,7 @@ def end_response( self._response.active_response_input_commit_seq = None self._response.active_response_awaits_input_commit = False self._clear_response_metrics() + self._clear_response_timing(drop_pending=False) self.turn_state = DuplexTurnState.IDLE self._restore_response_config() return message @@ -1336,6 +1421,7 @@ def barge_in(self) -> int: self._response.active_response_input_commit_seq = None self._response.active_response_awaits_input_commit = False self._clear_response_metrics() + self._clear_response_timing(drop_pending=True) self._restore_response_config() self.turn_state = DuplexTurnState.BARGE_IN return self.epoch @@ -1351,6 +1437,7 @@ def close(self) -> None: self._response.active_response_input_commit_seq = None self._response.active_response_awaits_input_commit = False self._clear_response_metrics() + self._clear_response_timing(drop_pending=True) self._restore_response_config() def signal_turn(self, event_type: str, payload: Mapping[str, object] | None = None) -> TurnEvent: diff --git a/vllm_omni/engine/duplex/session/model_channel.py b/vllm_omni/engine/duplex/session/model_channel.py index f244a68597f..fde9d8df9b9 100644 --- a/vllm_omni/engine/duplex/session/model_channel.py +++ b/vllm_omni/engine/duplex/session/model_channel.py @@ -165,6 +165,9 @@ async def append_runtime_input( # Anchor the submission time before the RPC; the acceptance callback # commits timing state only if the append actually submitted. submit_time = time.monotonic() + # Per-response TTFT/TTFP start here: the model-turn request begins + # executing before the client ever sees response.created. + session.mark_model_turn_request_started(helpers.append_fence(session, payload).turn_id, submit_time) try: result = await self._append_via_data_plane( payload, @@ -776,7 +779,16 @@ async def _send_one_model_output_event( if response_id is None: response_id = session.begin_response(turn_id=model_turn_id) response_created = True - self._out.emit(self.response_created_payload(response_id, epoch=session.epoch)) + response_request_metrics = session.mark_response_first_outputs( + observed_at_s=time.monotonic(), + has_text=has_text, + has_audio=has_audio, + ) + if response_created: + created_payload = self.response_created_payload(response_id, epoch=session.epoch) + if response_request_metrics: + created_payload["response_request_metrics"] = dict(response_request_metrics) + self._out.emit(created_payload) stage_metrics = model_result.get("stage_metrics") response_stage_metrics = session.accumulate_response_stage_metrics( stage_metrics if isinstance(stage_metrics, Mapping) else None @@ -791,7 +803,12 @@ async def _send_one_model_output_event( "end_of_turn": end_of_turn, "model_speak": True, } - self._attach_runtime_metadata(speak_payload, model_result, stage_metrics=response_stage_metrics) + self._attach_runtime_metadata( + speak_payload, + model_result, + stage_metrics=response_stage_metrics, + response_request_metrics=response_request_metrics, + ) self._out.emit(speak_payload) previous_sent_ms = session.playback.sent_ms text_chars_before_append = len("".join(session.assistant_text_buffer)) @@ -854,7 +871,12 @@ async def _send_one_model_output_event( sample_rate_hz = model_result.get("sample_rate_hz") or model_result.get("audio_sample_rate_hz") if isinstance(sample_rate_hz, int | float) and int(sample_rate_hz) > 0: payload["sample_rate_hz"] = int(sample_rate_hz) - self._attach_runtime_metadata(payload, model_result, stage_metrics=response_stage_metrics) + self._attach_runtime_metadata( + payload, + model_result, + stage_metrics=response_stage_metrics, + response_request_metrics=response_request_metrics, + ) self._out.emit(payload) if ( not end_of_turn @@ -950,6 +972,7 @@ def _attach_runtime_metadata( model_result: dict[str, object], *, stage_metrics: Mapping[str, object] | None = None, + response_request_metrics: Mapping[str, object] | None = None, ) -> None: metadata: dict[str, object] = {} runtime_impl = model_result.get("runtime_impl") @@ -972,6 +995,8 @@ def _attach_runtime_metadata( for stage_id, values in effective_stage_metrics.items() if isinstance(values, Mapping) } + if response_request_metrics: + metadata["response_request_metrics"] = dict(response_request_metrics) if metadata: payload["vllm_omni"] = metadata