Skip to content
Open
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
10 changes: 10 additions & 0 deletions docs/design/feature/host_weight_runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -313,6 +313,16 @@ 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.

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 filesystem lock acquisition. It does not
cancel synchronous validation, a producer that has already started, or atomic
publication. A hung in-process producer therefore blocks its owning process and
Expand Down
38 changes: 35 additions & 3 deletions docs/design/module/host_weight_runtime.md
Original file line number Diff line number Diff line change
Expand Up @@ -311,9 +311,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
305 changes: 305 additions & 0 deletions tests/host_weight_runtime/test_build_admission.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,305 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM-Omni project
"""CPU process and policy coverage for domain-wide producer admission."""

from __future__ import annotations

import errno
import json
import multiprocessing as mp
import os
import time
from multiprocessing.connection import Connection
from pathlib import Path
from typing import Literal

import pytest
import torch

from tests.host_weight_runtime.test_filesystem_store import FakeProducer, _identity, _make_store, _publish_test_artifact
from vllm_omni.host_weight_runtime import (
BuildRequest,
CapacityPolicy,
FailureCode,
HostWeightError,
HostWeightRuntime,
HostWeightRuntimeConfig,
ProductionMetadata,
ResolutionOutcome,
RuntimeMode,
StoreResult,
StoreStatus,
TensorWriteSpec,
ValidationLevel,
WaitPolicy,
WeightArtifactIdentity,
WeightProductionSpec,
)
from vllm_omni.host_weight_runtime.filesystem.locks import FileLock, lock_is_active
from vllm_omni.host_weight_runtime.protocols import ArtifactWriter

pytestmark = [pytest.mark.core_model, pytest.mark.cpu]


class AdmissionProducer:
def __init__(self, identity: WeightArtifactIdentity, control: Connection | None = None) -> None:
self._spec = FakeProducer(identity).spec
self.control = control

@property
def spec(self) -> WeightProductionSpec:
return self._spec

def produce(self, writer: ArtifactWriter) -> ProductionMetadata:
action = "produce"
if self.control is not None:
self.control.send("entered")
if not self.control.poll(60):
raise TimeoutError("test did not release producer")
action = self.control.recv()
spec = TensorWriteSpec("payload", (65536,), torch.bfloat16)
with writer.open_tensor_file("weights.safetensors", (spec,)) as output:
output.write_tensor("payload", torch.arange(65536, dtype=torch.float32).to(torch.bfloat16))
if action == "crash":
os._exit(27)
return ProductionMetadata("test-producer-v1", "test-restorer-v1")


def _admission_process(root: Path, policy: CapacityPolicy, rank: int, control: Connection) -> None:
try:
store = _make_store(root, capacity=policy)
identity = _identity(tp_rank=rank)
original_lookup = store.lookup
first_lookup = True

def notify_lookup(
identity: WeightArtifactIdentity, *, validation: ValidationLevel, deadline: float | None = None
) -> StoreResult:
nonlocal first_lookup
result = original_lookup(identity, validation=validation, deadline=deadline)
if first_lookup:
# Prove the caller observed a miss before the parent lets the
# active producer publish, independent of process scheduling.
control.send("looked_up")
first_lookup = False
return result

store.lookup = notify_lookup # type: ignore[method-assign]
control.send("ready")
assert control.poll(60) and control.recv() == "start"
result = store.get_or_build(
BuildRequest(identity),
AdmissionProducer(identity, control),
validation=ValidationLevel.FULL_CHECKSUM,
deadline=time.monotonic() + 30,
)
if result.lease is not None:
result.lease.close()
control.send((result.status.value, result.failure.code.value if result.failure is not None else None))
finally:
control.close()


