Skip to content

vllm_omni.engine.duplex.session

One engine-resident duplex session, and everything that runs it.

DuplexEngineSession (engine_session) is the state -- the input, response, playback and conversation ledgers, the lease and the identity fence. DuplexSessionRunner (runner) owns one of those on the orchestrator loop and is the only writer of it; DuplexSessionManager (manager) owns the runners and the admission slots.

The rest of the package is what the runner was split into, and each module names the thing it owns rather than a layer:

  • context -- the state a runner and its components share, and the only services a component may ask the runner for.
  • emitter -- everything the session sends: the Realtime projection, the epoch filter, the domain effects a terminal event applies.
  • model_channel -- the model side: submitting an append, projecting the stage output that comes back, continuing a turn with silence.
  • control -- server VAD and the events that reconfigure the session.
  • append_task -- one append in flight, and the rollback its failure owes.
  • helpers -- the small pure reads and payload builders over a session.
  • overlap_policy / commit_policy / playback_ledger -- the decisions a session makes about speech that arrives while the model is speaking, what a commit should do, and what the client has actually played.
  • lease -- the idle/activity lease that decides when a session expires.

Imports here are the package's public surface; a module of this package is fair game for anything inside it, but code outside engine.duplex should reach for these three names.

Modules:

Name Description
append_task

One append in flight, and what its failure owes the session.

commit_policy
context

Shared state for :class:DuplexSessionRunner and its components.

control

Control of one duplex session: server VAD, and the events that reconfigure it.

emitter

Everything one duplex session sends to its client.

engine_session

The one duplex session state, owned by the engine-side DuplexSessionRunner.

helpers

Small reads and payload builders over one duplex session.

lease

Generic engine lease primitives shared by optional session runtimes.

manager

Engine-side owner of duplex sessions: admission, lease reaping, command dispatch.

model_channel

The model side of one duplex session.

overlap_policy

How one duplex session reacts to input that arrives while it is speaking.

playback_ledger

Applying one playback.ack to a duplex session.

runner

Engine-resident session runner: one ordered mailbox and one session state per duplex session.

DuplexEngineSession dataclass

The one session state (see module docstring).

Owned and mutated only by its DuplexSessionRunner on the orchestrator loop.

accepted_fence class-attribute instance-attribute

accepted_fence: DuplexFence = field(
    default=None, repr=False
)

active_request_id property

active_request_id: str | None

active_response_id property

active_response_id: str | None

active_response_turn_id property

active_response_turn_id: int | None

assistant_audio_text_marks property

assistant_audio_text_marks: tuple[
    DuplexAssistantAudioTextMark, ...
]

assistant_text_buffer property

assistant_text_buffer: tuple[str, ...]

capabilities class-attribute instance-attribute

capabilities: DuplexCapabilities = field(
    default_factory=DuplexCapabilities
)

config instance-attribute

config_generation class-attribute instance-attribute

config_generation: int = 0

created_monotonic class-attribute instance-attribute

created_monotonic: float = field(
    default_factory=time.monotonic
)

epoch class-attribute instance-attribute

epoch: int = 0

fence property

fence: DuplexFence

history property

history: tuple[dict[str, object], ...]

history_item_ids property

history_item_ids: Mapping[str, dict[str, object]]

input_commit_seq property

input_commit_seq: int

input_seq class-attribute instance-attribute

input_seq: int = 0

input_turn_seq class-attribute instance-attribute

input_turn_seq: int = 0

last_assistant_audio_text_marks property

last_assistant_audio_text_marks: tuple[
    DuplexAssistantAudioTextMark, ...
]

last_assistant_full_message property

last_assistant_full_message: dict[str, object] | None

last_response_id property

last_response_id: str | None

lease class-attribute instance-attribute

lease: DuplexLeaseState = field(
    default_factory=_default_lease, repr=False
)

lease_generation property

lease_generation: int

log_stats class-attribute instance-attribute

log_stats: bool = False

model_state class-attribute instance-attribute

model_state: DuplexModelSessionState | None = field(
    default=None, repr=False
)

num_stages class-attribute instance-attribute

num_stages: int = 1

overlap_speech_ms property

overlap_speech_ms: int

pending_history_item_ids property

pending_history_item_ids: Mapping[str, dict[str, object]]

pending_history_truncations_ms property

