diff --git a/docs/features/kv_offloading_usage.md b/docs/features/kv_offloading_usage.md index 22a2350d8095..5b72abf762c1 100644 --- a/docs/features/kv_offloading_usage.md +++ b/docs/features/kv_offloading_usage.md @@ -108,6 +108,7 @@ vllm serve \ | `store_threshold` | no | `0` | single-tier | Min lookups before a block is offloaded. Values ≥ 2 are rejected by `TieringOffloadingSpec`. | | `max_tracker_size` | no | `64000` | single-tier | Max entries in the lookup tracker. | | `secondary_tiers` | no | `[]` | multi-tier | List of secondary tier configs (see below). | +| `tier_poll_interval_s` | no | `0.001` | multi-tier | Pause between control-plane rounds for tiers that must be serviced between engine steps (currently `p2p`). See [Threading model](#threading-model). Set to `0` to disable the thread and service those tiers once per engine step instead. | | `offload_prompt_only` | no | `true` | both | If `true`, only prompt (prefill) blocks are offloaded; decode blocks are skipped. | | `self_describing_kv_events` | no | `false` | both | Opt-in. When `true` *and* KV cache events are enabled (`--kv-events-config` with `enable_kv_cache_events`), the connector emits self-describing block-granular `BlockStored`/`BlockRemoved` payloads (constituent block hashes, whole-chunk `token_ids`, per-block `block_size`, parent hash, LoRA + group/cache-spec metadata) instead of the placeholder fallback, so external KV-event consumers can index offloaded blocks. Inert unless events are enabled. With `TieringOffloadingSpec`, a CPU promotion is self-describing when a local request observes its primary-tier `HIT` before event translation; otherwise its stored event may retain the placeholder, while a later `HIT` can backfill metadata for removal. Pending-removal/re-promotion races and externally initiated promotions may also produce placeholders, and consumers must ignore removals for unknown hashes. Partial recurrent tails emit the hash-aligned portion from the physical block start through the tail boundary. Other sliding-window/SSM chunks keep the placeholder fallback. In chunk mode (`block_size` > GPU block size, or `blocks_per_chunk` > 1), overlapping chunks re-announce shared per-block hashes, so consumers must reference-count (deduplicate) repeated store/remove announcements. | | `spec_module_path` | no | — | both | Python import path for a custom `OffloadingSpec` not in the built-in registry. Required only when `spec_name` is not built-in (advanced). | @@ -150,6 +151,19 @@ The filesystem and object-store tiers can publish hash-only `BlockStored` KV eve Set the optional `locality` tier field to `LOCAL` or `REMOTE` to describe the tier's storage location relative to the publishing vLLM instance. `LOCAL` marks storage local to that instance, while `REMOTE` marks storage that is not local to it. When the setting is omitted, locality is unspecified. vLLM does not infer it from the tier type, so an `obj` tier is not implicitly `REMOTE`. A KV event includes `locality` only when the tier explicitly configures it. This metadata describes the tier property without implying that a consumer can already route requests to its blocks. +### Threading model + +Tier methods run in the scheduler process, under the tiering manager's lock, so they are never entered concurrently. They are not always called from the same thread, though. + +Most tiers only ever see the scheduler thread, which holds the lock for the length of one engine step. A tier whose counterpart is not driven by the engine — the `p2p` tier, answering a remote peer — needs servicing while the engine is busy running the model, because otherwise a peer's lookup or fetch waits for the next step boundary. On a saturated rank that wait can be seconds, dwarfing the transfer it gates. Such tiers set `serves_external_requests`, and the tiering manager runs one thread that, once per `tier_poll_interval_s`, polls them for finished jobs and lets them serve external requests. A tier that implements `serve_external_requests()` without setting the flag is served once per engine step only. + +Two consequences worth knowing: + +- Any tier can be reached from that thread indirectly, because serving a peer looks keys up through the tiering manager, which fans out to the other tiers. Tier state must therefore not assume a particular thread, only that the manager lock is held. +- Work the thread deliberately leaves alone is batched per step: a promotion it initiates is submitted at the next `on_schedule_end()`, not immediately. + +Setting `tier_poll_interval_s` to `0` disables the thread; opted-in tiers then get serviced once per engine step, as they were before it existed. + ### Filesystem (FS) The filesystem tier (`type: "fs"`) writes blocks to a filesystem directory. diff --git a/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py b/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py index e71c3dbf18d4..48b6982d76fd 100644 --- a/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py +++ b/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py @@ -1017,6 +1017,129 @@ def test_scheduler_reports_lookup_sync_delay(request_runner): assert reduced[f"{_ConnectorMetricName.LOOKUP_SYNC_DELAY}_sum"] > 0 +def test_on_schedule_end_is_the_last_manager_call_of_a_step(request_runner): + """on_schedule_end must run after every other manager call of the step. + + Callers bound a step's exclusive access to the manager by releasing at + on_schedule_end, so anything issued after it — notably prepare_store, which + evicts — would fall outside that window. + """ + block_size = 4 + runner = request_runner( + block_size=block_size, + num_gpu_blocks=8, + async_scheduling=False, + ) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output([]) + ) + + connector_scheduler = runner.connector_scheduler + build_meta = connector_scheduler.build_connector_meta + windows: list[list[str]] = [] + + def spy(scheduler_output): + first = len(runner.manager.mock_calls) + meta = build_meta(scheduler_output) + # Entering and exiting manager.lock brackets the step; only the state + # operations it guards are ordered relative to on_schedule_end. + windows.append( + [ + name + for name, _, _ in runner.manager.mock_calls[first:] + if not name.startswith("lock.") + ] + ) + return meta + + connector_scheduler.build_connector_meta = spy + + runner.new_request(token_ids=[0] * (block_size * 2)) + runner.run(decoded_tokens=[EOS_TOKEN_ID]) + + assert windows, "build_connector_meta was never called" + assert any("prepare_store" in w for w in windows), ( + "step must exercise the calls that follow on_schedule_end today" + ) + assert any("on_request_finished" in w for w in windows) + for names in windows: + assert names[-1] == "on_schedule_end", ( + f"on_schedule_end must be the last manager call of a step, got {names}" + ) + + +def test_step_releases_the_manager_before_the_model_runs(request_runner): + """The step-scoped region must close when build_connector_meta returns. + + Everything after it -- the model future above all -- runs with the manager + free, which is the only window another thread has to use it. + """ + runner = request_runner( + block_size=4, + num_gpu_blocks=10, + async_scheduling=False, + ) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output([]) + ) + + runner.new_request(token_ids=[1] * 4) + runner.run(decoded_tokens=[EOS_TOKEN_ID]) + + scheduler = runner.connector_scheduler + assert scheduler._step_region is None + lock = runner.manager.lock + assert lock.__enter__.call_count > 0 + assert lock.__enter__.call_count == lock.__exit__.call_count + + +def test_has_pending_push_work_does_not_open_the_step_region(request_runner): + """It must take a short region, never the step-scoped one. + + The engine asks this before schedule() and returns early when the answer is + False, so a step-scoped acquire here would never reach the release in + build_connector_meta and would hold the manager for good. + """ + runner = request_runner( + block_size=4, + num_gpu_blocks=10, + async_scheduling=False, + ) + scheduler = runner.connector_scheduler + runner.manager.has_pending_work.return_value = False + + assert scheduler.has_pending_push_work() is False + assert scheduler._step_region is None + lock = runner.manager.lock + assert lock.__enter__.call_count == lock.__exit__.call_count + + +def test_take_events_drains_the_manager_inside_the_region(request_runner): + """take_events must not hold the manager across its own iteration. + + manager.take_events() is a generator whose backing list is cleared once + exhausted, so draining it lazily would both pin the manager for as long as + the consumer takes and lose events appended in the meantime. + """ + runner = request_runner( + block_size=4, + num_gpu_blocks=10, + async_scheduling=False, + ) + scheduler = runner.connector_scheduler + runner.manager.take_events.return_value = iter(()) + runner.manager.reset_mock() + + events = list(scheduler.take_events()) + + assert events == [] + assert runner.manager.take_events.call_count == 1 + assert scheduler._step_region is None + lock = runner.manager.lock + assert lock.__enter__.call_count == 1 + assert lock.__exit__.call_count == 1 + + def test_scheduler_reports_lookup_async_delay_on_resolve(request_runner): """A deferred lookup reports its async delay once it resolves.""" runner = request_runner( diff --git a/tests/v1/kv_offload/tiering/test_control_plane_thread.py b/tests/v1/kv_offload/tiering/test_control_plane_thread.py new file mode 100644 index 000000000000..7358565f53e9 --- /dev/null +++ b/tests/v1/kv_offload/tiering/test_control_plane_thread.py @@ -0,0 +1,406 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the vLLM project +"""Tests for OffloadingManager.lock and the tiering control-plane thread. + +The control plane of a tier that answers something the engine does not drive -- +a remote peer, say -- used to advance only from on_schedule_end(), so a peer's +request waited for a model-step boundary. These tests cover the thread that +services such tiers between steps and the exclusion that makes it safe. +""" + +import threading +import time +from collections.abc import Iterable +from contextlib import nullcontext +from typing import ClassVar +from unittest.mock import MagicMock + +import pytest + +from vllm.v1.kv_offload.base import ( + LookupResult, + Medium, + OffloadingManager, + OffloadKey, + ReqContext, + RequestOffloadingContext, + ScheduleEndContext, +) +from vllm.v1.kv_offload.cpu.manager import CPUOffloadingManager +from vllm.v1.kv_offload.tiering.base import ( + JobResult, + ParentManager, + SecondaryTierManager, + TransferJob, +) +from vllm.v1.kv_offload.tiering.manager import ( + CPUPrimaryTierOffloadingManager, + TieringOffloadingManager, +) + +from .test_tiering_offloading import _mock_mmap_region + +_CTX = ReqContext(req_id="test") +_EMPTY_SCHEDULE_END = ScheduleEndContext(new_req_ids=(), preempted_req_ids=()) + +# Long enough that a 1ms-interval thread gets many rounds, short enough to keep +# the suite quick. +_SETTLE_S = 0.5 + + +class _FakeTier(SecondaryTierManager): + """Tier that records which thread serviced it, and when.""" + + medium: ClassVar[Medium] = Medium.CPU + serves_external_requests: ClassVar[bool] = True + + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + # Counters are only ever incremented, and int writes are atomic under + # the GIL, so the test thread may read them without the manager lock. + self.serves = 0 + self.polls = 0 + self.serve_threads: set[int] = set() + self.served = threading.Event() + self.shutdown_saw_live_thread: bool | None = None + self.on_serve = None + + def lookup(self, key: OffloadKey, req_context: ReqContext) -> LookupResult: + return LookupResult.MISS + + def submit_store(self, job_metadata: TransferJob) -> None: + pass + + def submit_load(self, job_metadata: TransferJob) -> None: + pass + + def get_finished_jobs(self) -> Iterable[JobResult]: + self.polls += 1 + return () + + def on_new_request(self, req_context: ReqContext) -> RequestOffloadingContext: + return RequestOffloadingContext() + + def drain_jobs(self) -> None: + pass + + def serve_external_requests(self, parent: ParentManager) -> None: + self.serves += 1 + self.serve_threads.add(threading.get_ident()) + if self.on_serve is not None: + self.on_serve(parent) + self.served.set() + + def shutdown(self) -> None: + self.shutdown_saw_live_thread = any( + t.name == "vllm_offload_control_plane" and t.is_alive() + for t in threading.enumerate() + ) + + +class _QuietTier(_FakeTier): + """Tier that does not want servicing between steps.""" + + serves_external_requests: ClassVar[bool] = False + + +def _make_manager(tier_cls=_FakeTier, num_chunks: int = 8, **kwargs): + primary = CPUPrimaryTierOffloadingManager( + num_chunks=num_chunks, + mmap_region=_mock_mmap_region(num_chunks), + ) + tier = tier_cls( + offloading_spec=MagicMock(), + primary_kv_view=primary.get_kv_memoryview(), + tier_type="fake", + ) + manager = TieringOffloadingManager( + primary_tier=primary, + secondary_tiers=[tier], + **kwargs, + ) + return manager, tier + + +@pytest.fixture +def manager_and_tier(request): + kwargs = getattr(request, "param", {}) + manager, tier = _make_manager(**kwargs) + try: + yield manager, tier + finally: + manager.shutdown() + + +def _live_control_threads() -> list[threading.Thread]: + return [ + t + for t in threading.enumerate() + if t.name == "vllm_offload_control_plane" and t.is_alive() + ] + + +def test_control_thread_and_scheduler_never_overlap(manager_and_tier): + """The lock must keep the two threads out of the manager at the same time. + + Both sides must also make progress: a thread that never runs would pass a + mutual-exclusion check trivially, and so would one that starves the + scheduler. + """ + manager, tier = manager_and_tier + state = {"inside": False, "overlaps": 0} + + def critical_section(): + if state["inside"]: + state["overlaps"] += 1 + state["inside"] = True + time.sleep(0.002) + state["inside"] = False + + tier.on_serve = lambda parent: critical_section() + + steps = 0 + deadline = time.monotonic() + _SETTLE_S + while time.monotonic() < deadline: + with manager.lock: + critical_section() + steps += 1 + + assert state["overlaps"] == 0 + assert steps > 0, "scheduler thread was starved by the control thread" + assert tier.serves > 0, "control thread never ran" + + +def test_serves_while_the_scheduler_is_outside_the_lock(manager_and_tier): + """The point of the thread: serve during the model-execution window. + + The engine releases the manager at the end of a step and then blocks on the + model future. Here the "scheduler" does the same -- one step, then a sleep + with the lock free -- and the tier must be serviced during that sleep. + """ + manager, tier = manager_and_tier + + with manager.lock: + manager.on_schedule_end(_EMPTY_SCHEDULE_END) + + tier.served.clear() + serves_before = tier.serves + + # Stand in for future.result(): the engine thread is blocked, holding + # nothing. + assert tier.served.wait(_SETTLE_S), ( + "tier was not serviced while the scheduler held no lock" + ) + assert tier.serves > serves_before + assert tier.polls > 0 + + scheduler_thread_id = threading.get_ident() + assert scheduler_thread_id not in tier.serve_threads or len(tier.serve_threads) > 1 + + +@pytest.mark.parametrize( + "manager_and_tier", [{"tier_poll_interval_s": 0.0}], indirect=True +) +def test_zero_interval_disables_the_thread(manager_and_tier): + """The escape hatch back to per-step servicing. + + Without the thread the tier is serviced only by on_schedule_end, which is + the behaviour the control plane had before it existed. + """ + manager, tier = manager_and_tier + + assert not _live_control_threads() + assert not tier.served.wait(0.05) + assert tier.serves == 0 + + with manager.lock: + manager.on_schedule_end(_EMPTY_SCHEDULE_END) + + assert tier.serves == 1 + + +def test_no_thread_when_no_tier_asks_for_one(): + """A tier that does not opt in must not cost a thread.""" + manager, tier = _make_manager(tier_cls=_QuietTier) + try: + assert not _live_control_threads() + finally: + manager.shutdown() + + +def test_tier_that_did_not_opt_in_is_not_serviced_off_thread(): + """Only opted-in tiers are polled and served by the control thread. + + A tier that expects the scheduler thread must not be dragged onto another + one just because it shares a manager with a tier that opted in. Fan-out + through ParentManager (lookup, on_new_request, on_request_finished) can + still reach it from that thread; that is the documented exception. + """ + num_chunks = 8 + primary = CPUPrimaryTierOffloadingManager( + num_chunks=num_chunks, + mmap_region=_mock_mmap_region(num_chunks), + ) + common = { + "offloading_spec": MagicMock(), + "primary_kv_view": primary.get_kv_memoryview(), + } + opted_in = _FakeTier(tier_type="fake", **common) + quiet = _QuietTier(tier_type="quiet", **common) + manager = TieringOffloadingManager( + primary_tier=primary, + secondary_tiers=[opted_in, quiet], + ) + try: + assert opted_in.served.wait(_SETTLE_S) + assert quiet.serves == 0 + assert quiet.polls == 0 + + with manager.lock: + manager.on_schedule_end(_EMPTY_SCHEDULE_END) + + # on_schedule_end still serves every tier, so the thread dying degrades + # to per-step servicing rather than to none. + assert quiet.serves == 1 + finally: + manager.shutdown() + + +def test_per_step_gate_is_not_touched_by_the_control_thread(manager_and_tier): + """_processed_jobs_this_step stays owned by the scheduler thread. + + Sharing it would let a fast step skip a poll it would otherwise do, pushing + a completion out by a whole step. + """ + manager, tier = manager_and_tier + + assert tier.served.wait(_SETTLE_S) + polls_after_rounds = tier.polls + assert polls_after_rounds > 0 + assert manager._processed_jobs_this_step is False + + +class _LookingUpTier(_FakeTier): + """Opted-in tier that looks keys up through its parent while serving. + + That is what the p2p server role does with a peer's lookup, and the parent + lookup lands in TieringOffloadingManager.lookup(), which polls tiers for + finished jobs. + """ + + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self.in_serve = False + self.nested_polls = 0 + self.on_serve = self._look_up + + def _look_up(self, parent: ParentManager) -> None: + self.in_serve = True + try: + parent.lookup(OffloadKey(b"\x01" * 8), ReqContext(req_id="peer")) + finally: + self.in_serve = False + + def get_finished_jobs(self) -> Iterable[JobResult]: + if self.in_serve: + self.nested_polls += 1 + return super().get_finished_jobs() + + +def test_parent_lookup_while_serving_keeps_the_round_contained(): + """A lookup issued while serving must not escape the control round. + + Before the fix it ran the step's once-per-step poll from the control + thread: every tier got polled off-thread, the opted-in tier re-entrantly + from inside its own serve, and the gate was left set, so the next step + skipped its poll. + """ + num_chunks = 8 + primary = CPUPrimaryTierOffloadingManager( + num_chunks=num_chunks, + mmap_region=_mock_mmap_region(num_chunks), + ) + common = { + "offloading_spec": MagicMock(), + "primary_kv_view": primary.get_kv_memoryview(), + } + opted_in = _LookingUpTier(tier_type="fake", **common) + quiet = _QuietTier(tier_type="quiet", **common) + manager = TieringOffloadingManager( + primary_tier=primary, + secondary_tiers=[opted_in, quiet], + ) + try: + # Several rounds, each issuing a parent lookup between steps. + deadline = time.monotonic() + _SETTLE_S + while opted_in.serves < 3 and time.monotonic() < deadline: + time.sleep(0.01) + assert opted_in.serves >= 3, "control thread never served" + + with manager.lock: + assert manager._processed_jobs_this_step is False + assert opted_in.nested_polls == 0 + assert quiet.polls == 0, "quiet tier was polled off-thread" + + # The next step's first lookup still does its poll. + manager.lookup(OffloadKey(b"\x02" * 8), _CTX) + assert quiet.polls == 1 + finally: + manager.shutdown() + + +def test_shutdown_joins_the_control_thread_before_tier_teardown(manager_and_tier): + """The thread drives tier transports, so it must be stopped first. + + Closing a transport under a running round can take the process down rather + than raise. + """ + manager, tier = manager_and_tier + assert tier.served.wait(_SETTLE_S) + assert _live_control_threads() + + manager.shutdown() + + assert tier.shutdown_saw_live_thread is False + assert not _live_control_threads() + + +def test_lock_releases_on_error(manager_and_tier): + """Leaving a with-block on an exception must not leak the lock.""" + manager, _ = manager_and_tier + + with pytest.raises(RuntimeError), manager.lock: + raise RuntimeError("boom") + + assert manager.lock.acquire(timeout=0.5), "lock leaked after an error" + manager.lock.release() + + +def test_lock_excludes_other_threads(manager_and_tier): + """While one thread holds the lock, another cannot take it.""" + manager, _ = manager_and_tier + + with manager.lock: + acquired: list[bool] = [] + t = threading.Thread( + target=lambda: acquired.append(manager.lock.acquire(timeout=0.01)) + ) + t.start() + t.join(timeout=5.0) + assert acquired == [False] + + +def test_lock_defaults_to_no_op(): + """Managers entered from one thread only pay nothing for the lock. + + The base default must be a no-op context manager, while the CPU manager, + which runs under the plain offloading connector, keeps a real one. + """ + base_lock = OffloadingManager.lock.fget(MagicMock()) + assert isinstance(base_lock, nullcontext) + with base_lock, base_lock: + pass + + cpu = CPUOffloadingManager(num_chunks=4) + with cpu.lock: + assert not cpu.lock.acquire(blocking=False) diff --git a/vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py b/vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py index 806267dcf3dd..8b3ad7f2b3f1 100644 --- a/vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py +++ b/vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py @@ -2,7 +2,8 @@ # SPDX-FileCopyrightText: Copyright contributors to the vLLM project import time from collections import Counter -from collections.abc import Iterable, Sequence +from collections.abc import Iterable, Iterator, Sequence +from contextlib import ExitStack, contextmanager from dataclasses import dataclass, field, replace from itertools import chain, islice from typing import Any, NamedTuple @@ -636,6 +637,71 @@ def _sliding_window_sort_key(i: int) -> int: self._events_tracker = OffloadingEventsTracker(spec.kv_events_config) + # Holds the manager lock for the current schedule(), or None between steps. + # Read and written only by the scheduler thread, so it needs no + # synchronization of its own. Declared on the class so an instance built + # without __init__ still reads None rather than raising. + _step_region: ExitStack | None = None + + # ------------------------------------------------------------------ + # Manager exclusion + # + # The manager is entered from the scheduler thread and, for managers that + # run one, from a control-plane thread. Manager methods never take their + # own lock, so the exclusion boundary lives here. manager.lock may be a + # no-op for a manager that needs no exclusion; it is entered regardless. + # + # schedule() holds one region for the whole step: opened by the first + # manager-touching hook and closed at the end of build_connector_meta. + # That span is required, not tidiness -- get_num_new_matched_tokens + # establishes lookup() HITs that update_state_after_alloc later pins with + # prepare_load(), and prepare_store() evicts after both. Releasing between + # them would let another thread evict a block the step already treated as + # present. + # + # Hooks outside schedule() use a short region instead, so the manager is + # free while the engine runs the model. + # + # The step region outlives any single call, so it cannot be a with-block; + # an ExitStack carries the entered lock from open to close instead. + # ------------------------------------------------------------------ + + def _acquire_step(self) -> None: + """Open the step-scoped region. Idempotent within a step.""" + if self._step_region is not None: + return + region = ExitStack() + region.enter_context(self.manager.lock) + self._step_region = region + + def _release_step(self) -> None: + """Close the step-scoped region, letting other threads in.""" + region = self._step_region + if region is None: + return + self._step_region = None + region.close() + + @contextmanager + def _manager_locked(self) -> Iterator[None]: + """Hold the manager for the duration of one hook outside schedule(). + + None of these hooks can legitimately run inside the step-scoped region: + on_new_request comes from add_request, has_pending_push_work from the + engine's liveness check before schedule(), the rest from output + processing or a client RPC. So finding the region open here means a + previous schedule() raised between _acquire_step() and the release in + build_connector_meta. Recover instead of deadlocking. + """ + if self._step_region is not None: + logger.error( + "Offloading manager lock was left held by a previous " + "schedule(); releasing it before continuing." + ) + self._release_step() + with self.manager.lock: + yield + def _maybe_observe_lookup_async_delay( self, req_status: RequestOffloadState ) -> None: @@ -1027,7 +1093,8 @@ def _lookup( def on_new_request(self, request: Request) -> None: """Called when a new request is added to the scheduler.""" req_context = _create_req_context(request) - offloading_context = self.manager.on_new_request(req_context) + with self._manager_locked(): + offloading_context = self.manager.on_new_request(req_context) req_status = RequestOffloadState( config=self.config, req=request, @@ -1063,6 +1130,10 @@ def get_num_new_matched_tokens( (between scheduler steps). """ + # First manager-touching hook of a step: opens the step-scoped region, + # closed at the end of build_connector_meta. + self._acquire_step() + req_status = self._req_status[request.request_id] for group_state in req_status.group_states: group_state.block_ids.clear() @@ -1149,6 +1220,10 @@ def update_state_after_alloc( if num_external_tokens == 0: return + # prepare_load() below pins keys that get_num_new_matched_tokens() + # already saw as HIT, so both must sit inside the same region. + self._acquire_step() + req_status = self._req_status[request.request_id] num_locally_computed_tokens = req_status.num_locally_computed_tokens @@ -1845,13 +1920,24 @@ def _build_store_jobs( def build_connector_meta( self, scheduler_output: SchedulerOutput + ) -> KVConnectorMetadata: + # A step that scheduled nothing reaches build_connector_meta without + # any earlier manager touch, so open the region here too. + self._acquire_step() + try: + return self._build_connector_meta(scheduler_output) + finally: + # The step is over: hand the manager back. + self._release_step() + + def _build_connector_meta( + self, scheduler_output: SchedulerOutput ) -> KVConnectorMetadata: self._update_req_states(scheduler_output) schedule_end_context = ScheduleEndContext( new_req_ids=[req.req_id for req in scheduler_output.scheduled_new_reqs], preempted_req_ids=scheduler_output.preempted_req_ids or (), ) - self.manager.on_schedule_end(schedule_end_context) # Flush jobs for preempted requests. for req_id in scheduler_output.preempted_req_ids or (): @@ -1897,6 +1983,12 @@ def build_connector_meta( self._current_batch_load_jobs = {} self._current_batch_jobs_to_flush = set() self._current_batch_allocated_block_ids = set() + + # Last manager call of the step, after every prepare_store and + # on_request_finished above. on_schedule_end is documented as running at + # the end of the step, and callers rely on that to bound the step's + # exclusive access to the manager. + self.manager.on_schedule_end(schedule_end_context) return meta def has_pending_push_work(self) -> bool: @@ -1904,8 +1996,17 @@ def has_pending_push_work(self) -> bool: While True, build_connector_meta() and update_connector_output() continue to be called even when no requests are scheduled. + + Must use a short region, never the step-scoped one. The engine asks this + before schedule() and returns early when the answer is False, so a + step-scoped acquire here would never reach build_connector_meta's + release and would hold the manager for good -- which is precisely the + state an idle engine settles into. """ - return bool(self._jobs) or self.manager.has_pending_work() + if self._jobs: + return True + with self._manager_locked(): + return self.manager.has_pending_work() def update_connector_output(self, connector_output: KVConnectorOutput): """Update KVConnector state from worker-side connectors output. @@ -1915,6 +2016,10 @@ def update_connector_output(self, connector_output: KVConnectorOutput): connectors output. """ + with self._manager_locked(): + self._update_connector_output(connector_output) + + def _update_connector_output(self, connector_output: KVConnectorOutput): meta = connector_output.kv_connector_worker_meta if not isinstance(meta, OffloadingWorkerMetadata): assert meta is None @@ -1993,7 +2098,8 @@ def get_stats(self) -> OffloadingConnectorStats | None: stats = self._connector_stats self._connector_stats = OffloadingConnectorStats() - manager_stats = self.manager.get_stats() + with self._manager_locked(): + manager_stats = self.manager.get_stats() if manager_stats is not None: if stats is None: stats = manager_stats @@ -2022,8 +2128,9 @@ def request_finished( # Untracked request (offloading never started): no in-flight jobs, # nothing was deferred, so finalize immediately. req_context = _create_req_context(request) - self.manager.on_new_request(req_context) - self.manager.on_request_finished(req_context) + with self._manager_locked(): + self.manager.on_new_request(req_context) + self.manager.on_request_finished(req_context) return False, None self._maybe_observe_lookup_async_delay(req_status) @@ -2055,7 +2162,16 @@ def take_events(self) -> Iterable[KVCacheEvent]: the underlying :class:`OffloadingEvent` stream. """ - yield from self._events_tracker.take_events(self.manager.take_events()) + # Drain the manager eagerly, inside the region. manager.take_events() is + # itself a generator whose backing list is cleared once exhausted, so + # iterating it lazily would both hold the manager for as long as the + # consumer takes and let a concurrent writer append events that the + # clear then discards unpublished. + with self._manager_locked(): + manager_events = list(self.manager.take_events()) + # Still a generator, so callers testing the result for truthiness keep + # seeing a non-empty object the way they do today. + yield from self._events_tracker.take_events(manager_events) def reset_cache(self) -> None: """Reset the offloading manager cache, evicting all stored chunks.""" @@ -2064,6 +2180,13 @@ def reset_cache(self) -> None: assert not self._current_batch_jobs_to_flush assert not self._current_batch_allocated_block_ids + # Held for the whole reset. manager.reset_cache() drains tiers by + # polling them itself, so a control-plane thread locked out here loses + # nothing; it simply skips rounds until the reset completes. + with self._manager_locked(): + self._reset_cache() + + def _reset_cache(self) -> None: # Flush all in-flight jobs self._current_batch_jobs_to_flush.update(self._jobs.keys()) diff --git a/vllm/v1/kv_offload/base.py b/vllm/v1/kv_offload/base.py index 398bc269e6f2..54178bb0f091 100644 --- a/vllm/v1/kv_offload/base.py +++ b/vllm/v1/kv_offload/base.py @@ -4,6 +4,7 @@ from abc import ABC, abstractmethod from collections.abc import Collection, Iterable, Sequence +from contextlib import AbstractContextManager, nullcontext from dataclasses import dataclass, field from enum import Enum, auto from typing import TYPE_CHECKING, Any, ClassVar, NamedTuple, NewType, TypeVar @@ -198,6 +199,16 @@ class OffloadingEvent: as well as a list of blocks that were evicted as a result. complete_store() - marks a previous store as completed. Following this call, the given blocks will become loadable. + +Exclusion: + lock is a context manager guarding a manager's whole state subtree: the + manager itself, every tier it composes, and any front-end state those + tiers own. Manager methods never take the lock themselves, so a caller + holds it across as many operations as its invariants need -- notably a + lookup() HIT and the prepare_load() that pins the result, which must not + be separated by another thread's eviction. The default is a no-op; + managers that can be entered from more than one thread supply a real + lock. """ @@ -233,6 +244,26 @@ class OffloadingKVEventsConfig: class OffloadingManager(ABC): + @property + def lock(self) -> AbstractContextManager[Any]: + """Context manager granting exclusive access to this manager's state. + + Scope is the manager, every tier it composes, and any front-end state + those tiers own (e.g. a tier's async lookup front end). Every caller of + an OffloadingManager method must hold it. Manager methods never enter + it themselves, so a caller may hold it across an arbitrary sequence of + operations -- which is required wherever one call establishes a fact a + later call relies on, such as a lookup() HIT followed by the + prepare_load() that pins it. + + The default is a no-op, for managers that are only ever entered from + one thread. A manager that may be entered from another thread -- its + own background thread, say -- returns a real lock instead. Callers must + not assume either: always enter it, and never re-enter it while held, + since a real lock need not be reentrant. + """ + return nullcontext() + @abstractmethod def lookup(self, key: OffloadKey, req_context: ReqContext) -> LookupResult: """Checks whether a single block is offloaded and ready to be read. @@ -395,6 +426,10 @@ def take_events(self) -> Iterable[OffloadingEvent]: def on_schedule_end(self, context: ScheduleEndContext) -> None: """Called once at the end of each scheduler step. + This is the last manager call of a step: every lookup(), prepare_load(), + prepare_store() and on_request_finished() for the step has already been + issued. + Managers may override this to flush deferred work accumulated during the step (e.g., batched promotions). """ diff --git a/vllm/v1/kv_offload/cpu/manager.py b/vllm/v1/kv_offload/cpu/manager.py index 2c08ac897c6c..f2acf1d76ef7 100644 --- a/vllm/v1/kv_offload/cpu/manager.py +++ b/vllm/v1/kv_offload/cpu/manager.py @@ -1,5 +1,6 @@ # SPDX-License-Identifier: Apache-2.0 # SPDX-FileCopyrightText: Copyright contributors to the vLLM project +import threading from collections import OrderedDict from collections.abc import Collection, Iterable from dataclasses import dataclass, field @@ -64,6 +65,7 @@ def __init__( store_threshold: int = 1, max_tracker_size: int = 64_000, ): + self._lock = threading.Lock() self.medium: Medium = Medium.CPU self._num_chunks: int = num_chunks self._num_allocated_chunks: int = 0 @@ -89,6 +91,11 @@ def __init__( OrderedDict() if store_threshold >= 2 else None ) + @property + @override + def lock(self) -> threading.Lock: + return self._lock + # --- chunk pool --- def _get_num_free_chunks(self) -> int: diff --git a/vllm/v1/kv_offload/tiering/async_lookup.py b/vllm/v1/kv_offload/tiering/async_lookup.py index a6879d8c44e5..e57fb36bab5e 100644 --- a/vllm/v1/kv_offload/tiering/async_lookup.py +++ b/vllm/v1/kv_offload/tiering/async_lookup.py @@ -11,8 +11,12 @@ -------------- There is no explicit lock. Thread safety is achieved by ownership: -* _lookup_state and _lookup_batch are owned exclusively by the scheduler - thread. lookup(), flush(), and cleanup() read and write them directly. +* _lookup_state and _lookup_batch are owned by whichever thread holds the + OffloadingManager lock. lookup(), flush(), and cleanup() read and write + them directly. That is usually the scheduler thread, but a tier reached + through ParentManager -- a peer lookup fanning out from another tier, say -- + can arrive on the tiering manager's control-plane thread instead. Both hold + the manager lock, so the accesses stay serialized. * _lookup_queue is written by the scheduler (flush → put_nowait, one item per step) and read by the background thread (get). queue.Queue is diff --git a/vllm/v1/kv_offload/tiering/base.py b/vllm/v1/kv_offload/tiering/base.py index b855b6b000a5..e827d8232b19 100644 --- a/vllm/v1/kv_offload/tiering/base.py +++ b/vllm/v1/kv_offload/tiering/base.py @@ -139,6 +139,12 @@ class SecondaryTierManager(ABC): IMPORTANT: All methods run in the Scheduler process and must be lightweight and non-blocking. submit_load() and submit_store() submit async jobs; get_finished_jobs() polls for completion. + + Methods are called under the tiering manager's lock, so they are never + entered concurrently -- but not always from the same thread. A tier that + sets serves_external_requests is also serviced from the manager's + control-plane thread, and any tier can be reached from there via + ParentManager fan-out. Keep per-thread assumptions out of tier state. """ medium: ClassVar[Medium | None] = None @@ -154,6 +160,16 @@ def __init_subclass__(cls, **kwargs: Any) -> None: Medium.STORAGE: CacheHitSource.DISK, }[cls.medium] + # Whether this tier answers a counterpart the engine does not drive -- a + # remote peer, say -- whose requests would otherwise wait for a step + # boundary. The tiering manager then services the tier between engine + # steps too, not only from on_schedule_end(): on its polling thread, it + # calls both get_finished_jobs() and serve_external_requests(); see + # poll_tiers(). Overriding serve_external_requests() alone does not opt + # in: setting this asserts the tier tolerates those calls off the + # scheduler thread. + serves_external_requests: ClassVar[bool] = False + def __init__( self, offloading_spec: "OffloadingSpec", @@ -324,9 +340,20 @@ def on_request_finished(self, req_context: ReqContext) -> None: def serve_external_requests(self, parent: ParentManager) -> None: """Process remotely-originated requests using the parent manager. - Called once per scheduler step, BEFORE _flush_pending_promotions(). - The parent handle is valid only for the duration of this call. - Tiers that don't serve external requests leave this as a no-op. + Called once per scheduler step, BEFORE _flush_pending_promotions(), and + additionally once per control-plane round for tiers that set + serves_external_requests. The _flush_pending_promotions() ordering + holds only for the per-step call: a promotion this method initiates from + a control-plane round is submitted at the next on_schedule_end(). + + The parent handle is valid only for the duration of this call, and the + caller holds the manager lock throughout it -- so a lookup() HIT taken + here stays valid until this call pins it. An implementation must not + release that lock partway through. + + Tiers that don't serve external requests leave this as a no-op. A tier + that overrides it without setting serves_external_requests is served + once per step only. """ return diff --git a/vllm/v1/kv_offload/tiering/kvcr/manager.py b/vllm/v1/kv_offload/tiering/kvcr/manager.py index 4804f193e3e8..f5d30c916cdc 100644 --- a/vllm/v1/kv_offload/tiering/kvcr/manager.py +++ b/vllm/v1/kv_offload/tiering/kvcr/manager.py @@ -9,7 +9,7 @@ import uuid from collections.abc import Collection, Iterable from pathlib import Path -from typing import TYPE_CHECKING, Any, cast +from typing import TYPE_CHECKING, Any, ClassVar, cast import msgspec from kvcr import ( @@ -265,6 +265,12 @@ def decode(self, key: BlockKey) -> int | bytes: class KVCRSecondaryTierManager(SecondaryTierManager): """Secondary tier wrapper around the KVCR KV P2P API.""" + # Serves framework pin requests from serve_external_requests(), but only + # once per step for now: opting in would also poll the KVCR library from + # the tiering manager's polling thread, which it is not yet verified to + # tolerate. + serves_external_requests: ClassVar[bool] = False + @classmethod @override def build_metric_definitions( diff --git a/vllm/v1/kv_offload/tiering/manager.py b/vllm/v1/kv_offload/tiering/manager.py index 83ff2e2c1280..b136ebb11c65 100644 --- a/vllm/v1/kv_offload/tiering/manager.py +++ b/vllm/v1/kv_offload/tiering/manager.py @@ -19,6 +19,7 @@ protecting chunks from eviction until complete_read() is called """ +import threading import time from collections.abc import Collection, Iterable, Sequence from dataclasses import dataclass, field @@ -59,6 +60,26 @@ logger = init_logger(__name__) +# Default pause between control-plane rounds; the user-facing default of the +# ``tier_poll_interval_s`` extra-config key. A round is a non-blocking sweep, so +# this sets how long a peer's request can sit unserved; 1 ms is three orders of +# magnitude below the model step it hides, and matches the pause the p2p tier's +# own drain loops and EngineCoreProc's GIL-yield sleep already use. +DEFAULT_TIER_POLL_INTERVAL_S = 0.001 +# How long a round waits for the manager lock before giving up and retrying. +# Bounded so the thread stays responsive to shutdown even if a caller holds the +# lock for a long time, or leaked it. +_CONTROL_PLANE_LOCK_TIMEOUT_S = 1.0 +# Warn after this much continuous failure to acquire the lock. Above the 5 s +# warning drain_jobs() emits, so a slow reset_cache() does not warn twice. +_CONTROL_PLANE_STALL_WARN_S = 30.0 +# Give up after this many consecutive failed rounds. on_schedule_end() still +# serves the tiers, so the fallback is per-step servicing rather than none. +_CONTROL_PLANE_MAX_CONSECUTIVE_ERRORS = 100 +# How long shutdown() waits for the thread to finish its round. +_CONTROL_PLANE_JOIN_TIMEOUT_S = 5.0 + + @dataclass class PendingPromotion: """Accumulator for chunks awaiting submit_load() for one (tier, request).""" @@ -192,6 +213,7 @@ def __init__( self, primary_tier: CPUPrimaryTierOffloadingManager, secondary_tiers: list[SecondaryTierManager] | None = None, + tier_poll_interval_s: float = DEFAULT_TIER_POLL_INTERVAL_S, ): """Initialize the TieringOffloadingManager. @@ -199,8 +221,16 @@ def __init__( primary_tier: The primary tier manager (CPU-based). secondary_tiers: List of secondary tier managers (e.g., Storage, Network). Can be None or empty list. + tier_poll_interval_s: Pause between control-plane rounds for + tiers that set serves_external_requests. Zero or negative + disables the thread, leaving those tiers serviced once per + engine step as before. """ + # Guards this manager and every tier it composes, the primary tier + # included; the primary's own lock is not used underneath it. A real + # lock because the control-plane thread enters the manager too. + self._lock = threading.Lock() self.primary_tier: CPUPrimaryTierOffloadingManager = primary_tier self.secondary_tiers = secondary_tiers or [] @@ -248,6 +278,136 @@ def __init__( tier: i for i, tier in enumerate(self.secondary_tiers) } + # Tiers that serve external requests, and so must also be serviced + # between engine steps. + self._external_serving_tiers: list[tuple[int, SecondaryTierManager]] = [ + (i, tier) + for i, tier in enumerate(self.secondary_tiers) + if tier.serves_external_requests + ] + self._tier_poll_interval_s = tier_poll_interval_s + self._control_plane_stop = threading.Event() + self._control_plane_thread: threading.Thread | None = None + self._start_control_plane() + + @property + @override + def lock(self) -> threading.Lock: + return self._lock + + # ------------------------------------------------------------------ + # Control plane + # ------------------------------------------------------------------ + + def _start_control_plane(self) -> None: + """Start the control-plane thread, if any tier asked for one.""" + if not self._external_serving_tiers or self._tier_poll_interval_s <= 0: + return + self._control_plane_thread = threading.Thread( + target=self._control_plane_loop, + name="vllm_offload_control_plane", + daemon=True, + ) + self._control_plane_thread.start() + logger.info( + "KV offload control-plane thread started for tier(s) %s, interval %.3fs", + [tier.tier_type for _, tier in self._external_serving_tiers], + self._tier_poll_interval_s, + ) + + def _stop_control_plane(self) -> None: + """Stop and join the control-plane thread. + + Must run before any tier teardown: the thread drives tier transports, + and closing one underneath it can crash the process rather than raise. + """ + self._control_plane_stop.set() + thread = self._control_plane_thread + if thread is None: + return + thread.join(timeout=_CONTROL_PLANE_JOIN_TIMEOUT_S) + if thread.is_alive(): + logger.error( + "KV offload control-plane thread did not exit within %.1fs", + _CONTROL_PLANE_JOIN_TIMEOUT_S, + ) + else: + self._control_plane_thread = None + + def _control_plane_loop(self) -> None: + stalled_since: float | None = None + errors = 0 + while not self._control_plane_stop.is_set(): + if not self._lock.acquire(timeout=_CONTROL_PLANE_LOCK_TIMEOUT_S): + now = time.monotonic() + if stalled_since is None: + stalled_since = now + elif now - stalled_since > _CONTROL_PLANE_STALL_WARN_S: + logger.warning( + "KV offload control-plane thread has not acquired the " + "manager lock for %.0fs; peer requests are only being " + "served at engine-step boundaries.", + now - stalled_since, + ) + stalled_since = None + continue + stalled_since = None + try: + self.poll_tiers() + errors = 0 + except Exception: + # Never let a round kill the thread silently: on_schedule_end() + # still serves these tiers, so persistent failure degrades to + # per-step servicing instead of stopping the control plane. + errors += 1 + logger.exception("KV offload control-plane round failed") + finally: + self._lock.release() + if errors >= _CONTROL_PLANE_MAX_CONSECUTIVE_ERRORS: + logger.error( + "KV offload control-plane thread stopping after %d " + "consecutive failures; falling back to per-step servicing.", + errors, + ) + return + # Yield unconditionally, and only after release(). Python locks are + # not fair, so a release-then-reacquire loop could starve the + # scheduler thread; this pause is what bounds its wait to one round. + self._control_plane_stop.wait(self._tier_poll_interval_s) + + def poll_tiers(self) -> None: + """Advance the control plane of tiers that cannot wait for a step. + + Polls those tiers for finished jobs and lets them serve whatever their + counterpart has asked for -- the same two steps on_schedule_end() runs, + narrowed to the tiers that opted in. + + Excludes _flush_pending_promotions() and _flush_pending_cascades() by + design: both are per-step batching points, and running them here would + fragment batches the scheduler thread accumulated during schedule(). A + promotion that a lookup() from this path initiates therefore waits for + the next on_schedule_end() to be submitted. + + The caller must hold self.lock for the whole call. serve_external_requests() + relies on that: it establishes lookup() HITs and then pins them, and the + two must not be separated by an eviction. + + Serving can look keys up through the parent, and lookup() runs the + step's once-per-step poll. From here that poll would reach every tier + off-thread, re-enter the serving tier's get_finished_jobs() from inside + its own serve, and leave the gate set so the next step skips its poll. + The round has just polled its tiers itself, so it holds the gate shut + for its duration and then puts it back untouched. + """ + gate = self._processed_jobs_this_step + self._processed_jobs_this_step = True + try: + self._process_finished_jobs(self._external_serving_tiers) + for _, tier in self._external_serving_tiers: + tier.serve_external_requests(self._tier_parents[tier]) + finally: + self._processed_jobs_this_step = gate + @property def _transfer_jobs(self) -> dict[JobId, JobMetadata]: return self._jobs @@ -319,8 +479,11 @@ def _complete_promotion( False, ) - def _process_finished_jobs(self): - """Unconditionally poll all secondary tiers for completed jobs. + def _process_finished_jobs( + self, + tiers: Sequence[tuple[int, SecondaryTierManager]] | None = None, + ): + """Unconditionally poll secondary tiers for completed jobs. This method: 1. Calls get_finished_jobs() on each secondary tier @@ -328,8 +491,15 @@ def _process_finished_jobs(self): to decrement ref_cnt 3. For completed loads (secondary→primary): calls primary.complete_write() to make chunks available + + Args: + tiers: (index, tier) pairs to poll. Defaults to every secondary + tier. Pass a subset to poll only the tiers a caller is + responsible for, so a tier that expects to be polled on the + scheduler thread is not dragged elsewhere. + """ - for i, tier in enumerate(self.secondary_tiers): + for i, tier in tiers if tiers is not None else enumerate(self.secondary_tiers): for completed_job in tier.get_finished_jobs(): job_id = completed_job.job_id job_metadata = self._pop_job(job_id) @@ -620,9 +790,9 @@ def prepare_store( ) -> PrepareStoreOutput | None: """Prepare chunks to be stored from GPU to primary tier. - CRITICAL: This method calls _maybe_process_finished_jobs() FIRST to ensure - that any completed async transfers have their ref_cnt decremented - before the primary tier makes eviction decisions. + CRITICAL: This method polls for finished jobs FIRST to ensure that any + completed async transfers have their ref_cnt decremented before the + primary tier makes eviction decisions. For request-level tiers, chunks already present in the primary tier are immediately cascaded via submit_store(). @@ -646,7 +816,12 @@ def prepare_store( # not-yet-ready chunk's ref_cnt from -1 to 0 via complete_write(), # making it evictable for the first time. # Both must be accounted for before the eviction decision below. - self._maybe_process_finished_jobs() + # Unconditional, not _maybe_process_finished_jobs(): eviction must see + # the freshest completions, and on_schedule_end now runs at the very end + # of the step, so the once-per-step gate would otherwise still be set by + # this step's first lookup() and skip the poll. + self._processed_jobs_this_step = True + self._process_finished_jobs() # Step 2: Store to primary tier (new chunks only). # Cascading of these newly-stored chunks to ALL secondary tiers @@ -882,8 +1057,8 @@ def on_schedule_end(self, context: ScheduleEndContext) -> None: """End-of-schedule hook: process finished jobs, flush deferred promotions, and reset the per-step gate. - Called once per scheduler step from - OffloadingConnectorScheduler.build_connector_meta(). + Called once per scheduler step, as the last manager call of the step, + from OffloadingConnectorScheduler.build_connector_meta(). """ # Catch-all poll: guarantees jobs are processed even on steps where # lookup()/prepare_store() were never called (e.g. no requests @@ -898,6 +1073,10 @@ def on_schedule_end(self, context: ScheduleEndContext) -> None: self._processed_jobs_this_step = False self._flush_pending_promotions() + # Keys parked by this step's prepare_store cannot have completed yet, so + # they are simply re-parked here and retried on the next step. That is + # safe: _maybe_finalize_request() holds off while pending_cascade_keys is + # non-empty, so nothing finalizes early and the list still drains. self._flush_pending_cascades() for tier in self.secondary_tiers: tier.on_schedule_end(context) @@ -1011,9 +1190,32 @@ def get_stats(self) -> OffloadingConnectorStats | None: def shutdown(self) -> None: """Shut down secondary tiers before releasing primary resources. + Stops the control-plane thread first: it drives tier transports, and + tearing one down underneath it (closing a socket, destroying a ZMQ + context) can take the process down rather than raise. If the thread + will not exit, hold the manager lock across tier teardown so it cannot + be mid-round, and skip teardown entirely if even that is unavailable -- + a lingering daemon thread in an exiting process beats a crash. + Every secondary tier is given a shutdown attempt. If any shutdown fails, preserve the primary mmap because a failed tier may still use it. """ + self._stop_control_plane() + if self._control_plane_thread is not None: + if not self._lock.acquire(timeout=_CONTROL_PLANE_JOIN_TIMEOUT_S): + logger.error( + "KV offload control-plane thread is still running and the " + "manager lock is unavailable; skipping tier shutdown." + ) + return + try: + self._shutdown_tiers() + finally: + self._lock.release() + return + self._shutdown_tiers() + + def _shutdown_tiers(self) -> None: shutdown_error: Exception | None = None for tier_idx, tier in enumerate(self.secondary_tiers): try: diff --git a/vllm/v1/kv_offload/tiering/p2p/control/base.py b/vllm/v1/kv_offload/tiering/p2p/control/base.py index 5a7030737bdb..890bc16b2441 100644 --- a/vllm/v1/kv_offload/tiering/p2p/control/base.py +++ b/vllm/v1/kv_offload/tiering/p2p/control/base.py @@ -20,8 +20,11 @@ ├── mark_dead() — signal that the peer is gone └── close() — tear down the connection -Threading model: all I/O is driven by the caller invoking poll(). -No background threads. poll() must be called periodically to: +Threading model: all I/O is driven by the caller invoking poll(). No background +threads of its own. More than one caller thread may drive poll() over a +transport's lifetime, but never concurrently -- callers serialize themselves +(the tiering manager does so with its lock), because sockets here are not +thread-safe. poll() must be called periodically to: - receive messages (buffered per-connection) - accept new inbound peers - detect disconnections diff --git a/vllm/v1/kv_offload/tiering/p2p/data/base.py b/vllm/v1/kv_offload/tiering/p2p/data/base.py index 62304b7544a3..866450bf04a2 100644 --- a/vllm/v1/kv_offload/tiering/p2p/data/base.py +++ b/vllm/v1/kv_offload/tiering/p2p/data/base.py @@ -66,7 +66,11 @@ - close() releases all resources (memory registrations, handles). After close(), no other methods may be called. -Threading model: no background threads. All I/O driven by poll(). +Threading model: no background Python threads. All I/O driven by poll(). More +than one caller thread may drive a transport over its lifetime, but never +concurrently -- callers serialize themselves (the tiering manager does so with +its lock). Note that serialization alone is not obviously sufficient for a +backend that binds state to the thread that created it; see NixlTransport. """ from __future__ import annotations diff --git a/vllm/v1/kv_offload/tiering/p2p/data/nixl.py b/vllm/v1/kv_offload/tiering/p2p/data/nixl.py index e8fb66dd6e30..2c79acd52929 100644 --- a/vllm/v1/kv_offload/tiering/p2p/data/nixl.py +++ b/vllm/v1/kv_offload/tiering/p2p/data/nixl.py @@ -1,6 +1,37 @@ # SPDX-License-Identifier: Apache-2.0 # SPDX-FileCopyrightText: Copyright contributors to the vLLM project -"""NixlTransport: Data-plane transport for RDMA-based KV block transfers via NIXL.""" +"""NixlTransport: Data-plane transport for RDMA-based KV block transfers via NIXL. + +Threading: the agent is reached from the scheduler thread and, when the tiering +manager runs one, from its control-plane thread. The manager's lock serializes +them strictly -- one Python thread inside any agent call at a time, with a full +memory barrier between. + +That lock is the only protection, not a belt over NIXL's own. nixl_agent picks +NIXL_THREAD_SYNC_STRICT only when ``enable_listen`` is set, which this transport +does not set, so the agent runs with NIXL_THREAD_SYNC_NONE -- under which NIXL's +internal agent lock compiles down to no-ops, leaving every agent entry point +unguarded. So the manager lock must cover all of them, and the lifetime of every +handle passed between calls. The exception is peer registration moved onto a +worker thread, as NixlConnector's handshake executor does: it runs outside the +lock and relies, like that connector, on NIXL tolerating it beside transfers. + +Serialized use from a second thread is supported (checked against NIXL 1.4.1 +sources): NIXL never creates its UCX worker in UCS_THREAD_MODE_SINGLE, the only +creator-thread-affine mode. The weakest mode it ever requests is +UCS_THREAD_MODE_SERIALIZED, documented as "multiple threads can access, but only +one at a time", and because ``enable_prog_thread`` defaults on here it actually +requests UCS_THREAD_MODE_MULTI and fails construction if the UCX build cannot +provide it. + +Two caveats. NIXL binds a shared UCX worker per calling thread +(``tlsSharedWorkerMap``), which is harmless only while there is exactly one +shared worker -- true here, but it would stop being true if a ``num_workers`` +UCX backend param were ever set. And NIXL 1.4.1 does accept an explicit +``sync_mode`` (including NIXL_THREAD_SYNC_RW), which would make its own locking +real; that parameter does not exist before 1.x, so passing it would need a +version gate, the way ``capture_telemetry`` does. +""" from __future__ import annotations diff --git a/vllm/v1/kv_offload/tiering/p2p/manager.py b/vllm/v1/kv_offload/tiering/p2p/manager.py index f759a4481e5d..0c2220c1f871 100644 --- a/vllm/v1/kv_offload/tiering/p2p/manager.py +++ b/vllm/v1/kv_offload/tiering/p2p/manager.py @@ -201,13 +201,18 @@ class P2PSecondaryTierManager(SecondaryTierManager): blocks from the peer) and server-role (serving blocks to the peer) over the same control connection. - Single-threaded: every public method runs on the scheduler thread, and - the engine drives polling via ``get_finished_jobs()`` once per step. + Never entered concurrently: every public method runs under the tiering + manager's ``lock``, held either by the scheduler thread for the length of + a step or by the control-plane thread for the length of one round. Peers are + not driven by the engine, so this tier sets ``serves_external_requests`` + to opt into that thread; without it a peer's lookup or fetch waits + for a step boundary, which on a saturated rank can be seconds. ``has_pending_work()`` keeps the engine ticking so the control transport and existing sessions are polled even when no requests are scheduled. """ cache_hit_source: ClassVar[CacheHitSource] = CacheHitSource.P2P + serves_external_requests: ClassVar[bool] = True def __init__( self, @@ -607,9 +612,10 @@ def submit_load(self, job_metadata: TransferJob) -> None: @override def get_finished_jobs(self) -> Iterable[JobResult]: - # Drive one polling sweep on the scheduler thread, then hand off - # whatever has accumulated. The engine calls this once per step - # (and keeps stepping while has_pending_work() is True). + # Drive one polling sweep, then hand off whatever has accumulated. The + # engine calls this once per step (and keeps stepping while + # has_pending_work() is True); the control-plane thread calls it between + # steps, so a peer is not left waiting for a step boundary. self._poll_once() result = self._finished_jobs self._finished_jobs = [] @@ -836,7 +842,8 @@ def _poll_once(self) -> None: Drains the control transport, polls every session, accumulates their results into ``_finished_jobs``, and reaps any dead sessions. - Runs on the scheduler thread. + Runs under the tiering manager's lock, on the scheduler thread or on + the control-plane thread. """ new_connections = self._control.poll() if new_connections: diff --git a/vllm/v1/kv_offload/tiering/p2p/session/server.py b/vllm/v1/kv_offload/tiering/p2p/session/server.py index c78918323266..d0a193c56879 100644 --- a/vllm/v1/kv_offload/tiering/p2p/session/server.py +++ b/vllm/v1/kv_offload/tiering/p2p/session/server.py @@ -584,10 +584,13 @@ def _pin_and_register_hits( """Pin primary slots for HIT keys and park them as the lookup's round supply via ``add_stored_blocks``. - Caller has already confirmed every key is HIT (single-threaded - scheduler ⇒ no eviction race), so the JobMetadata returned by - ``parent.create_store_job`` carries parallel ``keys``/``block_ids`` - of length ``len(keys)``. + Caller has already confirmed every key is HIT, so the JobMetadata + returned by ``parent.create_store_job`` carries parallel + ``keys``/``block_ids`` of length ``len(keys)``. Nothing can evict + between that confirmation and the pin because the whole + ``serve_external_requests`` call runs inside one hold of the tiering + manager's lock -- which is why that call must never release it + partway through. """ meta = parent.create_store_job(keys, lookup.ctx) self.add_stored_blocks( diff --git a/vllm/v1/kv_offload/tiering/spec.py b/vllm/v1/kv_offload/tiering/spec.py index 93c6cbabcb85..73921f447571 100644 --- a/vllm/v1/kv_offload/tiering/spec.py +++ b/vllm/v1/kv_offload/tiering/spec.py @@ -73,6 +73,7 @@ from vllm.v1.kv_offload.tiering.base import TieringOffloadingMetrics from vllm.v1.kv_offload.tiering.factory import SecondaryTierFactory from vllm.v1.kv_offload.tiering.manager import ( + DEFAULT_TIER_POLL_INTERVAL_S, CPUPrimaryTierOffloadingManager, TieringOffloadingManager, ) @@ -287,6 +288,13 @@ def __init__(self, config: OffloadingConfig): if not isinstance(self.secondary_tier_configs, list): raise ValueError("secondary_tiers must be a list of tier configurations") + # Pause between control-plane rounds for tiers that need servicing + # between engine steps. Zero or negative disables the thread, which + # restores per-step servicing. + self.tier_poll_interval_s = float( + self.extra_config.get("tier_poll_interval_s", DEFAULT_TIER_POLL_INTERVAL_S) + ) + # Backpressure config is merged field-by-field in priority order # (highest first): # 1. Per-tier ``backpressure`` dict in the tier config @@ -387,6 +395,7 @@ def get_manager(self) -> OffloadingManager: tiering_manager = TieringOffloadingManager( primary_tier=primary_tier, secondary_tiers=secondary_tiers, + tier_poll_interval_s=self.tier_poll_interval_s, ) self._manager = tiering_manager except Exception: