vllm_omni.clients.duplex ¶
Public asynchronous client for the vLLM-Omni full-duplex Realtime API.
:class:DuplexClientBase holds the transport-agnostic session client (typed events, response demultiplexing, incremental playback acking, the session handshake). :class:DuplexClient connects it to /v1/realtime?duplex=1 (the normative duplex contract) over WebSocket with transparent session resume on transport drops; :class:vllm_omni.clients.inline_duplex.InlineDuplexClient drives an in-process :class:~vllm_omni.entrypoints.duplex_omni.DuplexOmni with exactly the same usage.
Example::
from vllm_omni.clients.duplex import DuplexClient
from vllm_omni.clients.minicpmo_4_5 import create_duplex_session_config
cfg = create_duplex_session_config(ref_audio=audio_data_url(wav_bytes))
async with DuplexClient("ws://localhost:8099", model=model, config=cfg) as client:
await client.stream_pcm(pcm16)
await client.commit()
async for response in client.responses():
if response.decision == "listen":
continue
async for chunk in response.audio():
play(chunk)
await client.ack_playback(response.played_ms)
break
Media capture/playback, VAD, and resampling stay in the application; this module moves bytes and events.
AudioDelta ¶
Bases: DuplexEvent
AudioFormat dataclass ¶
One PCM wire format: encoding name plus sample rate.
ConnectionResumed ¶
DuplexClient ¶
Bases: DuplexClientBase
Async client for one duplex session over /v1/realtime?duplex=1.
DuplexClientBase ¶
Bases: ABC
Transport-agnostic duplex session client.
Everything a caller touches (typed events, response demux, input helpers, playback acks, the session.update handshake) lives here. A subclass only supplies the transport: :meth:_open establishes it and must deliver the server's session.created payload through :meth:_dispatch, :meth:_send_command delivers one client event, :meth:_teardown releases the transport.
ack_playback async ¶
ack_playback(
played_ms: float,
*,
response_id: str | None = None,
item_id: str | None = None,
committed_ms: float | None = None,
) -> None
Report cumulative playback progress; call periodically while playing.
append_audio async ¶
append_audio(
pcm: bytes,
*,
is_speech: bool | None = None,
video_frames: Sequence[str] | None = None,
) -> None
Append one chunk of input audio in the session's input format.
A turn is ended with :meth:commit, never by a flag on the append.
video_frames takes base64 JPEG/PNG strings (one per ~1 s unit for omni models). Data URLs (:func:image_data_url) are accepted too; their prefix is stripped, since the wire contract carries bare base64.
cancel_response async ¶
cancel_response(response_id: str | None = None) -> None
Cancel the active (or a specific) response.
Together with :meth:clear_input this composes a client-forced barge-in for sessions whose capabilities advertise supports_barge_in; the currently bundled duplex models do not, so their overlap handling is model-owned.
close async ¶
close(*, timeout_s: float = 20.0) -> None
Send session.close and wait for the server to confirm.
events async ¶
events() -> AsyncIterator[DuplexEvent]
Iterate over every server event from now on (typed).
responses async ¶
responses() -> AsyncIterator[ResponseHandle]
Iterate over response lifecycles (single consumer).
Raises :class:DuplexProtocolError when the server rejects a request with an error event while waiting — a rejected send (an invalid video frame, an empty commit) produces no response, so without this the wait would never end. Re-enter the iterator to keep consuming after handling the error.
send async ¶
Send one raw client event; returns the (possibly stamped) event_id.
stream_pcm async ¶
stream_pcm(
pcm: bytes,
*,
chunk_ms: int = 200,
realtime: bool = True,
is_speech: bool | None = None,
video_frames: Sequence[str] | None = None,
stacked_video_frames: Sequence[str | None]
| None = None,
) -> int
Slice pcm into chunk_ms chunks, pace, append; interleave frames.
video_frames holds base64/data-URL JPEG frames in capture order, one per second of the clip. Frame k rides the append that closes model unit k (see :func:duplex_unit_boundary_ms), which reproduces the official streaming_prefill(audio_waveform=<1 s>, frame_list=[frame]) pairing: a second of audio and the picture captured during it enter the same unit. Sending a frame before its whole-second unit boundary would strand it on an append that cannot close the unit yet.
stacked_video_frames is the optional parallel track of composites (see vllm_omni.experimental.fullduplex.video_stacking): entry k tiles the sub-frames captured inside unit k and rides the same append right after the base frame. None entries send the base frame alone.
Returns the number of base frames actually sent (a clip shorter than the frame list leaves the tail unsent).
wait_for async ¶
wait_for(*types: str, timeout_s: float) -> DuplexEvent
Wait for the next event whose type is in types.
Raises :class:DuplexProtocolError if an error event arrives first, and :class:DuplexSessionClosedError if the session ends first.
DuplexConnectionError ¶
Bases: DuplexClientError
The WebSocket transport could not be established or failed.
DuplexEvent ¶
Thin typed wrapper over one server event dict.
DuplexProtocolError ¶
Bases: DuplexClientError
The server rejected a request with an error event.
DuplexSessionClosedError ¶
ErrorEvent ¶
Bases: DuplexEvent
EventCollector ¶
Accumulate events for assertions and latency summaries (tests/benchmarks).
Unbounded by design — attach one to a short-lived probe session, not to a production client that runs for hours.
consume async ¶
consume(client: DuplexClientBase) -> None
Subscribe to client and collect until the session ends.
global_timing_summary ¶
global_timing_summary(
*,
after_s: float,
window_started_at_s: float,
response_ids: list[str],
measurement_origin: dict[str, str],
) -> dict[str, object]
Summarize one client-observed timing window across responses.
Returns raw measurements only. Derived metrics such as RTF are computed by the caller — e.g. with vllm_omni.metrics.definitions.compute_audio_rtf.
response_id staticmethod ¶
Response identity of one raw event (response_id or response.id).
response_text ¶
Join all text/transcript deltas for one response identity.
ListenDecision ¶
Bases: DuplexEvent
ReconnectPolicy dataclass ¶
ResponseCreated ¶
Bases: DuplexEvent
ResponseDone ¶
Bases: DuplexEvent
ResponseHandle ¶
One response.created → terminal lifecycle, demultiplexed.
A handle reaches :meth:DuplexClient.responses only once its decision is known: when the response starts producing output (decision is "speak" or about to become it) or when a terminal listen finishes it (decision == "listen", no audio). Standalone listen decisions surface as already-finished handles too, so a turn loop stays uniform.
audio async ¶
audio() -> AsyncIterator[bytes]
Yield decoded output-audio chunks until the response finishes.
played_ms advances as chunks are yielded, so it reflects what the application has taken for playback — feed it to :meth:DuplexClient.ack_playback. A handle whose audio is never drained drops its oldest buffered chunks once the buffer fills; it never stalls the client's reader.
SessionClosed ¶
Bases: DuplexEvent
SessionConfig dataclass ¶
Model-agnostic duplex session configuration.
Model-specific knobs ride in extra_body; each duplex model ships a create_duplex_session_config preset in its client module (e.g. vllm_omni.clients.minicpmo_4_5), so the client itself stays model-neutral.
extra_body class-attribute instance-attribute ¶
input_audio class-attribute instance-attribute ¶
input_audio: AudioFormat = AudioFormat('pcm16', 16000)
output_audio class-attribute instance-attribute ¶
output_audio: AudioFormat = AudioFormat('pcm16', 24000)
SessionCreated ¶
SessionExpired ¶
Bases: DuplexEvent
SessionResumed ¶
Bases: SessionCreated
SessionUpdated ¶
Bases: DuplexEvent
An acknowledged replacement of the public session configuration.
SpeakDecision ¶
Bases: DuplexEvent
TextDelta ¶
Bases: DuplexEvent
TranscriptDelta ¶
Bases: DuplexEvent
WebSocketTransport ¶
acknowledge_collected_playback async ¶
acknowledge_collected_playback(
client: DuplexClientBase, collector: EventCollector
) -> None
Ack playback of every collected response's audio (probe shorthand).
A response that has produced no audio yet is still acked (0 ms) unless it already finished: the ack checkpoints the response's history position on the server, so a user input committed later cannot displace it.
audio_data_url ¶
Encode audio bytes (or a file path) as a data URL for ref_audio.
build_realtime_url ¶
build_realtime_url(
url: str,
model: str | None,
*,
autostart: bool | None = None,
extra_query: dict[str, str] | None = None,
) -> str
Add explicit duplex query parameters to a Realtime URL.
:class:DuplexClient builds its own URL; this helper is for drivers that speak the wire protocol directly. Model-specific query flags ride in extra_query. http(s) URLs are rewritten to ws(s).
chunk_period_ms ¶
Read the negotiated native-duplex model-unit duration from session events.
distribution_summary ¶
distribution_summary(
values: Sequence[float], *, digits: int = 3
) -> dict[str, float | int] | None
Summarize values as {count, mean, p50, p99} for duplex report fields.
duplex_unit_boundary_ms ¶
Cumulative appended audio, in ms, that closes model unit unit_index.
has_residual_model_unit ¶
True when pcm16 (16 kHz mono) does not end on a model-unit boundary.
image_data_url ¶
Encode an image (or a file path) as a data URL for video_frames.
metric_mean ¶
Read a scalar mean, or the mean field of a distribution summary.
read_pcm16_wav ¶
Read a mono, uncompressed PCM16 WAV file at the expected rate.
reference_audio_data_url ¶
Encode a local reference WAV for a Realtime session update.
summarize_session_request_metrics ¶
summarize_session_request_metrics(
request_metrics: list[dict[str, object]],
*,
session_id: str | None,
) -> dict[str, object]
Summarize client- and engine-observed metrics across audio responses.
request_metrics entries are caller-assembled dicts; keys that are absent or non-numeric in an entry are simply skipped. rtf is not produced by :meth:EventCollector.timing_summary (which reports raw data only) — callers that want session rtf add an rtf value per turn, e.g. via vllm_omni.metrics.definitions.compute_audio_rtf. Zero or missing tpot_ms values are omitted from session tpot_ms.
Aggregatable fields are nested as {count, mean, p50, p99}.
summarize_stage_metrics ¶
summarize_stage_metrics(
request_metrics: Sequence[Mapping[str, object]],
) -> dict[str, dict[str, object]] | None
Roll per-response engine stage blocks into {count, mean, p50, p99}.
Reads stages when present, otherwise stage0_tokens as stage "0". Zero or missing tpot_ms / tpop_ms values are omitted, matching session tpot_ms.
wait_for_condition async ¶
Poll a collector predicate without coupling to a scenario runner.