Skip to content

vllm_omni.engine.duplex.session.runner

Engine-resident session runner: one ordered mailbox and one session state per duplex session.

DuplexSessionRunner owns the whole lifecycle of one session on the orchestrator loop (the session is never touched from another thread):

  • every mutation of DuplexEngineSession happens on this loop, through the mailbox worker (commands, stage outputs, internal items) or through tracked tasks that re-validate (epoch, turn_id) after each await;
  • appends are planned with the model plugin and submitted to the stage port in wire order on the per-session append tail (no RPC hop);
  • stage outputs are pushed in by DuplexOrchestrator._intercept_stage_output instead of being polled through request queues;
  • everything the session says leaves as typed events through DuplexSessionManager.emit after terminal-acceptance / stale-epoch filtering, so a cancelled epoch can never speak again.

logger module-attribute

logger = init_logger(__name__)

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,
)

manager instance-attribute

manager = manager

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,
)

model_config instance-attribute

model_config = model_config

model_state instance-attribute

model_state: DuplexModelSessionState = session.model_state

out instance-attribute

out = SessionEmitter(
    self.ctx,
    promote_deferred_overlap=self._promote_deferred_overlap_later,
)

plugin instance-attribute

plugin = plugin

run instance-attribute

session instance-attribute

session = session

stage_port instance-attribute

stage_port = stage_port

tasks instance-attribute

close async

close(reason: str, *, emit_closed: bool = True) -> None

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

emit(payload: dict[str, object]) -> None

Apply an internal event to the session, project it, and send it.

expire async

expire(reason: str, *, emit_expired: bool = True) -> None

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

mark_closed_emitted() -> None

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

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

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.

shutdown async

shutdown() -> None

spawn

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

start

start() -> None

submit

submit(command: DuplexCommand) -> None

compute_silence_continuation_deadline

compute_silence_continuation_deadline(
    *,
    chunk_period_s: float,
    now: float,
    last_submit: float | None,
    current_deadline: float | None,
) -> tuple[float, float]

Return the silence continuation schedule (delay_s, next_silence_deadline).

The first continuation anchors to the latest accepted append's submission time (last_submit + chunk_period_s), falling back to now when no submission exists. Later continuations advance from the current deadline instead of from now, so pipeline processing time does not accumulate as timer drift. When the schedule is more than one period overdue it is stale: one continuation submits immediately and the schedule restarts from now (next_silence_deadline = now + chunk_period_s) instead of firing a burst of catch-ups.

delay_s is the wait before this continuation and next_silence_deadline the deadline for the following one.