Skip to content

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.

ConnectFn module-attribute

ConnectFn = Callable[[str], Awaitable["WebSocketTransport"]]

DUPLEX_FIRST_UNIT_MS module-attribute

DUPLEX_FIRST_UNIT_MS = 1000

DUPLEX_UNIT_MS module-attribute

DUPLEX_UNIT_MS = 1000

PCM16_BYTES_PER_SAMPLE module-attribute

PCM16_BYTES_PER_SAMPLE = 2

PCM16_SAMPLE_RATE module-attribute

PCM16_SAMPLE_RATE = 16000

AudioDelta

Bases: DuplexEvent

AudioFormat dataclass

One PCM wire format: encoding name plus sample rate.

bytes_per_sample property

bytes_per_sample: int

encoding instance-attribute

encoding: str

sample_rate_hz instance-attribute

sample_rate_hz: int

byte_count

byte_count(duration_ms: float) -> int

duration_ms

duration_ms(byte_count: int) -> float

ConnectionResumed

Bases: DuplexEvent

Synthetic client event: the transport dropped and was resumed.

DuplexClient

Bases: DuplexClientBase

Async client for one duplex session over /v1/realtime?duplex=1.

resume_token instance-attribute

resume_token: str | None = None

url instance-attribute

url = url

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.

config instance-attribute

config = config or SessionConfig()

model instance-attribute

model = model

session_id instance-attribute

session_id: str | None = None

session_info instance-attribute

session_info: dict[str, object] = {}

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.

clear_input async

clear_input() -> None

close async

close(*, timeout_s: float = 20.0) -> None

Send session.close and wait for the server to confirm.

commit async

commit(
    *,
    final: bool = True,
    create_response: bool | None = None,
) -> None

events async

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(event: dict[str, object]) -> str

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.

DuplexClientError

Bases: Exception

Base class for all duplex client errors.

DuplexConnectionError

Bases: DuplexClientError

The WebSocket transport could not be established or failed.

DuplexEvent

Thin typed wrapper over one server event dict.

audio property

audio: bytes | None

event_id property

event_id: str | None

item_id property

item_id: str | None

raw instance-attribute

raw = raw

response_id property

response_id: str | None

sample_rate_hz property

sample_rate_hz: int | None

server_event_seq property

server_event_seq: int | None

session_id property

session_id: str | None

text property

text: str | None

type property

type: str

DuplexProtocolError

Bases: DuplexClientError

The server rejected a request with an error event.

code instance-attribute

code = code

event_id instance-attribute

event_id = event_id

raw instance-attribute

raw = dict(raw or {})

DuplexSessionClosedError

Bases: DuplexClientError

The session ended and the client cannot continue.

reason instance-attribute

reason = reason

ErrorEvent

Bases: DuplexEvent

code property

code: str | None

message property

message: str

related_event_id property

related_event_id: str | None

to_exception

to_exception() -> DuplexProtocolError

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.

event_received_at_s instance-attribute

event_received_at_s: list[float] = []

events instance-attribute

events: list[dict[str, object]] = []

output_sample_rate_hz instance-attribute

output_sample_rate_hz = 24000

response_audio instance-attribute

response_audio: dict[str, list[bytes]] = {}

response_ids instance-attribute

response_ids: list[str] = []

add

add(
    event: dict[str, object] | DuplexEvent,
    *,
    received_at_s: float | None = None,
) -> None

audio_bytes

audio_bytes(response_id: str | None = None) -> bytes

consume async

consume(client: DuplexClientBase) -> None

Subscribe to client and collect until the session ends.

count

count(event_type: str) -> int

errors

errors() -> list[dict[str, object]]

first_received_at

first_received_at(
    *event_types: str, after_s: float = 0.0
) -> float | None

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.

last_received_at

last_received_at(event_type: str) -> float | None

response_id staticmethod

response_id(event: dict[str, object]) -> str | None

Response identity of one raw event (response_id or response.id).

response_is_done

response_is_done(response_id: str) -> bool

response_text

response_text(response_id: str) -> str

Join all text/transcript deltas for one response identity.

timing_summary

timing_summary(
    *,
    after_s: float,
    input_committed_at_s: float | None = None,
    response_id: str | None = None,
    measurement_origin: dict[str, str] | None = None,
) -> dict[str, object]

