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: |
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 |
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 |
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
)
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
DuplexSessionManager ¶
executor instance-attribute ¶
executor = (
executor
or concurrent.futures.ThreadPoolExecutor(
max_workers=4, thread_name_prefix="duplex-session"
)
)
vad_backend_provider instance-attribute ¶
vad_backend_provider = SileroVADBackendProvider(
model_path=getattr(
runtime_config, "server_vad_model_path", 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]]
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.
stage_request_id staticmethod ¶
stage_request_id(
fence: DuplexFence,
*,
stage_id: int,
resumable: bool = True,
) -> str
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,
)
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,
)
out instance-attribute ¶
out = SessionEmitter(
self.ctx,
promote_deferred_overlap=self._promote_deferred_overlap_later,
)
close async ¶
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 ¶
Apply an internal event to the session, project it, and send it.
expire async ¶
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 ¶
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 ¶
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.