Skip to content

vllm_omni.engine.async_omni_engine

AsyncOmniEngine: the turn-based engine (add_request / streaming updates / CFG companions / interaction).

logger module-attribute

logger = init_logger(__name__)

AsyncOmniEngine

Bases: OmniEngineBase

Turn-based engine used by Omni / AsyncOmni.

add_request

add_request(
    request_id: str,
    prompt: EngineCoreRequest | PromptType,
    prompt_text: str | None = None,
    sampling_params_list: Sequence[Any] | None = None,
    final_stage_id: int = 0,
    final_output_stage_ids: Sequence[int] | None = None,
    arrival_time: float | None = None,
    lora_request: Any = None,
    tokenization_kwargs: dict[str, Any] | None = None,
    trace_headers: Mapping[str, str] | None = None,
    priority: int = 0,
    data_parallel_rank: int | None = None,
    reasoning_ended: bool | None = None,
    *,
    resumable: bool = False,
) -> None

Process stage-0 input locally, then send to the Orchestrator.

Input processing and output processor registration happen here in the caller's thread, avoiding a queue + coroutine-switch round-trip. The Orchestrator receives a ready-to-submit OmniEngineCoreRequest.

add_request_async async

add_request_async(
    request_id: str,
    prompt: EngineCoreRequest | PromptType,
    prompt_text: str | None = None,
    sampling_params_list: Sequence[Any] | None = None,
    final_stage_id: int = 0,
    final_output_stage_ids: Sequence[int] | None = None,
    arrival_time: float | None = None,
    lora_request: Any = None,
    tokenization_kwargs: dict[str, Any] | None = None,
    trace_headers: Mapping[str, str] | None = None,
    priority: int = 0,
    data_parallel_rank: int | None = None,
    reasoning_ended: bool | None = None,
    *,
    resumable: bool = False,
) -> None

Async add_request API.

add_streaming_update

add_streaming_update(
    request_id: str,
    prompt: EngineCoreRequest | PromptType,
    prompt_text: str | None = None,
    sampling_params_list: Sequence[Any] | None = None,
    final_stage_id: int = 0,
    final_output_stage_ids: Sequence[int] | None = None,
    arrival_time: float | None = None,
    lora_request: Any = None,
    *,
    resumable: bool = True,
) -> None

Send an incremental streaming update for an existing request.

add_streaming_update_async async

add_streaming_update_async(
    request_id: str,
    prompt: EngineCoreRequest | PromptType,
    prompt_text: str | None = None,
    sampling_params_list: Sequence[Any] | None = None,
    final_stage_id: int = 0,
    final_output_stage_ids: Sequence[int] | None = None,
    arrival_time: float | None = None,
    lora_request: Any = None,
    *,
    resumable: bool = True,
) -> None

Async wrapper for add_streaming_update().

submit_interaction

submit_interaction(
    request_id: str, interaction: OmniInteractionPrompt
) -> None

Send an interaction control message to the Orchestrator.

submit_interaction_async async

submit_interaction_async(
    request_id: str, interaction: OmniInteractionPrompt
) -> None

Async interaction API.

StageRuntimeInfo dataclass

final_output instance-attribute

final_output: bool

final_output_type instance-attribute

final_output_type: FinalOutputModalityType | None

model_stage class-attribute instance-attribute

model_stage: str | None = None

stage_type instance-attribute

stage_type: str