@pytest.mark.parametrize("admission", ["concurrent", "serialized"])
@pytest.mark.parametrize("same_key", [False, True])
def test_cross_identity_admission_and_capacity(
tmp_path: Path, admission: Literal["concurrent", "serialized"], same_key: bool
) -> None:
policy = CapacityPolicy(max_store_bytes=192 * 1024, build_admission=admission)
store = _make_store(tmp_path / "store", capacity=policy)
ctx = mp.get_context("spawn")
pairs = [ctx.Pipe() for _ in range(2)]
processes = [
ctx.Process(target=_admission_process, args=(store.root, policy, 0 if same_key else rank, pair[1]))
for rank, pair in enumerate(pairs)
]
try:
for process, pair in zip(processes, pairs, strict=True):
process.start()
pair[1].close()
first, second = [pair[0] for pair in pairs]
for connection in (first, second):
assert connection.poll(60) and connection.recv() == "ready"
first.send("start")
assert first.poll(10) and first.recv() == "looked_up"
assert first.poll(10) and first.recv() == "entered"
second.send("start")
assert second.poll(10) and second.recv() == "looked_up"
if admission == "concurrent" and not same_key:
assert second.poll(10) and second.recv() == "entered"
else:
assert not second.poll(0.2), "a second producer entered while coordination should exclude it"
first.send("produce")
assert first.poll(30) and first.recv() == (StoreStatus.BUILT.value, None)
if same_key:
assert second.poll(30) and second.recv() == (StoreStatus.JOINED.value, None)
else:
if admission == "serialized":
assert second.poll(10) and second.recv() == "entered"
second.send("produce")
assert second.poll(30)
assert second.recv() == (StoreStatus.FAILED.value, FailureCode.STORE_LIMIT_EXCEEDED.value)
assert policy.max_store_bytes is not None
assert store.inspect_domain().store_bytes < policy.max_store_bytes
assert not list(store.tmp_dir.iterdir())
assert {path.name for path in store.artifacts_dir.iterdir()} == {_identity(tp_rank=0).key}
for process in processes:
process.join(10)
assert process.exitcode == 0
finally:
for process in processes:
if process.pid is not None:
if process.is_alive():
process.kill()
process.join(5)
process.close()
for pair in pairs:
for connection in pair:
connection.close()


def test_admission_timeout_preserves_owner_and_warm_hit_bypasses_lock(tmp_path: Path) -> None:
policy = CapacityPolicy(max_store_bytes=512 * 1024, build_admission="serialized")
store = _make_store(tmp_path / "store", capacity=policy)
identity, _ = _publish_test_artifact(store)
other = _identity(tp_rank=1)
lock_path = store.locks_dir / "domain-build.lock"
with FileLock(lock_path, exclusive=True, deadline=None):
hit = store.get_or_build(
BuildRequest(identity),
FakeProducer(identity),
validation=ValidationLevel.FULL_CHECKSUM,
deadline=time.monotonic() + 0.02,
)
assert hit.status is StoreStatus.HIT and hit.lease is not None
hit.lease.close()
timeout = store.get_or_build(
BuildRequest(other),
FakeProducer(other),
validation=ValidationLevel.FULL_CHECKSUM,
deadline=time.monotonic() + 0.02,
)
assert timeout.status is StoreStatus.TIMEOUT
assert timeout.failure is not None and timeout.failure.code is FailureCode.ACTIVE_BUILD_TIMEOUT
assert timeout.failure.retryable
assert lock_is_active(lock_path)
_publish_test_artifact(store, other)


def test_producer_failure_releases_admission(tmp_path: Path) -> None:
store = _make_store(
tmp_path / "store", capacity=CapacityPolicy(max_store_bytes=16384, build_admission="serialized")
)
identity = _identity()
failed = store.get_or_build(
BuildRequest(identity),
FakeProducer(identity, write_mode="incomplete"),
validation=ValidationLevel.FULL_CHECKSUM,
deadline=time.monotonic() + 5,
)
assert failed.status is StoreStatus.FAILED
assert not lock_is_active(store.locks_dir / "domain-build.lock")
_publish_test_artifact(store, identity)


@pytest.mark.parametrize(
("mode", "expected"),
[(RuntimeMode.PREFERRED, ResolutionOutcome.CANONICAL_FALLBACK), (RuntimeMode.REQUIRED, ResolutionOutcome.FAILED)],
)
def test_admission_timeout_obeys_runtime_mode(tmp_path: Path, mode: RuntimeMode, expected: ResolutionOutcome) -> None:
policy = CapacityPolicy(max_store_bytes=16384, build_admission="serialized")
store = _make_store(tmp_path / "store", capacity=policy)
runtime = HostWeightRuntime.from_config(
HostWeightRuntimeConfig(mode=mode, domain=store.domain_policy, capacity=policy, wait=WaitPolicy(0.01))
)
identity = _identity()
with FileLock(store.locks_dir / "domain-build.lock", exclusive=True, deadline=None):
resolution = runtime.resolve(identity, producer=FakeProducer(identity))
assert resolution.report.outcome is expected
failure = resolution.report.attempts[-1].failure
assert failure is not None and failure.code is FailureCode.ACTIVE_BUILD_TIMEOUT


