vllm_omni.diffusion.worker.diffusion_worker ¶
Diffusion Worker for vLLM-Omni.
Handles GPU infrastructure initialization and delegates model operations to DiffusionModelRunner.
CustomPipelineWorkerExtension ¶
DiffusionWorker ¶
A worker that manages GPU infrastructure and delegates to the model runner.
This class handles infrastructure initialization only: - Device setup (CUDA device selection) - Distributed environment (NCCL, model parallel) - Memory management (sleep/wake)
All model-related operations (loading, compilation, execution) are delegated to DiffusionModelRunner.
model_runner instance-attribute ¶
model_runner: DiffusionModelRunner | None = (
model_runner_cls(
vllm_config=self.vllm_config,
od_config=self.od_config,
device=self.device,
)
)
close_ar_diffusion_session ¶
Close runner-owned AR state through the collective RPC boundary.
determine_available_kv_memory ¶
determine_available_kv_memory(
profile_requests: list[OmniDiffusionRequest],
) -> list[int]
Profile and return each rank's safe Diffusion KV memory budget.
execute_model ¶
execute_model(
req: OmniDiffusionRequest | list[NewRequestData],
od_config: OmniDiffusionConfig,
kv_prefetch_job: KVPrefetchJob | None = None,
diffusion_kv_metadata: DiffusionKVMetadata
| None = None,
) -> DiffusionOutput
Execute a forward pass by delegating to the model runner.
If req is a list (DP multi-concurrency), each rank picks one complete NewRequestData envelope based on its distributed rank. AllGather in the layerwise offload only gathers weight shards (request-independent), so all ranks stay synchronised at each AllGather call while computing different activations. Selecting the envelope keeps Scheduler-issued KV metadata bound to the request that owns its block tables.
Each rank returns its OWN DiffusionOutput (no gather). The executor collects N responses via the per-worker result queues.
execute_model_batch ¶
execute_model_batch(
scheduler_output: DiffusionSchedulerOutput,
od_config: OmniDiffusionConfig,
) -> BatchRunnerOutput
Batch forward: LoRA activate once, delegate to model runner.
execute_stepwise ¶
execute_stepwise(
scheduler_output: DiffusionSchedulerOutput,
) -> BaseRunnerOutput
Execute one diffusion step by delegating to the model runner.
get_kv_cache_specs ¶
Return native rank-local specs for every diffusion Worker.
load_model ¶
load_model(
load_format: str = "default",
custom_pipeline_name: str | None = None,
**kwargs,
) -> None
Load the diffusion model using DiffusionModelRunner.
profile ¶
remove_diffusion_kv_requests ¶
Clear Worker-local rows without freeing Scheduler-owned blocks.
reset_ar_diffusion_session ¶
Reset runner-owned AR state through the collective RPC boundary.
set_kv_cache_configs ¶
Select this rank's config and initialize its physical KV pages.
sleep ¶
Put the worker to sleep, offloading model weights.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
level | int | Sleep level. Level 1 offloads weights, level 2 also saves buffers. | 1 |
submit_interaction ¶
submit_interaction(
request_id: str, interaction: OmniInteractionPrompt
) -> None
Apply a midway interaction to an active stepwise request.
synchronize_device ¶
synchronize_device(timeout: float | None = None) -> None
Wait until this rank has no device work left.
A KV prefetch still receiving on its background thread has queued no device work yet, so it is joined first. timeout bounds that join (and the multi-process worker's output drain before this call); the device wait itself is unbounded.
wake_up ¶
Wake up the worker from sleep mode.
Re-activates the memory allocator for the specified tags and restores model buffers from CPU back to GPU if they were saved during Level 2 sleep.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
tags | list[str] | None | List of memory pool tags to re-activate (e.g., ["weights"] to match Level 1 sleep). If None, all pools are re-activated. | None |
WorkerProc ¶
Wrapper that runs one Worker in a separate process.
result_mq instance-attribute ¶
worker instance-attribute ¶
drain_async_outputs ¶
Block until background D2H/SHM packing has no work left.
Returns False if outputs are still in flight when timeout expires.
recv_message ¶
recv_message() -> Any
Receive one complete broadcast message without dropping overflow data.
worker_main staticmethod ¶
worker_main(
rank: int,
od_config: OmniDiffusionConfig,
pipe_writer: Connection,
broadcast_handle,
wake_event: Event,
worker_extension_cls: str | None = None,
custom_pipeline_args: dict[str, Any] | None = None,
) -> None
Worker initialization and execution loops.
WorkerWrapperBase ¶
Wrapper base class that creates DiffusionWorker with optional worker_extension_cls support. This enables dynamic inheritance for DiffusionWorker to extend with custom functionality.
uses_custom_pipeline instance-attribute ¶
uses_custom_pipeline = bool(
custom_pipeline_args
and "pipeline_class" in custom_pipeline_args
)
execute_method ¶
Execute a method on the worker.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
method | str | bytes | Method name (str) or serialized callable (bytes) | required |
Returns:
| Type | Description |
|---|---|
Any | Result of the method execution (type depends on the method) |
Raises:
| Type | Description |
|---|---|
Exception | If method execution fails |
execute_model ¶
execute_model(
req: OmniDiffusionRequest | list[NewRequestData],
od_config: OmniDiffusionConfig,
kv_prefetch_job: KVPrefetchJob | None = None,
diffusion_kv_metadata: DiffusionKVMetadata
| None = None,
) -> DiffusionOutput
Execute a forward pass.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
req | OmniDiffusionRequest | list[NewRequestData] | Diffusion request. | required |
od_config | OmniDiffusionConfig | OmniDiffusionConfig configuration | required |
kv_prefetch_job | KVPrefetchJob | None | Optional next-request KV prefetch descriptor. | None |
Returns:
| Type | Description |
|---|---|
DiffusionOutput | DiffusionOutput with generated results |
execute_stepwise ¶
execute_stepwise(
scheduler_output: DiffusionSchedulerOutput,
) -> BaseRunnerOutput
Execute one diffusion step.
sleep ¶
Put the worker to sleep. The worker should not process any requests. The caller should guarantee that no requests are being processed during the sleep period, before wake_up is called.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
level | int | The sleep level. Level 1 sleep will offload the model weights and discard the kv cache. Level 2 also saves buffers. Currently only support level 1 and level 2. | 1 |
Returns:
| Type | Description |
|---|---|
int | Bytes held by the allocator before sleeping. |
wake_up ¶
Wake up the worker from sleep mode. See the sleep function method for more details.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
tags | list[str] | None | An optional list of tags to reallocate the worker memory for specific memory allocations. Values must be in | None |
Returns:
| Type | Description |
|---|---|
bool | True on success |