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.
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)
tts_max_instructions_length instance-attribute ¶
abort async ¶
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 ¶
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_supported_tasks async ¶
get_supported_tasks() -> tuple[SupportedTask, ...]
Return the task set exposed by the orchestrator-backed engine.
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.
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.
- Stop frontend admission (
_paused). - For AR/LLM stages, call EngineCore.pause_scheduler via the Orchestrator loop (abort/wait/keep + optional cache clear).
- 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).
remove_lora async ¶
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 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 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 ¶
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 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.