pending_history_truncations_ms: Mapping[str, int]

pending_input_bytes property

pending_input_bytes: int

pending_input_turns property

pending_input_turns: int

playback property

projector class-attribute instance-attribute

projector: RealtimeProjectionState | None = field(
    default=None, repr=False
)

request_resources class-attribute instance-attribute

request_resources: dict[
    tuple[int, str], DuplexRequestResource
] = field(default_factory=dict, repr=False)

response_config property

response_config: DuplexSessionConfig

Return the immutable-for-this-lifecycle response configuration.

config remains the session defaults. A response takes a deep snapshot when it begins so response.create overrides and concurrent session updates cannot mutate each other's ownership domains.

runtime_config property

runtime_config: Mapping[str, object]

session_id instance-attribute

session_id: str

state class-attribute instance-attribute

turn_id class-attribute instance-attribute

turn_id: int = 0

turn_state class-attribute instance-attribute

accept_fence

accept_fence(fence: DuplexFence) -> None

accumulate_overlap_speech

accumulate_overlap_speech(duration_ms: int) -> int

accumulate_response_stage_metrics

accumulate_response_stage_metrics(
    stage_metrics: Mapping[Any, Any] | None,
) -> dict[str, dict[str, object]]

acknowledge_playback

acknowledge_playback(
    played_ms: int,
    committed_ms: int | None = None,
    *,
    response_id: str | None = None,
) -> DuplexPlaybackView

active_response_accepts_model_turn

active_response_accepts_model_turn(
    turn_id: int | None,
) -> bool

append_assistant_text

append_assistant_text(text: str) -> None

append_draining_assistant_text

append_draining_assistant_text(
    response_id: str, text: str
) -> None

Append text onto a snapshotted draining response, not the active one.

append_history_message

append_history_message(message: dict[str, object]) -> None

as_public_dict

as_public_dict() -> dict[str, object]

assistant_transcript

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

Joined assistant text for the active response, or a draining snapshot.

barge_in

barge_in() -> int

begin_close

begin_close(*, reason: str) -> bool

Make close irreversible while retaining stage resources for cleanup retry.

begin_lease_operation

begin_lease_operation(
    fence: DuplexFence, operation_id: str
) -> None

begin_response

begin_response(*, turn_id: int | None = None) -> str

bind_draining_request

bind_draining_request(
    request_id: str, response_id: str
) -> None

Map a still-playing draining-stage request onto the response that owns it.

bind_request

bind_request(request_id: str | None) -> None

bind_response_turn

bind_response_turn(turn_id: int | None) -> None

bind_stage_request

bind_stage_request(
    stage_id: int, request_id: str, *, fence: DuplexFence
) -> None

cancel_fence

cancel_fence(
    cancelled_fence: DuplexFence, next_fence: DuplexFence
) -> list[str]

cancel_pending_input

cancel_pending_input() -> dict[str, int]

Drop every pending reservation; the PCM buffer itself is model state cleared by the runner.

clear_draining_for_response

clear_draining_for_response(
    response_id: str | None,
) -> None

Drop draining bindings owned by response_id; leave other responses intact.

clear_draining_requests

clear_draining_requests() -> None

clear_playback_cursor

clear_playback_cursor() -> None

clear_request

clear_request(
    expected_request_id: str | None = None,
) -> bool

close

close() -> None

commit_append

commit_append(
    reservation: DuplexAppendReservation,
) -> DuplexInputAppend

commit_audio_input

commit_audio_input(
    *,
    transcript: str | None = None,
    turn_id: int | None = None,
) -> DuplexCommittedInput

complete_model_turn

complete_model_turn(turn_id: int) -> None

Advance the model-owned output identity after its terminal signal.

delete_history_item

delete_history_item(item_id: str) -> bool

detach_lease

detach_lease() -> None

discard_response_options

discard_response_options() -> None

draining_request_ids

draining_request_ids() -> list[str]

Request ids whose TTS is still running under a previous response.

end_lease_operation

end_lease_operation(operation_id: str) -> None

end_response

end_response(
    *,
    commit_text: bool = True,
    playback_commit_policy: str | None = None,
    preserve_request: bool = False,
) -> dict[str, object] | None

has_assistant_response_item

has_assistant_response_item(
    response_id: str, item_id: str
) -> bool

is_draining_request

