Skip to content

vllm_omni.engine.async_engine_utils

Stateless request and shutdown helpers for :mod:async_omni_engine.

SHUTDOWN_ENQUEUE_TIMEOUT_S module-attribute

SHUTDOWN_ENQUEUE_TIMEOUT_S = 1.0

SHUTDOWN_JOIN_TIMEOUT_S module-attribute

SHUTDOWN_JOIN_TIMEOUT_S = 30.0

logger module-attribute

logger = init_logger(__name__)

apply_omni_final_stage_metadata

apply_omni_final_stage_metadata(
    request: EngineCoreRequest,
    final_stage_id: int,
    *,
    force_kv_transfer: bool = False,
) -> EngineCoreRequest

Tag a request with its final stage and optional KV-transfer override.

enqueue_orchestrator_shutdown

enqueue_orchestrator_shutdown(
    request_queue: Queue[EngineQueueMessage] | None,
    *,
    timeout: float,
) -> bool

Deliver shutdown without waiting forever behind a full request queue.

inject_global_id

inject_global_id(target: Any, request_id: str) -> None

Inject global_request_id into a prompt's additional information.

is_abort_transport_shutdown

is_abort_transport_shutdown(exc: Exception) -> bool

is_janus_sync_queue_shutdown

is_janus_sync_queue_shutdown(exc: Exception) -> bool

shutdown_runtime_after_orchestrator

shutdown_runtime_after_orchestrator(
    orchestrator_thread: Thread, runtime: StageRuntime
) -> None

Release a stage runtime once its orchestrator can no longer use it.

upgrade_to_omni_request

upgrade_to_omni_request(
    request: EngineCoreRequest, raw_prompt: Any
) -> EngineCoreRequest

Restore omni-only fields omitted by the upstream input processor.

weak_shutdown_async_omni_engine

weak_shutdown_async_omni_engine(
    orchestrator_thread: Thread | None,
    request_queue: Queue[EngineQueueMessage] | None,
    output_queue: Queue[EngineQueueMessage] | None,
    rpc_output_queue: Queue[EngineQueueMessage] | None,
    rpc_client: CorrelatedRpcClient | None,
) -> None

Best-effort orchestrator cleanup for garbage-collection finalization.