Skip to content

vllm_omni.diffusion.offloader

Modules:

Name Description
base
block_discovery

Block discovery for layerwise offload.

component_utils

Shared component-plan helpers for diffusion offload backends.

config

Public diffusion CPU-offload configuration helpers.

cuda_host_registration

CUDA implementation of read-only host-mapping registration.

distributed_layerwise_backend

Distributed Layerwise Offload backend with double-buffered H2D.

host_registration

Platform-neutral lifecycle for registering existing host mappings.

layerwise_backend
module_collector
module_residency

On-demand module staging backed by immutable pinned CPU storage.

offload_plan

Declarative model capabilities for component layerwise offload.

plan_resolver

One topology resolver shared by every diffusion offload backend.

sequential_backend
startup

Generic startup handoff between model loaders and offload backends.

tensor_utils

Shared tensor utilities for distributed layerwise offload.

logger module-attribute

logger = init_logger(__name__)

BlockStack dataclass

One ring of repeated blocks that a backend streams together.

attrs names the owner attributes that contributed blocks; backends use it to skip those children when they place the non-streamed remainder. It is filled for DiT stacks only — encoder stacks are hooked as whole stacks and their remainder is placed by tensor identity.

attrs class-attribute instance-attribute

attrs: tuple[str, ...] = ()

blocks instance-attribute

blocks: tuple[Module, ...]

resident property

resident: tuple[Module, ...]

Leading blocks that stay on the device instead of streaming.

resident_head class-attribute instance-attribute

resident_head: int = 0

streaming property

streaming: tuple[Module, ...]

Blocks a hook ring transfers on demand.

BoundedAllocatorCache

Retain reusable allocator blocks without monopolizing device memory.

Component offload normally calls empty_cache after every stage. That makes the next stage return to the device allocator even though PyTorch's cached blocks are immediately reusable. This policy keeps the cache while both of these bounds hold:

  • cached-but-unallocated memory is at most 25% of device capacity; and
  • at least 5% of device capacity is physically free.

Missing memory telemetry is handled conservatively by releasing the cache. Failure paths can force release; normal executor shutdown keeps its own unconditional device-cache cleanup because this policy is not global.

device instance-attribute

device = device

max_cached_fraction instance-attribute

max_cached_fraction = max_cached_fraction

min_free_fraction instance-attribute

min_free_fraction = min_free_fraction

release_if_needed

release_if_needed(*, force: bool = False) -> bool

Release cached blocks when a bound is crossed or release is forced.

DistributedLayerwiseOffloadBackend

Bases: OffloadBackend

Distributed layer-wise (block-level) offloading backend.

Supports both GPU (CUDA) and NPU (CANN) platforms. Device type is determined by the device passed to enable().

Each rank stores only a shard of each block's weights on host memory. Two device slots alternate: one holds current weights, one holds next weights. H2D and AllGather run asynchronously on dedicated streams, overlapped with computation.

comm_stream instance-attribute

comm_stream = current_omni_platform.Stream()

copy_stream instance-attribute

copy_stream = current_omni_platform.Stream()

dp_group instance-attribute

dp_group: ProcessGroup | None = None

dp_size instance-attribute

dp_size = config.dp_size

host_weight_plan instance-attribute

host_weight_plan = host_weight_plan

rank instance-attribute

rank = 0

disable

disable() -> None

enable

enable(pipeline: Module) -> None

Enable DLO and make partial startup failures transactional.

load_resident_layers

load_resident_layers() -> None

Load the model-declared leading blocks for the denoise stage.

offload_resident_layers

offload_resident_layers() -> None

Release leading blocks before VAE decode to bound peak HBM.

shutdown

shutdown() -> None

DistributedLayerwiseOffloadHook

Bases: ModelHook

Hook for distributed layerwise offloading with fixed double-buffer.

Each rank stores only a shard of each block's weights on host memory. Two device slots alternate: one holds current weights, one holds next weights. H2D and AllGather run asynchronously on dedicated streams, overlapped with computation.

Supports both NVIDIA GPU (CUDA) and Ascend NPU (CANN) platforms.

block_buffers instance-attribute

block_buffers: dict[str, Tensor] = {}

block_parameters instance-attribute

block_parameters: dict[str, Parameter] = {}

comm_stream instance-attribute

comm_stream = comm_stream or current_omni_platform.Stream()

copy_stream instance-attribute

copy_stream = copy_stream or current_omni_platform.Stream()

cpu_shards instance-attribute

