Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions .lil/changes/lmcache-106.json
Original file line number Diff line number Diff line change
@@ -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": []
}
8 changes: 8 additions & 0 deletions docs/source/mp/l2_storage/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
99 changes: 88 additions & 11 deletions lmcache/v1/multiprocess/checkpoint_storage.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
"""

# Standard
from collections import OrderedDict
from dataclasses import dataclass, field
from typing import TYPE_CHECKING
import hashlib
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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__(
Expand All @@ -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(
Expand All @@ -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):
Expand Down Expand Up @@ -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
]
Expand Down Expand Up @@ -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:
Expand All @@ -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,
}

Expand All @@ -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:
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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()
Loading
Loading