Skip to content

vllm_omni.engine.duplex.session.helpers

Small reads and payload builders over one duplex session.

None of these hold state or decide anything: they answer a question about the session, or shape one of the wire payloads its events carry. They were methods on the runner only because that is where the session reference lived, which made every component that needed one a reason to reach back into the runner.

Functions, not a component, because there is nothing here to own.

advance_barge_in_epoch

advance_barge_in_epoch(
    session: DuplexEngineSession,
) -> tuple[int, dict[str, int]]

Start a new epoch for a barge-in, returning it with the playback it cut off.

append_fence

append_fence(
    session: DuplexEngineSession,
    payload: object,
    *,
    epoch: int | None = None,
) -> DuplexFence

The fence one append is submitted under.

A model without a resumable core request gets one stage request per turn, so the turn this resolves to is part of the request id. The runner names that id when it queues the append and the model channel names it again when it submits, and the two must not disagree --- otherwise the session binds a request id it never submitted.

assistant_playback_active

assistant_playback_active(
    session: DuplexEngineSession,
) -> bool

Whether the client is still playing audio the session has sent.

audio_committed_payload

audio_committed_payload(
    session: DuplexEngineSession,
    *,
    committed: DuplexCommittedInput | None = None,
    realtime_item_id: object | None = None,
    transcript: object | None = None,
) -> dict[str, object]

audio_payload_size_bytes

audio_payload_size_bytes(
    payload: Mapping[str, object],
) -> int

barge_in_unsupported_error

barge_in_unsupported_error() -> ErrorEvent

commit_audio_input

commit_audio_input(
    session: DuplexEngineSession,
    *,
    realtime_item_id: object | None = None,
    transcript: object | None = None,
    turn_id: int | None = None,
) -> DuplexCommittedInput

commit_played_response_history

commit_played_response_history(
    session: DuplexEngineSession,
    response_id: str | None,
    committed_ms: int,
) -> None

Truncate an interrupted response in history to what the client actually heard.

input_committed_payload

input_committed_payload(
    session: DuplexEngineSession,
    committed: DuplexCommittedInput,
    *,
    realtime_item_id: object | None = None,
) -> dict[str, object]

next_commit_allowed

next_commit_allowed(
    session: DuplexEngineSession,
    tasks: DuplexSessionTasks,
    *,
    concurrent_turn_requests_released: bool,
) -> bool

Whether a new commit may flush/submit now.

Idle sessions always allow it. When supports_concurrent_turn_requests is on and the plugin has released the commit gate, a commit is allowed even though prior assistant TTS/playback still counts as response_in_progress. Barge-in remains the abort path; this gate does not cancel anything.

overlap_decision_event

overlap_decision_event(
    session: DuplexEngineSession,
    decision: dict[str, object],
) -> OverlapDecision

response_in_progress

response_in_progress(
    session: DuplexEngineSession, tasks: DuplexSessionTasks
) -> bool

Whether the model still owns the turn.

Broader than active_response_id: audio already sent but not yet played back, and an append still in flight, both mean the turn is not free.

stage0_request_id

stage0_request_id(
    session: DuplexEngineSession, epoch: int
) -> str

Stable stage0 or ephemeral stage0-turn{T} request id for Stage0.

task_is_cancelling

task_is_cancelling(task: Task[object] | None) -> bool

True when task has a pending cancellation (Python 3.11+ Task.cancelling).