def test_crashed_producer_releases_admission_and_same_key_retry_recovers(tmp_path: Path) -> None:
policy = CapacityPolicy(max_store_bytes=192 * 1024, build_admission="serialized")
store = _make_store(tmp_path / "store", capacity=policy)
ctx = mp.get_context("spawn")
parent, child = ctx.Pipe()
process = ctx.Process(target=_admission_process, args=(store.root, policy, 0, child))
process.start()
child.close()
try:
assert parent.poll(60) and parent.recv() == "ready"
parent.send("start")
assert parent.poll(10) and parent.recv() == "looked_up"
assert parent.poll(10) and parent.recv() == "entered"
parent.send("crash")
process.join(10)
assert process.exitcode == 27
finally:
if process.is_alive():
process.kill()
process.join(5)
process.close()
parent.close()
assert not lock_is_active(store.locks_dir / "domain-build.lock")
assert list(store.tmp_dir.iterdir())
identity = _identity(tp_rank=0)
recovered = store.get_or_build(
BuildRequest(identity),
AdmissionProducer(identity),
validation=ValidationLevel.FULL_CHECKSUM,
deadline=time.monotonic() + 10,
)
assert recovered.status is StoreStatus.BUILT and recovered.lease is not None
assert torch.equal(recovered.lease.tensors["payload"], torch.arange(65536, dtype=torch.float32).to(torch.bfloat16))
recovered.lease.close()
assert not list(store.tmp_dir.iterdir())


def test_admission_release_error_preserves_publication(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None:
store = _make_store(
tmp_path / "store", capacity=CapacityPolicy(max_store_bytes=16384, build_admission="serialized")
)
original_close = FileLock.close

def fail_admission_close(lock: FileLock) -> None:
original_close(lock)
if lock.path.name == "domain-build.lock":
raise OSError(errno.EIO, "injected admission release failure")

monkeypatch.setattr(FileLock, "close", fail_admission_close)
_publish_test_artifact(store)
assert not lock_is_active(store.locks_dir / "domain-build.lock")


def test_admission_policy_is_authoritative_and_legacy_document_is_unchanged(tmp_path: Path) -> None:
concurrent = CapacityPolicy(max_store_bytes=16384)
serialized = CapacityPolicy(max_store_bytes=16384, build_admission="serialized")
legacy = _make_store(tmp_path / "legacy", capacity=concurrent)
path = legacy.root / "domain-policy.json"
original = path.read_bytes()
assert json.loads(original) == {
"schema_version": 1,
"policy_version": 1,
"max_artifact_bytes": None,
"max_store_bytes": 16384,
"min_free_bytes": 0,
"eviction": "none",
}
with pytest.raises(HostWeightError, match="incompatible schema"):
_make_store(legacy.root, capacity=serialized)
assert path.read_bytes() == original
strict = _make_store(tmp_path / "serialized", capacity=serialized)
document = json.loads((strict.root / "domain-policy.json").read_bytes())
assert document["schema_version"] == 2 and document["build_admission"] == "serialized"
assert _make_store(strict.root, capacity=serialized).domain_uuid == strict.domain_uuid
with pytest.raises(HostWeightError, match="incompatible schema"):
_make_store(strict.root, capacity=concurrent)


def test_admission_policy_validation() -> None:
with pytest.raises(ValueError, match="requires max_store_bytes"):
CapacityPolicy(build_admission="serialized")
with pytest.raises(ValueError, match="must be concurrent or serialized"):
CapacityPolicy(build_admission="unknown") # type: ignore[arg-type]
6 changes: 6 additions & 0 deletions vllm_omni/host_weight_runtime/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from dataclasses import dataclass, field
from enum import Enum
from pathlib import Path
from typing import Literal


class RuntimeMode(str, Enum):
Expand Down Expand Up @@ -66,6 +67,7 @@ class CapacityPolicy:
max_store_bytes: int | None = None
min_free_bytes: int = 0
eviction: str = "none"
build_admission: Literal["concurrent", "serialized"] = "concurrent"

def __post_init__(self) -> None:
for name, value in (
Expand All @@ -78,6 +80,10 @@ def __post_init__(self) -> None:
raise ValueError("min_free_bytes must not be negative")
if not isinstance(self.eviction, str) or self.eviction != "none":
raise ValueError("automatic host weight eviction is not implemented")
if self.build_admission not in ("concurrent", "serialized"):
raise ValueError("build_admission must be concurrent or serialized")
if self.build_admission == "serialized" and self.max_store_bytes is None:
raise ValueError("serialized build admission requires max_store_bytes")


@dataclass(frozen=True)
Expand Down
Loading
Loading