Skip to content

vllm_omni.engine.duplex.session.model_channel

The model side of one duplex session.

Everything between the session and the model runtime: planning an append with the plugin and submitting it to the stage port, turning the stage output that comes back into session events, and -- when the model chooses to keep listening -- offering it another unit of silence so the turn can continue.

Those three read as separate concerns but they call each other in a cycle. An append submits and immediately projects its own first outputs; a projected output can schedule a continuation; a continuation is another append. Splitting them into three modules would mean three mutual imports and no fewer edges, so they are one component whose seam with the runner is narrow instead.

That seam is three callbacks, for the three things this component triggers but does not own: closing the session, scheduling the continuation append as a tracked task, and aborting a stage request in the background.

logger module-attribute

logger = init_logger(__name__)

ModelChannel

Appends out to the model, events back from it, for one session.

append_runtime_input async

append_runtime_input(
    payload: object,
    *,
    operation_id: str | None = None,
    final: bool,
    expected_epoch: int | None = None,
    on_append_accepted: Callable[[float], None]
    | None = None,
) -> tuple[bool, bool]

cancel_data_plane_stream

cancel_data_plane_stream() -> bool

decide_output

decide_output(
    stage_id: int,
    output: RequestOutput,
    context: DuplexOutputContext,
) -> DuplexOutputDecision | None

Pure plugin decision (no session mutation: this runs inline on the orchestrator loop).

maybe_continue_response async

maybe_continue_response(
    *,
    expected_epoch: int | None,
    expected_model_turn_id: int | None = None,
) -> None

on_stage_output_item async

on_stage_output_item(item: StageOutput) -> None

project_intermediate_output

project_intermediate_output(
    stage_id: int,
    output: RequestOutput,
    context: DuplexOutputContext,
) -> bool

Project an intermediate stage without short-circuiting the pipeline.

release_concurrent_turn_requests

release_concurrent_turn_requests(
    stage_id: int,
    output: RequestOutput,
    context: DuplexOutputContext,
) -> bool

Ask the plugin whether the next commit may start while TTS drains.

response_continuations_remaining

response_continuations_remaining(response_id: str) -> bool

response_created_payload

response_created_payload(
    response_id: str,
    *,
    epoch: int,
    request_id: str | None = None,
) -> dict[str, object]

send_runtime_error

send_runtime_error(code: str, exc: BaseException) -> None

should_commit_response_to_history staticmethod

should_commit_response_to_history(
    session: DuplexEngineSession, response_id: str | None
) -> bool

signal_cancel_fence async

signal_cancel_fence(cancelled_fence: DuplexFence) -> bool

In-process equivalent of the old barge_in control signal.

silence_continuation_is_stale

silence_continuation_is_stale(
    *,
    request_id: str,
    response_id: str | None,
    response_owned: bool,
    expected_epoch: int | None,
    expected_model_turn_id: int | None,
) -> bool

silence_unit_payload

silence_unit_payload() -> dict[str, object]

stage_metrics_snapshot

stage_metrics_snapshot(
    stage_id: int, metrics: object, output: object
) -> dict[str, dict[str, object]] | None

SilenceContinuationScheduler

Bases: Protocol

Runs one continuation append as a tracked task (see _schedule_silence_continuation).