Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
57 changes: 57 additions & 0 deletions benchmarks/host_weight_runtime/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
# Host weight dependency memory diagnostic

`safetensors_retention.py` isolates repeated CPU `get_tensor()` calls from
access to one cached tensor view. It imports neither vLLM nor HWR and uses a
synthetic float32 payload controlled by `--tensor-elements` (default: 64).
Tensor size can affect retention: a negative result for 16 elements does not
rule out growth for 64 elements. The selected size is recorded in the JSON
arguments and used for both file creation and the correctness check.
HWR itself calls `get_tensor()` when
acquiring a lease, then exposes cached views; this probe does not measure HWR
requests or repeated lease acquisition.

Run one mode per fresh process in the same Python environment and CPU affinity:

```bash
git rev-parse HEAD > revision.txt
git status --short > working-tree.txt
git diff > working-tree.patch
# Select an allowed CPU from your process affinity; 0 is only an example.
CUDA_VISIBLE_DEVICES='' taskset -c 0 timeout 120s python \
benchmarks/host_weight_runtime/safetensors_retention.py \
--mode get_tensor --iterations 10000 --sample-every 5000 > feasibility.json
```

For a size comparison, run 16 and 64 elements twice each in fresh processes,
with cached reuse as a control at both sizes:

```bash
for repetition in 1 2; do
for elements in 16 64; do
for mode in get_tensor reuse; do
CUDA_VISIBLE_DEVICES='' taskset -c 0 timeout 120s python \
benchmarks/host_weight_runtime/safetensors_retention.py \
--mode "$mode" --tensor-elements "$elements" \
> "$mode-$elements-$repetition.json"
done
done
done
```

The diagnostic sets Torch to one CPU thread, warms up before the baseline,
and reports preparation separately from timed loop work. Checkpoints run GC
and report Linux private/anonymous memory, RSS, and open descriptors. It checks
that transient tensor objects and the final cached view are released. Linux
`/proc/self/smaps_rollup` must be readable; invalid arguments and missing proc
support fail explicitly. Keep raw JSON and the repository snapshot with results.

Compare the final loop sample with the warmed baseline, then inspect the closed
sample. Private-memory growth after Python tensor objects disappear is evidence
of retention in the dependency/allocator path, not proof of a leak's root cause.
RSS can include file-backed pages; it is not a private-memory metric. Samples
also include small diagnostic bookkeeping allocations. Do not assert a portable
memory threshold or turn this synthetic loop into a per-request service estimate.

Do not drop shared caches, trim allocators, or add serving-time GC to make the
numbers look smaller. A dependency replacement or version change requires its
own reproducible comparison and correctness/lifecycle validation.
120 changes: 120 additions & 0 deletions benchmarks/host_weight_runtime/safetensors_retention.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,120 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM-Omni project
"""Linux CPU diagnostic: repeated safetensors materialization versus view reuse."""

from __future__ import annotations

import argparse
import gc
import json
import os
import platform
import tempfile
import time
import weakref
from pathlib import Path


def positive_int(value: str) -> int:
parsed = int(value)
if parsed <= 0:
raise argparse.ArgumentTypeError("must be positive")
return parsed


def sample(phase: str, iterations: int, loop_seconds: float) -> dict[str, object]:
gc.collect()
counters = {}
for line in Path("/proc/self/smaps_rollup").read_text().splitlines():
fields = line.split()
if len(fields) == 3 and fields[2] == "kB":
counters[fields[0].removesuffix(":")] = int(fields[1]) * 1024
return {
"phase": phase,
"iterations": iterations,
"loop_seconds": loop_seconds,
"rss_bytes": counters["Rss"],
"private_bytes": counters["Private_Clean"] + counters["Private_Dirty"],
"anonymous_bytes": counters["Anonymous"],
"fd_count": len(tuple(Path("/proc/self/fd").iterdir())),
}


def main() -> None:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--mode", choices=("get_tensor", "reuse"), required=True)
parser.add_argument(
"--tensor-elements", type=positive_int, default=64, help="Number of float32 elements in the synthetic tensor"
)
parser.add_argument("--iterations", type=positive_int, default=1_000_000)
parser.add_argument("--sample-every", type=positive_int, default=100_000)
parser.add_argument("--warmup", type=positive_int, default=1000)
args = parser.parse_args()
if not Path("/proc/self/smaps_rollup").is_file():
parser.error("requires Linux /proc/self/smaps_rollup")

