Skip to content

vllm_omni.entrypoints.duplex.realtime_input

Thin OpenAI Realtime wire envelope for one websocket connection.

Everything that needed session state now lives engine-side (vllm_omni.engine.duplex.realtime_commands for the command mapping, vllm_omni.engine.duplex.realtime_events for the projection state), and the model-agnostic parsing both of those build on lives in vllm_omni.protocol.realtime. What is left here is the per-connection handshake policy (query-param defaults, autostart / resume-only rules, session.resume parsing), the wire defaults used to translate appends, and error rendering through the typed :class:~vllm_omni.engine.duplex.events.ErrorEvent.

ENVELOPE_EVENT_TYPES module-attribute

ENVELOPE_EVENT_TYPES = frozenset(
    {"session.resume", "session.event_ack"}
)

RealtimeEnvelope dataclass

Per-connection Realtime wire policy (no session state).

autostarted_default_session class-attribute instance-attribute

autostarted_default_session: bool = False

default_model class-attribute instance-attribute

default_model: str | None = None

defaults class-attribute instance-attribute

defaults: RealtimeInputDefaults = field(
    default_factory=RealtimeInputDefaults
)

opened class-attribute instance-attribute

opened: bool = False

resume_only class-attribute instance-attribute

resume_only: bool = False

command_error_payload staticmethod

command_error_payload(
    exc: DuplexCommandError,
) -> dict[str, object]

default_session_payload

default_session_payload() -> dict[str, object]

error_payload staticmethod

error_payload(
    code: str,
    message: str,
    *,
    event_id: object | None = None,
    param: object | None = None,
) -> dict[str, object]

Wire JSON of a transport-level error (derived from the typed ErrorEvent).

first_message

first_message(
    payload: Mapping[str, object],
) -> RealtimeHandshake

Classify the first client message (call only while not opened).

session.update opens with its session object; session.resume resumes; any other event autostarts the default session and is then treated as a command (pending_command_payload).

from_query_params classmethod

from_query_params(
    query_params: Mapping[str, str] | WebSocket,
) -> RealtimeEnvelope

initial_open_payload

initial_open_payload() -> dict[str, object] | None

Session object to open with before any client message (?model= autostart), else None.

is_envelope_event

is_envelope_event(payload: Mapping[str, object]) -> bool

note_session_payload

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

Track wire defaults declared by a session object (open or session.update).

translate

translate(payload: Mapping[str, object]) -> DuplexCommand

Wire event -> command (raises :class:DuplexCommandError).

RealtimeHandshake dataclass

What the first client message asks for.

kind instance-attribute

kind: str

pending_command_payload class-attribute instance-attribute

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

resume_payload class-attribute instance-attribute

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

session_payload class-attribute instance-attribute

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

ResumeRequest dataclass

A validated session.resume client event.

last_received_server_event_seq instance-attribute

last_received_server_event_seq: int

resume_token instance-attribute

resume_token: str

session_id instance-attribute

session_id: str

parse_resume_request

parse_resume_request(
    event: Mapping[str, object],
) -> ResumeRequest | None

Validate the shape of session.resume; None when any field is missing or malformed.