cpu_shards: dict[dtype, Tensor] = {}

cpu_sources instance-attribute

cpu_sources: dict[dtype, list[dict[str, Any]]] = {}

cpu_staging_buffers instance-attribute

cpu_staging_buffers: list[dict[dtype, Tensor] | None] = [
    None,
    None,
]

cpu_staging_events instance-attribute

cpu_staging_events: list[Any | None] = [None, None]

current_slot instance-attribute

current_slot = 0

device instance-attribute

device = device

dp_group instance-attribute

dp_group = dp_group

dp_size instance-attribute

dp_size = dp_size

gpu_buffers instance-attribute

gpu_buffers: list[dict[dtype, Tensor] | None] = (
    shared_buffers
)

gpu_shard_buffers instance-attribute

gpu_shard_buffers: list[dict[dtype, Tensor] | None] = [
    None,
    None,
]

is_materialized property

is_materialized: bool

Check whether this block's parameters hold real data on device.

metadata instance-attribute

metadata: dict[dtype, list[dict[str, Any]]] = {}

next_block instance-attribute

next_block = next_block

next_block_buffers instance-attribute

next_block_buffers: dict[str, Tensor] = {}

next_block_parameters instance-attribute

next_block_parameters: dict[str, Parameter] = {}

pin_memory instance-attribute

pin_memory = pin_memory

rank instance-attribute

rank = rank

rank_local_mmap instance-attribute

rank_local_mmap = rank_local_mmap

ready_events instance-attribute

ready_events: list[Any | None] = [None, None]

registered_mmap instance-attribute

registered_mmap = False

tensor_transforms instance-attribute

tensor_transforms = tensor_transforms or {}

initialize_hook

initialize_hook(module: Module) -> Module

offload_layer

offload_layer() -> None

Free GPU memory for current block by replacing tensors with placeholders.

post_forward

post_forward(module: Module, output: Any) -> Any

pre_forward

pre_forward(
    module: Module, *args: Any, **kwargs: Any
) -> tuple[tuple, dict]

prefetch_layer

prefetch_layer(
    slot: int, non_blocking: bool = True
) -> None

Prepare next block's weights into the shared device buffer for slot.

Uses the pre-allocated self.gpu_buffers[slot] instead of allocating fresh tensors every layer. This enforces the fixed double-buffer memory bound (exactly 2 blocks on device).

restore_next_block_to_cpu

restore_next_block_to_cpu() -> None

Restore hook-owned master weights before removing the hook.

Every circular hook owns the host backing for its next_block. Dropping that hook while the module points at rotating device buffers or placeholders would make a later enable shard invalid tensors.

LayerWiseOffloadBackend

Bases: OffloadBackend

Layer-wise (block-level) offloading backend.

Implements sliding window offloading where only a small number of transformer blocks reside on GPU at a time. Blocks are prefetched asynchronously while previous blocks compute, and freed after use.

copy_stream instance-attribute

copy_stream = current_omni_platform.Stream()

disable

disable() -> None

enable

enable(pipeline: Module) -> None

ModelLevelOffloadBackend

Bases: OffloadBackend

Model-level (sequential) offloading backend.

Uses SequentialOffloadHook registered via HookRegistry for automatic module swapping.

disable

disable() -> None

enable

enable(pipeline: Module) -> None

OffloadBackend

Bases: ABC

Base class for CPU offload backends

config instance-attribute

config = config

device instance-attribute

device = device

enabled instance-attribute

enabled = False

disable abstractmethod

disable() -> None

Disable offloading and cleanup resources.

Removes all registered hooks. Does NOT move modules back to original devices (caller responsible for that).

enable abstractmethod

enable(pipeline: Module) -> None

Enable offloading on the pipeline.

Discovers modules, moves them to appropriate devices, and registers forward hooks for swapping/prefetching.

Parameters:

Name Type Description Default
pipeline Module

Diffusion pipeline model (e.g., Wan22Pipeline)

required

is_enabled

is_enabled() -> bool

shutdown

shutdown() -> None

Release offload resources at process exit.

Backends may skip work that only matters for a later enable.

OffloadConfig dataclass

components class-attribute instance-attribute

components: frozenset[str] | None = None

dlo_host_registration_limit_gib class-attribute instance-attribute

dlo_host_registration_limit_gib: float = 0.0

dlo_resident_layers class-attribute instance-attribute

dlo_resident_layers: int = 0

dlo_transfers class-attribute instance-attribute

dlo_transfers: dict[str, DLOTransfer] | None = None