is_draining_request(request_id: str | None) -> bool

mark_audio_sent

mark_audio_sent(
    duration_ms: int | None = None,
    *,
    text_chars: int | None = None,
    audio_text_marks: list[dict[str, object]] | None = None,
    text_requires_complete_audio: bool = False,
    audio_complete: bool = False,
    response_id: str | None = None,
) -> None

mark_closing

mark_closing() -> None

mark_model_turn_request_started

mark_model_turn_request_started(
    turn_id: int, started_at_s: float
) -> None

Record the native request start that can own one model turn.

Pending turns keep the latest accepted append. After begin_response the first accepted append wins; later appends must not rebind.

mark_response_first_outputs

mark_response_first_outputs(
    *, observed_at_s: float, has_text: bool, has_audio: bool
) -> dict[str, object]

Return server-monotonic TTF metrics newly observed for the active response.

mark_user_input_activity

mark_user_input_activity() -> None

notify_new_user_item

notify_new_user_item() -> None

Record a user item that a later response.create may answer.

observe_stage_request_stats

observe_stage_request_stats(
    stage_id: int, metrics: StageRequestStats
) -> None

Feed one engine StageRequestStats into the response's logger table.

playback_ack_is_too_late

playback_ack_is_too_late(
    response_id: str, item_id: str
) -> bool

playback_for_response

playback_for_response(
    response_id: str | None = None,
) -> DuplexPlaybackView

pop_draining_request

pop_draining_request(request_id: str | None) -> str | None

prepare_append

prepare_append(
    fence: DuplexFence,
) -> DuplexAppendReservation

prepare_cancel_fence

prepare_cancel_fence(
    cancelled_fence: DuplexFence, next_fence: DuplexFence
) -> list[str]

Advance the cancellation fence without dropping cleanup records.

register_history_item

register_history_item(
    item_id: str | None, message: dict[str, object] | None
) -> None

release_all_input_bytes

release_all_input_bytes() -> None

release_all_requests

release_all_requests() -> list[str]

Drop every reserved/submitted stage request; returns their ids for stage cleanup.

release_fence

release_fence(fence: DuplexFence) -> list[str]

release_finished_drain_response

release_finished_drain_response(
    response_id: str | None,
) -> None

Drop unused books for a response whose TTS just finished.

Overlap leaves the next turn active, so this must not call end_response. Audio that was already sent stays ACK-admissible: the client may playback.ack after the drain response.done. A response that sent no audio has nothing to acknowledge, so its snapshot, playback cursor, and unused history placeholder are dropped. The live response is left untouched.

release_input_bytes

release_input_bytes(size: int) -> None

release_pending_turn

release_pending_turn() -> None

release_resources_for_request_ids

release_resources_for_request_ids(
    request_ids: Iterable[str],
) -> list[str]

Drop bindings for these request ids, whichever fence they sit on.

cancel_fence only releases the fence being cancelled. Overlapped draining output stages belong to an older turn fence and would otherwise stay until the session closes.

release_response_history_snapshot

release_response_history_snapshot(
    response_id: str | None,
) -> None

Release a final response snapshot after its playback fully commits.

release_response_playback

release_response_playback(response_id: str | None) -> None

replace_config

replace_config(config: DuplexSessionConfig) -> None

replace_response_stage_metric_snapshots

replace_response_stage_metric_snapshots(
    stage_metrics: Mapping[Any, Any] | None,
) -> dict[str, dict[str, object]]

Merge cumulative chat snapshots by replacing each stage's latest value.

replace_runtime_config

replace_runtime_config(
    runtime_config: Mapping[str, object],
) -> None

reserve_history_item

reserve_history_item(item_id: str) -> None

Reserve the current history position for a response-owned item.

reserve_input_bytes

reserve_input_bytes(size: int, *, limit: int) -> bool

reserve_pending_turn

reserve_pending_turn(*, limit: int) -> bool

reserve_response_options

reserve_response_options(
    options: ResponseCreateOptions,
) -> None

reserve_stage_request

reserve_stage_request(
    stage_id: int, request_id: str, *, fence: DuplexFence
) -> None

reset_overlap_speech

reset_overlap_speech() -> int

reset_unanswered_user_items

reset_unanswered_user_items() -> None

A turn has started: the pending items are now its input.

resource_request_ids

