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
DuplexEngineSessionhappens on this loop, through the mailbox worker (commands, stage outputs, internal items) or through tracked tasks that re-validate(epoch, turn_id)after eachawait; - 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_outputinstead of being polled through request queues; - everything the session says leaves as typed events through
DuplexSessionManager.emitafter terminal-acceptance / stale-epoch filtering, so a cancelled epoch can never speak again.
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.
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.