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.
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]
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
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_created_payload ¶
response_created_payload(
response_id: str,
*,
epoch: int,
request_id: str | None = None,
) -> dict[str, object]
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.