resource_request_ids(
    fence: DuplexFence | None = None,
    *,
    submitted: bool | None = None,
) -> list[str]

response_has_draining_request

response_has_draining_request(
    response_id: str | None,
) -> bool

response_id_for_request

response_id_for_request(
    request_id: str | None,
) -> str | None

resume_lease

resume_lease(*, expected_lease_generation: int) -> int

signal_turn

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

Apply one external turn signal and return the typed turn.event.

snapshot_active_response_for_drain

snapshot_active_response_for_drain() -> None

Keep the active response ACK-admissible after a later response starts.

Overlapped commit calls begin_response while prior TTS is still draining. That clears the live text buffer and playback cursor. end_response is the normal snapshot, but it would also close the response. Copy the same snapshot (and reserve the history slot) first so a playback ACK for the draining response is not playback_item_not_found.

stage_request_submitted

stage_request_submitted(
    stage_id: int, request_id: str
) -> bool

stash_stage_metrics

stash_stage_metrics(
    stage_metrics: Mapping[Any, Any] | None,
) -> None

Hold a stage snapshot until a response exists to attribute it to.

A stage that hands its output to the next stage instead of the client (stage 0 feeding the TTS stage) reports its token metrics before the response carrying those tokens is created. Dropping them would leave every response without an engine-side token count; holding them keeps the count whole, at the cost of attributing a turn the model never spoke to the next response it does speak.

sync_fence

sync_fence() -> DuplexFence

Publish the current (epoch, turn_id) identity as the accepted fence.

touch_lease

touch_lease(activity: DuplexLeaseActivity) -> None

transition_session

transition_session(state: DuplexSessionState) -> None

transition_turn

transition_turn(state: DuplexTurnState) -> None

truncate_history_item

truncate_history_item(
    item_id: str,
    *,
    audio_end_ms: int,
    playback: DuplexPlaybackCursor
    | DuplexPlaybackView
    | None = None,
    hard: bool = False,
) -> bool

truncate_playback_commit

truncate_playback_commit(
    committed_ms: int, *, response_id: str | None = None
) -> DuplexPlaybackView

unanswered_user_items

unanswered_user_items() -> int

How many user items are waiting for a response.

DuplexSessionManager

executor instance-attribute

executor = (
    executor
    or concurrent.futures.ThreadPoolExecutor(
        max_workers=4, thread_name_prefix="duplex-session"
    )
)

log_stats instance-attribute

log_stats = bool(log_stats)

model_config instance-attribute

model_config = model_config

plugin instance-attribute

plugin = plugin

runners instance-attribute

runners: dict[str, DuplexSessionRunner] = {}

runtime_config instance-attribute

runtime_config = runtime_config

stage_port instance-attribute

stage_port = stage_port

vad_backend_provider instance-attribute

vad_backend_provider = SileroVADBackendProvider(
    model_path=getattr(
        runtime_config, "server_vad_model_path", None
    )
)

accepts

accepts(message: object) -> bool

active_count

active_count() -> int

close async

close(message: CloseDuplexSessionMessage) -> None

close_from_runner

close_from_runner(
    runner: DuplexSessionRunner, reason: str
) -> None

Finish a close the runner started from its command stream (session.close).

