Skip to content

vllm_omni.entrypoints.duplex_omni

DuplexOmni: the Python API for full-duplex models.

Sessions run inside the engine (DuplexSessionRunner on the orchestrator loop of DuplexOrchestrator). This class opens / resumes / closes them and pipes typed DuplexCommand objects in and typed DuplexEvent objects out through a DuplexSessionHandle. It holds no session state beyond the handle registry, so the websocket handler and InlineDuplexClient are both thin consumers of the same surface.

Example::

omni = DuplexOmni(model="openbmb/MiniCPM-o-4_5", trust_remote_code=True)
async with await omni.open_session({"modalities": ["audio", "text"]}) as session:
    async def consume():
        async for event in session.events():
            if isinstance(event, AudioDelta):
                play(event.audio)
    task = asyncio.create_task(consume())
    await session.submit(AppendAudio(audio=pcm_bytes, format="pcm16", sample_rate_hz=16000))
    await session.submit(Commit())

logger module-attribute

logger = init_logger(__name__)

DuplexOmni

Bases: AsyncOmni

Async Python API for full-duplex models (see module docstring).

Construct it like AsyncOmni. The pipeline must declare duplex_plugin and the deploy config session_mode: duplex; the engine raises at startup otherwise. There is no generate(): duplex serving is session-only.

duplex_capabilities property

duplex_capabilities: DuplexCapabilities

duplex_session_config property

duplex_session_config: DuplexSessionRuntimeConfig

engine instance-attribute

sessions property

active_session_count

active_session_count() -> int

close_all_sessions async

close_all_sessions(*, reason: str = 'shutdown') -> None

close_session async

close_session(
    session_id: str,
    *,
    reason: str = "client_close",
    timeout: float | None = _DEFAULT_CONTROL_TIMEOUT_S,
) -> None

detach_session async

detach_session(
    session_id: str,
    *,
    expected_lease_generation: int | None = None,
    timeout: float | None = _DEFAULT_CONTROL_TIMEOUT_S,
) -> None

Start the engine-owned disconnect grace; expiry arrives as session.expired.

expected_lease_generation names the lease the caller opened or resumed against; the engine refuses to detach a newer one, so a connection giving up its lease cannot start the grace for the lease a later resume owns.

get_session

get_session(session_id: str) -> DuplexSessionHandle | None

open_session async

open_session(
    config: DuplexSessionConfig
    | Mapping[str, object]
    | None = None,
    *,
    timeout: float | None = _DEFAULT_CONTROL_TIMEOUT_S,
) -> DuplexSessionHandle

Open an engine-resident session and return its handle.

The session id is always allocated here (duplex-<uuid4 hex>; never reused, so the id alone identifies a session) and the handle is registered before the open RPC, so the first event (session.created, which announces the id) can never arrive before it exists.

resume_session async

resume_session(
    session_id: str,
    *,
    expected_lease_generation: int,
    timeout: float | None = _DEFAULT_CONTROL_TIMEOUT_S,
    on_abandoned: Callable[[int], Awaitable[None]]
    | None = None,
) -> DuplexSessionHandle

Engine lease resume (CAS on the lease generation); returns the existing handle.

The RPC runs to completion in an executor thread whatever happens to this waiter, so a caller cancelled while it is in flight does not stop the engine from applying the resume (lease generation bumped, disconnect grace cleared) for a connection that will never serve it. The outcome is then observed off the cancelled task. If the RPC times out before answering, the outcome is unknown, so the same resume is replayed under its control id: the engine answers a replay with the generation the resume produced instead of resuming again. A resume that landed is adopted into the handle (a later reconnect resumes against the real generation) and settled through on_abandoned with that generation; without a callback it is detached again, fenced on the generation, so a resume that came after it is never touched.

on_abandoned is for callers that own an attachment concept: they decide whether the generation now belongs to a connection that is still serving (a takeover that never activated leaves the previous socket attached) or whether the lease goes back into disconnect grace.

shutdown

shutdown(timeout: float | None = None) -> None

touch_session async

touch_session(
    session_id: str,
    *,
    activity: str = "heartbeat",
    expected_lease_generation: int | None = None,
    timeout: float | None = _DEFAULT_CONTROL_TIMEOUT_S,
) -> None

DuplexSessionHandle

Client-side view of one engine-resident duplex session.

Connection-independent: a websocket may detach and a later one resume the same handle. events() is single-consumer at any given time but may be re-entered after the previous iterator was closed (resume).

capabilities instance-attribute

close_reason property

close_reason: str | None

closed property

closed: bool

lease_generation instance-attribute

lease_generation: int = 0

public_session instance-attribute

public_session: dict[str, object] = {}

session_id instance-attribute

session_id = session_id

ack_playback async

ack_playback(
    played_ms: int,
    *,
    response_id: str | None = None,
    item_id: str | None = None,
    committed_ms: int | None = None,
) -> None

append_audio async

append_audio(
    audio: bytes,
    *,
    format: str = "pcm16",
    sample_rate_hz: int | None = None,
    is_speech: bool | None = None,
    video_frames: Sequence[str] | None = None,
    duration_ms: int | None = None,
    audio_end_ms: int | None = None,
    hints: Mapping[str, object] | None = None,
) -> None

append_text async

append_text(text: str) -> None

barge_in async

barge_in() -> None

cancel_input async

cancel_input() -> None

cancel_response async

cancel_response(response_id: str | None = None) -> None

clear_input async

clear_input() -> None

clear_output_audio async

clear_output_audio(response_id: str | None = None) -> None

close async

close(
    *,
    reason: str = "client_close",
    timeout: float | None = _DEFAULT_CONTROL_TIMEOUT_S,
) -> None

commit async

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

create_item async

create_item(
    item: Mapping[str, object],
    *,
    previous_item_id: str | None = None,
) -> None

create_response async

create_response(
    options: ResponseCreateOptions
    | Mapping[str, object]
    | None = None,
) -> None

delete_item async

delete_item(item_id: str) -> None

events async

Ordered public events; ends after session.closed / session.expired.

heartbeat async

heartbeat() -> None

signal_turn async

signal_turn(
    event: str, payload: Mapping[str, object] | None = None
) -> None

submit async

submit(command: DuplexCommand) -> None

Enqueue one command in caller order; rejections arrive as ErrorEvent on events().

truncate_item async

truncate_item(
    item_id: str,
    *,
    audio_end_ms: int,
    content_index: int = 0,
) -> None

update async

update(session_patch: Mapping[str, object]) -> None

wait_closed async

wait_closed() -> str