dlo_use_allgather class-attribute instance-attribute

dlo_use_allgather: bool = True

dp_size class-attribute instance-attribute

dp_size: int = 1

pin_cpu_memory class-attribute instance-attribute

pin_cpu_memory: bool = True

strategy instance-attribute

strategy: OffloadStrategy

use_hsdp class-attribute instance-attribute

use_hsdp: bool = False

from_od_config classmethod

from_od_config(
    od_config: OmniDiffusionConfig,
) -> OffloadConfig

Extract and validate offload settings from OmniDiffusionConfig.

diffusion_offload_config is the canonical public selector. The historical enable_*_offload booleans remain compatibility aliases; ambiguous combinations fail instead of using silent precedence.

The dp_size is automatically derived from parallel_config — it is NOT a user-configurable parameter. The distributed layerwise offload works with whatever DP/SP parallelism is already set up.

Parameters:

Name Type Description Default
od_config OmniDiffusionConfig

OmniDiffusionConfig with offload settings

required

Returns:

Type Description
OffloadConfig

OffloadConfig with validated settings

offloads

offloads(component: str) -> bool

offloads_encoder

offloads_encoder(
    name: str, plan: OffloadPlan | None = None
) -> bool

Return whether the selector covers a discovered encoder path.

Plans declare non-standard encoder names explicitly. The name-based fallback preserves compatibility with pipelines that predate OffloadPlan.

should_offload_encoder

should_offload_encoder(
    name: str, plan: OffloadPlan | None = None
) -> bool

Apply explicit selection while preserving the legacy encoder topology.

transfer_for

transfer_for(component: str) -> DLOTransfer

uses_allgather

uses_allgather(component: str) -> bool

OffloadPlan dataclass

Optional declarative metadata for component layerwise offload.

Models declare this as a class attribute _offload_plan on the pipeline class. When present, both layerwise backends use it instead of model-specific backend branches, making new integrations data-driven.

If not declared, the offloader falls back to _layerwise_offload_blocks_attrs on each DiT module class.

:func:~vllm_omni.diffusion.offloader.plan_resolver.resolve_offload_plan is the only consumer; backends read the resolved artifact it returns.

Attributes:

Name Type Description
block_attrs dict[str, tuple[str, ...]]

Maps DiT path → tuple of block-list attribute names. e.g. {"transformer": ("gen_layers",), "transformer.language_model": ("layers",)}

offload_submodules dict[str, str]

Maps child name → block-list attribute name, for large non-DiT submodules within a DiT that should be independently offloaded with their own hooks. e.g. {"context_encoder": "layers"}

resident_dit_paths frozenset[str]

DiT paths whose leading blocks may be kept on the device when dlo_resident_layers is nonzero. Keeping this model-declared avoids applying a consumer-GPU tuning knob to auxiliary or dual DiTs unintentionally.

encoder_component_types dict[str, str]

Maps encoder paths to public selector types (currently text_encoder). This declaration is used before the compatibility name heuristic.

encoder_block_attrs dict[str, tuple[str, ...]]

Maps encoder paths to streamable block-list paths.

encoder_dlo_weight_replication frozenset[str]

Encoder paths whose loader-produced block tensors are identical across the DiT DLO group. Only these encoders may use multi-rank AllGather transfer; this must not be declared for encoder-TP shards.

block_attrs class-attribute instance-attribute

block_attrs: dict[str, tuple[str, ...]] = field(
    default_factory=dict
)

encoder_block_attrs class-attribute instance-attribute

encoder_block_attrs: dict[str, tuple[str, ...]] = field(
    default_factory=dict
)

encoder_component_types class-attribute instance-attribute

encoder_component_types: dict[str, str] = field(
    default_factory=dict
)

encoder_dlo_weight_replication class-attribute instance-attribute

encoder_dlo_weight_replication: frozenset[str] = field(
    default_factory=frozenset
)

offload_submodules class-attribute instance-attribute

offload_submodules: dict[str, str] = field(
    default_factory=dict
)

on_demand_component_paths class-attribute instance-attribute

on_demand_component_paths: frozenset[str] = field(
    default_factory=frozenset
)

resident_dit_paths class-attribute instance-attribute

resident_dit_paths: frozenset[str] = field(
    default_factory=frozenset
)

OffloadStartupState dataclass

Loader-owned state consumed by the offloader startup boundary.

allow_fresh_retry class-attribute instance-attribute

allow_fresh_retry: bool = False

fresh_model_loader class-attribute instance-attribute