preparation_start = time.monotonic()
import safetensors
import torch
from safetensors import safe_open
from safetensors.torch import save_file

torch.set_num_threads(1)
torch.set_num_interop_threads(1)
samples = [sample("before_file", 0, 0.0)]
with tempfile.TemporaryDirectory(prefix="safetensors-retention-") as directory:
path = Path(directory) / "synthetic.safetensors"
save_file({"weight": torch.arange(args.tensor_elements, dtype=torch.float32)}, str(path))
with safe_open(path, framework="pt", device="cpu") as reader:
cached = reader.get_tensor("weight")
assert torch.equal(cached, torch.arange(args.tensor_elements, dtype=torch.float32))
for _ in range(args.warmup):
tensor = reader.get_tensor("weight") if args.mode == "get_tensor" else cached
del tensor
samples.append(sample("baseline", 0, 0.0))
preparation_seconds = time.monotonic() - preparation_start
completed = 0
loop_seconds = 0.0
while completed < args.iterations:
count = min(args.sample_every, args.iterations - completed)
start = time.monotonic()
for _ in range(count):
tensor = reader.get_tensor("weight") if args.mode == "get_tensor" else cached
last_tensor = weakref.ref(tensor)
del tensor
loop_seconds += time.monotonic() - start
completed += count
snapshot = sample("loop", completed, loop_seconds)
snapshot["last_tensor_alive"] = last_tensor() is not None
if args.mode == "get_tensor":
assert last_tensor() is None, "diagnostic retained a transient tensor"
samples.append(snapshot)
del cached
del reader
samples.append(sample("closed", completed, loop_seconds))
assert last_tensor() is None, "last tensor survived reader/view teardown"

print(
json.dumps(
{
"schema_version": 1,
"arguments": vars(args),
"python": platform.python_version(),
"torch": torch.__version__,
"safetensors": safetensors.__version__,
"kernel": platform.release(),
"hostname": platform.node(),
"pid": os.getpid(),
"cpu_affinity": sorted(os.sched_getaffinity(0)),
"torch_threads": torch.get_num_threads(),
"pythonmalloc": os.environ.get("PYTHONMALLOC", "default"),
"preparation_seconds": preparation_seconds,
"samples": samples,
},
indent=2,
)
)


if __name__ == "__main__":
main()
46 changes: 44 additions & 2 deletions docs/design/feature/host_weight_runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,11 @@ V1 does not include:

## Motivation and use cases

For investigating CPU memory retained by dependency tensor materialization,
see the [standalone safetensors diagnostic](../../../benchmarks/host_weight_runtime/README.md).
It distinguishes repeated dependency calls from reuse of cached views and does
not establish per-request HWR leakage.

Model loading can create the same final host representation repeatedly. This
is especially expensive when loading performs checkpoint decoding, tensor
renaming, TP slicing, quantization, packing, or scale construction before GPU
Expand All @@ -67,6 +72,26 @@ views and mapped ranges. A separate transport still decides whether to use
registered mmap, private pinned staging, synchronous copies, or asynchronous
H2D transfer.

## Payload validation at startup

Eligible diffusion loaders accept `--host-weight-runtime-validation` (or the
`host_weight_runtime_validation` offline/stage configuration field):

- `manifest_and_metadata` is the existing default. It validates identity,
metadata, tensor structure, and sizes, but does not detect payload-only
corruption.
- `full_checksum` reads and hashes every payload on warm acquisition before
restoring weights. This adds startup I/O/CPU work proportional to artifact
bytes; it is not a per-inference check.

For example, add `--host-weight-runtime-validation full_checksum` to an enabled
HWR deployment. Preferred mode falls back to canonical loading and can publish
a replacement after detecting corruption; required mode fails startup. Readers
can use different validation levels in the same domain without changing artifact
identity. Neither level authenticates an untrusted writer or prevents mutation
after validation. Checksumming does not imply automatic recovery of a running
service.

## Resolution behavior

The loader resolves the immutable canonical source and computes the exact
Expand Down Expand Up @@ -313,8 +338,25 @@ One process per exact identity owns a build; other workers wait and then acquire
leases for the published artifact. Publication is invisible until all payloads
and metadata are validated, hashed, fsynced, and atomically renamed.

