diff --git a/.lil/changes/lmcache-106.json b/.lil/changes/lmcache-106.json new file mode 100644 index 00000000000..7e95ad9d635 --- /dev/null +++ b/.lil/changes/lmcache-106.json @@ -0,0 +1,20 @@ +{ + "schema": "local-inference-release-change/v1", + "id": "lmcache-106", + "category": "fix", + "summary": "A checkpoint restore that waits too long for the disk tier no longer stops the model; the request recomputes its prompt", + "models": [ + "GLM-5.3-Flash", + "Qwen3.8-Flash-Next" + ], + "compatibility": "No user action required.", + "details": [ + "A restore that waited 30 s for the disk tier (for example behind a large burst of on-evict writes) was cancelled, and if the disk still had not answered 30 s later both ranks raised UnsafeCheckpointCopyError ('Checkpoint lookup cancellation did not drain', 'Cancelled checkpoint lease ownership could not drain'). The model stopped, and the cache logged 'Checkpoint worker copy leases must drain before close' at shutdown.", + "A restore still waiting for the disk has handed no memory to the GPU. The model now waits at most 1 s for the cancellation and treats the restore as a miss: the request looks its checkpoint up again or recomputes its prompt, other requests keep being served, and later restores and stores work. The cache releases the cancelled lookup once the disk answers, and shutdown no longer waits for it.", + "Copies that did receive cache memory keep the existing fail-safe behavior." + ], + "pull_requests": [ + 106 + ], + "requires": [] +} diff --git a/docs/source/mp/l2_storage/index.rst b/docs/source/mp/l2_storage/index.rst index cfdac3e8eab..291e8d99915 100644 --- a/docs/source/mp/l2_storage/index.rst +++ b/docs/source/mp/l2_storage/index.rst @@ -509,6 +509,14 @@ SIGKILL; with a shorter stop timeout the most valuable pages are written first, and the checkpoints left incomplete are retired by lookups after the restart. +**Slow storage.** A restore waits up to 30 s for its pages to be looked up +and loaded, for example behind a burst of on-evict writes to a slow disk. +Then the engine cancels it, and the request looks the checkpoint up again or +recomputes its prompt; the engine keeps serving. A lookup still waiting for +the disk has handed no memory to the GPU, so the engine waits at most a +second for the cancellation and the cache releases the lookup once the disk +answers. A clean shutdown does not wait for such lookups. + **Sizing.** With write-through (``default``) the L2 retention time is about:: L2 capacity / (requests per minute x checkpoint bytes per request) diff --git a/lmcache/v1/multiprocess/checkpoint_storage.py b/lmcache/v1/multiprocess/checkpoint_storage.py index 9918aeb8562..fd1fb6dbd1e 100644 --- a/lmcache/v1/multiprocess/checkpoint_storage.py +++ b/lmcache/v1/multiprocess/checkpoint_storage.py @@ -7,6 +7,7 @@ """ # Standard +from collections import OrderedDict from dataclasses import dataclass, field from typing import TYPE_CHECKING import hashlib @@ -48,6 +49,9 @@ # seconds, as store admission does: a pass frees nothing while its victims are # still pinned by other restores or L2 writes. _ROOM_REQUEST_INTERVAL_SECONDS = 0.5 +# Released cancelled lookups remembered so that a later poll from their worker +# still gets the miss; a poll of an older one raises KeyError. +_MAX_RELEASED_CANCELLED_LOOKUPS = 4096 @dataclass(frozen=True) @@ -254,6 +258,8 @@ class CheckpointPayloadStore: The caller stages a manifest with ``index.begin`` before rank stores. Each successful store acknowledgement follows a drained worker D2H transfer. Retrieval completion similarly follows H2D completion, not MQ delivery. + A cancelled lookup belongs to this store: its worker may stop polling it, + and the next call on this store after its storage lookup ends releases it. """ def __init__( @@ -277,6 +283,10 @@ def __init__( self._store_ranks: set[tuple[str, int]] = set() self._retrieves: dict[str, _RetrieveLease] = {} self._retrieve_admissions: set[str] = set() + # Cancelled lookups still waiting for storage, and the last ones + # released without a poll from their worker. + self._cancelled: set[str] = set() + self._released_cancelled: OrderedDict[str, None] = OrderedDict() self._lock = threading.Lock() def prepare_store( @@ -298,6 +308,7 @@ def prepare_store( ValueError: For invalid layouts or a duplicate producer rank. """ self._maybe_reclaim() + self._release_cancelled() groups = checkpoint_page_groups(manifest) key_groups = checkpoint_object_keys(manifest, rank) if not self._index.is_pending(manifest): @@ -468,6 +479,7 @@ def begin_retrieve(self, manifest: CheckpointManifest, rank: int) -> str | None: ValueError: If the manifest layout or rank is invalid. """ self._maybe_reclaim() + self._release_cancelled() keys = [ key for group in checkpoint_object_keys(manifest, rank) for key in group ] @@ -552,17 +564,14 @@ def reclaim_abandoned(self) -> int: ] for lease_id in lookups: self._retrieves[lease_id].cancelled = True + self._cancelled.add(lease_id) for _lease_id, store_lease in stores: self._storage.abort_write(store_lease.keys) self._index.abort(store_lease.manifest.generation) for _lease_id, read_lease in ready: self._storage.finish_read_prefetched(read_lease.keys) - for lease_id in lookups: - try: - # A cancelled lookup releases its locks once its prefetch ends. - self.poll_retrieve(lease_id) - except KeyError: - pass + # A cancelled lookup releases its locks once its prefetch ends. + self._release_cancelled() stale = self._index.abort_stale(self._abandoned_after) reclaimed = len(stores) + len(ready) + len(lookups) + len(stale) if reclaimed: @@ -582,13 +591,22 @@ def report_status(self) -> dict[str, int]: Store counts include reservations being prepared. Retrieve counts include prefetch submissions that have not yet returned their handle. - A shutdown coordinator must drain these leases before closing storage. + ``retrieve_lookups`` counts the retrieve leases that exposed no slots + to a worker, and ``cancelled_lookups`` those of them that were + cancelled and wait for storage. A shutdown coordinator must drain the + store leases and the retrieve leases with exposed slots before closing + storage. """ with self._lock: + lookups = sum( + 1 for lease in self._retrieves.values() if lease.slots is None + ) return { "store_leases": len(self._store_ranks), "retrieve_leases": len(self._retrieves) + len(self._retrieve_admissions), + "retrieve_lookups": lookups + len(self._retrieve_admissions), + "cancelled_lookups": len(self._cancelled), "max_leases": self._max_leases, } @@ -611,10 +629,51 @@ def poll_retrieve(self, lease_id: str) -> CheckpointSlots | bool | None: at most the storage admission timeout; a lease that never gets room misses without invalidating its generation. + A cancelled lookup that this store already released returns False + once more. + Raises: KeyError: If the lease is unknown or already finished. ValueError: If stored payload byte layouts do not match the manifest. """ + try: + return self._poll(lease_id) + except KeyError: + with self._lock: + if lease_id not in self._released_cancelled: + raise + del self._released_cancelled[lease_id] + return False + finally: + self._release_cancelled() + + def _release_cancelled(self) -> None: + """Release the cancelled lookups whose storage lookup has ended. + + A worker may stop polling a lookup once it cancelled it; this store + releases the lookup's locks instead, on its next call. + """ + with self._lock: + if not self._cancelled: + return + lease_ids = list(self._cancelled) + for lease_id in lease_ids: + try: + self._poll(lease_id) + except KeyError: + pass + except Exception: + # reclaim_abandoned retries it once the lease is abandoned. + logger.exception("Could not release a cancelled checkpoint lookup") + with self._lock: + self._cancelled.discard(lease_id) + continue + with self._lock: + if lease_id not in self._retrieves: + self._cancelled.discard(lease_id) + + def _poll(self, lease_id: str) -> CheckpointSlots | bool | None: + """Advance one lease as ``poll_retrieve`` describes.""" with self._lock: lease = self._retrieves[lease_id] if lease.slots is not None: @@ -708,8 +767,17 @@ def _repeat_with_room(self, lease_id: str, lease: _RetrieveLease) -> bool | None return None def _miss(self, lease_id: str, lease: _RetrieveLease, outcome: str) -> bool: - """Forget a lease whose read locks are released, and say why.""" + """Forget a lease whose read locks are released, and say why. + + A cancelled lookup is remembered, so a later poll from its worker + still gets the miss. + """ del self._retrieves[lease_id] + if lease.cancelled: + self._cancelled.discard(lease_id) + self._released_cancelled[lease_id] = None + while len(self._released_cancelled) > _MAX_RELEASED_CANCELLED_LOOKUPS: + self._released_cancelled.popitem(last=False) logger.info( "Checkpoint retrieve of %d tokens for rank %d %s after %.1f s: " "%d of %d pages were readable", @@ -744,20 +812,29 @@ def finish_retrieve(self, lease_id: str) -> None: self._storage.finish_read_prefetched(lease.keys) def cancel_retrieve(self, lease_id: str) -> None: - """Mark a pending lookup for draining without exposing its SHM slots. + """Cancel a pending lookup without exposing its SHM slots. + + The store owns the cancelled lookup from then on: it releases the + lookup's locks on its first call after the storage lookup ends, even + if the worker stops polling. A worker that keeps polling gets False + once it is released. Cancelling an unknown or released lease does + nothing. Args: lease_id: Lookup whose consumer was cancelled before H2D submission. - The caller must keep polling until False releases lookup locks. Raises: ValueError: If slots were already exposed; their GPU copy must first drain and use ``finish_retrieve`` instead. """ with self._lock: - lease = self._retrieves[lease_id] + lease = self._retrieves.get(lease_id) + if lease is None: + return if lease.slots is not None: raise ValueError( "prepared checkpoint retrieval requires copy completion" ) lease.cancelled = True + self._cancelled.add(lease_id) + self._release_cancelled() diff --git a/lmcache/v1/multiprocess/checkpoint_transfer.py b/lmcache/v1/multiprocess/checkpoint_transfer.py index ea52bf3d263..9620506ec76 100644 --- a/lmcache/v1/multiprocess/checkpoint_transfer.py +++ b/lmcache/v1/multiprocess/checkpoint_transfer.py @@ -22,6 +22,11 @@ logger = init_logger(__name__) +# Seconds a cancelled lookup may take to drain before the worker stops waiting +# for it. It exposed no slots, and the server releases a cancelled lookup +# itself once storage answers, so its request need not wait any longer. +_CANCELLED_LOOKUP_DRAIN_SECONDS = 1.0 + class UnsafeCheckpointCopyError(RuntimeError): """A copy or SHM lease could not drain; its resources must remain owned. @@ -65,6 +70,10 @@ class CheckpointTransferWorker: rpc_timeout: Metadata reply deadline, in seconds. A timed-out lease acquisition or completion receives one additional deadline to reconcile ownership. Failure is fatal instead of abandoning its SHM reservation. + A retrieve lookup still pending after this long is cancelled and + misses. A cancelled lookup exposed no slots: if storage does not + answer within a second, the worker leaves it to the server, which + releases it once storage answers. Successful store completion acknowledges drained rank bytes, not all-rank manifest publication. The directory exclusively owns that publication. @@ -194,8 +203,19 @@ def _call(self, request: RequestType, *payloads: object) -> Any: def _discard_uncopied_lease( self, lease: CheckpointLeaseResponse, *, store: bool - ) -> None: - """Release a lease whose slots never reached the GPU copy callback.""" + ) -> bool: + """Release a lease whose slots never reached the GPU copy callback. + + A pending lookup is cancelled. It exposed no slots, so if it does not + drain within a second, or the cancellation fails, it is left to the + server, which releases a cancelled lookup once storage answers. + + Returns: + False if a pending lookup was left to the server, otherwise True. + + Raises: + Exception: If a lease with exposed slots could not be finished. + """ def call(request: RequestType, *payloads: object) -> Any: return self._client.submit_request(request, list(payloads)).result( @@ -203,24 +223,49 @@ def call(request: RequestType, *payloads: object) -> Any: ) if lease.status in ("miss", "busy"): - return + return True if store: call(RequestType.CHECKPOINT_FINISH_STORE, lease.lease_id, False) - return + return True if lease.status == "ready": call(RequestType.CHECKPOINT_FINISH_RETRIEVE, lease.lease_id) - return - call(RequestType.CHECKPOINT_CANCEL_RETRIEVE, lease.lease_id) - deadline = time.monotonic() + self._rpc_timeout - while lease.status == "pending": - if time.monotonic() >= deadline: - raise LMCacheTimeoutError( - "Cancelled checkpoint storage lookup did not drain" - ) - time.sleep(0.001) - lease = call(RequestType.CHECKPOINT_POLL_RETRIEVE, lease.lease_id) + return True + try: + call(RequestType.CHECKPOINT_CANCEL_RETRIEVE, lease.lease_id) + deadline = time.monotonic() + min( + self._rpc_timeout, _CANCELLED_LOOKUP_DRAIN_SECONDS + ) + while lease.status == "pending": + if time.monotonic() >= deadline: + return False + time.sleep(0.001) + lease = call(RequestType.CHECKPOINT_POLL_RETRIEVE, lease.lease_id) + except Exception: + logger.warning( + "Could not cancel a pending checkpoint lookup; the LMCache server " + "releases it once it is abandoned", + exc_info=True, + ) + return False if lease.status == "ready": call(RequestType.CHECKPOINT_FINISH_RETRIEVE, lease.lease_id) + return True + + def _discard_retrieve_lookup(self, lease: CheckpointLeaseResponse) -> bool: + """Discard an uncopied retrieve lease; see ``_discard_uncopied_lease``. + + Raises: + UnsafeCheckpointCopyError: If a lease with exposed slots could not + be finished; transfer admission stops. + """ + try: + return self._discard_uncopied_lease(lease, store=False) + except BaseException as drain_error: + with self._lock: + self._unsafe = True + raise UnsafeCheckpointCopyError( + "Cancelled checkpoint lease ownership could not drain" + ) from drain_error def _run(self, job: CheckpointTransferJob) -> bool: store = job.direction == "STORE" @@ -245,30 +290,10 @@ def _run(self, job: CheckpointTransferJob) -> bool: # requests. A pending lookup owns no worker-visible byte slots yet. started = time.monotonic() deadline = started + self._rpc_timeout - cancelled = False try: - while lease.status == "pending": - if not cancelled and ( - self._closing or time.monotonic() >= deadline - ): - if not self._closing: - logger.warning( - "Checkpoint retrieve of %d tokens for rank %d is " - "still waiting for storage after %.0f s; cancelling " - "it, so the request recomputes its prompt", - job.manifest.prefix.num_tokens, - job.rank, - time.monotonic() - started, - ) - self._call( - RequestType.CHECKPOINT_CANCEL_RETRIEVE, lease.lease_id - ) - cancelled = True - deadline = time.monotonic() + self._rpc_timeout - if cancelled and time.monotonic() >= deadline: - raise LMCacheTimeoutError( - "Checkpoint lookup cancellation did not drain" - ) + while lease.status == "pending" and not ( + self._closing or time.monotonic() >= deadline + ): time.sleep(0.001) lease = self._call( RequestType.CHECKPOINT_POLL_RETRIEVE, lease.lease_id @@ -277,16 +302,28 @@ def _run(self, job: CheckpointTransferJob) -> bool: raise except BaseException: # No slots have reached the copy callback. The server can - # safely cancel this lookup, but must drain its storage locks. - try: - self._discard_uncopied_lease(lease, store=False) - except BaseException as drain_error: - with self._lock: - self._unsafe = True - raise UnsafeCheckpointCopyError( - "Cancelled checkpoint lease ownership could not drain" - ) from drain_error + # safely cancel this lookup and release its storage locks. + self._discard_retrieve_lookup(lease) raise + if lease.status == "pending": + if not self._closing: + logger.warning( + "Checkpoint retrieve of %d tokens for rank %d is still " + "waiting for storage after %.0f s; cancelling it, so the " + "request recomputes its prompt", + job.manifest.prefix.num_tokens, + job.rank, + time.monotonic() - started, + ) + if not self._discard_retrieve_lookup(lease): + logger.info( + "Storage has not answered the cancelled checkpoint " + "retrieve of %d tokens for rank %d; the LMCache server " + "releases it once storage answers", + job.manifest.prefix.num_tokens, + job.rank, + ) + return False if lease.status == "miss": return False if lease.status != "ready": diff --git a/lmcache/v1/multiprocess/modules/checkpoint.py b/lmcache/v1/multiprocess/modules/checkpoint.py index 2aa06000ffc..522fd236734 100644 --- a/lmcache/v1/multiprocess/modules/checkpoint.py +++ b/lmcache/v1/multiprocess/modules/checkpoint.py @@ -102,8 +102,10 @@ class CheckpointModule: At shutdown, :meth:`drain_stores` gives workers a bounded time to finish or abort their copy leases while the message queue still serves them. - A lease left after that makes :meth:`close` raise; its buffers are never - recycled while a worker can still access their bytes. + A copy lease left after that makes :meth:`close` raise; its buffers are + never recycled while a worker can still access their bytes. A retrieve + whose lookup is still waiting for storage exposed no buffers to a worker + and does not count. """ def __init__( @@ -425,9 +427,11 @@ def finish_retrieve(self, lease_id: str) -> bool: return True def cancel_retrieve(self, lease_id: str) -> bool: - """Cancel unexposed lookup; callers must poll until its locks drain. + """Cancel an unexposed lookup; the server releases it once storage answers. - Once slots are exposed, use finish_retrieve after H2D completion instead. + The caller may keep polling until the lookup misses, or stop polling: + its locks are released either way. Once slots are exposed, use + finish_retrieve after H2D completion instead. """ self._payloads.cancel_retrieve(lease_id) return True @@ -518,14 +522,23 @@ def close(self) -> None: :meth:`drain_stores` gives workers a bounded time first, and ``MPCacheServer.close`` logs this error and still closes the storage manager, so its shutdown flush runs and its shared memory is released - without reusing a lease's buffers. + without reusing a lease's buffers. A retrieve whose lookup is still + waiting for storage exposed no buffer to a worker and does not block + closing; the storage manager ends its lookup. """ status = self._payloads.report_status() - if status["store_leases"] or status["retrieve_leases"]: + exposed = status["retrieve_leases"] - status["retrieve_lookups"] + if status["store_leases"] or exposed: raise RuntimeError( f"Checkpoint worker copy leases must drain before close " - f"({status['store_leases']} store, " - f"{status['retrieve_leases']} retrieve)" + f"({status['store_leases']} store, {exposed} retrieve)" + ) + if status["retrieve_lookups"]: + logger.info( + "Closing with %d checkpoint lookups still waiting for storage " + "(%d cancelled); no worker received their pages", + status["retrieve_lookups"], + status["cancelled_lookups"], ) self._ctx.storage_manager.checkpoint_retention.clear_retirement( self._retire_generations diff --git a/tests/v1/multiprocess/test_checkpoint_cancel_drain.py b/tests/v1/multiprocess/test_checkpoint_cancel_drain.py new file mode 100644 index 00000000000..df8e4e305c9 --- /dev/null +++ b/tests/v1/multiprocess/test_checkpoint_cancel_drain.py @@ -0,0 +1,256 @@ +# SPDX-License-Identifier: Apache-2.0 +"""A restore whose storage lookup never answers misses; it is never fatal. + +A pending retrieve lookup exposed no SHM slots to its worker, so no GPU copy +can touch its pages. When it is cancelled and storage still does not answer, +the worker reports a miss and keeps serving, and the server releases the +lookup once storage answers. +""" + +# Standard +from dataclasses import replace +from types import SimpleNamespace +from typing import Any, cast +import threading +import time +import uuid + +# Third Party +import pytest + +# First Party +from lmcache.v1.mp_observability.errors import LMCacheTimeoutError +from lmcache.v1.multiprocess.checkpoint_storage import ( + CheckpointSlots, + checkpoint_object_keys, +) +from lmcache.v1.multiprocess.checkpoint_transfer import ( + CheckpointTransferJob, + CheckpointTransferWorker, + UnsafeCheckpointCopyError, +) +from lmcache.v1.multiprocess.engine_context import MPCacheServerContext +from lmcache.v1.multiprocess.modules.checkpoint import CheckpointModule +from lmcache.v1.multiprocess.mq import MessageQueueClient +from lmcache.v1.multiprocess.protocol import RequestType +from lmcache.v1.multiprocess.protocols.checkpoint import CheckpointLeaseResponse +from tests.v1.multiprocess.test_checkpoint_storage import ( + make_manifest, + open_checkpoint_rpc, + open_store, + poll, + publish_all, +) + + +class _Reply: + def __init__(self, value: Any = None, error: BaseException | None = None): + self.value = value + self.error = error + + def result(self, timeout: float | None = None) -> Any: + if self.error is not None: + raise self.error + return self.value + + +class _StuckStorageClient: + """Storage never answers a lookup, not even after it is cancelled.""" + + def __init__( + self, + cancel_error: BaseException | None = None, + poll_error: BaseException | None = None, + ) -> None: + self.cancel_error = cancel_error + self.poll_error = poll_error + self.requests: list[RequestType] = [] + + def submit_request(self, kind: RequestType, payload: list[Any]) -> _Reply: + self.requests.append(kind) + if kind == RequestType.CHECKPOINT_CANCEL_RETRIEVE: + return _Reply(True, self.cancel_error) + if kind == RequestType.CHECKPOINT_POLL_RETRIEVE and self.poll_error: + return _Reply(error=self.poll_error) + return _Reply(CheckpointLeaseResponse("pending", "lease")) + + +def _stuck_worker(client: _StuckStorageClient) -> CheckpointTransferWorker: + return CheckpointTransferWorker( + cast(MessageQueueClient, client), + lambda job, lease: pytest.fail("a lookup without slots must not copy"), + rpc_timeout=0.1, + ) + + +def _retrieve(rank: int = 0) -> CheckpointTransferJob: + return CheckpointTransferJob( + replace(make_manifest(), world_size=1), rank, "RETRIEVE", () + ) + + +def test_lookup_that_never_drains_after_cancel_is_a_miss() -> None: + client = _StuckStorageClient() + worker = _stuck_worker(client) + try: + for _ in range(2): + started = time.monotonic() + completion = worker.submit(_retrieve()) + # The worker keeps admitting transfers after the first abandon. + assert completion is not None + assert completion.result(timeout=10) is False + assert time.monotonic() - started < 5 + finally: + worker.close() + assert client.requests.count(RequestType.CHECKPOINT_CANCEL_RETRIEVE) == 2 + + +@pytest.mark.parametrize( + "error", + [RuntimeError("cancel handler failed"), LMCacheTimeoutError("no reply")], + ids=["error", "timeout"], +) +def test_failed_cancel_of_a_pending_lookup_is_a_miss(error: BaseException) -> None: + client = _StuckStorageClient(cancel_error=error) + worker = _stuck_worker(client) + try: + completion = worker.submit(_retrieve()) + assert completion is not None + assert completion.result(timeout=10) is False + again = worker.submit(_retrieve()) + assert again is not None and again.result(timeout=10) is False + finally: + worker.close() + + +def test_failed_poll_of_a_pending_lookup_is_not_fatal() -> None: + client = _StuckStorageClient(poll_error=RuntimeError("poll handler failed")) + worker = _stuck_worker(client) + try: + completion = worker.submit(_retrieve()) + assert completion is not None + with pytest.raises(RuntimeError, match="poll handler failed") as failure: + completion.result(timeout=10) + assert not isinstance(failure.value, UnsafeCheckpointCopyError) + assert RequestType.CHECKPOINT_CANCEL_RETRIEVE in client.requests + assert worker.submit(_retrieve()) is not None + finally: + worker.close() + + +def _stall_lookups(monkeypatch: pytest.MonkeyPatch, storage: Any) -> threading.Event: + """Keep every prefetch pending until the returned event is set.""" + answer = threading.Event() + query = storage.query_prefetch_status + monkeypatch.setattr( + storage, + "query_prefetch_status", + lambda handle: query(handle) if answer.is_set() else None, + ) + return answer + + +def test_cancelled_lookup_is_released_when_storage_answers_late( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """The engine keeps serving and restores again once storage answers.""" + with open_checkpoint_rpc() as (client, module, _mapping, _name): + storage = module.context.storage_manager + entry = replace(make_manifest(), world_size=1) + assert module.begin(entry) + stored = module.prepare_store(entry, 0) + assert stored.status == "ready" + assert module.finish_store(stored.lease_id, True) + answer = _stall_lookups(monkeypatch, storage) + copies: list[str] = [] + worker = CheckpointTransferWorker( + client, lambda job, lease: copies.append(lease.lease_id), rpc_timeout=0.2 + ) + job = CheckpointTransferJob(entry, 0, "RETRIEVE", ()) + try: + stalled = worker.submit(job) + assert stalled is not None and stalled.result(timeout=10) is False + assert not copies + status = module.report_status()["recurrent_checkpoints"] + assert status["retrieve_leases"] == status["cancelled_lookups"] == 1 + answer.set() + restored = worker.submit(job) + assert restored is not None and restored.result(timeout=10) is True + assert len(copies) == 1 + status = module.report_status()["recurrent_checkpoints"] + assert status["retrieve_leases"] == status["cancelled_lookups"] == 0 + assert module.find((entry.prefix,)) == entry + keys = [key for group in checkpoint_object_keys(entry, 0) for key in group] + assert storage.delete_l1_keys(keys)[0] == len(keys) + finally: + answer.set() + worker.close() + + +def test_released_cancelled_lookup_answers_its_late_poll_with_a_miss( + monkeypatch: pytest.MonkeyPatch, +) -> None: + with open_store() as (service, index, storage, mapping): + entry = make_manifest() + publish_all(service, index, mapping, entry) + answer = _stall_lookups(monkeypatch, storage) + cancelled = service.begin_retrieve(entry, 0) + other = service.begin_retrieve(entry, 1) + assert cancelled is not None and other is not None + service.cancel_retrieve(cancelled) + assert service.poll_retrieve(cancelled) is None + assert service.report_status()["cancelled_lookups"] == 1 + answer.set() + # Polling another lease releases the cancelled lookup as well. + assert isinstance(poll(service, other), CheckpointSlots) + assert service.report_status()["cancelled_lookups"] == 0 + assert service.poll_retrieve(cancelled) is False + with pytest.raises(KeyError): + service.poll_retrieve(cancelled) + service.cancel_retrieve(cancelled) + service.cancel_retrieve(uuid.uuid4().hex) + service.finish_retrieve(other) + assert service.report_status()["retrieve_leases"] == 0 + assert index.find((entry.prefix,)) == entry + keys = [key for group in checkpoint_object_keys(entry, 0) for key in group] + assert storage.delete_l1_keys(keys)[0] == len(keys) + + +def test_close_waits_only_for_leases_that_exposed_slots( + monkeypatch: pytest.MonkeyPatch, +) -> None: + name = f"lmcache_l1_pool_checkpoint_close_{uuid.uuid4().hex}" + with open_store(shm_name=name) as (_service, _index, storage, _mapping): + module = CheckpointModule( + cast( + MPCacheServerContext, + SimpleNamespace( + storage_manager=storage, + shm_pool_info={"shm_name": name, "pool_size": 4 * 1024 * 1024}, + ), + ) + ) + entry = replace(make_manifest(), world_size=1) + assert module.begin(entry) + stored = module.prepare_store(entry, 0) + assert module.finish_store(stored.lease_id, True) + lookup = module.begin_retrieve(entry, 0) + deadline = time.monotonic() + 5 + ready = module.poll_retrieve(lookup.lease_id) + while ready.status == "pending" and time.monotonic() < deadline: + ready = module.poll_retrieve(lookup.lease_id) + assert ready.status == "ready" + with pytest.raises(RuntimeError, match="must drain"): + module.close() + assert module.finish_retrieve(ready.lease_id) + answer = _stall_lookups(monkeypatch, storage) + try: + cancelled = module.begin_retrieve(entry, 0) + assert module.cancel_retrieve(cancelled.lease_id) + assert module.begin_retrieve(entry, 0).status == "pending" + status = module.report_status()["recurrent_checkpoints"] + assert status["retrieve_lookups"] == 2 + assert status["cancelled_lookups"] == 1 + module.close() + finally: + answer.set()