Skip to content

vllm_omni.entrypoints.async_omni

AsyncOmni - Refactored async orchestrator using AsyncOmniEngine.

This is the new implementation that uses AsyncOmniEngine (which manages StageEngineCoreClient instances) instead of OmniStage with worker processes.

logger module-attribute

logger = init_logger(__name__)

AsyncOmni

Bases: AsyncOmniBase, EngineClient

Asynchronous unified entry point for multi-stage pipelines using AsyncOmniEngine.

This is the refactored version that uses AsyncOmniEngine instead of OmniStage workers. It provides the same interface as AsyncOmni but with a cleaner architecture.

Parameters:

Name Type Description Default
model str

Model name or path to load.

''
**kwargs Any

Additional keyword arguments. - deploy_config: Optional path to a deploy YAML. If None, configurations are resolved from the model pipeline factory. - log_stats: Whether to enable statistics logging. - stage_init_timeout: Timeout for per-stage initialization. - init_timeout: Total timeout for orchestrator startup. - async_chunk: Whether to enable async chunk mode. - output_modalities: Requested output modalities. - Additional keyword arguments passed to stage engines.

{}
Example

async_omni = AsyncOmni(model="Qwen/Qwen2.5-Omni-7B") async for output in async_omni.generate( ... prompt="Hello", ... request_id="req-1", ... sampling_params_list=[SamplingParams(), SamplingParams()] ... ): ... print(output)

engine instance-attribute

tts_max_instructions_length instance-attribute

tts_max_instructions_length = tts_max_instructions_length

abort async

abort(
    request_id: str | Iterable[str],
    *,
    timeout: float | None = None,
) -> None

Abort request(s) via the Orchestrator.

add_lora async

add_lora(lora_request: LoRARequest) -> bool

Load a new LoRA adapter into all stages.

Returns True only if all concretely-implemented stages report success.

collective_rpc async

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]

Execute a best-effort control RPC on selected stages.

Unsupported stages currently return a TODO-style result dict instead of failing the entire call. This keeps AsyncOmni usable while the orchestrator control plane is still being filled out.

do_log_stats async

do_log_stats() -> None

Log statistics.

TODO: Forward to Orchestrator process via message.

encode async

encode(
    prompt: Any,
    pooling_params: PoolingParams,
    request_id: str,
    lora_request: LoRARequest | None = None,
    trace_headers: dict[str, str] | None = None,
    priority: int = 0,
    tokenization_kwargs: dict[str, Any] | None = None,
    reasoning_ended: bool | None = None,
) -> AsyncGenerator[PoolingRequestOutput, None]

EngineClient.encode() stub.

Omni pipeline currently exposes only generate() API at orchestrator level.

finish_weight_update async

finish_weight_update(
    weight_version: str | None = None,
) -> None

Finish the current weight update.

Omni does not currently support weight transfer, so this is a no-op. weight_version is accepted for upstream EngineClient protocol compatibility (RLHF weight-transfer routers pass it positionally).

generate async

generate(
    prompt: OmniPromptType
    | AsyncGenerator[StreamingInput, None]
    | list[OmniPromptType],
    sampling_params: Any = None,
    request_id: str = "",
    *,
    prompt_text: str | None = None,
    lora_request: Any = None,
    tokenization_kwargs: dict[str, Any] | None = None,
    sampling_params_list: Sequence[OmniSamplingParams]
    | None = None,
    output_modalities: list[str] | None = None,
    trace_headers: Mapping[str, str] | None = None,
    priority: int = 0,
    data_parallel_rank: int | None = None,
    session_id: str | None = None,
    reasoning_ended: bool | None = None,
    reasoning_parser_kwargs: dict[str, Any] | None = None,
    arrival_time: float | None = None,
) -> AsyncGenerator[OmniRequestOutput, None]

Generate outputs for the given prompt(s) asynchronously.

Coordinates multi-stage pipeline execution. Processes the prompt through all stages in the pipeline and yields outputs as they become available.

session_id is accepted for EngineClient protocol compatibility and is not duplex-session plumbing.

Diffusion batching: Diffusion stages accept only a single prompt per request. Passing a list of prompts to a diffusion stage will raise ValueError. To batch multiple diffusion prompts, submit each as an independent request; the scheduler will automatically co-batch compatible requests.

Parameters:

Name Type Description Default
prompt OmniPromptType | AsyncGenerator[StreamingInput, None] | list[OmniPromptType]

A single prompt or a list of prompts. For diffusion stages, only a single prompt is accepted; a list will be rejected with an error.

required
request_id str

Unique identifier for this request. If one is not provided, a random one will be generated.

''
sampling_params_list Sequence[OmniSamplingParams] | None

List of SamplingParams, one per stage. Must have the same length as the number of stages. If None, uses default sampling params for each stage.

None
output_modalities list[str] | None

Optional list of output modalities.

None

Yields:

Type Description
AsyncGenerator[OmniRequestOutput, None]

OmniRequestOutput objects as they are produced by each stage.

Raises:

Type Description
ValueError

If sampling_params_list has incorrect length, or if a list prompt is submitted to a diffusion stage.

get_input_preprocessor async

get_input_preprocessor() -> InputProcessor

Get input preprocessor.

get_supported_tasks async

get_supported_tasks() -> tuple[SupportedTask, ...]

Return the task set exposed by the orchestrator-backed engine.

get_tokenizer async

get_tokenizer() -> TokenizerLike

Get tokenizer for the comprehension stage.

