Skip to content

vllm_omni.engine.duplex.realtime_events

Session-internal duplex events -> typed :class:DuplexEvent objects (stateful projection).

This is the former entrypoints/duplex/realtime_output.py projector plus the projection-relevant fields of realtime_state.py, now owned by the session runner through :class:RealtimeProjectionState. The projection consumes the state (response / item ids, content-part bookkeeping) when it constructs events; rendering an event to wire JSON (event.to_realtime()) is pure and lives on the event classes in vllm_omni.engine.duplex.events.

Besides the output projection (:func:project_internal_event) the state also carries the input-side bookkeeping the old input translator kept (input buffer flags, conversation items, response-id fallbacks); the resolve_* / note_* helpers give the runner the same behaviour for the corresponding commands.

What is not here is the model- and runtime-agnostic half of the codec --- audio format negotiation, conversation-item shape and truncation, transcript extraction, audio conversion. That lives in vllm_omni.protocol.realtime so a non-duplex Realtime surface can use it without the duplex session; this module is the duplex consumer of it (RFC #6592 P0a).

RealtimeProjectionState dataclass

Per-session Realtime projection state (lives on the engine-side runner).

active_input_item_id class-attribute instance-attribute

active_input_item_id: str | None = None

active_response_id class-attribute instance-attribute

active_response_id: str | None = None

conversation_items class-attribute instance-attribute

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

default_payload class-attribute instance-attribute

default_payload: Mapping[str, object] | None = None

defaults class-attribute instance-attribute

defaults: RealtimeInputDefaults = field(
    default_factory=RealtimeInputDefaults
)

initial_session_update class-attribute instance-attribute

initial_session_update: bool = False

input_audio_buffer_had_non_speech class-attribute instance-attribute

input_audio_buffer_had_non_speech: bool = False

input_audio_buffer_has_audio class-attribute instance-attribute

input_audio_buffer_has_audio: bool = False

input_audio_buffer_transcript_parts class-attribute instance-attribute

input_audio_buffer_transcript_parts: list[str] = field(
    default_factory=list
)

input_audio_format property

input_audio_format: str

input_sample_rate_hz property

input_sample_rate_hz: int

input_speech_started class-attribute instance-attribute

input_speech_started: bool = False

item_truncation_cursors class-attribute instance-attribute

item_truncation_cursors: dict[str, tuple[int, int]] = field(
    default_factory=dict
)

last_conversation_item_id class-attribute instance-attribute

last_conversation_item_id: str | None = None

last_response_id class-attribute instance-attribute

last_response_id: str | None = None

model class-attribute instance-attribute

model: str | None = None

output_audio_format property

output_audio_format: str

output_sample_rate_hz property

output_sample_rate_hz: int | None

overlap_silence_rms property

overlap_silence_rms: float

pending_commit_item_ids class-attribute instance-attribute

pending_commit_item_ids: list[str] = field(
    default_factory=list
)

response_states class-attribute instance-attribute

response_states: dict[str | int, _ResponseProjection] = (
    field(default_factory=dict)
)

session_id instance-attribute

session_id: str

apply_session_defaults

apply_session_defaults(
    session_payload: Mapping[str, object],
) -> None

ResolvedCommit dataclass

What a Commit means given the projected input buffer.

events instance-attribute

events: list[DuplexEvent]

payload instance-attribute

payload: dict[str, object] | None

reset_vad class-attribute instance-attribute

reset_vad: bool = False

ResolvedControl dataclass

Result of resolving a cancel/clear/delete/truncate against the projection.

events instance-attribute

events: list[DuplexEvent]

payloads instance-attribute

payloads: list[dict[str, object]]

append_audio_payload

append_audio_payload(
    command: AppendAudio,
) -> dict[str, object]

Internal input_audio_buffer.append payload for a command (convenience for the runner).

clear_input_buffer

clear_input_buffer(state: RealtimeProjectionState) -> None

Input-buffer projection reset for input_audio_buffer.clear.

discard_pending_input_audio

discard_pending_input_audio(
    state: RealtimeProjectionState,
    audio_end_ms: int | None = None,
) -> list[DuplexEvent]

Drop Realtime input-buffer state that was consumed as overlap (returns typed events).

emit_input_speech_started

emit_input_speech_started(
    state: RealtimeProjectionState,
    audio_start_ms: object = 0,
) -> list[DuplexEvent]

emit_input_speech_stopped

emit_input_speech_stopped(
    state: RealtimeProjectionState,
    *,
    item_id: str,
    audio_end_ms: object = 0,
) -> list[DuplexEvent]

note_input_append

note_input_append(
    state: RealtimeProjectionState,
    payload: dict[str, object],
    *,
    vad_result: TurnDetectionResult | None = None,
    allows_video_without_audio: bool = False,
) -> list[DuplexEvent]

Update the input-buffer projection for one appended chunk; returns typed events.

payload is the internal input_audio_buffer.append dictionary (AppendAudio.payload()), optionally after apply_turn_detection_result; vad_result is the TurnDetectionResult when server VAD is active. The returned events are input_audio_buffer.speech_started / speech_stopped exactly as the old translator produced them.

project_internal_event

project_internal_event(
    state: RealtimeProjectionState,
    event: Mapping[str, object],
) -> list[DuplexEvent]

Project one session-internal event onto 0..n typed public events.

realtime_session_payload

realtime_session_payload(
    state: RealtimeProjectionState, session: object
) -> dict[str, object]

Fill the Realtime session object defaults for session.created/updated.

register_user_item

register_user_item(
    state: RealtimeProjectionState,
    item: Mapping[str, object],
    *,
    previous_item_id: str | None = None,
) -> list[DuplexEvent]

Record a client-created user message item and return its ack events.

Replaces the old translator's _send_realtime_input_ack: call it when handling a CreateItem whose item role is user (before the internal conversation.item.created is emitted, which then projects to done).

resolve_cancel_response

resolve_cancel_response(
    state: RealtimeProjectionState, command: CancelResponse
) -> ResolvedControl

resolve_clear_output_audio

resolve_clear_output_audio(
    state: RealtimeProjectionState,
    command: ClearOutputAudio,
) -> ResolvedControl

resolve_commit

resolve_commit(
    state: RealtimeProjectionState, command: Commit
) -> ResolvedCommit

Apply the old translator's commit rules (empty buffer, non-speech commit, item id).

resolve_create_item

resolve_create_item(
    state: RealtimeProjectionState, command: CreateItem
) -> ResolvedControl

Expand conversation.item.create the way the old translator did.

Assistant/system/function items and text-only user items become one turn.signal conversation.item.create payload. User items carrying audio become one input_audio_buffer.append per audio part followed by an input_audio_buffer.commit (no response). The user-item ack events are returned in events and must be emitted before running the payloads.

resolve_delete_item

resolve_delete_item(
    state: RealtimeProjectionState, command: DeleteItem
) -> ResolvedControl

resolve_truncate_item

resolve_truncate_item(
    state: RealtimeProjectionState, command: TruncateItem
) -> ResolvedControl

response_is_done

response_is_done(
    state: RealtimeProjectionState, response_id: object
) -> bool

retrieve_item_events

retrieve_item_events(
    state: RealtimeProjectionState,
    payload: Mapping[str, object],
) -> list[DuplexEvent]

Answer conversation.item.retrieve from the projected conversation items.