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).
conversation_items class-attribute instance-attribute ¶
default_payload class-attribute instance-attribute ¶
defaults class-attribute instance-attribute ¶
defaults: RealtimeInputDefaults = field(
default_factory=RealtimeInputDefaults
)
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 ¶
item_truncation_cursors class-attribute instance-attribute ¶
last_conversation_item_id class-attribute instance-attribute ¶
last_conversation_item_id: str | None = None
pending_commit_item_ids class-attribute instance-attribute ¶
response_states class-attribute instance-attribute ¶
ResolvedCommit dataclass ¶
ResolvedControl dataclass ¶
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
retrieve_item_events ¶
retrieve_item_events(
state: RealtimeProjectionState,
payload: Mapping[str, object],
) -> list[DuplexEvent]
Answer conversation.item.retrieve from the projected conversation items.