fresh_model_loader: Callable[[], Module] | None = None

host_weight_plan class-attribute instance-attribute

host_weight_plan: HostWeightPlan | None = None

close_loader_ownership

close_loader_ownership() -> None

Release a plan that never reached a backend.

OffloadStrategy

Bases: str, Enum

Resolved internal backend strategy.

DISTRIBUTED_LAYER_WISE class-attribute instance-attribute

DISTRIBUTED_LAYER_WISE = 'distributed_layer_wise'

LAYER_WISE class-attribute instance-attribute

LAYER_WISE = 'layer_wise'

MODEL_LEVEL class-attribute instance-attribute

MODEL_LEVEL = 'model_level'

NONE class-attribute instance-attribute

NONE = 'none'

PinnedModuleStager

Stage immutable module groups without copying device weights to CPU.

nn.Module.to("cpu") performs a device-to-host copy for every parameter after each forward. In inference the weights are immutable, so retain one pinned CPU master instead. load materializes device storage from that master; offload only rebinds Parameters and buffers to the master.

A module iterable is treated as one staging group. It uses one copy stream and one reusable completion event. Tensors sharing storage keep their shapes, strides, offsets, dtypes, and aliases across every transition.

cache_retention instance-attribute

cache_retention = cache_retention

copy_stream instance-attribute

copy_stream = (
    copy_stream
    if copy_stream is not None
    else current_omni_platform.Stream()
)

device instance-attribute

device = device

loaded instance-attribute

loaded = False

load

load() -> None

offload

offload() -> None

set_cache_retention

set_cache_retention(
    cache_retention: BoundedAllocatorCache | None,
) -> None

ResolvedComponent dataclass

One pipeline component with its resolved, validated offload topology.

selected means the active selector covers this component. on_demand means the pipeline owns its residency through load_to_device / offload_to_cpu; a VAE staged by a legacy model plan is on_demand without being selectable through the public component grammar.

children class-attribute instance-attribute

children: tuple[ResolvedComponent, ...] = ()

module instance-attribute

module: Module

on_demand class-attribute instance-attribute

on_demand: bool = False

path instance-attribute

path: str

selected instance-attribute

selected: bool

stacks class-attribute instance-attribute

stacks: tuple[BlockStack, ...] = ()

ResolvedOffloadPlan dataclass

Backend-neutral topology for one pipeline under one offload config.

components property

components: tuple[ResolvedComponent, ...]

dits instance-attribute

encoders instance-attribute

encoders: tuple[ResolvedComponent, ...]

residents instance-attribute

residents: tuple[ResolvedComponent, ...]

skip_reason class-attribute instance-attribute

skip_reason: str | None = None

vaes instance-attribute

SupportsModelCpuOffload

Bases: Protocol

Pipeline-owned lifecycle for model-level CPU offload.

Pipelines with non-forward component entry points (for example VAE decode_latent methods) need to activate those stages explicitly, so generic forward-hook discovery cannot manage their full lifecycle.

disable_omni_model_cpu_offload

disable_omni_model_cpu_offload() -> None

enable_omni_model_cpu_offload

enable_omni_model_cpu_offload(
    *,
    device: device,
    pin_memory: bool,
    use_hsdp: bool,
    offload_components: frozenset[str] | None = None,
) -> None

apply_distributed_block_hook

apply_distributed_block_hook(
    module: Module,
    next_block: Module,
    device: device,
    dp_group: ProcessGroup | None,
    dp_size: int,
    rank: int,
    copy_stream: Any | None = None,
    comm_stream: Any | None = None,
    pin_memory: bool = True,
    shared_buffers: list[dict[dtype, Tensor] | None]
    | None = None,
    rank_local_mmap: bool = False,
    tensor_transforms: dict[int, Any] | None = None,
    materialization_probe_tensor: Tensor | None = None,
) -> DistributedLayerwiseOffloadHook

Register a DistributedLayerwiseOffloadHook on module.

apply_sequential_offload

apply_sequential_offload(
    dit_modules: list[Module],
    encoder_modules: list[Module],
    device: device,
    pin_memory: bool = True,
    use_hsdp: bool = False,
    offload_initial_dits: bool = False,
    offload_dit_modules: Collection[Module] | None = None,
    offload_encoder_modules: Collection[Module]
    | None = None,
) -> None

Apply sequential offloading hooks to DiT and encoder modules.

Registers hooks on modules to implement mutual-exclusion GPU allocation. - Before DiT runs, encoders are offloaded to CPU. - Before encoders run, DiT is offloaded to CPU.