Summarize engine token metrics and client-observed audio cadence.

ListenDecision

Bases: DuplexEvent

ReconnectPolicy dataclass

Jittered exponential backoff for session.resume reconnects.

backoff_s class-attribute instance-attribute

backoff_s: tuple[float, float] = (0.25, 4.0)

max_attempts class-attribute instance-attribute

max_attempts: int = 5

delay

delay(attempt: int) -> float

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.

created_event instance-attribute

created_event = created_event

decision instance-attribute

decision = decision

done_event instance-attribute

done_event: DuplexEvent | None = None

finished property

finished: bool

played_ms instance-attribute

played_ms = 0.0

response_id instance-attribute

response_id = response_id

text instance-attribute

text = ''

transcript instance-attribute

transcript = ''

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.

wait async

wait(timeout_s: float | None = None) -> None

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.

auto_response class-attribute instance-attribute

auto_response: bool = True

extra_body class-attribute instance-attribute

extra_body: dict[str, object] = field(default_factory=dict)

idle_timeout_s class-attribute instance-attribute

idle_timeout_s: float | None = None

input_audio class-attribute instance-attribute

input_audio: AudioFormat = AudioFormat('pcm16', 16000)

instructions class-attribute instance-attribute

instructions: str | None = None

modalities class-attribute instance-attribute

modalities: tuple[str, ...] = ('audio', 'text')

output_audio class-attribute instance-attribute

output_audio: AudioFormat = AudioFormat('pcm16', 24000)

overlap_policy class-attribute instance-attribute

overlap_policy: str | None = None

playback_commit_policy class-attribute instance-attribute

playback_commit_policy: str | None = None

ref_audio class-attribute instance-attribute

ref_audio: str | None = None

temperature class-attribute instance-attribute

temperature: float | None = None

turn_detection class-attribute instance-attribute

turn_detection: dict[str, object] | None = None

voice class-attribute instance-attribute

voice: str | None = None

to_session_payload

to_session_payload(*, model: str) -> dict[str, object]

SessionCreated

Bases: DuplexEvent

resume_token property

resume_token: str | None

session property

session: dict[str, object]

SessionExpired

Bases: DuplexEvent

SessionResumed

SessionUpdated

Bases: DuplexEvent

An acknowledged replacement of the public session configuration.

session property

session: dict[str, object]

SpeakDecision

Bases: DuplexEvent

TextDelta

Bases: DuplexEvent

TranscriptDelta

Bases: DuplexEvent

WebSocketTransport

Bases: ABC

Minimal WebSocket surface the client needs (satisfied by websockets).

close abstractmethod async

close() -> None

recv abstractmethod async

recv() -> str | bytes

send abstractmethod async

send(data: str) -> None

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

audio_data_url(
    data: bytes | Path, *, mime: str = "audio/wav"
) -> str

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

chunk_period_ms(
    events: Sequence[dict[str, object]],
    *,
    default: int = 1000,
) -> int

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

duplex_unit_boundary_ms(unit_index: int) -> int

Cumulative appended audio, in ms, that closes model unit unit_index.

has_residual_model_unit

has_residual_model_unit(
    pcm16: bytes, *, chunk_period_ms: int
) -> bool

True when pcm16 (16 kHz mono) does not end on a model-unit boundary.

image_data_url

image_data_url(
    data: bytes | Path, *, mime: str = "image/jpeg"
) -> str

Encode an image (or a file path) as a data URL for video_frames.

metric_mean

metric_mean(value: object) -> float | None

Read a scalar mean, or the mean field of a distribution summary.

read_pcm16_wav

read_pcm16_wav(
    path: Path, *, sample_rate_hz: int = 16000
) -> bytes

Read a mono, uncompressed PCM16 WAV file at the expected rate.

reference_audio_data_url

reference_audio_data_url(path: str | None) -> str | None

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

wait_for_condition(
    predicate: Callable[[], bool],
    *,
    timeout_s: float,
    label: str,
) -> None

Poll a collector predicate without coupling to a scenario runner.

wrap_event

wrap_event(raw: dict[str, object]) -> DuplexEvent

write_pcm16_wav

write_pcm16_wav(
    path: Path, pcm16: bytes, *, sample_rate_hz: int
) -> None

Write mono PCM16 bytes as a WAV file.