`coordination_timeout_seconds` bounds filesystem lock acquisition. It does not
cancel synchronous validation, a producer that has already started, or atomic
By default, different identities may build concurrently. A domain configured
with `CapacityPolicy(build_admission="serialized", max_store_bytes=...)` admits
only one producer/publication operation at a time. Warm hits remain concurrent,
and admission uses the existing coordination deadline. All workers must agree
on the serialized domain's schema-2 capacity policy; use a new root when opting
in, since in-place migration of an active domain is unsupported. This prevents
competing producers from racing allocation checks but does not provide a
filesystem-wide byte quota or automatic stale-data reclamation. See the
[module capacity contract](../module/host_weight_runtime.md#validation-and-capacity).

`coordination_timeout_seconds` bounds domain-initialization and lookup/build
lock acquisition. Store construction and each later resolution or publication
operation have separate budgets from the same wait policy, rather than one
end-to-end startup deadline. A domain-init timeout follows the retryable domain
failure policy below. After contention ends, a fresh runtime construction can
retry initialization; the timed-out runtime retains its original failure.

The coordination budget does not cancel filesystem I/O, synchronous validation,
a producer that has already started, or atomic
publication. A hung in-process producer therefore blocks its owning process and
must be handled by external process supervision. Enforceable producer
cancellation requires a future process-isolated producer contract.
Expand Down
15 changes: 15 additions & 0 deletions docs/design/feature/offloader/distributed_layerwise_offload.md
Original file line number Diff line number Diff line change
Expand Up @@ -289,13 +289,28 @@ store, tensor ownership, or H2D payload. On success, each tensor view copies
directly into the existing rotating HBM block buffers and the two private host
staging slots are not allocated.

`--dlo-host-registration-mode` defaults to `auto`, preserving this behavior.
Select `disabled` to bypass HWR mapping registration and use bounded host staging
regardless of mapping size or budget. This leaves the existing pin-memory policy
for staging buffers unchanged. The equivalent offline/stage configuration field
is `dlo_host_registration_mode`; an explicit disable requires enabled HWR with
no-AllGather DLO.

`--dlo-host-registration-limit-gib` is an optional per-worker preflight ceiling
over page-aligned registered bytes. Zero adds no ceiling. A disabled pinned
memory policy, unsupported platform/capability, over-budget mapping, or fully
rolled-back registration error selects the existing two-slot staging path. A
partial registration that cannot be rolled back aborts startup because closing
the lease would unmap memory still owned by the platform.

Automatic registration logs its entry with device and budget, then each CUDA
region's ordinal/total and byte size before calling `cudaHostRegister`. A
completion message follows each successful call. These identify the last
entered boundary if progress stops; they do not impose a timeout. A synchronous
CUDA call cannot be safely cancelled by a Python timer, and a stuck process
still requires external supervision. The explicit disable option avoids this
registration path; it does not repair the underlying driver hang.

Direct checkpoint mmap remains unchanged and continues to use staging. It may
require loader-owned per-block transforms, while the HWR artifact already
contains final runtime bytes. DLO AllGather never receives an HWR final-layout
Expand Down
87 changes: 82 additions & 5 deletions docs/design/module/host_weight_runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -136,8 +136,16 @@ so only when it is marked retryable; unsupported capabilities, identity
collisions, producer failures, and publication failures remain visible as
startup failures.

`coordination_timeout_seconds` bounds acquisition of lookup and build locks; it
is not a hard wall-clock deadline for in-process validation, production, or
`coordination_timeout_seconds` bounds acquisition of the domain-initialization
lock and lookup/build locks. Store construction and each later resolution or
publication operation receive separate coordination budgets from the same
`WaitPolicy`; they do not share an end-to-end startup deadline. A domain-init
lock timeout is a retryable domain failure: preferred mode may fall back and
required mode fails. Retrying initialization requires constructing a fresh
runtime after contention ends; a waiter never cancels or unlocks the owner.

The coordination budget is not a hard wall-clock deadline for filesystem I/O,
in-process validation, production, or
atomic publication. Once this V1 implementation becomes the producer, the
synchronous producer and publication run to completion. A hung producer blocks
its owning process and requires external process supervision; enforceable
Expand Down Expand Up @@ -299,6 +307,43 @@ Explicit cleanup uses the locked move and retries exact-key
`.cleanup.*` tombstones left by an interrupted removal; a later build performs
the same tombstone reconciliation before producing replacement content.

Failure-quarantined entries remain until explicitly selected for removal;
ordinary `cleanup(identity)` does not remove that diagnostic history.
`FilesystemHostWeightStore.cleanup_quarantined(storage_name)` accepts one
quarantined inventory basename and preserves the current artifact, deny marker,
and other failure-quarantined entries. It takes the same nonblocking build and
exclusive artifact locks, so an active builder or lease returns a typed refusal.
Malformed names are rejected before mutation, and symlinks are not followed.

Selected entries move through the existing `.cleanup.*` tombstone protocol.
Retrying the original name reconciles pending cleanup tombstones for that key,
even if the original entry disappeared during an earlier attempt. A missing
selection is otherwise an idempotent success. Existing cleanup intents for the
key may be completed along with the selection; other quarantine history is
retained. Parent synchronization is retried even when an earlier deletion
removed the last tombstone before its sync failed.

For a configured store, operators can inspect and deliberately select an entry:

```python
from vllm_omni.host_weight_runtime import ArtifactInventoryState, HostWeightError

for entry in store.inspect_domain().inventory:
if entry.state is ArtifactInventoryState.QUARANTINED:
print(entry.storage_name, entry.size_bytes)

selected_storage_name = input("Quarantined entry to remove: ")
failure = store.cleanup_quarantined(selected_storage_name)
if failure is not None:
raise HostWeightError(failure)
```

Use the same authoritative domain and capacity policy as the service. Review
inventory sizes and inspect again after cleanup before explicitly retrying a
build or startup. These are logical-byte, point-in-time observations; concurrent
activity can affect the difference and physical disk/RAM reclamation is not
implied. This operation adds no automatic retention policy or service retry.

An operational failure while inspecting a noncooperative competing publication
is a retryable storage failure. The competing entry remains authoritative and
is not mislabeled as corrupt or quarantined merely because its manifest could
Expand All @@ -311,9 +356,41 @@ point-in-time evidence only; kernel locks are not a persistent owner registry.
Capacity policy covers all store-owned bytes, including ready, temporary, and
quarantined data. The local writer preflights `max_artifact_bytes`,
`max_store_bytes`, and `min_free_bytes`, and preallocates payloads where the
filesystem supports it. `ENOSPC` remains a normal typed store failure because
concurrent preflight is inherently racy. Automatic eviction and strict
concurrent reservations are not implemented.
filesystem supports it. `ENOSPC` remains a normal typed store failure.

`CapacityPolicy.build_admission` defaults to `concurrent`, preserving parallel
production for different identities and best-effort capacity preflight. An
explicit `serialized` policy, which requires `max_store_bytes`, admits one
cooperative build/publication operation per domain through
`locks/domain-build.lock`. Production lock order becomes domain admission,
per-key build, then artifact. Validated warm hits bypass admission; queued
same-key callers still recheck and join a completed artifact. Admission waiting
uses the caller's existing coordination deadline and returns a retryable build
timeout without interrupting its owner.

The kernel releases admission on process exit. An explicit same-key retry uses
the existing stale-temp recovery; files from unrelated dead builds continue to
count toward capacity. A synchronous producer that hangs still requires
external process supervision. Admission release errors are logged without
revising an already determined publication result.

Concurrent domains retain their exact schema-1 capacity document. Serialized
domains use a schema-2 capacity document with `build_admission=serialized`, so
mixed policies and old readers fail initialization rather than bypass admission.
Opt into serialization using a new root with an agreed policy; in-place policy
migration while old workers may exist is unsupported. For example:

```python
from vllm_omni.host_weight_runtime import CapacityPolicy

capacity = CapacityPolicy(max_store_bytes=192 * 1024**3, build_admission="serialized")
```

Serialization prevents cooperative producers from racing each other's
allocation checks. It is not a filesystem quota: control metadata, finalization,
and outside writers can still affect total usage. Automatic eviction and
parallel byte reservations are not implemented. This policy trades parallel
cold-build throughput for predictable producer admission.

For tmpfs, artifact bytes are host-memory consumption and may consume swap;
they must not be reported operationally as ordinary disk capacity. The store
Expand Down
Loading
Loading