Skip to content

vllm_omni.worker_v2.omni_data_plane

logger module-attribute

logger = init_logger(__name__)

OmniRunnerDataPlane

Bases: OmniConnectorModelRunnerMixin

MRv2-owned stage transport and request-side payload state.

model_config instance-attribute

model_config = model_config

vllm_config instance-attribute

vllm_config = vllm_config

abort_requests

abort_requests(req_ids: set[str]) -> int

Cancel deferred outputs and terminate each live request once.

close

close() -> None

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.

drain_outputs

drain_outputs() -> None

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

pop_local_stage_payload(req_id: str) -> Any

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.

register_receivers

register_receivers(handles: list[Any]) -> None

register_request

register_request(request_data: Any) -> None

request_terminal

request_terminal(req_ids: set[str]) -> int

Emit terminal markers only after all deferred outputs are enqueued.

reserve_outputs

reserve_outputs(req_ids: list[str]) -> None

Keep request state alive until deferred runner outputs are consumed.

shutdown_omni_connectors

shutdown_omni_connectors() -> None

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.