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())
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.
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.
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.
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).
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
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_response async ¶
create_response(
options: ResponseCreateOptions
| Mapping[str, object]
| None = None,
) -> None
events async ¶
events() -> AsyncIterator[DuplexEvent]
Ordered public events; ends after session.closed / session.expired.
submit async ¶
submit(command: DuplexCommand) -> None
Enqueue one command in caller order; rejections arrive as ErrorEvent on events().