Runs on the orchestrator loop as a tracked task (the runner's worker cannot await its own teardown): stage requests are aborted, the admission slot is released and session.closed is emitted last.

close_sessions_for_request_ids

close_sessions_for_request_ids(
    request_ids: list[str],
    *,
    abort: bool = False,
    cleanup_in_progress: bool = False,
) -> dict[str, list[str]]

defer_request_cleanups

defer_request_cleanups(session_ids: Iterable[str]) -> None

dispatch

dispatch(message: object) -> None

Route one engine request-queue message without blocking the request handler.

emit

emit(
    session: DuplexEngineSession, event: DuplexEvent
) -> None

Bind the session identity to one typed event and push it to the engine output queue.

ensure_stage_request

ensure_stage_request(
    session: DuplexEngineSession,
    *,
    stage_id: int,
    fence: DuplexFence | None = None,
) -> DuplexStageRequestContext | None

Reserve the session's stage request id and register it with the stage port.

finalize_closed_sessions

finalize_closed_sessions(
    session_ids: Iterable[str],
) -> None

get

get(session_id: str) -> DuplexEngineSession | None

handle async

handle(message: object) -> None

open async

open(message: OpenDuplexSessionMessage) -> None

reap_expired async

reap_expired(now: float | None = None) -> int

reaper_loop async

reaper_loop(shutdown_event: Event) -> None

register_request

register_request(request_id: str, session_id: str) -> None

resume async

resume(message: ResumeDuplexSessionMessage) -> None

runner_for_request_id

runner_for_request_id(
    request_id: str,
) -> DuplexSessionRunner | None

sampling_params_for

sampling_params_for(
    session: DuplexEngineSession,
) -> tuple[object, ...]

sessions

sessions() -> dict[str, DuplexEngineSession]

shutdown async

shutdown() -> None

stage_request_id staticmethod

stage_request_id(
    fence: DuplexFence,
    *,
    stage_id: int,
    resumable: bool = True,
) -> str

touch async

touch(message: TouchDuplexSessionMessage) -> None

unregister_request

unregister_request(request_id: str) -> None

DuplexSessionRunner

Owns one DuplexEngineSession on the orchestrator loop (see module docstring).

closed_emitted property

closed_emitted: bool

Whether session.closed / session.expired already left this runner.

closing property

closing: bool

Whether an irreversible close has begun (commands and control ops are refused).

control instance-attribute

control = SessionControl(
    self.ctx,
    self.out,
    self.model,
    wait_for_append_tail=self._wait_for_append_tail,
)

ctx instance-attribute

ctx = DuplexSessionContext(
    session=session,
    model_state=self.model_state,
    plugin=plugin,
    stage_port=stage_port,
    manager=manager,
    tasks=self.tasks,
    run=self.run,
    services=self,
)

manager instance-attribute

manager = manager

model instance-attribute

model = ModelChannel(
    self.ctx,
    self.out,
    close_from_runtime=self._close_from_runtime,
    schedule_silence_continuation=self._schedule_silence_continuation,
    abort_request=self._abort_request_background,
)

model_config instance-attribute

model_config = model_config

model_state instance-attribute

model_state: DuplexModelSessionState = session.model_state

out instance-attribute

out = SessionEmitter(
    self.ctx,
    promote_deferred_overlap=self._promote_deferred_overlap_later,
)

plugin instance-attribute

plugin = plugin

run instance-attribute

session instance-attribute

session = session

stage_port instance-attribute

stage_port = stage_port

tasks instance-attribute

close async

close(reason: str, *, emit_closed: bool = True) -> None

Graceful close: cancel work, release the data plane, emit session.closed.

With emit_closed=False the manager emits session.closed itself once the stage resources are released, so the event also means "the admission slot is free again".

emit

emit(payload: dict[str, object]) -> None

Apply an internal event to the session, project it, and send it.

expire async

expire(reason: str, *, emit_expired: bool = True) -> None

Lease expiry / runtime cleanup: emit session.expired and tear down.

With emit_expired=False the manager emits the event after the stage cleanup (see close).

mark_closed_emitted

mark_closed_emitted() -> None

Claim the session's one terminal event for the caller.

The manager emits the deferred terminal itself, after the stage cleanup; recording it here stops a late runtime close emitting a second one.

offload async

offload(
    fn: Callable[..., _OffloadT],
    *args: object,
    **kwargs: object,
) -> _OffloadT

on_stage_failure

on_stage_failure(
    stage_id: int,
    exc: BaseException,
    *,
    request_id: str | None = None,
) -> None

A stage rejected this session's request: fail the owning response.

Under concurrent turn requests the failing request may belong to a draining older response; resolve via response_id_for_request before falling back to active_response_id.

Runs synchronously on the loop (no mailbox hop): the orchestrator expires the session right after this call, so a queued item could be cancelled with the worker and the client would only see session.expired. Emitting here keeps the order error -> failed response.done -> session.expired.

on_stage_output

on_stage_output(
    stage_id: int,
    output: RequestOutput,
    metrics: StageRequestStats | None,
    *,
    request_id: str,
    context: DuplexOutputContext,
) -> bool

Accept one stage output (orchestrator loop); return True when it must not be forwarded.

shutdown async

shutdown() -> None

spawn

spawn(coro: Awaitable[None], *, name: str) -> None

start

start() -> None

submit

submit(command: DuplexCommand) -> None