From 1a37624be6cd56d3d59bfab3750aee685a048974 Mon Sep 17 00:00:00 2001 From: Costa Tsaousis Date: Sun, 27 Sep 2026 22:28:30 +0000 Subject: [PATCH 1/5] Wait for recurrent checkpoint restore admission under capacity pressure Signed-off-by: Costa Tsaousis --- .lil/changes/lmcache-restore-admission.json | 15 ++ .../integration/vllm/checkpoint_scheduler.py | 106 +++++++--- .../test_vllm_semantic_checkpoint_transfer.py | 185 ++++++++++++++++-- 3 files changed, 262 insertions(+), 44 deletions(-) create mode 100644 .lil/changes/lmcache-restore-admission.json diff --git a/.lil/changes/lmcache-restore-admission.json b/.lil/changes/lmcache-restore-admission.json new file mode 100644 index 00000000000..8d39483c904 --- /dev/null +++ b/.lil/changes/lmcache-restore-admission.json @@ -0,0 +1,15 @@ +{ + "schema": "local-inference-release-change/v1", + "id": "lmcache-restore-admission", + "category": "fix", + "summary": "Recurrent checkpoint restores wait for GPU and copy capacity instead of admitting requests to recompute cached prefixes.", + "models": ["all"], + "compatibility": "Requires the paired vLLM restore-admission change; deploy vLLM before enabling this LMCache change.", + "details": [ + "Capacity waits retain the answered checkpoint lookup and do not consume the timeout for an unanswered directory request.", + "A local checkpoint evicted before admission triggers a fresh external lookup, while genuine misses and exhausted transfer retries retain recomputation fallback.", + "This change does not repair missing checkpoint pages in cache storage." + ], + "authors": ["ktsaou"], + "requires": ["vllm-restore-admission"] +} diff --git a/lmcache/integration/vllm/checkpoint_scheduler.py b/lmcache/integration/vllm/checkpoint_scheduler.py index 1038455515d..fa05ca574c4 100644 --- a/lmcache/integration/vllm/checkpoint_scheduler.py +++ b/lmcache/integration/vllm/checkpoint_scheduler.py @@ -76,9 +76,11 @@ class _Lookup: done: bool = False task_id: str | None = None checkpoint_id: int | None = None + local_checkpoint_id: int | None = None attempts: int = 1 started: float = field(default_factory=time.monotonic) capacity_refusal_logged: bool = False + task_refusal_logged: bool = False class CheckpointSchedulerBridge: @@ -91,10 +93,10 @@ class CheckpointSchedulerBridge: Worker layout is added after all ranks report identical descriptors. world_size: Number of engine ranks contributing to one atomic generation. max_tasks: Admission limit for collective stores and restores. - lookup_timeout: Seconds a request may wait for directory replies, - including shorter-checkpoint retries, and a store for its begin - reply. A request whose lookup has not been answered by then is - admitted to recompute its prompt; an unanswered store is aborted. + lookup_timeout: Seconds allowed for each directory reply (including + retries) and for a store's begin reply. Capacity waiting does not + consume this deadline. An unanswered lookup permits prompt + recomputation; an unanswered store is aborted. A started worker copy is never abandoned. The scheduler must keep issuing connector-only steps while has_pending is @@ -175,9 +177,8 @@ def poll_prefix(self, request: "Request") -> bool: Returns: False while layout negotiation, lookup or H2D is pending. True when - ordinary GPU admission may proceed, including when an external - import cannot reserve capacity. An unadmitted request may poll - again to retry that reservation without another directory lookup. + ordinary GPU admission may proceed. Capacity and copy-task pressure + keep an available checkpoint waiting without another lookup. Reusing a finished request's ID waits for its admitted copies to drain, so their completion cannot cancel or erase another lookup. """ @@ -188,11 +189,15 @@ def poll_prefix(self, request: "Request") -> bool: ): return True local = self._cache.find(request, request.num_tokens) - if local is not None and local.num_tokens == request.num_tokens: + state = self._lookups.get(request.request_id) + if ( + state is None + and local is not None + and local.num_tokens == request.num_tokens + ): return True if self._layout is None: return False - state = self._lookups.get(request.request_id) if state is None: roots = self._roots(request) state = _Lookup( @@ -202,6 +207,27 @@ def poll_prefix(self, request: "Request") -> bool: self._lookups[request.request_id] = state return False if state.done: + # Local reuse can also disappear while ordinary admission waits. + selected_id = ( + state.local_checkpoint_id + if state.local_checkpoint_id is not None + else state.checkpoint_id + ) + if ( + selected_id is not None + and not ( + self._manager.external_boundary_admission_ready(request.request_id) + ) + and (local is None or local.checkpoint_id != selected_id) + ): + logger.info( + "Recurrent checkpoint invalidated before admission for " + "request %s; looking up its prefix again", + request.request_id, + ) + self._manager.release_external_boundary_admission(request.request_id) + del self._lookups[request.request_id] + return self.poll_prefix(request) return True if state.task_id is not None: return False @@ -227,21 +253,25 @@ def poll_prefix(self, request: "Request") -> bool: ) state.done = True return True - if manifest is None or ( - local is not None and manifest.prefix.num_tokens <= local.num_tokens - ): + if manifest is None: state.done = True return True - if len(self._tasks) >= self._max_tasks: - logger.info( - "Recurrent checkpoint restore of %d tokens skipped for request %s: " - "%d checkpoint copies in flight; recomputing its prompt", - manifest.prefix.num_tokens, - request.request_id, - len(self._tasks), - ) + if local is not None and manifest.prefix.num_tokens <= local.num_tokens: + # Keep local selection separate from external cache-hit accounting. + state.local_checkpoint_id = local.checkpoint_id state.done = True return True + if len(self._tasks) >= self._max_tasks: + if not state.task_refusal_logged: + state.task_refusal_logged = True + logger.info( + "Recurrent checkpoint restore of %d tokens waiting for " + "copy capacity for request %s (%d copies in flight)", + manifest.prefix.num_tokens, + request.request_id, + len(self._tasks), + ) + return False try: payload, positions = self._validate_manifest(manifest, state.roots) checkpoint = self._manager.reserve_external_boundary_checkpoint( @@ -251,6 +281,7 @@ def poll_prefix(self, request: "Request") -> bool: draft_prefix_len=payload["draft_prefix_len"], kind=payload["kind"], num_ranks=self._world_size, + reserve_admission=True, ) except (ValueError, KeyError, TypeError): logger.warning( @@ -263,21 +294,25 @@ def poll_prefix(self, request: "Request") -> bool: state.done = True return True if checkpoint is None: - # Insufficient GPU capacity is not an external-cache miss. Permit - # ordinary admission, but retain the manifest so an unadmitted - # request can retry its import after other owners release pages. + # Resource pressure is not a cache miss. Keep the answered manifest + # while runnable owners release capacity. if not state.capacity_refusal_logged: state.capacity_refusal_logged = True logger.info( "Recurrent checkpoint restore of %d tokens deferred for " - "request %s: not enough free GPU blocks (%d free); it is " - "admitted without the restore unless capacity returns first", + "request %s: waiting for restore and execution capacity " + "(%d free GPU blocks)", manifest.prefix.num_tokens, request.request_id, self._manager.block_pool.get_num_free_blocks(), ) - return True + return False task = self._make_task(manifest, checkpoint, "RETRIEVE") + logger.debug( + "Recurrent checkpoint restore of %d tokens copying for request %s", + manifest.prefix.num_tokens, + request.request_id, + ) self._tasks[task.task_id] = _PendingTask(task, checkpoint, request.request_id) state.task_id = task.task_id return False @@ -523,6 +558,22 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: ) if published and state is not None: state.checkpoint_id = pending.checkpoint.checkpoint_id + logger.debug( + "Recurrent checkpoint restore of %d tokens ready for " + "request %s; ownership retained until admission", + pending.checkpoint.num_tokens, + pending.request_id, + ) + elif state is not None: + retry = state.attempts < _MAX_LOOKUP_ATTEMPTS + logger.info( + "Recurrent checkpoint publication invalidated for " + "request %s%s", + pending.request_id, + "; retrying lookup" + if retry + else "; recomputing its prompt", + ) elapsed = time.monotonic() - pending.created if elapsed > _SLOW_RESTORE_SECONDS: logger.info( @@ -540,7 +591,6 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: state is not None and pending.request_id not in self._cancelled and state.attempts < _MAX_LOOKUP_ATTEMPTS - and time.monotonic() - state.started < self._lookup_timeout ) logger.info( "Recurrent checkpoint restore of %d tokens failed for " @@ -567,6 +617,7 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: ) state.task_id = None state.attempts += 1 + state.started = time.monotonic() state.future = self._client.submit_request( RequestType.CHECKPOINT_FIND, [state.roots.roots] ) @@ -578,6 +629,7 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: def finish_request(self, request_id: str) -> None: """Forget lookup state, retaining copy pins until every admitted task drains.""" + self._manager.release_external_boundary_admission(request_id) self._lookups.pop(request_id, None) if any(task.request_id == request_id for task in self._tasks.values()): self._cancelled.add(request_id) diff --git a/tests/v1/test_vllm_semantic_checkpoint_transfer.py b/tests/v1/test_vllm_semantic_checkpoint_transfer.py index 44c063222a5..ffa64844a84 100644 --- a/tests/v1/test_vllm_semantic_checkpoint_transfer.py +++ b/tests/v1/test_vllm_semantic_checkpoint_transfer.py @@ -26,6 +26,7 @@ from vllm.lora.request import LoRARequest # noqa: E402 from vllm.sampling_params import SamplingParams # noqa: E402 from vllm.utils.hashing import sha256 # noqa: E402 +from vllm.v1.core.boundary_checkpoint import BoundaryCheckpoint # noqa: E402 from vllm.v1.core.kv_cache_manager import KVCacheManager # noqa: E402 from vllm.v1.core.kv_cache_utils import ( # noqa: E402 get_request_block_hasher, @@ -485,6 +486,11 @@ def test_worker_metadata_requires_bound_rank() -> None: "rejected-store-reused-request-id", "inflight-store-reused-request-id", "capacity-retry", + "capacity-timeout", + "task-capacity", + "local-hit-during-copy", + "invalidated-before-use", + "publication-invalidated", ], ) def test_semantic_roundtrip_collective_visibility_and_cancellation( @@ -503,6 +509,8 @@ def test_semantic_roundtrip_collective_visibility_and_cancellation( "parallel": {"tp": 4, "dcp": 1}, }, 4, + max_tasks=1 if outcome == "task-capacity" else 32, + lookup_timeout=0.05 if outcome == "capacity-timeout" else 60.0, ) layout = {"schema_version": 1, "page_bytes": 128} bridge.accept_layouts({rank: layout for rank in range(4)}) @@ -591,30 +599,28 @@ def record(kind: RequestType, args: list[Any]) -> Any: assert roots and all(root.namespace == prefix.namespace for root in roots) assert sum(len(root.tail_tokens) for root in roots) >= prefix.num_tokens consumer = make_request("consumer") - if outcome == "capacity-retry": - # Ordinary admission may also fail while another request owns - # the pool. An available manifest must remain retryable until - # this consumer is admitted or cancelled. + if outcome in ("capacity-retry", "capacity-timeout"): + # Resource pressure must not authorize cold admission. pressure = manager.block_pool.get_new_blocks( manager.block_pool.get_num_free_blocks() ) try: - deadline = time.monotonic() + 5 - while not bridge.poll_prefix(consumer): - assert time.monotonic() < deadline + for _ in range(100): + assert not bridge.poll_prefix(consumer) time.sleep(0.001) assert bridge.take_tasks() == [] - # Persistent pressure still permits ordinary admission; - # retaining the manifest must not introduce a wait loop. - assert bridge.poll_prefix(consumer) + assert not bridge.poll_prefix(consumer) finally: manager.block_pool.free_blocks(pressure) inflight_store = None - if outcome == "inflight-store-reused-request-id": + if outcome in ("inflight-store-reused-request-id", "task-capacity"): # This cancelled producer has different tokens but the same # public ID as the waiting consumer. Its admitted copy is still # live while the consumer looks up the first producer's data. - predecessor = make_request("consumer", first_token=20) + predecessor = make_request( + "blocker" if outcome == "task-capacity" else "consumer", + first_token=20, + ) predecessor_checkpoint = manager.reserve_external_boundary_checkpoint( predecessor, 11, @@ -675,8 +681,28 @@ def drain_predecessor() -> None: assert len(tasks) == 1 restore_task = tasks[0] assert manager.get_computed_blocks(consumer)[1] == 0 + if outcome == "local-hit-during-copy": + local = make_request("local-producer") + local_checkpoint = manager.reserve_external_boundary_checkpoint( + local, + 11, + manager.boundary_checkpoint_page_positions(11), + draft_prefix_len=11, + kind="prompt", + num_ranks=1, + ) + assert local_checkpoint is not None + assert manager.acknowledge_external_boundary_checkpoint( + local_checkpoint.checkpoint_id, 0 + ) + assert not bridge.poll_prefix(consumer) if outcome == "cancelled": bridge.finish_request(consumer.request_id) + if outcome == "publication-invalidated": + assert manager.boundary_checkpoints is not None + manager.boundary_checkpoints.invalidate_block( + restore_task.block_ids[-1][0] + ) for rank in range(4): future = worker.submit( CheckpointTransferJob( @@ -695,12 +721,32 @@ def drain_predecessor() -> None: } ) if rank < 3: - assert manager.get_computed_blocks(consumer)[1] == 0 + if outcome != "local-hit-during-copy": + assert manager.get_computed_blocks(consumer)[1] == 0 assert manager.block_pool.get_num_free_blocks() < before - assert manager.block_pool.get_num_free_blocks() == before - expected_tokens = 0 if outcome in ("rank-miss", "cancelled") else 11 + expected_tokens = ( + 0 + if outcome in ("rank-miss", "cancelled", "publication-invalidated") + else 11 + ) + if expected_tokens: + assert manager.block_pool.get_num_free_blocks() < before + assert not manager.reset_prefix_cache() + else: + assert manager.block_pool.get_num_free_blocks() == before assert manager.get_computed_blocks(consumer)[1] == expected_tokens assert bridge.external_tokens(consumer) == expected_tokens + if outcome == "invalidated-before-use": + # Preemption releases request ownership while its connector + # lookup state can still remember the previous successful import. + checkpoint = consumer.boundary_checkpoint + assert checkpoint is not None + manager.free(consumer) + assert manager.boundary_checkpoints is not None + manager.boundary_checkpoints.invalidate(checkpoint.checkpoint_id) + assert not bridge.poll_prefix(consumer) + bridge.finish_request(consumer.request_id) + assert manager.block_pool.get_num_free_blocks() == before if inflight_store is not None: drain_predecessor() assert not bridge.has_pending @@ -846,7 +892,108 @@ def make_bridge(client, manager, lookup_timeout: float) -> CheckpointSchedulerBr return bridge -@pytest.mark.parametrize("reply", ["never", "error"]) +@pytest.mark.parametrize("prefix, already_local", [(11, False), (8, True)]) +def test_local_reuse_retries_external_restore_after_capacity_eviction( + prefix: int, already_local: bool +) -> None: + """Losing a local selection must not turn an external hit into cold prefill.""" + manager = make_manager() + manifest = None + lookups = 0 + + def submit(kind: RequestType, args: list[Any]) -> SimpleNamespace: + nonlocal manifest, lookups + if kind == RequestType.CHECKPOINT_BEGIN: + manifest = args[0] + if kind == RequestType.CHECKPOINT_FIND: + lookups += 1 + result = manifest if kind == RequestType.CHECKPOINT_FIND else True + return SimpleNamespace(query=lambda: True, result=lambda: result) + + bridge = make_bridge(SimpleNamespace(submit_request=submit), manager, 60) + producer = make_request("producer") + + def publish() -> BoundaryCheckpoint: + checkpoint = manager.reserve_external_boundary_checkpoint( + producer, + prefix, + manager.boundary_checkpoint_page_positions(prefix), + draft_prefix_len=prefix, + kind="prompt", + num_ranks=1, + ) + assert checkpoint is not None + assert manager.acknowledge_external_boundary_checkpoint( + checkpoint.checkpoint_id, 0 + ) + return checkpoint + + checkpoint = publish() + bridge.store(producer, checkpoint) + (store,) = bridge.take_tasks() + bridge.complete({store.task_id: {rank: True for rank in range(4)}}) + bridge.finish_request(producer.request_id) + assert manager.reset_prefix_cache() + consumer = make_request("consumer") + if already_local: + checkpoint = publish() + assert not bridge.poll_prefix(consumer) + if not already_local: + checkpoint = publish() + # Other requests occupy everything except the unpinned cached bundle. + pressure = manager.block_pool.get_new_blocks( + manager.block_pool.get_num_free_blocks() - len(checkpoint.dependencies) + ) + assert bridge.poll_prefix(consumer) + blocks, tokens, _ = manager.get_computed_blocks(consumer) + assert tokens == prefix + assert bridge.external_tokens(consumer) == 0 + assert bridge.poll_prefix(consumer) + assert lookups == 1 + assert bridge.take_tasks() == [] + assert ( + manager.allocate_slots( + consumer, + 1, + num_new_computed_tokens=tokens, + new_computed_blocks=blocks, + full_sequence_must_fit=True, + ) + is None + ) + # A running request grows into the cached bundle before capacity returns. + growth = manager.block_pool.get_new_blocks(1) + manager.block_pool.free_blocks(pressure + growth) + assert manager.get_computed_blocks(consumer)[1] == 0 + assert not bridge.poll_prefix(consumer) + assert lookups == 2 + assert not bridge.poll_prefix(consumer) + (restore,) = bridge.take_tasks() + bridge.complete({restore.task_id: {rank: True for rank in range(4)}}) + assert bridge.poll_prefix(consumer) + blocks, tokens, _ = manager.get_computed_blocks(consumer) + assert tokens == prefix + assert bridge.external_tokens(consumer) == prefix + assert ( + manager.allocate_slots( + consumer, + 1, + num_new_computed_tokens=tokens, + new_computed_blocks=blocks, + full_sequence_must_fit=True, + ) + is not None + ) + _, copies = manager.take_kv_cache_block_copies() + manager.block_pool.free_blocks(copies) + bridge.finish_request(consumer.request_id) + manager.free(consumer) + assert not bridge.has_pending + assert manager.external_boundary_reserved_blocks() == 0 + assert manager.block_pool.get_num_free_blocks() == 63 + + +@pytest.mark.parametrize("reply", ["never", "error", "miss"]) def test_lookup_without_a_usable_reply_admits_the_request( reply: str, monkeypatch: pytest.MonkeyPatch ) -> None: @@ -859,7 +1006,10 @@ def test_lookup_without_a_usable_reply_admits_the_request( future = ( SimpleNamespace(query=lambda: False, result=failed.result) if reply == "never" - else SimpleNamespace(query=lambda: True, result=failed.result) + else SimpleNamespace( + query=lambda: True, + result=(lambda: None) if reply == "miss" else failed.result, + ) ) consumer = make_request("consumer") with monkeypatch.context() as patch: @@ -869,6 +1019,7 @@ def test_lookup_without_a_usable_reply_admits_the_request( assert not bridge.poll_prefix(consumer) time.sleep(0.3) assert bridge.poll_prefix(consumer) + assert bridge.poll_prefix(consumer) assert bridge.take_tasks() == [] assert bridge.external_tokens(consumer) == 0 assert not bridge.has_pending From 65254d20cc1a05e4a196a0efe998a6dd3b6e4922 Mon Sep 17 00:00:00 2001 From: Costa Tsaousis Date: Sun, 27 Sep 2026 22:35:21 +0000 Subject: [PATCH 2/5] Link restore capacity release note to PR 100 Signed-off-by: Costa Tsaousis --- .lil/changes/lmcache-restore-admission.json | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/.lil/changes/lmcache-restore-admission.json b/.lil/changes/lmcache-restore-admission.json index 8d39483c904..2a460446ba8 100644 --- a/.lil/changes/lmcache-restore-admission.json +++ b/.lil/changes/lmcache-restore-admission.json @@ -3,13 +3,19 @@ "id": "lmcache-restore-admission", "category": "fix", "summary": "Recurrent checkpoint restores wait for GPU and copy capacity instead of admitting requests to recompute cached prefixes.", - "models": ["all"], + "models": [ + "all" + ], "compatibility": "Requires the paired vLLM restore-admission change; deploy vLLM before enabling this LMCache change.", "details": [ "Capacity waits retain the answered checkpoint lookup and do not consume the timeout for an unanswered directory request.", "A local checkpoint evicted before admission triggers a fresh external lookup, while genuine misses and exhausted transfer retries retain recomputation fallback.", "This change does not repair missing checkpoint pages in cache storage." ], - "authors": ["ktsaou"], - "requires": ["vllm-restore-admission"] + "requires": [ + "vllm-restore-admission" + ], + "pull_requests": [ + 100 + ] } From ef929e16abd87495c9e29d6afcf31e8cacbde5af Mon Sep 17 00:00:00 2001 From: Martin Vit Date: Mon, 28 Sep 2026 14:55:19 +0000 Subject: [PATCH 3/5] fix(checkpoint): gate restore admission on the vLLM API and bound waits Follow-up review fixes for #100: - Detect vLLM restore admission reservations once, when the bridge is created (the admission methods and the reserve_admission parameter). Without them, log one warning and keep the previous fallback: a restore without free GPU blocks is admitted and retried while unadmitted, one without a copy slot recomputes, and no missing method is ever called (finish_request no longer raises on every finished request). - Look a selected checkpoint up again only when no local checkpoint at least as long is cached, not whenever a different one replaced it. - Log a capacity wait when it starts and every 30 s with its age, the free GPU blocks and the blocks the checkpoint needs (or the copies in flight); report_status() counts current waits, the longest wait and the waits started. Every way out of a wait releases vLLM's waiter entry. - A restore that retains its answer for half the lookup timeout sends the lookup again without waiting for it, keeping its pages recent; the reply replaces the retained answer, and an empty one recomputes the prompt. - A complete local hit again bypasses a pending directory lookup, except behind the request's own reserved copy, and settles the lookup so a later eviction looks the prefix up again. - Validate an answered manifest once per answer, and log an invalidated publication without the "missed; looking up a shorter checkpoint" line. - Tests run on vLLM with and without the admission API: reservation-only cases skip with a reason, the fallback is tested through an allocator without the API, and exits assert that no restore stays queued. Co-Authored-By: Claude Opus 5.5 --- .../integration/vllm/checkpoint_scheduler.py | 489 ++++++++++++--- .../test_vllm_semantic_checkpoint_transfer.py | 577 ++++++++++++++++-- 2 files changed, 932 insertions(+), 134 deletions(-) diff --git a/lmcache/integration/vllm/checkpoint_scheduler.py b/lmcache/integration/vllm/checkpoint_scheduler.py index fa05ca574c4..cf2c1e2908d 100644 --- a/lmcache/integration/vllm/checkpoint_scheduler.py +++ b/lmcache/integration/vllm/checkpoint_scheduler.py @@ -4,6 +4,7 @@ # Standard from dataclasses import dataclass, field from typing import TYPE_CHECKING, Any, Literal +import inspect import json import math import os @@ -34,6 +35,18 @@ # the request waits in the scheduler, deferred, for the whole restore. _SLOW_RESTORE_SECONDS = 10.0 +# A restore waiting for GPU or copy capacity is logged when the wait starts +# and again at this interval, so a starved restore stays visible. +_WAIT_LOG_SECONDS = 30.0 + +# Allocator methods of vLLM's restore admission reservations. The scheduler +# gates admissions with the first; this bridge calls the other two. +_ADMISSION_METHODS = ( + "can_admit_external_boundary_request", + "external_boundary_admission_ready", + "release_external_boundary_admission", +) + # Settings that change encoder outputs, and so the KV of multimodal spans, # without changing the processor's content hash. They salt multimodal roots # only, so text namespaces are unaffected. @@ -46,6 +59,26 @@ from vllm.v1.request import Request +def _supports_restore_admission(manager: "KVCacheManager") -> bool: + """Whether vLLM can hold GPU capacity for a restore until its admission. + + Args: + manager: vLLM allocator bound to the bridge. + + Returns: + True if the allocator has the admission methods and reservations + accept ``reserve_admission``. Allocators older than the paired vLLM + restore-admission change have neither. + """ + if not all(callable(getattr(manager, name, None)) for name in _ADMISSION_METHODS): + return False + try: + signature = inspect.signature(manager.reserve_external_boundary_checkpoint) + except (TypeError, ValueError): + return False + return "reserve_admission" in signature.parameters + + @dataclass(frozen=True) class CheckpointEngineTask: """Ephemeral scheduler-to-worker transfer command, never persisted to storage.""" @@ -72,15 +105,29 @@ class _PendingTask: @dataclass class _Lookup: roots: CheckpointTokenRoots + # Directory lookup of the current attempt, sent at `started`. future: MessagingFuture[CheckpointManifest | None] done: bool = False + # Restore copy in flight for this request, or None. task_id: str | None = None + # Published import, counted as an external cache hit. checkpoint_id: int | None = None - local_checkpoint_id: int | None = None + # Prefix length selected for reuse, local or imported; 0 before selection. + selected_tokens: int = 0 attempts: int = 1 started: float = field(default_factory=time.monotonic) + # Answered candidate of the current attempt, retained until its copy + # starts, and its validated payload and page positions. + manifest: CheckpointManifest | None = None + plan: tuple[dict[str, Any], tuple[tuple[int, ...], ...]] | None = None + # When the directory last answered or was asked again for `manifest`, and + # the unanswered re-sent lookup, if any. + refreshed: float = 0.0 + refresh: MessagingFuture[CheckpointManifest | None] | None = None + # Start of the current wait for GPU or copy capacity, and its last log. + waiting_since: float | None = None + wait_logged: float = 0.0 capacity_refusal_logged: bool = False - task_refusal_logged: bool = False class CheckpointSchedulerBridge: @@ -88,15 +135,20 @@ class CheckpointSchedulerBridge: Args: manager: vLLM allocator with request-boundary checkpoint support enabled. + Restores wait for GPU capacity only if it supports restore + admission reservations; otherwise a restore that cannot reserve + its pages permits admission and the prompt is recomputed. client: Shared thread-safe LMCache queue client; no ownership transfer. identity: Immutable target/draft/source revisions and parallel geometry. Worker layout is added after all ranks report identical descriptors. world_size: Number of engine ranks contributing to one atomic generation. max_tasks: Admission limit for collective stores and restores. - lookup_timeout: Seconds allowed for each directory reply (including - retries) and for a store's begin reply. Capacity waiting does not - consume this deadline. An unanswered lookup permits prompt - recomputation; an unanswered store is aborted. + lookup_timeout: Seconds allowed for each directory reply, separately + for every retry after a failed restore, and for a store's begin + reply. Capacity waits and copies do not consume it. An unanswered + lookup permits prompt recomputation; an unanswered store is + aborted. A restore waiting longer than half of it sends its lookup + again, without waiting for the reply, to keep its pages recent. A started worker copy is never abandoned. The scheduler must keep issuing connector-only steps while has_pending is @@ -123,6 +175,17 @@ def __init__( self._lookup_timeout = lookup_timeout self._manager = manager self._cache = manager.boundary_checkpoints + # Fixed for the allocator's lifetime; never call a missing method. + self._reserve_admission = _supports_restore_admission(manager) + if not self._reserve_admission: + logger.warning( + "vLLM cannot reserve admission capacity for recurrent checkpoint " + "restores; a restore without enough free GPU blocks or copy slots " + "is admitted to recompute its prompt. Deploy the paired vLLM " + "restore-admission change to make restores wait for capacity." + ) + # Restores that had to wait for GPU or copy capacity, since creation. + self._waits_started = 0 self._client = client self._identity = dict(identity) self._multimodal_salt = json.dumps( @@ -176,145 +239,205 @@ def poll_prefix(self, request: "Request") -> bool: request: Waiting request whose prefix has not been admitted yet. Returns: - False while layout negotiation, lookup or H2D is pending. True when - ordinary GPU admission may proceed. Capacity and copy-task pressure - keep an available checkpoint waiting without another lookup. + False while layout negotiation, lookup or H2D is pending, and, if + vLLM supports restore admission reservations, while an answered + checkpoint waits for GPU or copy capacity. True when ordinary GPU + admission may proceed. Without those reservations, a restore that + cannot reserve GPU blocks permits admission; an unadmitted request + may poll again to retry it without another blocking lookup. + A complete local hit proceeds without a directory reply, except + behind this request's own reserved copy. Reusing a finished request's ID waits for its admitted copies to drain, so their completion cannot cancel or erase another lookup. """ - if request.request_id in self._cancelled: + request_id = request.request_id + if request_id in self._cancelled: return False if not self.handles(request) or not self._manager.prefix_cache_lookup_enabled( request ): return True local = self._cache.find(request, request.num_tokens) - state = self._lookups.get(request.request_id) + state = self._lookups.get(request_id) if ( - state is None - and local is not None + local is not None and local.num_tokens == request.num_tokens + and (state is None or state.task_id is None or not self._reserve_admission) ): + # A complete local hit needs no directory reply or capacity wait. A + # reserved copy must still finish: its publication hands the + # reservation to this request, which must still be waiting then. + if self._reserve_admission and state is not None and not state.done: + # Settle the lookup as a local selection. If this checkpoint + # leaves the cache before admission, the prefix is looked up + # again rather than resuming an old reply or its deadline. + state.selected_tokens = local.num_tokens + self._stop_waiting(request_id, state) + state.done = True return True if self._layout is None: return False if state is None: roots = self._roots(request) - state = _Lookup( + self._lookups[request_id] = _Lookup( roots, self._client.submit_request(RequestType.CHECKPOINT_FIND, [roots.roots]), ) - self._lookups[request.request_id] = state return False if state.done: - # Local reuse can also disappear while ordinary admission waits. - selected_id = ( - state.local_checkpoint_id - if state.local_checkpoint_id is not None - else state.checkpoint_id - ) + # A selected prefix can leave the GPU cache while ordinary admission + # waits or after preemption; only a shorter or missing local + # checkpoint needs another lookup. Without reservations a restored + # bundle is not held for its consumer, so another lookup could + # repeat the restore under the same pressure. if ( - selected_id is not None - and not ( - self._manager.external_boundary_admission_ready(request.request_id) - ) - and (local is None or local.checkpoint_id != selected_id) + self._reserve_admission + and state.selected_tokens > 0 + and (local is None or local.num_tokens < state.selected_tokens) + and not self._manager.external_boundary_admission_ready(request_id) ): logger.info( - "Recurrent checkpoint invalidated before admission for " - "request %s; looking up its prefix again", - request.request_id, + "Recurrent checkpoint of %d tokens selected for request %s is " + "no longer cached; looking up its prefix again", + state.selected_tokens, + request_id, ) - self._manager.release_external_boundary_admission(request.request_id) - del self._lookups[request.request_id] + self._manager.release_external_boundary_admission(request_id) + del self._lookups[request_id] return self.poll_prefix(request) return True if state.task_id is not None: return False - if not state.future.query(): - if time.monotonic() - state.started < self._lookup_timeout: - return False - # An unanswered lookup must not park the request indefinitely. - logger.warning( - "Recurrent checkpoint lookup for request %s got no reply within " - "%.0f s; recomputing its prompt", - request.request_id, - self._lookup_timeout, - ) - state.done = True - return True - try: - manifest = state.future.result() - except Exception: - logger.exception( - "Recurrent checkpoint lookup failed for request %s; recomputing " - "its prompt", - request.request_id, - ) - state.done = True - return True + manifest = state.manifest if manifest is None: - state.done = True - return True + if not state.future.query(): + if time.monotonic() - state.started < self._lookup_timeout: + return False + # An unanswered lookup must not park the request indefinitely. + logger.warning( + "Recurrent checkpoint lookup for request %s got no reply within " + "%.0f s; recomputing its prompt", + request_id, + self._lookup_timeout, + ) + state.done = True + return True + try: + manifest = state.future.result() + except Exception: + logger.exception( + "Recurrent checkpoint lookup failed for request %s; recomputing " + "its prompt", + request_id, + ) + state.done = True + return True + if manifest is None: + state.done = True + return True + state.manifest = manifest + state.refreshed = time.monotonic() + else: + manifest = self._refresh_manifest(request_id, state, manifest) + if manifest is None: + self._stop_waiting(request_id, state) + state.done = True + return True if local is not None and manifest.prefix.num_tokens <= local.num_tokens: - # Keep local selection separate from external cache-hit accounting. - state.local_checkpoint_id = local.checkpoint_id + # Reuse the local checkpoint; external_tokens counts imports only. + state.selected_tokens = local.num_tokens + self._stop_waiting(request_id, state) state.done = True return True if len(self._tasks) >= self._max_tasks: - if not state.task_refusal_logged: - state.task_refusal_logged = True + if not self._reserve_admission: logger.info( - "Recurrent checkpoint restore of %d tokens waiting for " - "copy capacity for request %s (%d copies in flight)", + "Recurrent checkpoint restore of %d tokens skipped for request " + "%s: %d checkpoint copies in flight; recomputing its prompt", manifest.prefix.num_tokens, - request.request_id, + request_id, len(self._tasks), ) + state.done = True + return True + self._wait( + request_id, + state, + manifest.prefix.num_tokens, + "a copy slot (%d copies in flight)", + len(self._tasks), + ) return False try: - payload, positions = self._validate_manifest(manifest, state.roots) - checkpoint = self._manager.reserve_external_boundary_checkpoint( - request, - manifest.prefix.num_tokens, - positions, - draft_prefix_len=payload["draft_prefix_len"], - kind=payload["kind"], - num_ranks=self._world_size, - reserve_admission=True, + if state.plan is None: + # Validated once per answer, not on every step of a wait. + state.plan = self._validate_manifest(manifest, state.roots) + payload, positions = state.plan + checkpoint = self._reserve( + request, manifest.prefix.num_tokens, payload, positions ) except (ValueError, KeyError, TypeError): logger.warning( "Recurrent checkpoint restore of %d tokens rejected for request " "%s; recomputing its prompt", manifest.prefix.num_tokens, - request.request_id, + request_id, exc_info=True, ) + self._stop_waiting(request_id, state) state.done = True return True if checkpoint is None: - # Resource pressure is not a cache miss. Keep the answered manifest - # while runnable owners release capacity. - if not state.capacity_refusal_logged: - state.capacity_refusal_logged = True + if not self._reserve_admission: + # Insufficient GPU capacity is not an external-cache miss. Permit + # ordinary admission, but retain the manifest so an unadmitted + # request can retry its import after other owners release pages. + if not state.capacity_refusal_logged: + state.capacity_refusal_logged = True + logger.info( + "Recurrent checkpoint restore of %d tokens deferred for " + "request %s: not enough free GPU blocks (%d free); it is " + "admitted without the restore unless capacity returns first", + manifest.prefix.num_tokens, + request_id, + self._manager.block_pool.get_num_free_blocks(), + ) + return True + # Resource pressure is not a cache miss. vLLM queues this request + # for the next restore reservation; keep the answered manifest. + self._wait( + request_id, + state, + manifest.prefix.num_tokens, + "GPU capacity (%d free blocks; the checkpoint needs %d plus " + "execution headroom)", + self._manager.block_pool.get_num_free_blocks(), + sum(len(group) for group in positions) + 1, + ) + return False + if state.waiting_since is not None: + waited = time.monotonic() - state.waiting_since + state.waiting_since = None + if waited >= _WAIT_LOG_SECONDS: logger.info( - "Recurrent checkpoint restore of %d tokens deferred for " - "request %s: waiting for restore and execution capacity " - "(%d free GPU blocks)", + "Recurrent checkpoint restore of %d tokens for request %s " + "starts after waiting %.0f s for capacity", manifest.prefix.num_tokens, - request.request_id, - self._manager.block_pool.get_num_free_blocks(), + request_id, + waited, ) - return False task = self._make_task(manifest, checkpoint, "RETRIEVE") logger.debug( "Recurrent checkpoint restore of %d tokens copying for request %s", manifest.prefix.num_tokens, - request.request_id, + request_id, ) - self._tasks[task.task_id] = _PendingTask(task, checkpoint, request.request_id) + self._tasks[task.task_id] = _PendingTask(task, checkpoint, request_id) state.task_id = task.task_id + # The copy carries the manifest; a retry looks the prefix up again. + state.manifest = None + state.plan = None + state.refresh = None return False def external_tokens(self, request: "Request") -> int: @@ -545,7 +668,10 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: self._cache.release(pending.checkpoint) else: state = self._lookups.get(pending.request_id) + if state is not None: + state.task_id = None retry = False + missed = False if ( all(pending.acknowledgements.values()) and pending.request_id not in self._cancelled @@ -558,6 +684,7 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: ) if published and state is not None: state.checkpoint_id = pending.checkpoint.checkpoint_id + state.selected_tokens = pending.checkpoint.num_tokens logger.debug( "Recurrent checkpoint restore of %d tokens ready for " "request %s; ownership retained until admission", @@ -565,14 +692,18 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: pending.request_id, ) elif state is not None: + # The GPU cache evicted a destination page before + # publication. The stored checkpoint is intact, so the + # retry may receive the same one. retry = state.attempts < _MAX_LOOKUP_ATTEMPTS logger.info( - "Recurrent checkpoint publication invalidated for " - "request %s%s", + "Recurrent checkpoint restore of %d tokens for request " + "%s was invalidated before publication; %s", + pending.task.manifest.prefix.num_tokens, pending.request_id, - "; retrying lookup" + "looking up its prefix again" if retry - else "; recomputing its prompt", + else "recomputing its prompt", ) elapsed = time.monotonic() - pending.created if elapsed > _SLOW_RESTORE_SECONDS: @@ -592,6 +723,7 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: and pending.request_id not in self._cancelled and state.attempts < _MAX_LOOKUP_ATTEMPTS ) + missed = True logger.info( "Recurrent checkpoint restore of %d tokens failed for " "request %s on ranks %s after %.1f s%s", @@ -606,16 +738,19 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: "" if retry else "; recomputing its prompt", ) if state is not None and retry: - # A shorter checkpoint can survive eviction or a restart - # when the longest one did not; look it up again. - logger.info( - "Recurrent checkpoint restore of %d tokens missed for " - "request %s; looking up a shorter checkpoint (attempt %d)", - pending.task.manifest.prefix.num_tokens, - pending.request_id, - state.attempts + 1, - ) - state.task_id = None + if missed: + # A shorter checkpoint can survive eviction or a restart + # when the longest one did not; look it up again. + logger.info( + "Recurrent checkpoint restore of %d tokens missed for " + "request %s; looking up a shorter checkpoint " + "(attempt %d)", + pending.task.manifest.prefix.num_tokens, + pending.request_id, + state.attempts + 1, + ) + # Each attempt has its own reply deadline; capacity waits + # and copies of earlier attempts do not consume it. state.attempts += 1 state.started = time.monotonic() state.future = self._client.submit_request( @@ -628,14 +763,41 @@ def complete(self, results: dict[str, dict[int, bool]]) -> None: self.finish_request(pending.request_id) def finish_request(self, request_id: str) -> None: - """Forget lookup state, retaining copy pins until every admitted task drains.""" - self._manager.release_external_boundary_admission(request_id) + """Forget lookup state, retaining copy pins until every admitted task drains. + + With vLLM restore admission reservations, also release the request's + reservation, or its place among restores waiting for capacity. + """ + if self._reserve_admission: + self._manager.release_external_boundary_admission(request_id) self._lookups.pop(request_id, None) if any(task.request_id == request_id for task in self._tasks.values()): self._cancelled.add(request_id) else: self._cancelled.discard(request_id) + def report_status(self) -> dict[str, float]: + """Report restores waiting for GPU or copy capacity, without advancing them. + + Returns: + ``capacity_waits``: requests whose answered restore is waiting now. + ``longest_capacity_wait_seconds``: age of the oldest such wait, or + 0.0 without one. ``capacity_waits_started``: restores that had to + wait since this bridge was created. Requests that vLLM holds back + before they reach this bridge are not counted. + """ + now = time.monotonic() + waits = [ + now - state.waiting_since + for state in self._lookups.values() + if state.waiting_since is not None + ] + return { + "capacity_waits": len(waits), + "longest_capacity_wait_seconds": max(waits, default=0.0), + "capacity_waits_started": self._waits_started, + } + def _roots(self, request: "Request") -> CheckpointTokenRoots: tokens = request.all_token_ids if request.mm_features: @@ -690,3 +852,130 @@ def _make_task( ) + (checkpoint.auxiliary_block_ids,), ) + + def _reserve( + self, + request: "Request", + num_tokens: int, + payload: dict[str, Any], + positions: tuple[tuple[int, ...], ...], + ) -> "BoundaryCheckpoint | None": + """Reserve restore destinations, returning None under resource pressure.""" + if self._reserve_admission: + # Also reserve a running slot and continuation blocks; vLLM then + # keeps the published restore owned until this request is admitted. + return self._manager.reserve_external_boundary_checkpoint( + request, + num_tokens, + positions, + draft_prefix_len=payload["draft_prefix_len"], + kind=payload["kind"], + num_ranks=self._world_size, + reserve_admission=True, + ) + return self._manager.reserve_external_boundary_checkpoint( + request, + num_tokens, + positions, + draft_prefix_len=payload["draft_prefix_len"], + kind=payload["kind"], + num_ranks=self._world_size, + ) + + def _refresh_manifest( + self, request_id: str, state: _Lookup, manifest: CheckpointManifest + ) -> CheckpointManifest | None: + """Return the retained answer, asking the directory again while it waits. + + A lookup refreshes the eviction recency of its candidate's pages. A + restore that has retained its answer for half the reply deadline sends + the lookup again, one at a time. It keeps using the retained answer + until the reply arrives, which then replaces it. + + Args: + request_id: Request that owns the lookup. + state: Lookup of the current attempt, retaining ``manifest``. + manifest: Answered candidate of the current attempt. + + Returns: + The current candidate, or None once the directory lists no + checkpoint of the prefix. + """ + now = time.monotonic() + if state.refresh is None: + if now - state.refreshed >= self._lookup_timeout / 2: + state.refresh = self._client.submit_request( + RequestType.CHECKPOINT_FIND, [state.roots.roots] + ) + state.refreshed = now + return manifest + if not state.refresh.query(): + return manifest + refresh, state.refresh = state.refresh, None + try: + answer = refresh.result() + except Exception: + logger.warning( + "Recurrent checkpoint lookup refresh failed for request %s; " + "keeping its answered checkpoint", + request_id, + exc_info=True, + ) + return manifest + if answer is None: + logger.info( + "Recurrent checkpoint of %d tokens for request %s is no longer " + "listed; recomputing its prompt", + manifest.prefix.num_tokens, + request_id, + ) + return None + if answer.generation != manifest.generation: + state.manifest = answer + state.plan = None + return answer + + def _wait( + self, request_id: str, state: _Lookup, num_tokens: int, reason: str, *args: int + ) -> None: + """Keep an answered restore waiting; log its start, then periodically. + + Args: + request_id: Request whose restore cannot start yet. + state: Its lookup, retaining the answered checkpoint. + num_tokens: Length of that checkpoint. + reason: Log format naming the missing capacity, filled from args. + args: Current values for ``reason``. + """ + now = time.monotonic() + if state.waiting_since is None: + state.waiting_since = now + state.wait_logged = now + self._waits_started += 1 + logger.info( + "Recurrent checkpoint restore of %d tokens for request %s is " + "waiting for " + reason, + num_tokens, + request_id, + *args, + ) + elif now - state.wait_logged >= _WAIT_LOG_SECONDS: + state.wait_logged = now + logger.info( + "Recurrent checkpoint restore of %d tokens for request %s has " + "waited %.0f s for " + reason, + num_tokens, + request_id, + now - state.waiting_since, + *args, + ) + + def _stop_waiting(self, request_id: str, state: _Lookup) -> None: + """End a capacity wait of a request that proceeds without this restore.""" + if state.waiting_since is None: + return + state.waiting_since = None + if self._reserve_admission: + # vLLM would otherwise still count the request as waiting for a + # restore reservation and hold other admissions back for it. + self._manager.release_external_boundary_admission(request_id) diff --git a/tests/v1/test_vllm_semantic_checkpoint_transfer.py b/tests/v1/test_vllm_semantic_checkpoint_transfer.py index ffa64844a84..28e924643b3 100644 --- a/tests/v1/test_vllm_semantic_checkpoint_transfer.py +++ b/tests/v1/test_vllm_semantic_checkpoint_transfer.py @@ -8,6 +8,7 @@ from types import SimpleNamespace from typing import Any, cast from unittest.mock import patch +import inspect import json import os import time @@ -52,6 +53,7 @@ RecurrentCheckpointMetadata, RecurrentCheckpointWorkerMetadata, ) +from lmcache.v1.multiprocess.checkpoint_index import CheckpointManifest # noqa: E402 from lmcache.v1.multiprocess.checkpoint_storage import ( # noqa: E402 checkpoint_object_keys, ) @@ -74,6 +76,20 @@ reason="vLLM requires the atomic boundary import allocator API", ) +# vLLM before the paired restore-admission change cannot reserve admission +# capacity for a restore; the bridge then keeps its earlier fallback, which +# admits a restore without enough free GPU blocks to recompute its prompt. +_reserve = getattr(KVCacheManager, "reserve_external_boundary_checkpoint", None) +RESTORE_ADMISSION = ( + _reserve is not None + and hasattr(KVCacheManager, "release_external_boundary_admission") + and "reserve_admission" in inspect.signature(_reserve).parameters +) +requires_restore_admission = pytest.mark.skipif( + not RESTORE_ADMISSION, + reason="vLLM cannot reserve admission capacity for checkpoint restores", +) + def test_connector_capability_rejects_mutable_revision_names() -> None: identity = { @@ -294,13 +310,13 @@ def test_lora_requests_do_not_read_or_publish_the_base_weight_namespace() -> Non ) -def make_request(name: str, *, first_token: int = 0) -> Request: - """Construct an exact eleven-token prompt with ordinary cache authentication.""" +def make_request(name: str, *, first_token: int = 0, num_tokens: int = 11) -> Request: + """Construct an exact prompt, eleven tokens by default, with cache hashes.""" params = SamplingParams(max_tokens=1) params.update_from_generation_config({}, eos_token_id=100) return Request( request_id=name, - prompt_token_ids=list(range(first_token, first_token + 11)), + prompt_token_ids=list(range(first_token, first_token + num_tokens)), sampling_params=params, pooling_params=None, block_hasher=get_request_block_hasher(4, sha256), @@ -486,10 +502,10 @@ def test_worker_metadata_requires_bound_rank() -> None: "rejected-store-reused-request-id", "inflight-store-reused-request-id", "capacity-retry", - "capacity-timeout", - "task-capacity", - "local-hit-during-copy", - "invalidated-before-use", + pytest.param("capacity-timeout", marks=requires_restore_admission), + pytest.param("task-capacity", marks=requires_restore_admission), + pytest.param("local-hit-during-copy", marks=requires_restore_admission), + pytest.param("invalidated-before-use", marks=requires_restore_admission), "publication-invalidated", ], ) @@ -600,16 +616,28 @@ def record(kind: RequestType, args: list[Any]) -> Any: assert sum(len(root.tail_tokens) for root in roots) >= prefix.num_tokens consumer = make_request("consumer") if outcome in ("capacity-retry", "capacity-timeout"): - # Resource pressure must not authorize cold admission. pressure = manager.block_pool.get_new_blocks( manager.block_pool.get_num_free_blocks() ) try: - for _ in range(100): + if RESTORE_ADMISSION: + # Resource pressure must not authorize cold admission, + # even after the lookup deadline has passed. + for _ in range(100): + assert not bridge.poll_prefix(consumer) + time.sleep(0.001) + assert bridge.take_tasks() == [] assert not bridge.poll_prefix(consumer) - time.sleep(0.001) - assert bridge.take_tasks() == [] - assert not bridge.poll_prefix(consumer) + else: + # Without reservations, ordinary admission may proceed, + # and the manifest stays retryable until this consumer + # is admitted or cancelled. + deadline = time.monotonic() + 5 + while not bridge.poll_prefix(consumer): + assert time.monotonic() < deadline + time.sleep(0.001) + assert bridge.take_tasks() == [] + assert bridge.poll_prefix(consumer) finally: manager.block_pool.free_blocks(pressure) inflight_store = None @@ -703,33 +731,46 @@ def drain_predecessor() -> None: manager.boundary_checkpoints.invalidate_block( restore_task.block_ids[-1][0] ) - for rank in range(4): - future = worker.submit( - CheckpointTransferJob( - restore_task.manifest, - rank, - "RETRIEVE", - restore_task.block_ids, - ) + messages: list[str] = [] + with monkeypatch.context() as logs: + logs.setattr( + checkpoint_scheduler.logger, + "info", + lambda message, *args: messages.append(message), ) - assert future is not None and future.result(timeout=5) - bridge.complete( - { - restore_task.task_id: { - rank: not (outcome == "rank-miss" and rank == 2) + for rank in range(4): + future = worker.submit( + CheckpointTransferJob( + restore_task.manifest, + rank, + "RETRIEVE", + restore_task.block_ids, + ) + ) + assert future is not None and future.result(timeout=5) + bridge.complete( + { + restore_task.task_id: { + rank: not (outcome == "rank-miss" and rank == 2) + } } - } - ) - if rank < 3: - if outcome != "local-hit-during-copy": - assert manager.get_computed_blocks(consumer)[1] == 0 - assert manager.block_pool.get_num_free_blocks() < before + ) + if rank < 3: + if outcome != "local-hit-during-copy": + assert manager.get_computed_blocks(consumer)[1] == 0 + assert manager.block_pool.get_num_free_blocks() < before + if outcome == "publication-invalidated": + # Local eviction is not a directory miss; the retry may find + # the same checkpoint rather than a shorter one. + assert any("invalidated before publication" in m for m in messages) + assert not any("missed" in m for m in messages) expected_tokens = ( 0 if outcome in ("rank-miss", "cancelled", "publication-invalidated") else 11 ) - if expected_tokens: + if expected_tokens and RESTORE_ADMISSION: + # The consumer owns its published restore until admission. assert manager.block_pool.get_num_free_blocks() < before assert not manager.reset_prefix_cache() else: @@ -747,6 +788,7 @@ def drain_predecessor() -> None: assert not bridge.poll_prefix(consumer) bridge.finish_request(consumer.request_id) assert manager.block_pool.get_num_free_blocks() == before + assert_restore_queue_empty(manager) if inflight_store is not None: drain_predecessor() assert not bridge.has_pending @@ -868,11 +910,15 @@ def run(task: CheckpointEngineTask) -> dict[int, bool]: assert manager.get_computed_blocks(consumer)[1] == 8 assert bridge.external_tokens(consumer) == 8 assert not bridge.has_pending + bridge.finish_request(consumer.request_id) + assert_restore_queue_empty(manager) finally: worker.close() -def make_bridge(client, manager, lookup_timeout: float) -> CheckpointSchedulerBridge: +def make_bridge( + client, manager, lookup_timeout: float, *, max_tasks: int = 32 +) -> CheckpointSchedulerBridge: """Bridge with negotiated layouts and a short directory reply deadline.""" bridge = CheckpointSchedulerBridge( manager, @@ -884,6 +930,7 @@ def make_bridge(client, manager, lookup_timeout: float) -> CheckpointSchedulerBr "parallel": {"tp": 4, "dcp": 1}, }, 4, + max_tasks=max_tasks, lookup_timeout=lookup_timeout, ) bridge.accept_layouts( @@ -892,6 +939,114 @@ def make_bridge(client, manager, lookup_timeout: float) -> CheckpointSchedulerBr return bridge +def assert_restore_queue_empty(manager: KVCacheManager) -> None: + """No restore reservation or queued restore holds back other admissions.""" + if RESTORE_ADMISSION: + assert manager.can_admit_external_boundary_request("unrelated-request") + assert manager.external_boundary_reserved_blocks() == 0 + + +class FakeDirectory: + """Checkpoint RPC double that lists the last stored manifest. + + ``listed`` and ``replies`` apply to lookups sent after they change. + """ + + def __init__(self) -> None: + self.listed: CheckpointManifest | None = None + self.replies = True + self.finds = 0 + + def submit_request(self, kind: RequestType, args: list[Any]) -> SimpleNamespace: + if kind == RequestType.CHECKPOINT_BEGIN: + self.listed = args[0] + if kind != RequestType.CHECKPOINT_FIND: + return SimpleNamespace(query=lambda: True, result=lambda: True) + self.finds += 1 + answered, listed = self.replies, self.listed + return SimpleNamespace(query=lambda: answered, result=lambda: listed) + + +class LegacyAllocator: + """A real allocator seen through vLLM's API before restore admission. + + Reservations lack ``reserve_admission`` and the admission methods are + missing, so any use of them fails the test. + """ + + missing = ( + "can_admit_external_boundary_request", + "external_boundary_admission_ready", + "external_boundary_reserved_blocks", + "has_external_boundary_admission", + "has_pending_external_boundary_admissions", + "release_external_boundary_admission", + "set_external_boundary_admission_context", + ) + + def __init__(self, manager: KVCacheManager) -> None: + self.manager = manager + + def __getattr__(self, name: str) -> Any: + if name in LegacyAllocator.missing: + raise AttributeError(name) + return getattr(self.manager, name) + + def reserve_external_boundary_checkpoint( + self, + request: Request, + num_tokens: int, + page_positions: tuple[tuple[int, ...], ...], + *, + draft_prefix_len: int, + kind: str, + num_ranks: int, + ) -> BoundaryCheckpoint | None: + return self.manager.reserve_external_boundary_checkpoint( + request, + num_tokens, + page_positions, + draft_prefix_len=draft_prefix_len, + kind=kind, + num_ranks=num_ranks, + ) + + +def publish_local( + manager: KVCacheManager, request: Request, prefix: int +) -> BoundaryCheckpoint: + """Publish a GPU-cached checkpoint of the request's first tokens.""" + checkpoint = manager.reserve_external_boundary_checkpoint( + request, + prefix, + manager.boundary_checkpoint_page_positions(prefix), + draft_prefix_len=prefix, + kind="prompt", + num_ranks=1, + ) + assert checkpoint is not None + assert manager.acknowledge_external_boundary_checkpoint(checkpoint.checkpoint_id, 0) + return checkpoint + + +def publish_external( + bridge: CheckpointSchedulerBridge, + manager: KVCacheManager, + prefix: int, + *, + num_tokens: int = 11, +) -> Request: + """Store a producer's checkpoint and drop every GPU copy of it.""" + producer = make_request("producer", num_tokens=num_tokens) + bridge.store(producer, publish_local(manager, producer, prefix)) + (store,) = bridge.take_tasks() + bridge.complete({store.task_id: {rank: True for rank in range(4)}}) + bridge.finish_request(producer.request_id) + assert manager.reset_prefix_cache() + return producer + + +@requires_restore_admission @pytest.mark.parametrize("prefix, already_local", [(11, False), (8, True)]) def test_local_reuse_retries_external_restore_after_capacity_eviction( prefix: int, already_local: bool @@ -989,7 +1144,7 @@ def publish() -> BoundaryCheckpoint: bridge.finish_request(consumer.request_id) manager.free(consumer) assert not bridge.has_pending - assert manager.external_boundary_reserved_blocks() == 0 + assert_restore_queue_empty(manager) assert manager.block_pool.get_num_free_blocks() == 63 @@ -1023,6 +1178,8 @@ def test_lookup_without_a_usable_reply_admits_the_request( assert bridge.take_tasks() == [] assert bridge.external_tokens(consumer) == 0 assert not bridge.has_pending + bridge.finish_request(consumer.request_id) + assert_restore_queue_empty(manager) def test_unanswered_store_begin_is_aborted_after_the_timeout( @@ -1064,3 +1221,355 @@ def submit(kind: RequestType, *_args) -> SimpleNamespace: assert submitted == [RequestType.CHECKPOINT_BEGIN, RequestType.CHECKPOINT_ABORT] assert not bridge.has_pending assert manager.reset_prefix_cache() + + +@pytest.mark.parametrize("api", ["complete", "without-release", "without-parameter"]) +def test_restore_admission_api_is_detected_before_use(api: str) -> None: + """Reservation calls reach only a vLLM that has every reservation method.""" + + def reserve( + request: Request, + num_tokens: int, + page_positions: tuple[tuple[int, ...], ...], + *, + draft_prefix_len: int, + kind: str, + num_ranks: int, + reserve_admission: bool = False, + ) -> None: + return None + + def reserve_before_admission( + request: Request, + num_tokens: int, + page_positions: tuple[tuple[int, ...], ...], + *, + draft_prefix_len: int, + kind: str, + num_ranks: int, + ) -> None: + return None + + released: list[str] = [] + manager = SimpleNamespace( + boundary_checkpoints=SimpleNamespace(), + reserve_external_boundary_checkpoint=( + reserve_before_admission if api == "without-parameter" else reserve + ), + can_admit_external_boundary_request=lambda request_id: True, + external_boundary_admission_ready=lambda request_id: False, + ) + if api != "without-release": + manager.release_external_boundary_admission = released.append + with patch.object(checkpoint_scheduler.logger, "warning") as warning: + bridge = make_bridge(FakeDirectory(), manager, 60) + bridge.finish_request("finished") + supported = api == "complete" + assert released == (["finished"] if supported else []) + assert warning.call_count == (0 if supported else 1) + + +@pytest.mark.parametrize("pressure", ["gpu-blocks", "copy-slots"]) +def test_restores_without_admission_reservations_keep_the_earlier_fallback( + pressure: str, +) -> None: + """An older vLLM receives no reservation calls and keeps the prior fallback. + + A restore that finds no free GPU blocks or copy slots is admitted to + recompute its prompt; a request still waiting for ordinary admission + retries its GPU reservation with the answer it already has. + """ + manager = make_manager() + directory = FakeDirectory() + with patch.object(checkpoint_scheduler.logger, "warning") as warning: + bridge = make_bridge( + directory, + LegacyAllocator(manager), + 60, + max_tasks=1 if pressure == "copy-slots" else 32, + ) + assert warning.call_count == 1 + publish_external(bridge, manager, 11) + consumer = make_request("consumer") + assert not bridge.poll_prefix(consumer) + if pressure == "gpu-blocks": + blocked = manager.block_pool.get_new_blocks( + manager.block_pool.get_num_free_blocks() + ) + else: + blocker = make_request("blocker", first_token=20) + bridge.store(blocker, publish_local(manager, blocker, 11)) + (store,) = bridge.take_tasks() + with ( + patch.object(checkpoint_scheduler.logger, "info") as info, + patch.object(checkpoint_scheduler.logger, "warning") as warning, + ): + assert bridge.poll_prefix(consumer) + assert bridge.poll_prefix(consumer) + # A reservation passing reserve_admission would have been rejected. + assert not warning.called + (message,) = (call.args[0] for call in info.call_args_list) + assert ("deferred" if pressure == "gpu-blocks" else "skipped") in message + assert bridge.take_tasks() == [] + if pressure == "gpu-blocks": + manager.block_pool.free_blocks(blocked) + assert not bridge.poll_prefix(consumer) + (restore,) = bridge.take_tasks() + bridge.complete({restore.task_id: {rank: True for rank in range(4)}}) + else: + bridge.complete({store.task_id: {rank: True for rank in range(4)}}) + bridge.finish_request(blocker.request_id) + assert bridge.take_tasks() == [] + assert bridge.poll_prefix(consumer) + assert directory.finds == 1 + expected = 11 if pressure == "gpu-blocks" else 0 + assert manager.get_computed_blocks(consumer)[1] == expected + assert bridge.external_tokens(consumer) == expected + bridge.finish_request(consumer.request_id) + assert not bridge.has_pending + + +def test_complete_local_hit_does_not_wait_for_the_directory() -> None: + """A full-prompt checkpoint in the GPU cache needs no directory reply.""" + manager = make_manager() + directory = FakeDirectory() + directory.replies = False + bridge = make_bridge(directory, manager, 60) + consumer = make_request("consumer") + assert not bridge.poll_prefix(consumer) + # A sibling request with the same prompt publishes its checkpoint. + publish_local(manager, make_request("sibling"), 11) + assert bridge.poll_prefix(consumer) + assert manager.get_computed_blocks(consumer)[1] == 11 + assert bridge.external_tokens(consumer) == 0 + assert directory.finds == 1 + # Losing that checkpoint before admission starts a new lookup when vLLM + # can reserve restores, instead of resuming the unanswered one. + assert manager.reset_prefix_cache() + assert not bridge.poll_prefix(consumer) + assert directory.finds == (2 if RESTORE_ADMISSION else 1) + bridge.finish_request(consumer.request_id) + assert_restore_queue_empty(manager) + + +@requires_restore_admission +def test_longer_local_checkpoint_keeps_the_selection() -> None: + """A selected checkpoint superseded by a longer one needs no new lookup.""" + manager = make_manager() + directory = FakeDirectory() + bridge = make_bridge(directory, manager, 60) + producer = publish_external(bridge, manager, 8, num_tokens=15) + consumer = make_request("consumer", num_tokens=15) + publish_local(manager, producer, 8) + assert not bridge.poll_prefix(consumer) + # The local checkpoint is as long as the stored one, so it is selected. + assert bridge.poll_prefix(consumer) + assert manager.get_computed_blocks(consumer)[1] == 8 + # Before ordinary admission succeeds, a sibling publishes more tokens. + publish_local(manager, producer, 12) + with patch.object(checkpoint_scheduler.logger, "info") as info: + assert bridge.poll_prefix(consumer) + assert not info.called + assert directory.finds == 1 + assert manager.get_computed_blocks(consumer)[1] == 12 + assert bridge.external_tokens(consumer) == 0 + bridge.finish_request(consumer.request_id) + assert_restore_queue_empty(manager) + + +@requires_restore_admission +def test_preempted_import_reuses_a_longer_local_checkpoint() -> None: + """A preempted import whose prompt now has a longer checkpoint resumes at once. + + The directory is unresponsive, so another lookup would hold the request. + """ + manager = make_manager() + directory = FakeDirectory() + bridge = make_bridge(directory, manager, 60) + producer = publish_external(bridge, manager, 8, num_tokens=15) + consumer = make_request("consumer", num_tokens=15) + assert not bridge.poll_prefix(consumer) + assert not bridge.poll_prefix(consumer) + (restore,) = bridge.take_tasks() + bridge.complete({restore.task_id: {rank: True for rank in range(4)}}) + assert bridge.poll_prefix(consumer) + blocks, tokens, _ = manager.get_computed_blocks(consumer) + assert tokens == 8 and bridge.external_tokens(consumer) == 8 + assert ( + manager.allocate_slots( + consumer, + consumer.num_tokens - tokens, + num_new_computed_tokens=tokens, + new_computed_blocks=blocks, + ) + is not None + ) + _, copies = manager.take_kv_cache_block_copies() + manager.block_pool.free_blocks(copies) + # While running, the prompt gains a longer checkpoint (as the request's + # own prompt capture would publish); then the request is preempted. + publish_local(manager, producer, 12) + manager.free(consumer) + directory.replies = False + assert bridge.poll_prefix(consumer) + assert directory.finds == 1 + assert manager.get_computed_blocks(consumer)[1] == 12 + bridge.finish_request(consumer.request_id) + assert_restore_queue_empty(manager) + + +@requires_restore_admission +def test_capacity_wait_is_logged_periodically_and_reported( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """A restore starved of GPU blocks stays visible while it waits. + + Its answer is validated once, not on every scheduler step of the wait. + """ + monkeypatch.setattr(checkpoint_scheduler, "_WAIT_LOG_SECONDS", 0.05) + parsed: list[CheckpointManifest] = [] + page_groups = checkpoint_scheduler.checkpoint_page_groups + + def parse(manifest: CheckpointManifest) -> Any: + parsed.append(manifest) + return page_groups(manifest) + + monkeypatch.setattr(checkpoint_scheduler, "checkpoint_page_groups", parse) + manager = make_manager() + bridge = make_bridge(FakeDirectory(), manager, 60) + publish_external(bridge, manager, 11) + consumer = make_request("consumer") + assert not bridge.poll_prefix(consumer) + pressure = manager.block_pool.get_new_blocks( + manager.block_pool.get_num_free_blocks() + ) + with patch.object(checkpoint_scheduler.logger, "info") as info: + for _ in range(10): + assert not bridge.poll_prefix(consumer) + time.sleep(0.06) + assert not bridge.poll_prefix(consumer) + started, repeated = (call.args for call in info.call_args_list) + assert "is waiting for GPU capacity" in started[0] + # Tokens, request, free blocks, and blocks the checkpoint alone needs. + assert started[1:4] == (11, consumer.request_id, 0) + assert started[4] > 0 + assert "has waited" in repeated[0] + assert repeated[1:3] == (11, consumer.request_id) + assert repeated[3] >= 0.05 + assert repeated[4:] == started[3:] + status = bridge.report_status() + assert status["capacity_waits"] == 1 + assert status["longest_capacity_wait_seconds"] >= 0.05 + assert status["capacity_waits_started"] == 1 + manager.block_pool.free_blocks(pressure) + with patch.object(checkpoint_scheduler.logger, "info") as info: + assert not bridge.poll_prefix(consumer) + (resumed,) = (call.args for call in info.call_args_list) + assert "starts after waiting" in resumed[0] + assert bridge.report_status()["capacity_waits"] == 0 + (restore,) = bridge.take_tasks() + bridge.complete({restore.task_id: {rank: True for rank in range(4)}}) + assert bridge.poll_prefix(consumer) + assert bridge.external_tokens(consumer) == 11 + assert len(parsed) == 1 + bridge.finish_request(consumer.request_id) + assert_restore_queue_empty(manager) + assert bridge.report_status() == { + "capacity_waits": 0, + "longest_capacity_wait_seconds": 0.0, + "capacity_waits_started": 1, + } + + +@requires_restore_admission +@pytest.mark.parametrize( + "exit_path", + ["finished", "restored", "local-covers", "complete-local", "unlisted", "rejected"], +) +def test_capacity_waiter_leaves_the_restore_queue(exit_path: str) -> None: + """Every way out of a capacity wait lets other requests be admitted.""" + manager = make_manager() + directory = FakeDirectory() + bridge = make_bridge(directory, manager, 0.2) + producer = publish_external(bridge, manager, 8) + consumer = make_request("consumer") + assert not bridge.poll_prefix(consumer) + pressure = manager.block_pool.get_new_blocks( + manager.block_pool.get_num_free_blocks() + ) + assert not bridge.poll_prefix(consumer) + assert bridge.report_status()["capacity_waits"] == 1 + if exit_path == "finished": + bridge.finish_request(consumer.request_id) + elif exit_path == "restored": + manager.block_pool.free_blocks(pressure) + pressure = [] + assert not bridge.poll_prefix(consumer) + (restore,) = bridge.take_tasks() + bridge.complete({restore.task_id: {rank: True for rank in range(4)}}) + assert bridge.poll_prefix(consumer) + elif exit_path in ("local-covers", "complete-local"): + prefix = 8 if exit_path == "local-covers" else 11 + pages = sum(map(len, manager.boundary_checkpoint_page_positions(prefix))) + 1 + manager.block_pool.free_blocks(pressure[:pages]) + pressure = pressure[pages:] + publish_local(manager, producer, prefix) + assert bridge.poll_prefix(consumer) + assert manager.get_computed_blocks(consumer)[1] == prefix + else: + # A lookup re-sent during the wait finds the checkpoint gone, or + # finds one this engine cannot restore. + assert directory.listed is not None + directory.listed = ( + None + if exit_path == "unlisted" + else replace(directory.listed, generation="other-layout", world_size=2) + ) + time.sleep(0.11) + assert not bridge.poll_prefix(consumer) + assert directory.finds == 2 + assert bridge.poll_prefix(consumer) + assert bridge.take_tasks() == [] + assert manager.can_admit_external_boundary_request("unrelated-request") + assert bridge.report_status()["capacity_waits"] == 0 + manager.block_pool.free_blocks(pressure) + bridge.finish_request(consumer.request_id) + assert_restore_queue_empty(manager) + assert not bridge.has_pending + + +@requires_restore_admission +def test_waiting_restore_refreshes_its_lookup_without_blocking() -> None: + """A long capacity wait sends its lookup again but never waits for it. + + The lookup refreshes the eviction recency of the checkpoint's pages. + """ + manager = make_manager() + directory = FakeDirectory() + bridge = make_bridge(directory, manager, 0.2) + publish_external(bridge, manager, 11) + consumer = make_request("consumer") + assert not bridge.poll_prefix(consumer) + pressure = manager.block_pool.get_new_blocks( + manager.block_pool.get_num_free_blocks() + ) + assert not bridge.poll_prefix(consumer) + assert not bridge.poll_prefix(consumer) + assert directory.finds == 1 + time.sleep(0.11) + # Half the reply deadline has passed; this re-sent lookup is never answered. + directory.replies = False + assert not bridge.poll_prefix(consumer) + assert directory.finds == 2 + time.sleep(0.11) + assert not bridge.poll_prefix(consumer) + assert directory.finds == 2 + # Capacity starts the restore from the retained answer, beyond the deadline. + manager.block_pool.free_blocks(pressure) + assert not bridge.poll_prefix(consumer) + (restore,) = bridge.take_tasks() + bridge.complete({restore.task_id: {rank: True for rank in range(4)}}) + assert bridge.poll_prefix(consumer) + assert bridge.external_tokens(consumer) == 11 + bridge.finish_request(consumer.request_id) + assert_restore_queue_empty(manager) From 9750f48c2e0f992efbf15f606741fcb73cd7537d Mon Sep 17 00:00:00 2001 From: Martin Vit Date: Mon, 28 Sep 2026 14:55:19 +0000 Subject: [PATCH 4/5] docs(release): describe restore-admission fallback, waits and timeouts Scope the fragment to the request-boundary recurrent models, document the fallback with an older vLLM, the per-reply lookup timeout (up to four replies for a request whose restores fail), wait logging, lookup refresh and the selection and local-hit rules. It still requires the paired vLLM fragment vllm-restore-admission. Co-Authored-By: Claude Opus 5.5 --- .lil/changes/lmcache-restore-admission.json | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/.lil/changes/lmcache-restore-admission.json b/.lil/changes/lmcache-restore-admission.json index 2a460446ba8..7d5ca4e8164 100644 --- a/.lil/changes/lmcache-restore-admission.json +++ b/.lil/changes/lmcache-restore-admission.json @@ -4,13 +4,17 @@ "category": "fix", "summary": "Recurrent checkpoint restores wait for GPU and copy capacity instead of admitting requests to recompute cached prefixes.", "models": [ - "all" + "Qwen3.8-Flash-Next", + "GLM-5.3-Flash" ], - "compatibility": "Requires the paired vLLM restore-admission change; deploy vLLM before enabling this LMCache change.", + "compatibility": "Waiting needs the paired vLLM restore-admission change. With an older vLLM, LMCache logs one warning at startup and keeps the previous behavior: a restore without enough free GPU blocks or copy slots is admitted and its prompt is recomputed.", "details": [ - "Capacity waits retain the answered checkpoint lookup and do not consume the timeout for an unanswered directory request.", - "A local checkpoint evicted before admission triggers a fresh external lookup, while genuine misses and exhausted transfer retries retain recomputation fallback.", - "This change does not repair missing checkpoint pages in cache storage." + "A waiting restore keeps its directory answer; the wait does not count against lmcache.mp.mq_timeout. It is logged when it starts and every 30 s with its age and what it waits for: free GPU blocks against the blocks its checkpoint needs, or copies in flight.", + "lmcache.mp.mq_timeout now limits each directory reply separately, including each retry after a failed restore, instead of all lookups of a request together: a request whose restores fail three times can wait up to four times mq_timeout for replies.", + "A restore waiting longer than half of lmcache.mp.mq_timeout sends its lookup again, which keeps the checkpoint's pages recent in RAM and disk eviction order. Its answer replaces the kept one; if the directory lists no checkpoint of the prompt any more, the prompt is recomputed.", + "A checkpoint chosen before admission is looked up again only if the GPU cache no longer holds one at least as long; a longer local checkpoint that appears meanwhile is used without a lookup.", + "A request whose complete prompt is cached on the GPU is admitted without waiting for the directory, unless its own restore copy is in flight.", + "Genuine misses and exhausted restore retries still recompute the prompt. This change does not repair missing checkpoint pages in cache storage." ], "requires": [ "vllm-restore-admission" From 4ac0b48e75b4cb8a1d065c9dfc7421df911e98b1 Mon Sep 17 00:00:00 2001 From: Martin Vit Date: Mon, 28 Sep 2026 15:47:28 +0000 Subject: [PATCH 5/5] docs(release): link the restore-admission fragment to PR 103 Co-Authored-By: Claude Opus 5.5 --- .lil/changes/lmcache-restore-admission.json | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/.lil/changes/lmcache-restore-admission.json b/.lil/changes/lmcache-restore-admission.json index 7d5ca4e8164..a5e95b50ec5 100644 --- a/.lil/changes/lmcache-restore-admission.json +++ b/.lil/changes/lmcache-restore-admission.json @@ -20,6 +20,11 @@ "vllm-restore-admission" ], "pull_requests": [ - 100 + 100, + 103 + ], + "authors": [ + "ktsaou", + "Local Inference Lab" ] }