is_paused async

is_paused() -> bool

Check if frontend admission is paused.

is_sleeping async

is_sleeping() -> bool

Return whether all stages are sleeping.

TODO(AsyncOmni): query the orchestrator once all stage backends expose a real sleeping-state RPC. For now we track the requested state locally.

is_tracing_enabled async

is_tracing_enabled() -> bool

Check if tracing is enabled.

list_loras async

list_loras() -> list[int]

List all loaded LoRA adapter IDs across stages.

notify_kv_transfer_request_rejected async

notify_kv_transfer_request_rejected(
    request_id: str,
    kv_transfer_params: dict[str, Any],
    *,
    data_parallel_rank: int | None = None,
) -> None

Notify engine that a KV-transfer request was rejected before admission.

Omni does not currently use KV-transfer pre-admission resources, so this is a no-op.

pause_generation async

pause_generation(
    *,
    mode: PauseMode = "abort",
    wait_for_inflight_requests: bool = False,
    clear_cache: bool = True,
    stage_ids: list[int] | None = None,
) -> None

Pause generation, mirroring vLLM AsyncLLM.pause_generation.

  1. Stop frontend admission (_paused).
  2. For AR/LLM stages, call EngineCore.pause_scheduler via the Orchestrator loop (abort/wait/keep + optional cache clear).
  3. For diffusion stages, mode="keep" pauses the DiffusionEngine scheduler and returns once the batch that was running has finished on every worker; that batch is delivered before any control RPC issued after this call runs, so the documented pause -> sleep order is safe. Queued requests stay queued until :meth:resume_generation. Other modes pause frontend admission only.

Note: sleep() already pauses the AR scheduler internally (same as vLLM EngineCore.sleep). Call this API when you need pause without freeing GPU memory (e.g. weight sync).

pin_lora async

pin_lora(adapter_id: int) -> bool

Pin a LoRA adapter across stages.

remove_lora async

remove_lora(adapter_id: int) -> bool

Remove a LoRA adapter from all stages.

TODO(AsyncOmni): add richer per-stage error reporting to the public API.

reset_encoder_cache async

reset_encoder_cache() -> None

Reset the encoder cache for all stages.

TODO: Forward to Orchestrator process via message.

reset_mm_cache async

reset_mm_cache() -> None

Reset the frontend (P0) multimodal processor cache.

EngineCore.sleep(level>=1) already clears the P1 receiver cache. Clearing P0 avoids hash-only follow-up requests after that reset.

reset_prefix_cache async

reset_prefix_cache(
    reset_running_requests: bool = False,
    reset_connector: bool = False,
) -> bool

Reset the prefix cache for all stages.

TODO: Forward to Orchestrator process via message.

resume_generation async

resume_generation(
    stage_ids: list[int] | None = None,
) -> None

Resume generation after :meth:pause_generation.

sleep async

sleep(
    stage_ids: list[int] | None = None,
    level: int = 2,
    mode: PauseMode = "abort",
) -> list[OmniACK]

Put stages to sleep.

AR/LLM stages use EngineCore.sleep (pause scheduler, wait idle, then offload/discard memory) — matching vLLM AsyncLLM.sleep.

Diffusion stages keep the worker-level handle_sleep_task RPC, which does not stop the DiffusionEngine scheduler; quiesce a busy diffusion stage first with pause_generation(mode="keep") (or abort it).

Frontend admission is blocked at the start of this call (_paused) so pipelined :meth:generate cannot race into stages while sleep is in flight. This does not invoke EngineCore.pause_scheduler again (sleep already pauses the AR scheduler).

For AR / mixed engines, wake_up does not clear _paused; callers must :meth:resume_generation when ready (typical trainer order: pause → abort → sleep → train → wake → resume). Diffusion-only engines have no EngineCore pause to hold, so wake_up restores admission and sleep → wake → generate keeps working.

start_profile async

start_profile(
    profile_prefix: str | None = None,
    stages: list[int] | None = None,
) -> list[Any]

Start profiling specified stages.

Uses vLLM-compatible profile(is_start=True, profile_prefix) interface.

Parameters:

Name Type Description Default
profile_prefix str | None

Optional prefix for the trace file names.

None
stages list[int] | None

List of stage IDs to profile. If None, profiles all stages.

None

start_weight_update async

start_weight_update(
    is_checkpoint_format: bool = True,
) -> None

Start a new weight update.

Omni does not currently support weight transfer, so this is a no-op.

stop_profile async

stop_profile(stages: list[int] | None = None) -> list[Any]

Stop profiling specified stages.

Uses vLLM-compatible profile(is_start=False) interface.

Parameters:

Name Type Description Default
stages list[int] | None

List of stage IDs to profile. If None, stops all stages.

None

submit_interaction_async async

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

Apply a midway interaction to an active streaming diffusion request.

request_id is the external id created by the server-side session, matching the value passed to :meth:generate.

wake_up async

wake_up(
    stage_ids: list[int] | None = None,
    tags: list[str] | None = None,
) -> list[OmniACK]

Wake stages after sleep.

AR/LLM stages use EngineCore.wake_up (restore memory, auto-resume scheduler). Diffusion stages keep the worker-level wake RPC.

Does not clear the frontend _paused admission gate when :meth:pause_generation ran or AR stages were slept — call :meth:resume_generation when the trainer is ready to admit new requests. Diffusion-only sleep uses _paused only as a race guard; this method restores admission after a successful wake.