vllm_omni.engine.duplex.session.manager ¶
Engine-side owner of duplex sessions: admission, lease reaping, command dispatch.
DuplexOrchestrator hosts exactly one manager. The whole session (DuplexEngineSession) lives here, owned by one DuplexSessionRunner per session on the orchestrator loop; nothing outside the engine holds session state.
Control operations (open / close / resume / touch) arrive as correlated RPC messages and answer through result_sink; session commands arrive one-way as DuplexSessionCommandMessage and are pushed onto the runner's ordered mailbox after backpressure admission; everything a session emits leaves through output_sink as DuplexSessionEventMessage.
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