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,
    *,
    timeout: float | None = _DEFAULT_CONTROL_TIMEOUT_S,
) -> None

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

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,
) -> DuplexSessionHandle

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

shutdown

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

touch_session async

touch_session(
    session_id: str,
    *,
    activity: str = "heartbeat",
    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