Skip to content

vllm_omni.protocol.duplex

vLLM-Omni's full-duplex extension of the OpenAI Realtime wire protocol.

OpenAI's Realtime API assumes one turn at a time: the client commits input, the server answers. A full-duplex session does not --- the model listens while it speaks, it decides when to take a turn, the client reports playback progress, and a dropped socket re-attaches to a session that kept running. Those need events and commands the GA vocabulary does not have.

This package is that delta. It builds on vllm_omni.protocol.realtime (same base classes, same pure wire rendering, no duplication) and, like it, owns no session and imports no runtime: the duplex engine depends on this package, not the other way round.

The one door

It is also the only door a duplex consumer uses. protocol.duplex.commands and protocol.duplex.events carry the complete vocabulary --- Tier 1 classes re-exported unchanged alongside the Tier 2 and Tier 3 ones --- and this module re-exports the Tier 1 helper functions. So the layering is a chain, not a mesh::

protocol.realtime   (OpenAI's surface)
    ^
    | imports
    |
protocol.duplex     (this package: re-exports Tier 1, adds Tier 2 and 3)
    ^
    | imports
    |
duplex engine / entrypoints / clients

tests/protocol/duplex/test_duplex_protocol_facade.py asserts that the last arrow is the only one: no duplex consumer imports protocol.realtime directly. The payoff is that a helper which later needs a duplex-specific version --- convert_input_audio_with_rate is the standing example, it resamples to MiniCPM-o's 16 kHz rather than the client's rate --- is overridden in one file instead of at every call site.

What stays outside

The mailbox half of a command does not live here. A duplex command carries two different things: the client event it decodes from (wire, this package) and the dictionary the session runner consumes (engine). They genuinely differ --- session.update, conversation.item.create / .delete / .truncate all travel on the runner's turn.signal channel --- so payload() and its type stay in vllm_omni.engine.duplex.commands. Events have no such half, which is why :data:~vllm_omni.protocol.duplex.events.DuplexEvent can be a plain alias of RealtimeEvent while DuplexCommand cannot.

Modules:

Name Description
commands

The client-event vocabulary a full-duplex session accepts.

errors

vLLM-Omni's Realtime error-code vocabulary.

events

vLLM-Omni's half of the Realtime server-event vocabulary.

MAX_INPUT_SAMPLE_RATE_HZ module-attribute

MAX_INPUT_SAMPLE_RATE_HZ = 192000

MIN_INPUT_SAMPLE_RATE_HZ module-attribute

MIN_INPUT_SAMPLE_RATE_HZ = 8000

REALTIME_ERROR_TYPES_BY_CODE module-attribute

REALTIME_ERROR_TYPES_BY_CODE: dict[str, str] = {
    "bad_event": "invalid_request_error",
    "bad_audio": "invalid_request_error",
    "config_timeout": "invalid_request_error",
    "invalid_json": "invalid_request_error",
    "event_too_large": "invalid_request_error",
    "unknown_event": "invalid_request_error",
    "internal_error": "server_error",
    "runtime_append_failed": "server_error",
    "runtime_append_task_failed": "server_error",
    "runtime_signal_failed": "server_error",
    "runtime_abort_failed": "server_error",
    "runtime_data_plane_stream_failed": "server_error",
    "runtime_data_plane_text_without_audio": "server_error",
    "resource_exhausted": "rate_limit_error",
    "session_exists": "invalid_request_error",
    "session_closed": "invalid_request_error",
    "unknown_session": "invalid_request_error",
    "invalid_duplex_runtime_config": "invalid_request_error",
    "instructions_update_unsupported": "invalid_request_error",
    "persona_update_unsupported": "invalid_request_error",
    "voice_update_unsupported": "invalid_request_error",
    "unsupported_nemotron_duplex_mode": "invalid_request_error",
    "unsupported_native_response_options": "invalid_request_error",
    "runtime_touch_failed": "server_error",
    "engine_error": "server_error",
    "input_backpressure": "rate_limit_error",
    "response_already_active": "invalid_request_error",
    "response_not_active": "invalid_request_error",
    "response_create_without_input": "invalid_request_error",
    "text_only_turn_unsupported": "invalid_request_error",
    "commit_aborted": "server_error",
    "input_audio_buffer_empty": "invalid_request_error",
    "missing_item_id": "invalid_request_error",
    "item_not_found": "invalid_request_error",
    "playback_item_mismatch": "invalid_request_error",
    "playback_item_not_found": "invalid_request_error",
    "playback_ack_too_late": "invalid_request_error",
    "unsupported_audio_format": "invalid_request_error",
    "unsupported_turn_detection": "invalid_request_error",
    "unsupported_ref_audio_path": "invalid_request_error",
    "ref_audio_required": "invalid_request_error",
    "model_update_unsupported": "invalid_request_error",
    "voice_update_after_audio_unsupported": "invalid_request_error",
    "ref_audio_update_unsupported": "invalid_request_error",
    "native_text_append_unsupported": "invalid_request_error",
    "invalid_video_frames": "invalid_request_error",
    "invalid_input_modality": "invalid_request_error",
    "invalid_function_call_output": "invalid_request_error",
    "server_vad_unavailable": "server_error",
}

REALTIME_INPUT_AUDIO_FORMATS module-attribute

REALTIME_INPUT_AUDIO_FORMATS = {
    "pcm16",
    "pcm_s16le",
    "s16le",
    "pcm_f32le",
    "g711_ulaw",
    "g711_alaw",
}

REALTIME_INPUT_HINT_KEYS module-attribute

REALTIME_INPUT_HINT_KEYS = (
    "duration_ms",
    "audio_duration_ms",
    "audio_start_ms",
    "audio_end_ms",
    "is_speech",
    "speech",
    "speech_probability",
    "vad",
    "overlap_action",
    "overlap",
    "force_barge_in",
    "force_listen",
    "text",
    "transcript",
)

REALTIME_OUTPUT_AUDIO_FORMATS module-attribute

REALTIME_OUTPUT_AUDIO_FORMATS = {
    "pcm16",
    "pcm_s16le",
    "s16le",
    "wav",
    "pcm",
    "g711_ulaw",
    "g711_alaw",
}

RealtimeAudioAppend dataclass

One decoded input_audio_buffer.append.

audio is raw bytes in format at sample_rate_hz --- after :func:decode_audio_append that is 16 kHz pcm_f32le for every input format the codec converts.

audio instance-attribute

audio: bytes

audio_end_ms class-attribute instance-attribute

audio_end_ms: int | None = None

duration_ms class-attribute instance-attribute

duration_ms: int | None = None

event_id class-attribute instance-attribute

event_id: str | None = None

format instance-attribute

format: str

hints class-attribute instance-attribute

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

is_speech class-attribute instance-attribute

is_speech: bool | None = None

sample_rate_hz class-attribute instance-attribute

sample_rate_hz: int | None = None

video_frames class-attribute instance-attribute

video_frames: tuple[str, ...] = ()

RealtimeInputDefaults dataclass

Session-level wire defaults an append may omit (derived from the session object).

input_audio_format class-attribute instance-attribute

input_audio_format: str = 'pcm16'

input_sample_rate_hz class-attribute instance-attribute

input_sample_rate_hz: int = 16000

output_audio_format class-attribute instance-attribute

output_audio_format: str = 'pcm16'

output_sample_rate_hz class-attribute instance-attribute

output_sample_rate_hz: int | None = None

overlap_silence_rms class-attribute instance-attribute

overlap_silence_rms: float = 0.003

with_session_payload

with_session_payload(
    session_payload: Mapping[str, object],
) -> RealtimeInputDefaults

Return defaults updated from a Realtime session object (session.update).

RealtimeProtocolCapabilities dataclass

One consumer's answer to what a session object may ask for.

The defaults are permissive, not restrictive: every format the codec can decode, and no turn-detection validation at all. A consumer that leaves validate_turn_detection at None therefore accepts whatever turn_detection the session object carries --- including values it cannot actually serve, such as semantic_vad. A consumer that supports only some turn-detection modes (or none) must supply its own validator; see vllm_omni.engine.duplex.realtime_commands.DUPLEX_REALTIME_CAPABILITIES.

input_audio_formats class-attribute instance-attribute

input_audio_formats: frozenset[str] = field(
    default_factory=lambda: frozenset(
        REALTIME_INPUT_AUDIO_FORMATS
    )
)

output_audio_formats class-attribute instance-attribute

output_audio_formats: frozenset[str] = field(
    default_factory=lambda: frozenset(
        REALTIME_OUTPUT_AUDIO_FORMATS
    )
)

validate_turn_detection class-attribute instance-attribute

validate_turn_detection: TurnDetectionValidator | None = (
    None
)

RealtimeProtocolError

Bases: ValueError

A client payload could not be decoded into a valid Realtime intent.

code is the internal error code that :data:REALTIME_ERROR_TYPES_BY_CODE maps to an OpenAI error.type; event_id is the client event id the error answers. vllm_omni.engine.duplex.commands.DuplexCommandError is the duplex specialization, so a consumer catching either sees the same three attributes.

code instance-attribute

code = code

event_id instance-attribute

event_id = event_id

RealtimeSessionRejection dataclass

Why a session object was refused: the internal error code and the message.

param names the offending session field where there is one, for the error.param slot of the OpenAI error envelope.

code instance-attribute

code: str

message instance-attribute

message: str

param class-attribute instance-attribute

param: str | None = None

apply_realtime_session_defaults

apply_realtime_session_defaults(
    defaults: RealtimeInputDefaults,
    session_payload: Mapping[str, object],
) -> RealtimeInputDefaults

Derive the wire defaults a Realtime session object declares.

convert_input_audio_with_rate

convert_input_audio_with_rate(
    audio: object,
    fmt: object,
    *,
    sample_rate_hz: int | float | None = None,
    target_sample_rate_hz: int = 16000,
) -> tuple[object, object, int | float | None]

convert_output_audio

convert_output_audio(
    audio: str,
    *,
    source_fmt: str,
    target_fmt: str,
    source_sample_rate_hz: int | None = None,
    target_sample_rate_hz: int | None = None,
) -> tuple[str, str, int | None]

copy_realtime_input_hints

copy_realtime_input_hints(
    source: Mapping[str, object], target: dict[str, object]
) -> None

decode_audio_append

decode_audio_append(
    event: Mapping[str, object],
    *,
    defaults: RealtimeInputDefaults,
    hints_source: Mapping[str, object] | None = None,
) -> RealtimeAudioAppend

Validate and pack one input append (audio and/or video frames).

Capability checks (required/optional modalities) run later on the session. Audio path converts to 16 kHz pcm_f32le.

hints_source is the enclosing payload when the audio arrives inside something larger than a bare append --- a conversation.item.create audio part, say --- so its hints apply unless the part overrides them.

Raises :class:RealtimeProtocolError for an unsupported format, undecodable audio or invalid camera frames.

decode_g711_alaw

decode_g711_alaw(raw: bytes) -> bytes

decode_g711_ulaw

decode_g711_ulaw(raw: bytes) -> bytes

encode_float32_mono_wav_base64

encode_float32_mono_wav_base64(
    samples: ndarray, *, sample_rate_hz: int
) -> str

Package normalized mono float32 samples as a PCM16 WAV data payload.

encode_g711_alaw

encode_g711_alaw(raw: bytes) -> bytes

encode_g711_ulaw

encode_g711_ulaw(raw: bytes) -> bytes

input_audio_transcription_config

input_audio_transcription_config(
    session_payload: Mapping[str, object],
) -> dict[str, object] | None

input_explicitly_non_speech

input_explicitly_non_speech(
    event: Mapping[str, object],
) -> bool

input_looks_like_speech

input_looks_like_speech(
    event: Mapping[str, object],
    *,
    audio: object,
    fmt: object,
    overlap_silence_rms: float,
) -> bool

input_transcript_from_item

input_transcript_from_item(
    item: Mapping[str, object],
) -> str

is_supported_realtime_input_format

is_supported_realtime_input_format(fmt: object) -> bool

json_safe_realtime_payload

json_safe_realtime_payload(
    payload: Mapping[str, object],
) -> dict[str, object]

normalize_conversation_item

normalize_conversation_item(
    item: Mapping[str, object],
) -> dict[str, object]

parse_realtime_audio_format

parse_realtime_audio_format(
    raw_format: object,
) -> tuple[object, int | None]

realtime_audio_format_object

realtime_audio_format_object(
    fmt: object, *, sample_rate_hz: int | None = None
) -> dict[str, object]

realtime_error_type

realtime_error_type(code: str) -> str

The OpenAI error.type bucket for an internal error code.

realtime_max_output_tokens

realtime_max_output_tokens(value: object) -> int | None

Normalize Realtime max output tokens ("inf" -> None).

realtime_output_format

realtime_output_format(duplex_format: object) -> str

realtime_overlap_fields

realtime_overlap_fields(
    session_payload: Mapping[str, object],
) -> dict[str, object]

resample_pcm16_mono

resample_pcm16_mono(
    raw: bytes, *, source_rate_hz: int, target_rate_hz: int
) -> bytes

text_chars_for_audio_ms_from_marks

text_chars_for_audio_ms_from_marks(
    audio_end_ms: int,
    text_len: int,
    marks: list[object],
    *,
    final_ms: object | None = None,
) -> int

truncate_realtime_item_content

truncate_realtime_item_content(
    item: dict[str, object],
    *,
    content_index: int,
    audio_end_ms: int,
) -> None

validate_conversation_item_audio_formats

validate_conversation_item_audio_formats(
    item: object,
) -> str | None

validate_input_sample_rate_hz

validate_input_sample_rate_hz(
    sample_rate_hz: object,
) -> int

validate_realtime_item_truncate

validate_realtime_item_truncate(
    item: Mapping[str, object],
    *,
    content_index: int,
    audio_end_ms: int,
) -> str | None

validate_realtime_response_audio_formats

validate_realtime_response_audio_formats(
    response_payload: Mapping[str, object],
) -> str | None

validate_realtime_session_audio_formats

validate_realtime_session_audio_formats(
    session_payload: Mapping[str, object],
    *,
    input_audio_formats: Collection[str] | None = None,
    output_audio_formats: Collection[str] | None = None,
) -> str | None

Reject a session object that declares an audio format we cannot serve.

The format sets default to everything the codec can decode; a consumer that serves a narrower set passes its own (see vllm_omni.protocol.realtime.capabilities).

validate_realtime_video_frames

validate_realtime_video_frames(
    video_frames: object, max_slice_nums: object
) -> str | None

Validate omni-duplex camera frames on input_audio_buffer.append.

Wire contract matches the official MiniCPM-o duplex loop: one base base64 JPEG per ~1 s audio chunk, optionally followed by that unit's stacked composite tiling the sub-frames captured inside it (at most 2 images either way). A caller-supplied max_slice_nums is rejected rather than silently ignored: slicing is Stage 0's decision here, and Stage 0 already applies the official HD suggestion for a stacked unit (max_slice_nums=[2, 1]). The wire simply does not let the client choose it.

validate_session_payload

validate_session_payload(
    session_payload: Mapping[str, object],
    *,
    capabilities: RealtimeProtocolCapabilities,
) -> RealtimeSessionRejection | None

Check a session object against one consumer; None when it is acceptable.

Audio formats are checked before turn detection, so a session object that is wrong in both ways reports the format problem --- keep that order, it is what clients already see.

wav_payload_to_pcm16

wav_payload_to_pcm16(
    raw: bytes,
) -> tuple[bytes | None, int | None]