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.

CACHE_RESET_TIMEOUT_S module-attribute

CACHE_RESET_TIMEOUT_S = 60.0

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(
    *,
    stage_ids: list[int] | None = None,
    timeout: float = CACHE_RESET_TIMEOUT_S,
) -> None

Reset encoder caches on selected AR stages (all by default).

Diffusion stages are skipped. RPC failures are raised to the caller.

reset_mm_cache async

reset_mm_cache(
    *,
    stage_ids: list[int] | None = None,
    timeout: float = CACHE_RESET_TIMEOUT_S,
) -> None

Reset sender and receiver MM caches on selected AR stages.

By default all AR stages are selected; diffusion stages are skipped. Stage 0's sender is the frontend renderer. Downstream senders are cleared by the orchestrator before resetting their engine cores. Call while generation is paused to avoid racing new inputs.

reset_prefix_cache async

reset_prefix_cache(
    reset_running_requests: bool = False,
    reset_connector: bool = False,
    *,
    stage_ids: list[int] | None = None,
    timeout: float = CACHE_RESET_TIMEOUT_S,
) -> bool

Reset prefix caches on selected AR stages (all by default).

Diffusion stages are skipped. Return False if a stage cannot reset its cache; unsupported operations and RPC failures raise instead of silently retaining KV computed under previous model weights.

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.