Skip to content

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.

logger module-attribute

logger = init_logger(__name__)

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