Parameters:

Name Type Description Default
dit_modules list[Module]

DiT/transformer modules to register hooks on

required
encoder_modules list[Module]

Encoder modules to register hooks on

required
device device

Target GPU device for loading

required
pin_memory bool

Whether to pin CPU memory for faster transfers

True
use_hsdp bool

Whether HSDP is enabled (affects non_blocking behavior)

False
offload_initial_dits bool

Whether to begin with all DiT modules on CPU.

False
offload_dit_modules Collection[Module] | None

DiT modules allowed to move to CPU. None selects every DiT for backward compatibility.

None
offload_encoder_modules Collection[Module] | None

Encoder/stage modules allowed to move to CPU. None selects every supplied module for backward compatibility.

None
Example

apply_sequential_offload( ... dit_modules=[pipeline.transformer], ... encoder_modules=[pipeline.text_encoder, pipeline.vae], ... device=torch.device("cuda:0"), ... )

Modules of pipeline now automatically swap between CPU and GPU

dtype_size

dtype_size(dtype: dtype) -> int

Return element size in bytes for a torch.dtype.

enable_offload_backend

enable_offload_backend(
    od_config: OmniDiffusionConfig,
    pipeline: Module,
    device: device | None = None,
) -> tuple[Module, OffloadBackend | None]

Create and enable the loader-selected backend transactionally.

The model runner only calls this generic offloader boundary. Loader-owned host plans and optional fresh-model recovery callbacks stay inside the startup state consumed here.

get_blocks_attr_names

get_blocks_attr_names(model: Module) -> list[str]

Get block attribute names from model class.

get_blocks_from_dit

get_blocks_from_dit(
    model: Module,
    block_attrs: tuple[str, ...] | None = None,
) -> tuple[list[str], list[Module]]

Retrieve blocks from an explicit plan or the DiT class metadata.

get_offload_backend

get_offload_backend(
    od_config: OmniDiffusionConfig,
    device: device | None = None,
    host_weight_plan: HostWeightPlan | None = None,
) -> OffloadBackend | None

Create appropriate offload backend based on configuration.

Parameters:

Name Type Description Default
od_config OmniDiffusionConfig

OmniDiffusionConfig with offload settings

required
device device | None

Target device (auto-detected if None)

None
host_weight_plan HostWeightPlan | None

Exact loader-produced backing plan, if ordinary weight materialization was skipped.

None

Returns:

Type Description
OffloadBackend | None

OffloadBackend instance or None if offloading disabled

Example

backend = get_offload_backend(od_config, device=torch.device("cuda:0")) if backend: ... backend.enable(pipeline)

get_offload_plan

get_offload_plan(pipeline: Module) -> OffloadPlan | None

Retrieve the OffloadPlan declared by the pipeline, if any.

is_dtensor

is_dtensor(t: Tensor) -> bool

Check if tensor is a DTensor.

is_materialized_tensor

is_materialized_tensor(t: Tensor) -> bool

Check if tensor holds real data (not meta or empty placeholder).

make_offload_placeholder

make_offload_placeholder(tensor: Tensor) -> Tensor

Create a zero-element placeholder to free GPU memory.

remove_distributed_block_hook

remove_distributed_block_hook(module: Module) -> None

Remove the distributed layerwise offload hook from module.

remove_sequential_offload

remove_sequential_offload(modules: list[Module]) -> None

Remove sequential offloading hooks from modules.

Parameters:

Name Type Description Default
modules list[Module]

Modules to remove hooks from

required
Example

all_modules = [dit_modules, encoder_modules] remove_sequential_offload(all_modules)

resolve_offload_plan

resolve_offload_plan(
    pipeline: Module, config: OffloadConfig
) -> ResolvedOffloadPlan

Resolve and validate one pipeline's offload topology.

Raises ValueError for every topology that the requested configuration cannot serve, before the caller places a module or installs a hook.

sequential_offload_component

sequential_offload_component(
    module: Module,
) -> Iterator[None]

Activate and release a hooked component called outside forward.

set_blocks_attr_names

set_blocks_attr_names(
    model: Module, names: list[str]
) -> None

set_tensor_storage

set_tensor_storage(target: Tensor, value: Tensor) -> None

Replace target's underlying storage with value (zero-copy).

take_offload_startup_state

take_offload_startup_state(
    model: Module,
) -> OffloadStartupState | None

Take and remove the loader handoff from a pipeline exactly once.