vllm_omni.worker_v2.omni_data_plane ¶
OmniRunnerDataPlane ¶
Bases: OmniConnectorModelRunnerMixin
MRv2-owned stage transport and request-side payload state.
abort_requests ¶
Cancel deferred outputs and terminate each live request once.
complete_outputs ¶
complete_outputs(
*,
req_ids: list[str],
inter_stage_outputs: list[Any | None] | None,
sampled_token_ids: list[list[int]] | None,
) -> int
Commit one deferred output batch, then release its lifecycle holds.
emit_chunks ¶
emit_chunks(
*,
req_ids: list[str],
inter_stage_outputs: list[Any | None] | None,
sampled_token_ids: list[list[int]] | None,
terminal_req_ids: set[str],
) -> int
enqueue_outputs ¶
enqueue_outputs(
*,
req_ids: list[str],
inter_stage_outputs: list[Any | None] | None,
sampled_token_ids: list[list[int]] | None,
) -> None
Transfer one completed model batch to the ordered output worker.
pop_local_stage_payload ¶
Hand one accumulated delta to the model and acknowledge its rows.
Connector accumulation and model_intermediate_buffer are separate ownership domains. Once decode rows are handed to the model, remove them from connector accumulation so the next chunk cannot replay the same rows into Talker's cached_decode.
request_terminal ¶
Emit terminal markers only after all deferred outputs are enqueued.
reserve_outputs ¶
Keep request state alive until deferred runner outputs are consumed.
shutdown_omni_connectors ¶
Bound MRv2 shutdown even if a connector call cannot be cancelled.
Connector close() runs concurrently with the I/O threads so a backend that supports cancellation can release a blocked put(). Backends that do not support it are quarantined and left only on daemon threads after the bounded deadline; V1 keeps its existing shutdown implementation in the shared mixin.