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. |
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.
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.
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.
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.
cpu_staging_buffers instance-attribute ¶
gpu_shard_buffers instance-attribute ¶
is_materialized property ¶
is_materialized: bool
Check whether this block's parameters hold real data on device.
offload_layer ¶
Free GPU memory for current block by replacing tensors with placeholders.
prefetch_layer ¶
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 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.
ModelLevelOffloadBackend ¶
Bases: OffloadBackend
Model-level (sequential) offloading backend.
Uses SequentialOffloadHook registered via HookRegistry for automatic module swapping.
OffloadBackend ¶
Bases: ABC
Base class for CPU offload backends
disable abstractmethod ¶
Disable offloading and cleanup resources.
Removes all registered hooks. Does NOT move modules back to original devices (caller responsible for that).
enable abstractmethod ¶
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 |
shutdown ¶
Release offload resources at process exit.
Backends may skip work that only matters for a later enable.
OffloadConfig dataclass ¶
dlo_host_registration_limit_gib class-attribute instance-attribute ¶
dlo_host_registration_limit_gib: float = 0.0
dlo_transfers class-attribute instance-attribute ¶
dlo_transfers: dict[str, DLOTransfer] | None = None
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_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.
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. |
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. |
resident_dit_paths | frozenset[str] | DiT paths whose leading blocks may be kept on the device when |
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 ¶
encoder_block_attrs class-attribute instance-attribute ¶
encoder_component_types class-attribute instance-attribute ¶
encoder_dlo_weight_replication class-attribute instance-attribute ¶
offload_submodules class-attribute instance-attribute ¶
on_demand_component_paths class-attribute instance-attribute ¶
OffloadStartupState dataclass ¶
Loader-owned state consumed by the offloader startup boundary.
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 ¶
Release a plan that never reached a backend.
OffloadStrategy ¶
Resolved internal backend strategy.
DISTRIBUTED_LAYER_WISE class-attribute instance-attribute ¶
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.
copy_stream instance-attribute ¶
copy_stream = (
copy_stream
if copy_stream is not None
else current_omni_platform.Stream()
)
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.
ResolvedOffloadPlan dataclass ¶
Backend-neutral topology for one pipeline under one offload config.
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.
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¶
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 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_materialized_tensor ¶
is_materialized_tensor(t: Tensor) -> bool
Check if tensor holds real data (not meta or empty placeholder).
make_offload_placeholder ¶
Create a zero-element placeholder to free GPU memory.
remove_distributed_block_hook ¶
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_tensor_storage ¶
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.