Repository navigation
feat(offload): support native DSv4 checkpoints with LMCache MP - #2250
Conversation
|
While tracing the native MP path, could we add In the pinned LMCache 05fc77a, I also noticed the PR description still points to mp/README.md, while the guide is now in the parent offload README. Would a one-line recipe correction plus updating that pointer be the right scope? I’d keep the topology, checkpoint and lifetime design unchanged. AI-assisted source review and drafting; no GPU results claimed. |
|
Follow-up with a one-line patch and real-parser validation. Against ATOM The check uses the unmodified The patch is just the README line for your existing branch—none of the downstream experiment machinery. AI-assisted preparation and evidence review; no GPU results claimed. |
Save and restore DSv4 PAGE KV together with the matching recurrent STATE checkpoint through a standalone LMCache multiprocess server. The lmcache_mp connector selects its implementation from backend capabilities: backends that publish a PagedStateCheckpointSpec and execute_paged_state_copies use the native PAGE/STATE path, others keep the PAGE-only transport. Saves lease immutable READY checkpoint PAGE units instead of snapshotting the live SLOT; restores reserve fresh units and adopt them as a local READY checkpoint. Also: - allow single-host DP and DP-attention, scoping MP sessions per replica; - release source PAGE blocks per chunk and reacquire a finished request's still-resident prefix at save admission; - allow partial release of recurrent-state requests only when the connector guarantees an independent state lease; - read offload env knobs through atom.utils.envs; - refuse PAGE-copied state backends that lack the native contract. Requires LMCache dev@05fc77a (LMCache/LMCache#5132), pinned on main by #2395. Start the server with --null-block-id -1 --separate-object-groups. Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
9c64bea to
62893d7
Compare
The pinned LMCache (05fc77a) requires --eviction-policy on `lmcache server`; without it the documented command exits 2 during argument parsing. Reported-by: kvnloo Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Native MP stored every READY checkpoint boundary, including whole short prompts, although a prefix below OFFLOAD_MIN_LOAD_TOKENS (8192 by default, equal to OFFLOAD_MIN_SAVE_TOKENS) can never be loaded back. On DSv4-Pro TP8 with a 1K/1K workload at concurrency 64, offload stored 257 such prefixes and cost 12.4% output throughput; skipping them stores none and leaves 1.6-2.2%, within run-to-run noise. 16K prompts still save as before. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Review — native DSv4 checkpoints with LMCache MPReviewed The native-state path is a real capability and the checkpoint/lease accounting around it is careful work — The first is a boundary:
The second is a shape that recurs:
Provenance: [verified] means I read the deciding lines in the PR-head worktree myself. [reported] means the shape matches the code but I did not trace it end to end. 1. The late-save path cannot run on the vLLM plugin, and is silently wrong there instead of loudly [verified]These are two findings that only make sense together, because the second is currently masking the first.
So the whole safety argument rests on the reacquire at if getattr(seq, "_offload_finished", False) and not seq.block_table:
late_source = self._late_save_source(seq, saved, aligned)
else:
block_ids = list(seq.block_table)On the plugin path Now the masking. If that gate were fixed, the request would reach if self._block_manager is None:
raise RuntimeError("late offload save requires a bound block manager")
This is not a missing wire. Worth deciding before merge: either the plugin keeps a lease-based protection (what was deleted), or 2. Both new hash-chain walks seed from a hardcoded
|
| site | seed |
|---|---|
native_state_scheduler.py:168 |
chain[-1] if chain else **-1** |
block_manager.py:2572 |
parent_hash = **-1** |
block_manager.py:864 |
chain[-1] if chain else **seq.cache_seed** |
block_manager.py:1584 |
return **seq.cache_seed** |
The comment directly above the first one says "Use its exact public hashing algorithm and token slices" — the algorithm and the slices were copied; the seed was not.
The blast radius is precise: sequence.py:268 sets cache_seed = -1 and only overwrites it when multimodal_data is not None, and compute_hash mixes the prefix in only when prefix != -1. So the two are equivalent for text requests and diverge from block 0 onward for multimodal ones. On a multimodal DSv4 request self._checkpoints.contains(self._boundary_hash(...)) never matches, every native save is silently skipped, and adopt_transfer_units(op, prefix_hash) publishes a restored checkpoint under a key no loader will ever look up. In acquire_offload_prefix the same seed makes kv.lookup(block_hash) miss block 0, so the late save degrades to "nothing resident".
This is the live path, not a corner: page_unit_checkpoint.py:849 declares readable_midstep = False, and block_manager.py:1967 is if not self.state.readable_midstep or ...: seq.block_hashes = [] — for exactly the model family this PR targets, the empty-chain fallback is the normal case. _boundary_hash's two branches also use different seeds for the same boundary and are only accidentally equal today. No offload test uses a multimodal sequence, and acquire_offload_prefix's one test reference (test_offload_early_block_release.py:385) is a lambda stub.
3. A failed submit() wedges the engine permanently, and every recovery path is switched off [reported]
native_state_worker.py:255 catches, logs, and returns — leaving the pending entry holding _UncertainSubmission, whose query() is return False. The operation can never become terminal, so _complete_native_save — the only caller of _refund_state_image and settle_offload_store — is never reached.
What makes this worse than a leak is that all four escapes are deliberately disabled: abandon_save returns None (:511), reclaim_stale_leases returns [] (:516), the pin was taken timeout_reclaimable=False (page_unit_checkpoint.py:600) so reclaim_stale_offload_pins skips it, and _reconcile_stalled_deferred_saves exits on its first line because save_abandon_timeout_s() is 0.0 (backend.py:1137, scheduler.py:1276).
One ZMQ hiccup or LMCache server restart therefore costs, permanently: units_per_checkpoint PAGE units pinned, _pinned_state_bytes charged (after _max_pending_saves such events _has_state_budget() is False and all native saves and loads stop being admitted), the request's blocks never freed because should_defer_free stays True, and has_pending_work() stuck True so EngineCore busy-loops over idle GPUs. len(pending) also never drops, so the worker's own bound at :219 starts rejecting every subsequent save.
tests/test_lmcache_mp_native_worker.py:208 asserts the lease is retained — so the wedge is the tested contract. Retaining the lease is right; what is missing is any bounded path out of it.
4. OFFLOAD_MIN_SAVE_TOKENS is read in two different units [reported]
The base scheduler compares it against the save increment; the native scheduler compares it against the absolute boundary:
# chunked_scheduler.py:453
if target - saved < self._min_save_tokens: # increment
# native_state_scheduler.py
floor = max(exhausted, self._min_save_tokens - 1, 0) # absoluteBoth cite the same justification, and it is only true for the absolute reading. At the default 8192, a 40k-token request that has already saved 32k finishes with a sub-8k remainder; the base scheduler frees the reacquired blocks and pops the request from _save_tracker, so the tail of every long request is silently never persisted and later hits stop 32k in — while the native scheduler on the same knob would have accepted that boundary. Pre-PR the base had no _min_save_tokens at all and emitted whenever aligned > saved, so this also silently changes dense/hybrid behaviour; the README still documents the threshold only under "Native MP" while it now lives in the shared base.
The same knob has a second failure at the other end. With OFFLOAD_MIN_SAVE_TOKENS=0 (reachable — envs.py clamps with max(0, ...)) and a fully-evicted prefix, acquire_offload_prefix breaks on the first block, so target == saved, the early-out is skipped, and a zero-length save is emitted forever: entry[1] = aligned leaves the watermark unchanged and the retirement _save_tracker.pop(sid) at :585 only runs when late_source is None, which now never happens. On the native path the same shape is worse — _build_save_request calls _boundary_hash(seq, aligned=0), which raises ValueError("native checkpoint boundary must align to hash blocks") straight out of build_connector_meta.
Both new fixtures set OFFLOAD_MIN_SAVE_TOKENS=0 (test_offload_early_block_release.py:47, test_lmcache_mp_native_scheduler.py:83), so the default-valued behaviour in the first half is never exercised and the fully-evicted case in the second half is never reached.
5. Three sweeps that stopped short
descriptor_slot was added to two of three implementations [reported]. backends.py:341 declares execute_paged_state_copies(self, store_ops, restore_ops, descriptor_slot: int = 0); deepseek_v4_attn.py:968 and gdn_attn.py:815 were updated; deepseek_v41/backend.py:263 still has def execute_paged_state_copies(self, stores, restores). The native worker always calls it with the keyword (native_state_worker.py:285-295), so the moment DSv4.1 is reached the TypeError lands inside _begin_restore's except Exception and is reported as "LMCache MP native restore failed" rather than failing loudly.
The uint8 fix was applied to one of two halves of the same registration [reported]. backend.py:433 now registers view.view(torch.uint8) with a comment explaining that LMCache's ROCm raw-pointer fallback "cannot express FP8 through the CUDA array interface". native_state_layout.py:162 hands tensors = list(page_views) to the same register_kv_caches in native dtype — while the STATE aliases built eleven lines below do retype via page_view.view(torch.uint8). A native-state model publishing FP8 block views hits the bug the sibling path just fixed: wrong bit patterns after a cache hit, not a crash.
The DP guard was relaxed without the namespace that made it safe [reported]. The blanket if dp_size != 1 or enable_dp_attention: raise NotImplementedError("TP-only") became a multi-node-only if dp_size_local != dp_size:. But _model_namespace hashes only model/layout/checkpoint with no DP rank, _parallel_strategy returns worker_id = rank_in_group // replication so every replica registers worker_id 0, and _server_urls returns the one configured URL for all of them. Four processes register four different sets of GPU KV tensors under an identical (model_name, worker_id, world_size) identity. Only _mp_session_id was DP-scoped, which fixes per-request routing but not the registration namespace. The deleted cases are the tell: ({"dp": 2}, "TP-only") and ({"enable_dp_attention": True}, "TP-only") were removed from tests/test_lmcache_mp.py, and enable_dp_attention no longer appears anywhere under offload/mp/.
6. Two selectors that encode the wrong predicate [reported]
is not None is standing in for "enabled". mp/connector.py:144 picks the native-state scheduler whenever block_manager.paged_state_checkpoints is not None. But block_manager.py:240 constructs that coordinator with enabled=self.enable_prefix_caching and self.num_state_slots > 0 — the attribute is non-None regardless, and enabled gates only applies() (page_unit_checkpoint.py:893); acquire_checkpoint_source and reserve_transfer_units never consult it. DSv4 served with prefix caching off therefore binds the native scheduler, whose _save_frontier queries a store that can never hold a READY image, and no offload happens at all where the PAGE-only scheduler would have worked. tests/test_lmcache_mp_shell.py only ever passes object() or None, so the disabled case has no coverage.
can_partially_deallocate_state is fail-open while the safety net it cites is fail-closed. multi_connector.py:656 returns True if any sub says True; protected_block_ids, sixteen lines above, returns None if any sub cannot answer. With a native-state MP sub plus a dense PAGE sub both deferring on the same finished sequence, the composite says True and the state slot is recycled while the non-declaring sub is still deferring. The docstring's claimed mitigation does not reach it: protected_block_ids narrows PAGE blocks only, and a dense sub legitimately returns frozenset() there while still needing the request alive. tests/test_multi_connector.py:291 passes under either ANY or ALL, so it has no discriminating power here.
7. Lifetimes and ordering [reported]
The MP session is ended before the late save is submitted. This PR sets _supports_early_block_release = True on the MP scheduler, where the base had it False precisely so no save could outlive request_finished. Now: request_finished → end_session(_mp_session_id(config, seq.id)); then a later build_connector_meta reacquires the prefix and emits the final save; then the worker submits it under that same, already-ended session id. Either the server drops the store — losing the end-of-prompt checkpoint, the most valuable one and the one a later full-prompt lookup depends on — or it resurrects a torn-down session. The native scheduler's own request_finished (:465) defers only when a dispatched load lease exists, not for a pending save.
The restore side stream is never fenced against the compute stream. grep -E 'wait_stream|wait_event|record_stream' across atom/kv_transfer/offload/mp/ and deepseek_v4_attn.py returns nothing. The same execute_paged_state_copies runs on the compute stream from AttentionMetadataBuilder.build() and on self._restore_stream from _begin_restore on the connector thread, with no handshake. Read-after-write is safe (the worker host-polls restore_event.query()), but write-after-read is not: relocate_state_slots runs on the compute stream from batch.state_maintenance_ops.relocations and can move slots while a parked restore is in flight, and the previous occupant's still-enqueued decode work is unfenced. Only suspend_queued_restore removes one competitor. Symptom: non-deterministic garbage after a cache hit that vanishes under HIP_LAUNCH_BLOCKING=1. tests/test_lmcache_mp_native_worker.py:360 monkeypatches torch.cuda.stream to a no-op, so it pins the descriptor-slot plumbing and enshrines the missing sync.
DSv4 + MTP cannot start. deepseek_v4_attn.py:1986 validates sum(region.unit_bytes) == checkpoint_spec.page_unit_bytes over the target's regions; model_runner.py:2107 then extends block_regions and block_tensor_views with the draft's; native_state_layout.py:159 re-sums over all regions and raises ValueError("PAGE regions do not cover the native PAGE unit") inside register_kv_caches. Were the check looser, the per-ordinal image loop would fold draft KV rows into the checkpoint image. Separately, line 2117 gcds tp_replication_factor down to 1 but leaves native_state_tp_replication_factor at tp_size, so two declarations describing the same physical object disagree and are validated independently. The PR's e2e was DSv4-Pro TP8, presumably with no draft builder.
A pin and a budget charge leak on the exception path. native_state_scheduler.py:258 charges the state budget and takes the checkpoint PAGE lease, then calls super()._build_save_request(...) and builds NativeStateTransfer with no try/except around the remainder. If anything after the acquire raises, the caller at chunked_scheduler.py:594 frees only the late-acquired KV blocks and re-raises; nothing ever reaches release_offload_store_source / settle_offload_store / _refund_state_image, because _native_saves[operation] was never written. The pin is timeout_reclaimable=False, so the reclaimer skips it forever: one whole checkpoint image leaves the pool for the process lifetime. _charge_state_image is carefully exception-safe for the acquire itself, but its unwind scope ends at acquire(). Same shape at _decide_load_after_alloc:350-369, where reserve_transfer_units precedes a suspend_queued_restore that can assert and a _boundary_hash that can raise.
8. Cleanliness, cost, docs
_save_frontiernow runs an O(prompt/chunk) hash-and-lookup scan per tracked request per scheduler step.- Descriptor slots ≥1 lazily allocate pinned memory on the connector thread mid-serving, while
warmup_per_req_cachewarms only slot 0 — contradicting that function's own docstring. _begin_restore's exception path leaks a descriptor slot.- 15 new env vars are absent from
docs/environment_variables.md. chunked_scheduler.py:1162still cites the "frozen block table" this PR deleted.
9. The tests encode several of these defects as the contract
This is the part I would act on first, because it is why the PR is green:
| test | what it pins |
|---|---|
test_lmcache_mp_native_worker.py:208 |
asserts the lease is retained — makes §3's permanent wedge the contract |
test_lmcache_mp_native_worker.py:360 |
torch.cuda.stream monkeypatched to a no-op — pins §7's missing sync |
test_multi_connector.py:291 |
passes under ANY or ALL — no discriminating power for §6 |
test_lmcache_mp_shell.py |
only object() or None — §6's disabled case uncovered |
| both new fixtures | OFFLOAD_MIN_SAVE_TOKENS=0 — §4's default behaviour never runs |
run_unit_tests.sh --ignore=tests/plugin |
§1 is unreachable by CI by construction |
None of these are wrong tests. Each asserts something true about the code as written; the gap is that the thing asserted is the symptom.
Checked and not upheld
Recorded so they are not re-raised: a suspected KeyError in _new_load_operation and a suspected stride collapse in _build_cache_views are both properly guarded, and a suspected late-save lease leak with _early_release off is unreachable — _offload_finished is only stamped under early release.
…e threshold Review follow-ups (#2250 §1, §4). - Unbound schedulers (the vLLM plugin, where vLLM owns the blocks) cannot reacquire a finished request's prefix, so teardown again freezes the block table and leases the unemitted suffix, and the final save reads exactly those leased blocks. The late-save reacquire path now runs only when a BlockManager is bound, instead of keying off an empty block table that the plugin never clears. - The shared base no longer reads OFFLOAD_MIN_SAVE_TOKENS, so dense and hybrid behave as before this PR. A late save with nothing savable resident now retires the request instead of emitting an empty save forever. - Native MP applies OFFLOAD_MIN_SAVE_TOKENS to the absolute boundary for normal and late saves alike, so a long request's short tail is still stored, and the floor can never reach boundary 0. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Review follow-ups (#2250 §3, §7). - A submission that raises, a restore that raises, or a restore event whose query raises may still be touching engine memory, so its lease is kept -- but no longer forever. _UncertainSubmission now fails the transfer after lmcache.mp.uncertain_transfer_timeout_s (default: twice lmcache.mp.mq_timeout) with a warning, so the save or load settles and its lease, budget, pending slot and descriptor slot are released. Before, one ZMQ hiccup could stop all native saves and loads and keep EngineCore busy-looping. - PAGE-only MP treats a raising save submission the same way instead of as a definite failure. The failure reported no quiescence, so under #2339 the save was never retried and its lease was never released. The now-dead _immediate_save_failures set is removed. - Native restores are fenced both ways: the restore stream waits for work already on the compute stream (write-after-read on the SLOT's previous occupant), and every step's compute stream waits on in-flight restore events (read-after-write, including SLOT relocations in build()). The fence is on the GPU, so the host does not block. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ssions Review follow-ups (#2250 §7). - The MP scheduler ended a request's session in request_finished, but with early release its final save is emitted and submitted under that session afterwards. Session ends are now deferred until no save of the request is tracked or in flight. They are swept on retirement and at every metadata build. - Native save admission returns the checkpoint lease and budget charge if anything after the acquire raises; the pin is never timeout-reclaimable, so it used to leak for the process lifetime. Load admission computes the boundary hash before reserving units, and releases the units, the budget and any suspended local restore if the rest raises. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Review follow-ups (#2250 §5, §7). - DSv4 now publishes paged_state_region_count, the number of its own PAGE regions. A draft with a pool of its own appends regions after them, and the native layout validates checkpoint coverage and builds STATE aliases from the leading regions only. Draft rows are registered as ordinary PAGE KV and never folded into the checkpoint image. Before, DSv4 with such a draft failed registration with "PAGE regions do not cover the native PAGE unit". - Auto rank collapse is off when a DSpark draft's backend owns a KV pool: that draft publishes per-rank PAGE regions, so the complete PAGE object is not replicated and registration would otherwise reject the requested collapse. - The native layout registers its PAGE group as uint8 views, like the PAGE-only path and like its own STATE aliases. LMCache's ROCm raw-pointer fallback cannot express FP8 through the CUDA array interface. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ts early Review follow-ups (#2250 §5, §8). - Native _save_frontier walked the prompt checkpoint by checkpoint for every tracked request on every scheduler step. PageUnitCheckpointStore now has a generation that is bumped whenever the READY set can change, and the answer is cached per request on (frontier, floor, generation). - Restore descriptor slots >= 1 were first allocated as pinned memory on the connector thread mid-serving. The native worker now reserves them at registration through the new builder hook reserve_checkpoint_descriptors (base builders allocate their descriptor buffers, DSv4.1 its per-slot staging). - Document why single-host DP replicas share one (model_name, worker_id, world_size) identity. It is the content-addressed storage namespace, while the server registers GPU memory per instance_id, refcounts layout descriptors per (model_name, world_size), and _mp_session_id scopes sessions per replica. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Review follow-ups (#2250 §8, §9). - docs/environment_variables.md listed only OFFLOAD_MAX_PENDING_SAVES and the LMCache pin timeout, and still said they bypass atom.utils.envs. All 16 ATOM-owned offload knobs are now documented with type, default, precedence and invalid-value behavior. A new test fails if an OFFLOAD_*/LMCACHE_* var registered in envs.py is missing from the reference. - README and envs.py describe OFFLOAD_MIN_SAVE_TOKENS as the native absolute boundary for normal and late saves, and the README documents the uncertain-transfer bound. - Test that a restore raising after it took a descriptor slot fails the load and returns the slot once the uncertainty bound expires. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
|
Thanks for the thorough review. Every point is addressed in §1 Late save without a BlockManager — §2 Hash-chain seed — §3 Wedge on an unprovable submission — §4 One knob, two units — §5
§6 —
§7
§8
§9 Each test in your table now asserts the property rather than the symptom (see above). About 20 tests were added. Validation
|
On a CPU-only torch, torch.cuda.Event is a dummy class that raises when it is built, which happens before the restore takes its descriptor slot, so the test's "slot is held until the bound" assertion failed in CI. Stub the Event so the restore fails deterministically at the stream fence, after the slot is taken, on every runner. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Second pass —
|
…at a deadline The MP path answered "what settles a transfer whose outcome is unprovable" three ways: never (native checkpoint pins, a zeroed abandon timeout), after 600 s (the uncertainty bound, which then freed load destinations and PAGE sources a server might still be using), and never-unbounded (a future whose poll keeps raising). One rule now: memory under an MP transfer is released only on a terminal report, and a transfer still not terminal after `lmcache.mp.transfer_deadline_s` (default 1200 s) raises `LMCacheTransferUnprovable`, stopping the engine. - Worker: every pending PAGE-only and native load/save carries its start time and is checked against the deadline, which bounds a raising submission, a future whose poll keeps raising, and a restore whose event cannot be queried. A descriptor that cannot be built is provably unsent and still fails immediately; only the transport call is unprovable. - Scheduler: a watchdog over dispatched saves, loads and native checkpoint sources catches a completion that never arrives (lost report, silent TP rank), with a margin so the worker fails first. - `save_abandon_timeout_s` is inherited again, so the engine-wide stalled save, state-pin and orphan-load-slot reclaimers stay on; MP opts out through its own `abandon_save` / `reclaim_stale_leases` instead. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The late final save reacquired a finished request's prefix through the hash index whenever a BlockManager was bound. With prefix caching off nothing is indexed, so the lookup found nothing and the save was dropped without a warning -- the combination the offload README ships for `lmcache_offload`. Hash reacquire now requires a bound manager with prefix caching; otherwise teardown leases the unemitted suffix of the frozen table, as on main. A late save that finds part of its prefix evicted is counted in `truncated_late_saves`. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
`protected_block_ids` marked the request finished before the scheduler decided whether it could release it partially. When per-request state made that unsafe, the request was deferred whole and kept its table, but the next metadata build still took the late-save path and claimed every block a second time through the hash index (dropping the save if a hash was missing). `activate_block_leases`, which only the partial-release branch calls, now marks the release, and only that mark selects the reacquire path. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…erminal The worker released a save's checkpoint image as soon as any PAGE chunk milestone reached the boundary. Chunk milestones report PAGE token ranges; nothing in LMCache promises that a chunk completes only after every engine group registered at it, STATE groups included, has been read. The pinned LMCache emits no chunk events, so this path never ran, but it would have been unsound on the first server that did. The STATE source-safe channel is removed: the image pin is released at the STORE terminal, and chunk milestones release PAGE leases only. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
`OFFLOAD_PUBLICATION_TIMEOUT_S=` made `float('')` raise out of connector
init, and `OFFLOAD_COPY_WORKERS` / `OFFLOAD_LOAD_WORKERS` died with a bare
`invalid literal for int()`. The offload section now states one policy:
unset or empty is the default; a set but unusable value warns and falls back
for knobs that only tune reuse, and is rejected at startup, naming the
variable, for widths, sizes and timeouts. The transfer-mode comment no
longer lists `engine_driven` as valid.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…invariants - `adopt_transfer_units` popped the transfer record before adopting; if adoption raised, the units were reserved but unreachable by any release path. The record is now removed only after adoption returns. - A native restore that loses the publish race to an identical image was dropped silently; it is now logged as a deduplication. - The worker's own pending-save refusal is unreachable while the scheduler enforces the same bound; if it ever fires it now logs an error instead of passing as an ordinary save failure. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…y doubles The restore-copy doubles were positional-only, so the production call with `descriptor_slot=` would have raised inside `_begin_restore` and been reported as a failed restore; the success tests replaced `_begin_restore` outright, so no test ran the happy path. The doubles now take the production signature, and a new test drives a terminal retrieve through `get_finished` into the real `_begin_restore`, then asserts the copy's descriptor slot, that the load finishes only once the restore event is done, and that the slot is returned. The CUDA stand-ins are shared with the fence test. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…n checkpoints The per-request `_save_frontier` memo was keyed on the store-wide generation, so any request's publish or eviction forced every tracked request to rescan its prompt on the next step. The store now keeps a bounded log of which prefix hash each generation bump touched; a request rescans only when a change hits one of the boundaries it scanned (its answer or anything above it), or when the log cannot tell. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…int copier DSv4.1 kept its own `_DescriptorStaging`, a third copy of the per-slot pinned descriptor pool the base builder shares with DSv4 and GDN, and the only one that fenced reuse: a non-blocking H2D reads its pinned rows when the stream reaches it, so refilling a buffer whose last upload is still queued rewrites the descriptor under an earlier copy. `DescriptorStaging` in `pool_layout/paged_state_copy.py` now owns the per-slot buffers and the fence, and the base builder (DSv4, GDN) and DSv4.1's `StateCopies` both use it, so the base path gains the fence too. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
`build_native_state_mp_layout` and `_build_cache_views` each validated the backend-published PAGE views, and had already drifted: PAGE-only required full contiguity and compared total bytes against the view, native checked inner contiguity plus the block stride against the region. Both now call `validate_page_views` (new `mp/page_views.py`), which applies the stricter union -- shape, tight block-major stride, unit and total bytes, aliasing, forward indexing, one device -- and, for native, that a block's physical slots divide the block size. Each builder keeps only what is its own: PAGE-only rejects stateful fields; native checks the image coverage. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
`get_finished` was 91 lines at six levels of nesting, with three near-identical recovery blocks. It is now a loop over `_poll_native_save` and `_poll_native_load`; the load's retrieve-then-restore progression lives in `_advance_native_load`, and both unprovable-restore paths share `_hold_unprovable_restore`. No behaviour change. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
`mp/backend.py` had grown to 1332 lines holding configuration, PAGE view validation, transfer bookkeeping, lookups and both connector halves. It is split without behaviour change into: - `deployment.py`: config validation, TP/DP topology and rank collapse, model namespace, server adapters - `transfer.py`: operation identity, terminal detection, the transfer deadline and `LMCacheTransferUnprovable` - `lookup.py`: the scheduler's lookup client and read-lock bookkeeping - `page_views.py`: PAGE-only `_build_cache_views`, beside the shared view validation it uses - `worker.py` / `scheduler.py`: the PAGE-only connector halves The native modules, the public shell and the tests import from the new homes; tests patch each name where it is looked up. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ate saves - With the abandon window inherited again, a slow but legitimate MP save deferred past it was "abandoned" (a no-op for MP) and, a minute later, logged as a wedged P/D send -- both long before MP's own transfer deadline. A connector can now answer `waits_for_transfer_report(seq)`; the stalled-save reclaim skips such requests. LMCache MP answers True, the in-process shell forwards (default False), and the composite answers True only if every sub still deferring the request does, so P/D sends and in-process saves keep their abandon path. - A late save that finds part of its prefix evicted now logs a warning with the running count, not a debug line: early release returns the tail's blocks at teardown, and this is the cost of that trade. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
|
Thanks for the second pass. You were right that the three answers to "what settles an unprovable transfer" were the root; the fix starts there. Everything is in The policy —
§5 — §6 — §7
§8
Validation
|
…ether Every backend kept `block_regions` and `block_tensor_views` as two lists paired by hand (DSV4 even collected a third list of sources to zip later). `KVTransferTensors.add_block_region(tensor, semantic_role=...)` now appends both from one tensor: the region's addresses and a zero-copy `uint8 [num_units, 1, unit_bytes]` alias of exactly those bytes, cut from a larger allocation when `total_bytes` says so. DSV4, MHA, the MHA draft and MLA all publish through it. MHA and the draft already published this shape. MLA's views change from `[n, block_size, width]` to the same byte form; LMCache stores the same bytes in the same order either way, and neither its object key nor ATOM's model namespace depends on view shape, so existing cache entries stay valid. The MLA builder test's module stubs had drifted since #2399 (`atom.utils.block_tables`); they are fixed so the test runs again. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…w together `add_block_region` paired a region and its view at one call site, but the pair was only a convention: both lists stayed public and mutable, the constructor still took them separately, DSV4 and MLA used a throwaway `KVTransferTensors` as a scratch builder, the draft merge extended the two lists by hand, and tensor code lived in the torch-free `types` contract. - `PageRegion(region, view)` is a frozen value; `KVTransferTensors.pages` holds them. `block_regions` and `block_tensor_views` are read-only tuples derived from it, so no code can add to one without the other. - `page_region(tensor, ...)` in the new `disaggregation/page_region.py` builds one from the owning tensor (zero-copy byte view, contiguity and size checks); `types.py` no longer touches torch. - Builders pass `pages=[...]` straight to the final object: no scratch instance. The address-only producer (Qwen4 exp) publishes `PageRegion(region)` with no view, which LMCache MP refuses as before. - `merge_pages(other)` replaces the draft merge in `ModelRunner`, carrying the gcd replication rule with it. `validate_page_views` stays on the LMCache MP side: it is the consumer checking what it is handed, not a second copy of construction. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Resolves docs/environment_variables.md: #2414's ATOM_KV_OFFLOAD and ATOM_KV_OFFLOAD_EXTRA_CONFIG rows join the rewritten LMCache offload table. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Summary
Add native DSv4 PAGE/STATE checkpoint support through a standalone LMCache multiprocess server.
The integration saves and restores PAGE KV together with the matching recurrent STATE checkpoint while active requests retain fixed SLOTs. State sources are represented by immutable checkpoint PAGE units, and each emitted save acquires an independent checkpoint lease before request teardown. Restore reserves fresh PAGE units and copies the native checkpoint image back into the request SLOT before the request is allowed to resume.
Sparse STATE lookup requires a complete checkpoint at the same boundary as PAGE KV. Full-prompt lookup truncates tokens before querying. ATOM uses the same server-wide absent marker for every registered group: PAGE IDs are nonnegative, including valid block
0, while missing STATE entries are-1. Start the matching LMCache server with--null-block-id -1 --separate-object-groups. Native layout/model revision namespaces and exact operation generations reject incompatible or stale transfers. Uncertain DMA completion retains its lease.The partial PAGE-release path remains fail-closed for recurrent-state requests. Only the native-state MP scheduler explicitly opts in after it is bound to the PAGE checkpoint coordinator; connectors without that guarantee keep the whole request deferred.
Scope
This PR is focused on the LMCache MP path. The existing standalone DSv4 and Kimi LMCache connector behavior is intentionally unchanged.
Native MP keeps a simple
max_pending_savessafety bound. Value-ranked admission, PAGE pin-ratio budgets, candidate replacement, and the related policy metrics are not included. If those policies are needed later, they should be designed specifically around the MP path instead of changing legacy connectors.Dependencies
rocm/atom-dev:lmcache-v0.5.6.dev98-g05fc77a0-rocm-torch210wheel image (ROCm 7.2.4 / torch 2.10, wheel sha256a5fe8f3f…). LMCache's own nightly ROCm wheel targets torch 2.11 and does not load in the ATOM image. This PR no longer changes.github/ordocker/; the release-based pin and the pin-bump PR it used to carry (ci: open a pin-bump PR after publishing an LMCache wheel #2391) were superseded by ci(lmcache): ship LMCache wheels as Docker Hub images, not GitHub Releases #2395, which deleted that release._fresh_tier_lookupqueries_lookup_token_ids, and fix(offload): memoize the external-tier lookup per HBM frontier (native + plugin) #2305's timeout test uses the DP-scoped MP session ids._offload_finishedflag this branch sets on finished requests, and its retired-failure test expects a teardown lease of only the emitted save's source blocks: the unemitted suffix is reacquired at late-save admission.Main implementation
Start with the "LMCache Multiprocess (
lmcache_mp)" section ofatom/kv_transfer/offload/README.mdfor launch configuration, lifetime rules, and supported scope. The primary implementation is in:atom/kv_transfer/offload/mp/connector.pyatom/kv_transfer/offload/mp/native_state_scheduler.pyatom/kv_transfer/offload/mp/native_state_worker.pyatom/kv_transfer/offload/mp/native_state_layout.pyCleanups on this branch:
OFFLOAD_*,LMCACHE_MP_TRANSFER_MODE) is read throughatom/utils/envs.py;kv_connector_extra_configoverrides still take precedence.PageUnitCheckpointStore.adopt_units; chunk ranges and PAGE source-safe completions come from shared helpers inmp/backend.py;_checkpoint_descriptor_bufferlives once onAttentionMetadataBuilder; native-state budget charges and refunds are paired; finished requests carry an_offload_finishedflag instead of a copied block table.state_transfer()copies PAGE-backed checkpoints but publishes nopaged_state_checkpoint_spec/execute_paged_state_copies(Kimi K3 today) is refused atlmcache_mpregistration, instead of running a PAGE-only worker under a native-state scheduler whose lookups never hit.--eviction-policy LRU, which the pinned LMCache requires (reported by @kvnloo).Validation
rocm/atom-dev:nightly_202609161445: 3452 passed; no test fails only on this branch (one main failure, theSeqViewslot scan, passes here).0:17920across the 8192 and 16384 checkpoints. GSM8K (64-shot, 50 samples) scored 0.94 on all three passes, and cold-vs-restore answers matched 47/50, the same as the cold-vs-control floor (FP8 with dynamic batching is not bit-reproducible run to run).--worker-registration-grace-seconds 30, the restarted workers were reaped before their first heartbeat, and every restore silently fell back to prefill while accuracy still looked normal.Logit-level equivalence and performance measurements remain outstanding.