Skip to content

vllm_omni.engine.omni_engine_base

OmniEngineBase: process/thread ownership and transport shared by every Omni engine.

logger module-attribute

logger = init_logger(__name__)

OmniEngineBase

Generic engine: launches an orchestrator in a background thread.

Shared by AsyncOmniEngine (turn-based requests) and DuplexOmniEngine (duplex sessions). Subclasses construct their orchestrator in _create_orchestrator.

All stage clients, input/output processors, and stage-to-stage transfer logic live inside the Orchestrator coroutine (running in its own thread with a dedicated asyncio event loop). This class communicates with it via janus queues (sync side for callers, async side for orchestrator).

Parameters:

Name Type Description Default
model str

Model name or path

required
init_timeout int

Total timeout waiting for orchestrator startup (seconds).

600
stage_init_timeout int

Timeout for stage initialization (seconds)

300
**kwargs Any

Additional arguments

{}

async_chunk instance-attribute

async_chunk = any(
    bool(
        getattr(
            getattr(stage, "connector_config", None),
            "async_chunk",
            getattr(
                getattr(stage, "engine_args", None),
                "async_chunk",
                False,
            ),
        )
    )
    for stage in self.stage_configs
)

default_sampling_params_list instance-attribute

default_sampling_params_list: list[OmniSamplingParams] = []

deploy_config instance-attribute

deploy_config = load_deploy_config(deploy_config_source)

endpoint_restrictions instance-attribute

endpoint_restrictions = (
    pipeline_config.endpoint_restrictions
    if pipeline_config is not None
    else ()
)

input_processor instance-attribute

input_processor: InputProcessor | None = None

model instance-attribute

model = model

num_stages instance-attribute

num_stages = len(self.stage_configs)

orchestrator_thread instance-attribute

orchestrator_thread = threading.Thread(
    target=self._bootstrap_orchestrator,
    args=(stage_init_timeout, startup_future),
    daemon=True,
    name="orchestrator",
)

output_queue instance-attribute

output_queue: Queue[EngineQueueMessage] = janus.Queue()

pipeline_config instance-attribute

pipeline_config = pipeline_config

prompt_expand_func instance-attribute

prompt_expand_func: Any | None = None

prompt_transform_func instance-attribute

prompt_transform_func: Any | None = None

request_queue instance-attribute

request_queue: Queue[EngineQueueMessage] = janus.Queue(
    maxsize=_REQUEST_QUEUE_MAXSIZE
)

rpc_client property

rpc_client: CorrelatedRpcClient

The correlated RPC transport (available once the orchestrator is up).

rpc_output_queue instance-attribute

rpc_output_queue: Queue[EngineQueueMessage] = janus.Queue()

single_stage_mode instance-attribute

single_stage_mode: bool = single_stage_mode

stage_clients instance-attribute

stage_clients: list[StageClient] = []

stage_metadata instance-attribute

stage_metadata: list[StageRuntimeInfo] = []

stage_pools instance-attribute

stage_pools: list[StagePool] = []

supported_tasks instance-attribute

supported_tasks: tuple[str, ...] = ('generate',)

tokenizer instance-attribute

tokenizer = tokenizer

abort

abort(request_ids: list[str]) -> None

Fire-and-forget abort: enqueue and return without waiting.

Prefer :meth:abort_async when the caller needs acknowledgment that stage aborts, binding release, and orchestrator request cleanup finished.

abort_async async

abort_async(
    request_ids: list[str], timeout: float | None = None
) -> list[OutputMessage]

Abort requests and wait for orchestrator acknowledgment.

Unlike :meth:abort, this generates an rpc_id, correlates the :class:AbortResultMessage via :class:CorrelatedRpcClient, and raises if the orchestrator reports failure or times out.

Returns:

Type Description
list[OutputMessage]

Final-stage AR abort OutputMessage list carrying partial

list[OutputMessage]

tokens generated before abort (empty for diffusion / no OP state).

collective_rpc

collective_rpc(
    method: str,
    timeout: float | None = None,
    args: tuple[Any, ...] = (),
    kwargs: dict[str, Any] | None = None,
    stage_ids: list[int] | None = None,
) -> list[Any]

Send a control RPC to the Orchestrator and wait for aggregated results.

This uses a dedicated RPC output queue so control-plane messages do not race with the normal request output polling loop.

collective_rpc_async async

collective_rpc_async(
    method: str,
    timeout: float | None = None,
    args: tuple[Any, ...] = (),
    kwargs: dict[str, Any] | None = None,
    stage_ids: list[int] | None = None,
) -> list[Any]

Async wrapper around collective_rpc().

get_diffusion_od_config

get_diffusion_od_config() -> Any

Expose the diffusion model_class_name to client-side model-extras.

The worker holds the full config; here we just resolve the pipeline class name from the model config (cached). model_class_name may be None.

get_output_blocking_async async

get_output_blocking_async(
    timeout: float = 1.0,
) -> EngineQueueMessage | None

Blocking-wait read from the Orchestrator output queue.

Waits up to timeout seconds in a dedicated drain thread for the next message (condition-variable wakeup instead of a poll cadence); returns None on timeout so the caller keeps its liveness check, mirroring try_get_output_async's contract. Used by the serving final-output drain when VLLM_OMNI_EVENT_DRIVEN_ORCH is on.

get_stage_metadata

get_stage_metadata(stage_id: int) -> StageRuntimeInfo

Get cached metadata for a stage.

is_alive

is_alive() -> bool

Whether the orchestrator thread is alive.

shutdown

shutdown() -> None

Send shutdown message and wait for the Orchestrator thread to exit.

try_get_output

try_get_output(
    timeout: float = 0.001,
) -> EngineQueueMessage | None

Read one output message from the Orchestrator output queue.

try_get_output_async async

try_get_output_async() -> EngineQueueMessage | None

Async read from the Orchestrator output queue.

load_and_resolve_stage_configs

load_and_resolve_stage_configs(
    model: str,
    kwargs: dict[str, Any],
    *,
    trust_remote_code: bool | None,
    deploy_config_path: str | None,
    stage_overrides: Mapping[str, Mapping[str, Any]] | None,
    strategy_config_path: str | None,
) -> OmniConfigResolution

Compatibility seam delegating to the single config resolver.