Skip to content

vllm_omni.engine.duplex.session.context

Shared state for :class:DuplexSessionRunner and its components.

The runner is a state machine that several components read and write. Splitting it into modules only helps if the shared part stops being implicit: a component holding a back-reference to the runner is a mixin with extra steps, and one holding its own copy of a flag diverges from the others.

So the mutable part is named. DuplexRunState is the small set of flags that more than one component touches -- grep "\.run\." finds every mutation -- and DuplexSessionContext is the read-mostly collaborator bundle everything is constructed with. Anything a component needs from the runner that is neither of those is infrastructure, and goes through :class:RunnerServices.

DuplexAppendTaskMeta dataclass

epoch instance-attribute

epoch: int

final instance-attribute

final: bool

response_bound instance-attribute

response_bound: bool

DuplexRunState dataclass

Mutable per-run flags shared by the runner and its components.

Every field here is written by one component and read by another; that is the criterion for being in this object rather than private to a component.

close_reason class-attribute instance-attribute

close_reason: str | None = None

closed_deferred class-attribute instance-attribute

closed_deferred: bool = False

closed_emitted class-attribute instance-attribute

closed_emitted: bool = False

closing class-attribute instance-attribute

closing: bool = False

concurrent_turn_requests_released class-attribute instance-attribute

concurrent_turn_requests_released: bool = False

runtime_closed class-attribute instance-attribute

runtime_closed: bool = False

stream_request_id class-attribute instance-attribute

stream_request_id: str | None = None

DuplexSessionContext dataclass

Collaborators one session's components are built with.

Read-mostly: the session and the model state are mutated through their own APIs, not by rebinding these fields. The mutable flags live in run.

manager instance-attribute

model_state instance-attribute

plugin instance-attribute

run instance-attribute

services instance-attribute

services: RunnerServices

session instance-attribute

stage_port instance-attribute

stage_port: DuplexStagePort

tasks instance-attribute

DuplexSessionTasks dataclass

Tracked task handles of one session (append tail, active response, pending silence).

active_response_task class-attribute instance-attribute

active_response_task: Task[None] | None = None

append_tail class-attribute instance-attribute

append_tail: Task[bool] | None = None

append_tasks class-attribute instance-attribute

append_tasks: dict[Task[bool], DuplexAppendTaskMeta] = (
    field(default_factory=dict)
)

cancel_append_tasks async

cancel_append_tasks(
    timeout_s: float = 0.25,
    *,
    response_bound_only: bool = False,
) -> bool

has_response_bound_append_tasks

has_response_bound_append_tasks() -> bool

track_append_task

track_append_task(
    task: Task[bool],
    *,
    epoch: int,
    final: bool,
    response_bound: bool,
) -> None

RunnerServices

Bases: Protocol

The only things a component may ask the runner for.

Deliberately two methods. A component that needs more than task scheduling from the runner is reaching for orchestration that belongs in the runner.

offload async

offload(
    fn: Callable[..., _OffloadT],
    *args: object,
    **kwargs: object,
) -> _OffloadT

Run a blocking call off the orchestrator loop.

spawn

spawn(coro: Awaitable[None], *, name: str) -> None

Run coro as a tracked background task on the session's loop.

StageOutput dataclass

One stage result, queued for the session that owns the request.

context instance-attribute

decision instance-attribute

decision: DuplexOutputDecision | None

metrics instance-attribute

metrics: StageRequestStats | None

output instance-attribute

output: RequestOutput

request_id instance-attribute

request_id: str

stage_id instance-attribute

stage_id: int