vllm_omni.engine.omni_engine_base ¶
OmniEngineBase: process/thread ownership and transport shared by every Omni engine.
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] = []
endpoint_restrictions instance-attribute ¶
endpoint_restrictions = (
pipeline_config.endpoint_restrictions
if pipeline_config is not None
else ()
)
orchestrator_thread instance-attribute ¶
orchestrator_thread = threading.Thread(
target=self._bootstrap_orchestrator,
args=(stage_init_timeout, startup_future),
daemon=True,
name="orchestrator",
)
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).
abort ¶
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 |
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.
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.