Skip to content

vllm_omni.diffusion.worker.diffusion_worker

Diffusion Worker for vLLM-Omni.

Handles GPU infrastructure initialization and delegates model operations to DiffusionModelRunner.

logger module-attribute

logger = init_logger(__name__)

CustomPipelineWorkerExtension

re_init_pipeline

re_init_pipeline(
    custom_pipeline_args: dict[str, Any],
) -> None

Re-initialize the pipeline with custom arguments.

Parameters:

Name Type Description Default
custom_pipeline_args dict[str, Any]

Dictionary of arguments for custom pipeline initialization

required

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.

device instance-attribute

device: device | None = None

init_snapshot instance-attribute

init_snapshot: MemorySnapshot | None = None

local_rank instance-attribute

local_rank = local_rank

lora_manager instance-attribute

lora_manager: DiffusionLoRAManager | None = None

model_runner instance-attribute

model_runner: DiffusionModelRunner | None = (
    model_runner_cls(
        vllm_config=self.vllm_config,
        od_config=self.od_config,
        device=self.device,
    )
)

od_config instance-attribute

od_config = od_config

profiler instance-attribute

profiler: WorkerProfiler | None = self._create_profiler()

rank instance-attribute

rank = rank

requested_memory instance-attribute

requested_memory: int | None = None

stage_id instance-attribute

stage_id = getattr(od_config, 'stage_id', 0)

vllm_config instance-attribute

vllm_config: VllmConfig | None = None

add_lora

add_lora(lora_request: LoRARequest) -> bool

close_ar_diffusion_session

close_ar_diffusion_session(session_id: str) -> bool

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

get_kv_cache_specs() -> list[dict[str, KVCacheSpec]]

Return native rank-local specs for every diffusion Worker.

handle_sleep_task

handle_sleep_task(
    task: OmniSleepTask | dict,
) -> OmniACK | None

handle_wake_task

handle_wake_task(
    task: OmniWakeTask | dict,
) -> OmniACK | None

init_device

init_device() -> None

Initialize the device and distributed environment.

init_lora_manager

init_lora_manager() -> None

Initialize the LoRA manager for this worker.

list_loras

list_loras() -> list[int]

load_model

load_model(
    load_format: str = "default",
    custom_pipeline_name: str | None = None,
    **kwargs,
) -> None

Load the diffusion model using DiffusionModelRunner.

pin_lora

pin_lora(adapter_id: int) -> bool

prepare_kv_for_forward

prepare_kv_for_forward(
    scheduler_output: DiffusionSchedulerOutput,
)

profile

profile(
    is_start: bool = True, profile_prefix: str | None = None
) -> None

Start or stop profiling for this GPU worker.

Parameters:

Name Type Description Default
is_start bool

True to start profiling, False to stop.

True
profile_prefix str | None

Optional prefix for trace filename.

None

remove_diffusion_kv_requests

remove_diffusion_kv_requests(
    request_ids: list[str | tuple[str, int]],
) -> int

Clear Worker-local rows without freeing Scheduler-owned blocks.

remove_lora

remove_lora(adapter_id: int) -> bool

reset_ar_diffusion_session

reset_ar_diffusion_session(session_id: str) -> bool

Reset runner-owned AR state through the collective RPC boundary.

set_kv_cache_configs

set_kv_cache_configs(
    kv_cache_configs: list[KVCacheConfig],
    resolved_max_model_len: int,
) -> None

Select this rank's config and initialize its physical KV pages.

shutdown

shutdown() -> None

Shutdown the worker and release process-global resources.

sleep

sleep(level: int = 1) -> int

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(tags: list[str] | None = None) -> bool

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.

context instance-attribute

context = zmq.Context(io_threads=2)

gpu_id instance-attribute

gpu_id = gpu_id

mq instance-attribute

mq = MessageQueue.create_from_handle(
    broadcast_handle, gpu_id
)

od_config instance-attribute

od_config = od_config

result_mq instance-attribute

result_mq = MessageQueue(
    n_reader=1, n_local_reader=1, local_reader_ranks=[0]
)

result_mq_handle instance-attribute

result_mq_handle = self.result_mq.export_handle()

wake_event instance-attribute

wake_event = wake_event

worker instance-attribute

worker = self._create_worker(
    gpu_id,
    od_config,
    worker_extension_cls,
    custom_pipeline_args,
)

drain_async_outputs

drain_async_outputs(timeout: float | None = None) -> bool

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.

shutdown

shutdown() -> None

Stop background work and release worker-owned IPC resources.

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.

base_worker_class instance-attribute

base_worker_class = base_worker_class

custom_pipeline_args instance-attribute

custom_pipeline_args = custom_pipeline_args

gpu_id instance-attribute

gpu_id = gpu_id

od_config instance-attribute

od_config = od_config

uses_custom_pipeline instance-attribute

uses_custom_pipeline = bool(
    custom_pipeline_args
    and "pipeline_class" in custom_pipeline_args
)

worker instance-attribute

worker = worker

worker_extension_cls instance-attribute

worker_extension_cls = worker_extension_cls

execute_method

execute_method(method: str | bytes, *args, **kwargs) -> Any

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.

handle_sleep_task

handle_sleep_task(
    task: OmniSleepTask | dict,
) -> OmniACK | None

handle_wake_task

handle_wake_task(
    task: OmniWakeTask | dict,
) -> OmniACK | None

shutdown

shutdown() -> None

Shutdown the worker and cleanup resources.

sleep

sleep(level: int = 1) -> int

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(tags: list[str] | None = None) -> bool

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 ("weights"). If None, all memory is reallocated. wake_up should be called with all tags (or None) before the worker is used again.

None

Returns:

Type Description
bool

True on success