From 90815cc86632562192156a79446beaf250e66021 Mon Sep 17 00:00:00 2001 From: Hongsheng Liu Date: Sun, 6 Sep 2026 16:44:52 +0800 Subject: [PATCH] [Core] Add opt-in serialized HWR build admission Signed-off-by: Hongsheng Liu --- docs/design/feature/host_weight_runtime.md | 10 + docs/design/module/host_weight_runtime.md | 38 ++- .../test_build_admission.py | 305 ++++++++++++++++++ vllm_omni/host_weight_runtime/config.py | 6 + .../host_weight_runtime/filesystem/store.py | 39 ++- 5 files changed, 394 insertions(+), 4 deletions(-) create mode 100644 tests/host_weight_runtime/test_build_admission.py diff --git a/docs/design/feature/host_weight_runtime.md b/docs/design/feature/host_weight_runtime.md index a9bb1b0b838..955f72cf5e2 100644 --- a/docs/design/feature/host_weight_runtime.md +++ b/docs/design/feature/host_weight_runtime.md @@ -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 diff --git a/docs/design/module/host_weight_runtime.md b/docs/design/module/host_weight_runtime.md index ac92711fc4d..6107f93be27 100644 --- a/docs/design/module/host_weight_runtime.md +++ b/docs/design/module/host_weight_runtime.md @@ -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 diff --git a/tests/host_weight_runtime/test_build_admission.py b/tests/host_weight_runtime/test_build_admission.py new file mode 100644 index 00000000000..8bd06bbb6cf --- /dev/null +++ b/tests/host_weight_runtime/test_build_admission.py @@ -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] diff --git a/vllm_omni/host_weight_runtime/config.py b/vllm_omni/host_weight_runtime/config.py index 0bacd69d200..6ab6c7de771 100644 --- a/vllm_omni/host_weight_runtime/config.py +++ b/vllm_omni/host_weight_runtime/config.py @@ -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): @@ -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 ( @@ -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) diff --git a/vllm_omni/host_weight_runtime/filesystem/store.py b/vllm_omni/host_weight_runtime/filesystem/store.py index 0f01e36dac5..ae98a1ca859 100644 --- a/vllm_omni/host_weight_runtime/filesystem/store.py +++ b/vllm_omni/host_weight_runtime/filesystem/store.py @@ -627,7 +627,7 @@ def _initialize_domain(self) -> None: self.policy_version = 1 def _capacity_document(self, *, policy_version: int) -> dict[str, object]: - return { + document: dict[str, object] = { "schema_version": 1, "policy_version": policy_version, "max_artifact_bytes": self.capacity_policy.max_artifact_bytes, @@ -635,6 +635,9 @@ def _capacity_document(self, *, policy_version: int) -> dict[str, object]: "min_free_bytes": self.capacity_policy.min_free_bytes, "eviction": self.capacity_policy.eviction, } + if self.capacity_policy.build_admission == "serialized": + document.update(schema_version=2, build_admission="serialized") + return document def _build_lock_path(self, key: str) -> Path: return self.locks_dir / f"{key}.build.lock" @@ -1013,6 +1016,40 @@ def get_or_build( if initial.status in {StoreStatus.TIMEOUT, StoreStatus.FAILED}: return initial + if self.capacity_policy.build_admission == "concurrent": + return self._build_after_miss(identity, producer, validation=validation, deadline=deadline) + try: + admission = FileLock(self.locks_dir / "domain-build.lock", exclusive=True, deadline=deadline) + except FileLockTimeoutError as exc: + return StoreResult( + StoreStatus.TIMEOUT, + failure=_failure( + ResolutionStage.PRODUCTION, FailureCode.ACTIVE_BUILD_TIMEOUT, str(exc), retryable=True + ), + ) + except OSError as exc: + return StoreResult( + StoreStatus.FAILED, + failure=_failure(ResolutionStage.DOMAIN, FailureCode.DOMAIN_UNAVAILABLE, str(exc), retryable=True), + ) + try: + return self._build_after_miss(identity, producer, validation=validation, deadline=deadline) + finally: + try: + admission.close() + except OSError: + # Releasing coordination must not revise the completed build's + # publication outcome or replace its primary failure. + logger.warning("Failed to release domain build admission under %s", self.root, exc_info=True) + + def _build_after_miss( + self, + identity: WeightArtifactIdentity, + producer: WeightProducer, + *, + validation: ValidationLevel, + deadline: float, + ) -> StoreResult: try: build_lock = FileLock(self._build_lock_path(identity.key), exclusive=True, deadline=deadline) except FileLockTimeoutError as exc: