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 the encoder cache for all stages.
TODO: Forward to Orchestrator process via message.
reset_mm_cache async ¶
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 the prefix cache for all stages.
TODO: Forward to Orchestrator process via message.
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.