vllm_omni.engine.duplex.session.engine_session ¶
The one duplex session state, owned by the engine-side DuplexSessionRunner.
DuplexEngineSession is the single session object: the input / response / playback / conversation ledgers, the lease, the identity fence, the stage request resources, the append sequencing, the model plugin's per-session state (model_state) and the Realtime projection state (projector). There is no cross-boundary fence protocol; session.fence is derived from the session's own epoch/turn and accepted_fence is the monotonic high-water mark of fences accepted for stage requests (sync_fence() publishes the current identity there).
RESPONSE_REQUEST_MEASUREMENT_ORIGIN module-attribute ¶
RESPONSE_REQUEST_MEASUREMENT_ORIGIN: dict[str, str] = {
"ttft": "accepted native-append start to first non-empty text output; pending turns keep the latest append before first output, so overlapping user speech shortens TTFT",
"ttfp": "accepted native-append start to first audio output; pending turns keep the latest append before first output, so overlapping user speech shortens TTFP",
}
ConversationHistory dataclass ¶
assistant_response_snapshots class-attribute instance-attribute ¶
assistant_response_snapshots: dict[
str,
tuple[
dict[str, object],
tuple[DuplexAssistantAudioTextMark, ...],
int,
],
] = field(default_factory=dict)
hard_truncations_ms class-attribute instance-attribute ¶
history_item_placeholders class-attribute instance-attribute ¶
item_audio_text_marks class-attribute instance-attribute ¶
item_audio_text_marks: dict[
str, list[DuplexAssistantAudioTextMark]
] = field(default_factory=dict)
item_ids class-attribute instance-attribute ¶
last_assistant_audio_text_marks class-attribute instance-attribute ¶
last_assistant_audio_text_marks: list[
DuplexAssistantAudioTextMark
] = field(default_factory=list)
last_assistant_full_message class-attribute instance-attribute ¶
messages class-attribute instance-attribute ¶
pending_item_audio_text_marks class-attribute instance-attribute ¶
pending_item_audio_text_marks: dict[
str, list[DuplexAssistantAudioTextMark]
] = field(default_factory=dict)
pending_item_ids class-attribute instance-attribute ¶
pending_item_input_commit_seqs class-attribute instance-attribute ¶
pending_truncations_ms class-attribute instance-attribute ¶
DuplexAppendReservation dataclass ¶
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
)
assistant_audio_text_marks property ¶
assistant_audio_text_marks: tuple[
DuplexAssistantAudioTextMark, ...
]
capabilities class-attribute instance-attribute ¶
capabilities: DuplexCapabilities = field(
default_factory=DuplexCapabilities
)
created_monotonic class-attribute instance-attribute ¶
last_assistant_audio_text_marks property ¶
last_assistant_audio_text_marks: tuple[
DuplexAssistantAudioTextMark, ...
]
lease class-attribute instance-attribute ¶
lease: DuplexLeaseState = field(
default_factory=_default_lease, repr=False
)
model_state class-attribute instance-attribute ¶
model_state: DuplexModelSessionState | None = field(
default=None, repr=False
)
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.
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 ¶
append_draining_assistant_text ¶
Append text onto a snapshotted draining response, not the active one.
assistant_transcript ¶
Joined assistant text for the active response, or a draining snapshot.
begin_close ¶
Make close irreversible while retaining stage resources for cleanup retry.
bind_draining_request ¶
Map a still-playing draining-stage request onto the response that owns it.
bind_stage_request ¶
bind_stage_request(
stage_id: int, request_id: str, *, fence: DuplexFence
) -> None
cancel_pending_input ¶
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.
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.
draining_request_ids ¶
Request ids whose TTS is still running under a previous response.
end_response ¶
end_response(
*,
commit_text: bool = True,
playback_commit_policy: str | None = None,
preserve_request: bool = False,
) -> dict[str, object] | None
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_model_turn_request_started ¶
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.
notify_new_user_item ¶
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_for_response ¶
playback_for_response(
response_id: str | None = None,
) -> DuplexPlaybackView
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 ¶
release_all_requests ¶
Drop every reserved/submitted stage request; returns their ids for stage cleanup.
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_resources_for_request_ids ¶
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.
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.
reserve_history_item ¶
reserve_history_item(item_id: str) -> None
Reserve the current history position for a response-owned item.
reserve_stage_request ¶
reserve_stage_request(
stage_id: int, request_id: str, *, fence: DuplexFence
) -> None
reset_unanswered_user_items ¶
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]
signal_turn ¶
Apply one external turn signal and return the typed turn.event.
snapshot_active_response_for_drain ¶
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.
stash_stage_metrics ¶
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.
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
DuplexFenceMismatchError ¶
Bases: RuntimeError
DuplexInputAppend dataclass ¶
DuplexRequestResource dataclass ¶
InputBufferState dataclass ¶
PlaybackLedger dataclass ¶
by_response class-attribute instance-attribute ¶
by_response: dict[str, DuplexPlaybackCursor] = field(
default_factory=dict
)
current class-attribute instance-attribute ¶
current: DuplexPlaybackCursor = field(
default_factory=DuplexPlaybackCursor
)
ResponseState dataclass ¶
active_options class-attribute instance-attribute ¶
active_options: ResponseCreateOptions | None = None
active_response_awaits_input_commit class-attribute instance-attribute ¶
active_response_awaits_input_commit: bool = False
active_response_input_commit_seq class-attribute instance-attribute ¶
active_response_input_commit_seq: int | None = None
active_response_request_started_at_s class-attribute instance-attribute ¶
active_response_request_started_at_s: float | None = None
active_response_ttfp_ms class-attribute instance-attribute ¶
active_response_ttfp_ms: float | None = None
active_response_ttft_ms class-attribute instance-attribute ¶
active_response_ttft_ms: float | None = None
active_response_turn_id class-attribute instance-attribute ¶
active_response_turn_id: int | None = None
assistant_audio_text_marks class-attribute instance-attribute ¶
assistant_audio_text_marks: list[
DuplexAssistantAudioTextMark
] = field(default_factory=list)
assistant_text_buffer class-attribute instance-attribute ¶
draining_response_by_request class-attribute instance-attribute ¶
pending_options class-attribute instance-attribute ¶
pending_options: ResponseCreateOptions | None = None