Skip to content

fix(offload): fence dense LMCache saves - #2339

Merged
valarLip merged 5 commits into
ROCm:mainfrom
NidhoggD1:fix/lmcache-dense-save-fence
Sep 25, 2026
Merged

valarLip merged 5 commits into
ROCm:mainfrom
NidhoggD1:fix/lmcache-dense-save-fence

Conversation

@NidhoggD1

@NidhoggD1 NidhoggD1 commented Sep 21, 2026 •

Copy link
Copy Markdown
Contributor

Motivation

Dense LMCache saves run on a background executor. The save thread packs KV on
its thread-local GPU stream, but current main does not order that stream after
the model/RPC stream that produced the KV. A background save can therefore
persist stale or incomplete KV and later return incorrect model output on a
host-tier hit.

Technical Details

  • Record one producer event when a metadata step contains dense save work.
  • Share that event with every background save submitted for the step.
  • Forward the event through LMCacheEngine.store() and consume it inside
    BlockGPUConnector.batched_from_gpu(), so the wait is enqueued on the pack
    stream of the thread that actually reads KV.
  • Keep the dependency stream-to-stream so the save worker does not block on
    Event.synchronize() in the normal GPU path.
  • Fall back to event synchronization only when the connector has no pack
    stream.

The event stays owned by each submitted save job until that save has
consumed the dependency. batched_from_gpu enqueues the wait before any
other pack-stream work and reports producer_fenced=1 in the transfer
stats, so a store that reaches the connector without the event is visible in
[OFFLOAD-SAVE-PROF] rather than silently unfenced.

Fence creation is isolated inside the per-request dispatch loop. A runtime
error while recording the event rejects this step's saves without dropping
loads from the same metadata step; no unfenced save is submitted. A rejected
or failed save no longer advances the scheduler's save watermark: the
scheduler restores the previous frontier and re-emits the save once every TP
rank is quiescent (a pre-submit rejection reports quiescence directly; a
post-submit failure becomes retryable when all of its source groups are
source-safe, or is reclaimed with its lease). A save that fails after its
request was retired is dropped rather than retried against recycled blocks.

Related: #2332 contains a similar correctness fix inside a much larger Kimi-K3
P/D change. This PR isolates the dense-save fence against main and uses a
nonblocking pack-stream wait instead of host synchronization.

Validation

The same producer-event to pack-stream dependency was validated on an ATOM
dense LMCache stack using MI350X, LMCache 0.4.5, FP8 KV, and TP4:

  • Connector-level delayed-write reproduction: 38/40 stale reads without the
    dependency, 0/40 with pack_stream.wait_event.
  • Model-level host-tier reads: 41/41 correct with the fence.
  • Adjacent candidate/control/candidate gate: corrupted host-tier reads
    0/4 -> 4/4 -> 0/4.
  • Integrated coverage included single TP4, two routed TP4 engines, concurrent
    and reused prefixes, and forced GPU-prefix eviction followed by host recall.

Checks run on the current branch:

Checks run on the current head (e9be221f, rebased onto main at
a5ad0a50; the rebase folds the failed-save cleanup into #2167's
_drop_finished_save_state, so both request_finished and
source_blocks_released clear a parked failure):

  • Linux targeted unit tests: pytest -q tests/test_dense_offload_connector.py
    (55 passed); offload/connector suites
    (tests/test_offload_*.py tests/test_lmcache_*.py tests/test_kv_*.py tests/test_multi_connector.py tests/test_dense_offload_connector.py):
    700 passed, 8 skipped.
  • black --check on the repository and ruff check on the changed files:
    pass.
  • Real-GPU smoke on 4a8d76be (the pre-rebase head; the rebase changed
    only scheduler bookkeeping) (MI350X, ROCm 7.2, LMCache 0.4.5,
    one GPU, no model): DenseOffloadConnector.start_load_kv() saved 16/16
    tokens to LMCache LocalCPUBackend; the save-worker pack path consumed
    exactly one producer event, on the pack thread with a live pack stream, and
    reported producer_fenced=1; the store terminal and the source-quiescent
    completion were both emitted for the save operation; after the GPU KV
    tensors were cleared, retrieval restored all K/V and scale tensors with
    exact equality. The same smoke passed earlier on c6af1092.

GitHub Pre Checkin is pending repository approval for Actions from this fork.

Test Plan

  • Verify one producer event is recorded per save step and shared by multiple
    save requests.
  • Verify each save waits before calling store().
  • Verify the GPU connector enqueues wait_event on the pack stream before
    its block-ID upload and does not host-synchronize in the normal GPU path.
  • Verify a rejected or failed save restores the scheduler frontier and is
    re-emitted only after every rank is quiescent.
  • Run the existing non-GPU unit-test suite in CI.

Submission Checklist

  • The change is limited to dense LMCache save correctness and tests.
  • Static style and formatting checks pass.
  • Targeted dense connector unit tests pass on Linux.
  • Real-GPU LMCache save/retrieve smoke passes on the PR head.
  • ATOM CI unit tests pass.

@github-actions

Copy link
Copy Markdown
Contributor

🏷️ CI Guide

Runs automatically on every eligible PR before approval:

  • ✅ Pre Checkin: Black, Ruff, catalog schema validation, non-GPU unit tests

Heavy model tests:

  • ✅ Run after the PR is approved and Pre Checkin passes
  • ✅ Run immediately when an approval review is submitted
  • ✅ Can be requested before approval with labels
Label Tests
ci:full Run all heavy PR model tests: native ATOM, vLLM, and SGLang
ci:atom Run native ATOM model accuracy tests
ci:vllm Run ATOM vLLM OOT model accuracy tests
ci:sglang Run ATOM SGLang model accuracy tests

Heavy jobs are skipped when the PR is not approved and no matching ci:* label is present.
Add labels via the sidebar or gh pr edit 2339 --add-label <label>

@NidhoggD1
NidhoggD1 marked this pull request as draft September 21, 2026 17:58
@NidhoggD1
NidhoggD1 marked this pull request as ready for review September 22, 2026 07:39
@NidhoggD1

Copy link
Copy Markdown
Contributor Author

Ready for review. The validation section now includes the Linux targeted unit
suite and a real-GPU LMCache save/retrieve smoke on the current PR head.

The GitHub Pre Checkin run is currently action_required because this PR comes
from a fork. Could a maintainer please approve the workflow run and review the
change when convenient? Thank you.

@IvanShan177 IvanShan177 left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed at 54749d5.

The core premise holds, and I checked it before looking for problems: start_load_kv is dispatched before forward (process_kvconnector_output, atom/model_engine/model_runner.py:3333), and the dense scheduler only emits saves for tokens whose forward has already completed, so recording on the current stream at dispatch time really does fence the producing writes. pack_stream is also the right stream to fence, since batched_from_gpu stage A is the only leg that touches KV. The partial(...) change is compatible with OffloadWorkerMixin._guard(kind, fn, req), and DSV4's 3-arg _do_save_req is a separate class rather than a subclass.

Three comments below, all on error / None paths rather than on the fence itself. Concurrency first.

Comment on lines +183 to +191
save_ready_event = None
if self._do_save and any(
req.save_spec is not None for req in metadata.requests
):
# Save metadata is dispatched after the producing forward. Record
# that stream here so the background pack stream cannot read KV
# blocks before their writes are complete.
save_ready_event = torch.cuda.Event()
save_ready_event.record(torch.cuda.current_stream())

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The event record sits outside the dispatch loop, so a failure here drops the step's loads as well as its saves, and the requests hang.

Lines 183-191 run ahead of the for req in metadata.requests loop on line 192, which submits both load and save jobs. The raising call is really line 191 — torch.cuda.Event() is lazy in PyTorch (the hipEvent is created on first record), while current_stream() + .record() is a real runtime call, and on this path it is the first runtime call of the step. That makes it exactly where a sticky async HIP error from a previous step's kernel surfaces (async faults are reported at the next runtime call, not at the offending kernel).

The interleaving: .record() raises on the worker RPC thread -> the exception escapes start_load_kv -> ModelRunner.process_kvconnector_output (model_runner.py:3333-3338) calls it bare, with no handler -> the loop on 192 never runs, so no load future is ever submitted. Meanwhile the scheduler has already pinned those loads and is waiting on finished_loading / failed_loading from get_finished().

Neither ever arrives, because the only thing that reports a load failure is _guard (_offload_common.py:361-376), and _guard is the submitted job — if submission never happens, self._failed_load is never filled. The _lookup_unpin loop on 180-182 has already run at that point, so the step's pins are gone too.

This is worse under KimiK3OffloadConnector.start_load_kv (hybrid/kimi_k3/connector.py:259-285), which calls super().start_load_kv(metadata) from a finally after _arm_joint_loads. Its own comment states the invariant this breaks:

ThreadPoolExecutor.submit raises RuntimeError if the executor was shut down by a racing close() [...] and ModelRunner.process_kvconnector_output has no handler -- so each submit is isolated inside its helper [...] and super() runs in a finally so a KV load or save is never dropped.

So every raise site was deliberately isolated per item precisely because there is no handler upstream. A raise at 191 leaves the park armed with the KV leg never reporting (permanently parked), and the exception from the finally also masks whatever the try raised.

The sibling DSV4OffloadConnector already does this correctly for the same fence (hybrid/dsv4/connector.py:745-762): the create+record is inside the per-request loop, wrapped in try/except, and on failure it rejects just that one save and continues. Suggest matching that shape here — or at minimum wrapping 190-191 and leaving save_ready_event = None so the saves degrade instead of taking the loads down with them.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in c6af109. Fence creation is now lazy inside the per-save dispatch branch and wrapped in try/except. A record failure calls _record_save_failure(req) and continues the loop, so loads from the same metadata step are still submitted while no unfenced save is launched. Added a mixed load/save regression test that injects a record failure and verifies the load reaches finished_loading.

Comment on lines +437 to +451
def wait_for_save_source(self, producer_event) -> None:
"""Order this thread's pack stream after the KV producer stream.

Dense saves run on a background executor while the KV writes they read
were issued on the model/RPC stream. Enqueueing the event dependency on
the thread-local pack stream preserves that ordering without blocking
the CPU save worker. The following ``batched_from_gpu`` call runs on
the same thread and therefore reuses this stream state.
"""

state = self._thread_state()
if state.pack_stream is None:
producer_event.synchronize()
return
state.pack_stream.wait_event(producer_event)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This fence is thread-affine, and it silently degrades to a no-op if the pack ever runs off the calling thread.

wait_for_save_source enqueues the dependency on the calling thread's _ThreadTransferState.pack_stream, and the docstring makes the assumption explicit ("The following batched_from_gpu call runs on the same thread").

The interleaving that breaks it: save worker T1 enqueues wait_event(E) on T1's pack stream, then calls LMCacheEngine.store(). If store() dispatches the GPU pack on any other thread T2 — an LMCache storage-manager worker, an async-store configuration, the save_unfull_chunk path — then _thread_state() on T2 lazily constructs a fresh _ThreadTransferState (atom_lmcache_staging.py:98-111) with brand-new pack/copy streams that carry no dependency on E. batched_from_gpu then reads KV concurrently with the model stream's still-pending writes: precisely the stale-KV persistence this PR exists to prevent, with no exception and no log line.

DSV4's equivalent uses a host-side producer_event.synchronize() (hybrid/dsv4/connector.py:1523), which is thread-independent. The existing track_save_source leans on the same same-thread assumption, but its failure mode is a benign degradation (no source-safe callback), not silent data corruption — so I don't think the precedent carries over.

Worth either recording the owning thread on _ThreadTransferState and asserting it here, or falling back to producer_event.synchronize() when the state was not created by the current thread.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in c6af109 by removing the same-thread assumption. DenseOffloadConnector now forwards producer_event through LMCacheEngine.store(...); LMCache forwards the extra kwarg to BlockGPUConnector.batched_from_gpu(), which waits only after obtaining the actual calling thread state and pack stream. The GPU smoke confirms the event is consumed once on the real save-worker pack path without host synchronization.

Comment on lines +354 to +355
if producer_event is not None:
gpu_connector.wait_for_save_source(producer_event)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is now the only unguarded gpu_connector dereference in _do_save_req.

Every neighbouring access is defensive: _reset_gpu_connector_transfer_stats and _last_gpu_connector_transfer_stats both early-return when gpu_connector is None (_offload_common.py:336-360), and the very next line uses getattr(gpu_connector, "track_save_source", None) plus a callable(...) check. tests/test_dense_offload_connector.py:307-310 constructs SimpleNamespace(gpu_connector=None, store=...), so None is a contemplated state rather than a theoretical one.

Because production now always passes a non-None producer_event, any engine whose gpu_connector is None — or which predates wait_for_save_source — turns every save into an AttributeError. That gets swallowed by _guard into _record_save_failure, i.e. 100% silent save failure (offload effectively off, visible only as an exception log), rather than a degraded-but-working save.

getattr(gpu_connector, "wait_for_save_source", None) with a producer_event.synchronize() fallback would keep both cases correct, and would match the style of the line directly below.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in c6af109 as part of the event handoff change. _do_save_req no longer dereferences gpu_connector.wait_for_save_source; it passes producer_event through engine.store, while the adjacent track_save_source access remains guarded with getattr/callable. The existing gpu_connector=None-compatible unit path remains intact. The targeted Linux suite is now 44 passed.

@valarLip
valarLip requested a review from yhl-amd September 22, 2026 15:06
yhl-amd
yhl-amd previously approved these changes Sep 22, 2026
@valarLip

Copy link
Copy Markdown
Collaborator

Review — fence dense LMCache saves

Reviewed at the PR head (3 commits, +180/−8) in a detached worktree; the main checkouts were never touched.

The premise checks out, so this is not a no-op. I verified independently that the save frontier only covers already-forwarded chunks on both dispatch paths, that disagg-side streams rejoin the default stream, and that LMCache does forward store(**kwargs) through to batched_from_gpu in both shipped versions (the 0.5.5rc3 wheel and the v0.4.5 source build pinned in docker/). The race is real and this closes it.

What follows is about the invariants it leaves unasserted, the legs it did not sweep, and the failure path. Two statements:

The fence is carried on a channel that can disappear silently — producer_event=None means "no fence" with no log, counter or assertion, and LMCacheEngine.store() has four early returns that never reach the connector at all. The comment justifying that channel is contradicted by the mechanism twelve lines above it.

The failure path converts "one missing fence" into a permanent silent cache hole plus an unbounded traceback flood on the RPC thread — because it reports a never-attempted save as finished after the scheduler watermark has already advanced.

And the channel this wants already exists in the tree, owned by the producer: record_kv_cache_ready.

Provenance: [verified] means the deciding code was read at this head or the behavior was confirmed against LMCache/torch on this box. [reported] means the shape matches the code but was not traced end to end.


1. The fence's arrival is unverifiable [verified]

dense/connector.py:373. batched_from_gpu(..., producer_event=None) skips _wait_for_save_source with zero signal — no log, no counter, no assertion. And LMCacheEngine.store() returns before gpu_connector.batched_from_gpu on four paths: not is_healthy(), _is_passive(), is_frozen(), and empty memory_objs.

So any LMCache bump that builds an explicit connector-kwargs dict, or a layerwise/remote store route, drops the kwarg and the race returns — observable only as an intermittent accuracy drop under offload hits, which this repo's own noise floor (~9pp at 100 GSM8K questions) makes effectively invisible.

The connector already carries the exact idiom for this class of question — stats['batch_block_ids_enabled'] = 1 (_block_gpu_connector.py:666). Record a producer_fenced stat in _capture_transfer_stats and surface it in [OFFLOAD-SAVE-PROF].

The comment defending the kwarg route is contradicted in the same with block (:374). It reads "even if LMCache changes its dispatch model" — but twelve lines above, track_save_source (_block_gpu_connector.py:792) carries save identity through a thread-local, and its docstring says it does so "without changing LMCache's public API". The same hypothetical dispatch change silently breaks that one: if LMCache ever runs batched_from_gpu on its own thread, the fence survives while early block release stops — _source_safe_identities (:808-821) returns () for a non-SaveOperationId, _queue_source_safe_group returns without publishing, no DENSE_PAGE_SOURCE_SAFE_CHANNEL completion is ever emitted, and every save's source blocks stay leased until reclaim_stale_leases. Either delete the overclaim, or put the event on the same thread-local channel — self._engine.gpu_connector is already reachable from _do_save_req (see _last_gpu_connector_transfer_stats, _offload_common.py:337).

2. The failure path is worse than having no fence [verified]

dense/connector.py:204. _record_save_failure(req) → _record_store_terminal(req, False) adds the operation to _done_save (i.e. finished_saving) in both branches. Four consequences compound:

  1. A save that was never attempted is reported as finished.
  2. chunked_scheduler.py:435 already did entry[1] = aligned at emit time and nothing resets it, so those chunks are never retried — a permanent silent cache hole.
  3. _store_finished(succeeded=False) (chunked_scheduler.py:751-753) deliberately does not call _release_operation_lease — "Failed stores keep unsafe ranges leased until abandon timeout" — so the block lease is held to the timeout.
  4. Because save_ready_event stays None, a sticky HIP error — the stated motivation for this change — retries torch.cuda.Event() and emits a full logger.exception traceback for every save-bearing request, every step, forever, on the pre-forward RPC thread.

The comment "rejects only this save" is inaccurate: it rejects every save for the life of the error. Caching a per-step sentinel instead of None breaks item 4.

No CUDA precondition on the producer side (:200) [verified]. The consumer side explicitly contemplates non-CUDA (state.pack_stream is None), but start_load_kv calls torch.cuda.Event() / torch.cuda.current_stream() unconditionally for every step that has a save, with no _use_cuda()-style predicate and no once-only degradation. On any worker whose torch has no usable CUDA context, every save in every step takes the except branch: a stack-trace flood on the RPC thread — the thread this whole connector design exists to keep free — plus a terminal "finished" report per save, with the watermark already advanced so nothing is ever retried. Offload is silently dead rather than loudly rejected. One check at register_kv_caches is louder and cheaper.

3. Three legs not swept, one of them a subclass of the class being fixed

KimiK3's state-tier stores have the identical hazard, by their own docstring (hybrid/kimi_k3/connector.py:399) [verified]. They are dispatched from the same start_load_kv, before super(), and explicitly carry no producer fence:

"No producer fence and no staging copy: the source is the checkpoint's PAGE units, reserved out of the KV pool and pinned by the engine."

Pinning answers "can these blocks be reallocated", not "have the producing writes retired". start_load_kv (:283) submits these to self._state_tier's own executor threads before super().start_load_kv(metadata) (:285) creates the new fence. Under deferred output the host has not synced on forward N−1, so the tier packer gathers a torn mix of pre/post-write bytes into a recurrent-state checkpoint that restores cleanly and corrupts output on reuse. The sibling docstring at :355 (_start_state_loads: "No producer fence, unlike the save path") is now self-contradictory.

DSV4 is the other _engine.store( call site and was not swept (hybrid/dsv4/connector.py:1603) [verified]. grep -n '_engine.store(' returns exactly two hits; one was fixed. dsv4 does producer_event.synchronize() at :1525 (and :945) then calls store(tokens, mask=, block_ids=, req_id=) with no kwarg — so the new _wait_for_save_source is dead code for the entire DSV4 PAGE path, which keeps the host stall this PR exists to remove, while the PR's own test is named ..._without_host_sync. Note it cannot be converted mechanically: dsv4's synchronize() also gates the SLOT D2H, and the _SlotStagingSyncError quarantine depends on it raising. That makes it a design decision the PR should state, not omit.

Loads get nothing, and the in-repo precedent fences both directions (dense/connector.py:188) [verified]. mp/backend.py:723-729 records one torch.cuda.Event(interprocess=True) and hands it to _submit_load and _submit_save. Here the load submit at :188-193 runs before the event is even created, and batched_to_gpu has no parameter for it. A block freed by a request that finished or was preempted in step N's schedule() and reallocated to a loading request can take the H2D restore while forward N−1 (deferred output, not host-synced) is still reading or writing it; wait_for_requests fences at request granularity, not against the producing forward. Either the load direction genuinely does not need it — then say so, because the mp backend's symmetry now reads as an inconsistency someone will "fix" — or it is the same bug one line up.

Under PP the invariant does not hold (dense/connector.py:201) [verified, latent]. Scheduler.advance_on_schedule is pipeline_parallel_size > 1 (scheduler.py:598), and _advance_prefill_on_schedule (:2567) does seq.num_cached_tokens += num_scheduled_tokens[i] inside schedule(), before the forward. chunked_scheduler.build_connector_meta reads that same field (:397-408; _offload_finished_cached_tokens is only set for retiring requests) and emits a SaveSpec covering the in-flight chunk. So candidate_event.record(...) runs before that forward is enqueued and pack_stream.wait_event orders the pack against nothing: the save packs pre-forward KV and stores a corrupt chunk under a valid content key, restored silently on a later prefix hit. Latent only because PP is currently broken on this box. Needs an assertion (save frontier ≤ previously-forwarded frontier) or a comment naming the chunked_scheduler invariant it rests on.

4. This is the third per-step producer event in the package, and the producer-owned channel already exists [verified]

dense/connector.py:195. ModelRunner._record_kv_cache_ready (model_runner.py:3199, fired at :3290/:3312/:4659 at the end of forward, per final prefill chunk) already calls connector.record_kv_cache_ready(req_ids), and offload/connector.py:133 already forwards it to the impl. That docstring still says "no offload impl needs it now: the event exists only for the Mooncake producer" — which this PR makes false without updating it. Mooncake is the reference consumer (mooncake_connector.py:1059-1064: a torch.cuda.Event() on the post-forward stream, keyed per req_id, with explicit pop cleanup at three terminal points). The other two copies are mp/backend.py:723 and dsv4/connector.py:743-762.

The right-depth fix is to implement record_kv_cache_ready on OffloadWorkerMixin and wait it in _do_save_req: one fence, owned by the producer, on a channel LMCache cannot drop, shared by dense/dsv4/kimi_k3/m3/mp, and symmetric for loads. That collapses §1, §3 and this section together.

5. Four mechanical issues [verified]

  • The wait is enqueued after work is already on the same stream (_block_gpu_connector.py:915). :910-914 runs _prepare_block_id_stage, whose prepared path calls prepare([...], device=self.device, stream=state.pack_stream) (:659-663) — real stream work — and only then does :915-916 enqueue the wait. So "nothing this call submits runs before the producer event" is already false. Benign today because that work is a host-side block-id H2D, but the fence's validity is now an implicit, unstated property of every codec's prepare_block_id_groups. Swap the two lines (cost: nothing) and state the invariant in _wait_for_save_source's docstring.
  • producer_event is a positional-or-keyword 4th parameter on a method a third-party library calls (:897). The abstract contract is three positionals plus **kwargs, so any wrapper drift binds a memory-obj-ish value to producer_event and _wait_for_save_source calls .wait_event(<not an event>) inside the save thread — where _guard turns it into a save reported as finished (§2). Make it keyword-only. Separately, batched_to_gpu (:934) and to_gpu (:889) have no such parameter, so a producer_event= added to a retrieve "by symmetry" later falls into **kwargs → _prepare_transfer → _ranges_to_block_ids, which reads only block_ids and drops the rest: it would compile, pass, and fence nothing.
  • functools.partial silently degrades the failure log (dense/connector.py:216). _offload_common.py:368 logs getattr(fn, "__name__", kind), and a partial has no __name__ — so every dense save failure now logs offload save failed for <req> instead of offload _do_save_req failed for <req>, losing the one field that distinguishes a save-path failure from a submit-path one, on exactly the path where a new failure class was just introduced. Latent trap: dsv4/connector.py:1078-1084 has its own _guard whose handler does a bare fn.__name__, so copying this idiom there raises AttributeError inside the except block. Use functools.update_wrapper, or widen _guard to (kind, fn, req, **fn_kwargs) — dsv4 already passes producer_event as a plain argument rather than binding it into a partial.
  • The pack_stream is None → producer_event.synchronize() branch is unreachable, untested, and would busy-spin (_block_gpu_connector.py:441). pack_stream is None only when _ThreadTransferState was built with use_cuda=False (atom_lmcache_staging.py:105-110), but the only source of a producer_event is dense/connector.py:201, which requires a live CUDA context — and _run_staged_pipeline → _assert_fused_chunk_major_available() raises three lines later on a non-CUDA device anyway. Meanwhile this torch.cuda.Event() takes default flags, unlike every sibling in the package which writes blocking=False explicitly (:839), and this is the first event in the package ever host-synchronize()d — so the wait would spin a full core on the save daemon against a co-resident model. Delete the branch (let _assert_fused_chunk_major_available own the non-CUDA case) or raise; if it is kept, the event needs blocking=True.

6. All three new tests pass for weaker claims than their names assert [verified]

tests/test_dense_offload_connector.py:333.

(a) The one-event-per-step property is unfalsifiable. torch.cuda.Event is patched to lambda: event — a shared singleton — and current_stream globally to rpc_stream. An implementation that created a fresh event per request, or created and recorded it inside _do_save_req on the save thread (the exact anti-design this shape exists to prevent), yields an identical trace and stays green.

(b) The failure test cannot distinguish failure from success. assert output.finished_saving == {save_operation} (:427) is identical either way — dense sets _early_release=True, so the real verdict lives in ConnectorCompletion(DENSE_PAGE_STORE_CHANNEL, op, succeeded), which the test never inspects. Changing :209 to _record_store_terminal(req, True) keeps it green. Its load is also degenerate: hbm=lmc=0 early-returns before retrieve is ever called.

(c) The ordering test asserts only its own wiring. It stubs five private methods, and _prepare_block_id_stage's stub appends nothing to trace — so assert trace == [("wait", ...), ("pipeline", ...)] is true by construction and blind to §5's first item.

(d) The pack_stream is None branch has zero coverage; the only test that touches it wires synchronize to pytest.fail.

(e) The zero-arg stubs pin the exact signatures under review: adding blocking=True or a device argument makes both tests TypeError, so the fixes above will read as "the fix broke the tests".

7. Docs-match-code: every document describing this call chain is stale, including the comment that justifies the change [verified]

  • README.md:352 — the save-flow diagram is still engine.store(tokens, mask, block_ids) → batched_from_gpu with no fence.
  • README.md:273 — frames the pre-return current-stream touch as DSV4-only.
  • README.md:616-626 — the "Invariants enforced in code" table gains no row for the invariant the PR's own test treats as hard.
  • README.md:808-823 — documents batched_from_gpu/batched_to_gpu with no mention of the new parameter or the LMCache-kwargs dependency. That is the single most fragile cross-repo assumption in the subsystem, documented only in a three-line inline comment.
  • batched_from_gpu's own docstring (:900) is unchanged.
  • dense/connector.py:16 still claims start_load_kv "only submits ... so the worker RPC thread is free for forward", which the new CUDA-event record and _record_save_failure call break.
  • offload/connector.py:136 still says no offload impl needs record_kv_cache_ready (§4).

Finally, the new comment at dense/connector.py:196 ("Save metadata is dispatched after the producing forward") reads as false against engine_core.py:418 and contradicts dsv4's own comment at :722 ("Metadata is dispatched before this batch's forward"). It should name the chunked_scheduler frontier invariant it actually depends on — which is also what §3's PP item needs.


Of the above, §2 is the one I would hold on: it is the only item that turns a transient failure into permanent silent data loss, and it does so on the exact error class named as the motivation. §1 is the one worth fixing even though the fence works today — an invariant whose presence cannot be observed will be removed by accident. §3's dsv4 and kimi_k3 legs need at minimum a sentence each saying whether they were considered; right now the PR reads as if dense were the only producer.

@kvnloo

kvnloo commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

Thanks for documenting the real-GPU smoke and the follow-up work on the save path. I'm working through the ATOM–LMCache integration downstream and don't have the target AMD hardware locally.

Would it be useful if I turned your existing LocalCPUBackend smoke into a small reproducible validation recipe, using whichever script and revision you think are appropriate?

I can handle the source/CPU checks and result write-up downstream. I'm mainly trying to preserve one narrow claim: actual production-path storage and reload with K/V + scale equality, separately from model accuracy or performance. I don't want to duplicate the fence work you're already doing here.
I also found the existing GPU/disk round-trip test useful as an earlier check, although because it explicitly synchronizes the producer I'd treat it as narrower evidence rather than validation of this async path.

Does that decomposition make sense, or is there an existing test/script you'd prefer we build the qualification around?

AI-assisted source review and drafting; I don't have GPU results of my own to report yet.

@NidhoggD1

Copy link
Copy Markdown
Contributor Author

@valarLip thanks for the thorough pass. Addressed in 6aefe2ee (scheduler retry) and 4a8d76be (fence observability, ordering, tests, docs). Point by point:

§1 fence arrival unverifiable — done. batched_from_gpu now records producer_fenced in the transfer stats (default 0, set to 1 when the wait is enqueued) and [OFFLOAD-SAVE-PROF] prints it, so a store that reaches the connector without the event is visible per save. The "even if LMCache changes its dispatch model" comment is gone; the replacement names the store(**kwargs) forwarding as the one cross-repo assumption, and the README's _block_gpu_connector.py section says the same. I kept the kwargs route rather than the thread-local: if LMCache ever ran the pack on its own thread, the thread-local would lose the event silently, whereas kwargs survive that, and the stat now makes loss visible either way.

§2 failure path — done in 6aefe2ee. A rejected or failed save no longer advances the watermark: the scheduler keeps the pre-emit frontier per operation, restores it on failure and re-emits the save once every TP rank is quiescent. A pre-submit rejection (fence record raised) reports quiescence directly; a post-submit failure becomes retryable when all of its source groups have been reported source-safe, otherwise the lease is reclaimed on the abandon timeout. A save that fails after its request was retired is dropped rather than retried against recycled blocks. Fence failure is logged once per failure streak, and one sentinel per step rejects the remaining saves of that step without re-recording. On the CUDA precondition: _assert_fused_chunk_major_available already refuses a non-CUDA device on the first pipeline run, so offload is rejected loudly there; I did not add a second check at registration.

§3 legs not swept —

  • DSV4: deliberate, not converted. Its producer_event.synchronize() also gates the SLOT D2H and the _SlotStagingSyncError quarantine depends on it raising, so replacing it is a separate change with its own test plan. The PR description now says the PAGE fence is dense-only.
  • KimiK3 state stores: not changed here. Whether the tier packer needs the same fence depends on whether the step that marks a checkpoint ready has host-observed the producing forward; I have not traced that and would rather fix it in a follow-up than assert it is safe.
  • Loads: the hazard you describe (a partial-prefill request preempted in schedule(), its blocks reallocated to a load in the same step) is a block-reuse ordering question that exists independently of the save fence, and batched_to_gpu has no parameter for it today. Out of scope for this PR; worth its own issue.
  • PP: the chunked scheduler now warns at init when pipeline_parallel_size > 1 and saves are enabled, with a comment naming advance_on_schedule as the reason the frontier invariant does not hold there.

§4 record_kv_cache_ready — not adopted. ModelRunner._record_kv_cache_ready fires only for requests whose is_final_chunk is true, but the racing saves are exactly the intermediate prefill chunks (no token is produced, so nothing host-synchronizes before the next dispatch). Using that channel would require it to fire per chunk, which is a model-runner change I would keep separate. The offload/connector.py docstring stays accurate because dense still does not consume the callback.

§5 mechanical — all four done: the wait is enqueued before _prepare_block_id_stage and the docstring states that nothing else may precede it; producer_event is keyword-only; _guard takes **fn_kwargs so the partial (and the lost fn.__name__) is gone; the pack_stream is None branch raises instead of host-synchronizing.

§6 tests — (a) torch.cuda.Event is now a counting class: exactly one instance, recorded on the RPC thread, shared by both stores. (b) The failure test asserts ConnectorCompletion(DENSE_PAGE_STORE_CHANNEL, op, False) and the quiescent report. (c) The ordering test traces _prepare_block_id_stage and asserts wait < block_ids < pipeline, plus producer_fenced == 1. (d) New tests cover the keyword-only signature and the no-pack-stream raise. (e) The stubs were rewritten with the new signatures.

§7 docs — README save-flow diagram, start_load_kv bullet, invariants table and _block_gpu_connector.py section updated; the dense module docstring no longer claims start_load_kv only submits; the fence comment now states the frontier invariant (metadata dispatched before the forward; the save frontier covers only already-forwarded chunks).

The real-GPU smoke in the description was run on c6af1092; the GPU-facing change since then is the wait moving ahead of the block-ID upload. I will re-run the smoke on the current head and update the description.

@NidhoggD1

Copy link
Copy Markdown
Contributor Author

Re-ran the real-GPU smoke on the current head 4a8d76be (MI350X, ROCm 7.2.53211, torch 2.10.0+rocm7.2.4, LMCache 0.4.5, single GPU, LocalCPUBackend, no model). Result: passed. Exactly one producer-event wait, executed on the save-worker pack thread against a live pack stream; the connector reported producer_fenced=1; two source-safe groups, one store terminal (succeeded) and one source-quiescent completion for the save; 16/16 tokens hit on the host tier; after zeroing the GPU KV, retrieval restored every K/V and scale tensor with exact equality. The PR description now reflects this run.

NidhoggD1 and others added 5 commits September 25, 2026 04:32
Treat a save whose every source group is source-safe on every rank as
quiescent: no rank still reads its source, so a store failure that arrives
after teardown leased nothing may retry or retire at once instead of being
parked in `_save_retry_blocked` with no reclaim path.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Follow-ups from review:

- `batched_from_gpu` enqueues the producer wait before the block-ID upload,
  so no pack-stream work of a transfer precedes the dependency, and records
  `producer_fenced=1` in the transfer stats. `[OFFLOAD-SAVE-PROF]` prints it,
  so a store that reaches the connector without the event is visible instead
  of silently unfenced.
- `producer_event` is keyword-only; a wrapper passing extra positionals can no
  longer bind one to it.
- A staging state without a pack stream now raises instead of spinning on
  `Event.synchronize()`; the fused staging pipeline already rejects that
  device.
- `OffloadWorkerMixin._guard` forwards keyword arguments, replacing the
  `functools.partial` that hid `fn.__name__` from the failure log.
- The fence comment names the scheduler frontier invariant it rests on
  (metadata is dispatched before the forward; the save frontier covers only
  already-forwarded chunks), and the chunked scheduler warns under PP, where
  `advance_on_schedule` breaks that invariant.
- Tests now falsify the one-event-per-step and wait-before-upload claims,
  cover the keyword-only and no-pack-stream paths, and check `_guard`'s
  kwargs and log name. README and the connector module docstring describe
  the fence and its LMCache-kwargs dependency.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
@NidhoggD1
NidhoggD1 force-pushed the fix/lmcache-dense-save-fence branch from 4a8d76b to e9be221 Compare September 24, 2026 20:44
@NidhoggD1

NidhoggD1 commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor Author

@kvnloo yes, that decomposition makes sense, and building the qualification around the LocalCPUBackend smoke is the right call. Two pointers:

  • Script and revision: use the version of the smoke that ran on 4a8d76be (results and description above), against the current head e9be221f. It is the earlier script plus three checks the head now supports: the producer wait must land on the pack thread with a live pack stream, last_transfer_stats()["producer_fenced"] must be 1 for the store, and both dense.page.store (succeeded) and dense.page.source_quiescent completions must be emitted for the save operation. It needs one GPU, LMCache 0.4.5 and the PR checkout on PYTHONPATH; no model. Script and runner: https://gist.github.com/NidhoggD1/4768a59519b2e4a9c83e80037a21db79
  • Agreed on test_lmcache_offload_gpu_disk_e2e.py: it records and forwards a producer event but the test itself is not a race, so treat it as coverage of the plumbing, not evidence for the async ordering.

One caveat for the write-up: the smoke proves "the shipped store path waits on the producer and reloads K/V + scales exactly". It does not exercise TP quorum or the failed-save retry path; those are unit-tested only.

@kvnloo

kvnloo commented Sep 24, 2026

Copy link
Copy Markdown
Contributor

Thanks, this gives me a concrete starting point. I’ll build the downstream recipe around the linked smoke, using e9be221 as the target and recording the reported 4a8d76b run separately.

I’ll preserve the pack-thread/live-stream check, producer_fenced=1, both save completions, and exact K/V + scale restoration. The write-up will explicitly exclude TP quorum and failed-save retry coverage, and describe the disk test as plumbing coverage rather than an ordering-race test.

I’ll review the script and prepare the recipe before requesting GPU time. No additional testing needed from you at this stage, thanks for sharing the runner.

kvnloo commented Sep 24, 2026

Copy link
Copy Markdown
Contributor

Following up with the pinned downstream recipe. It reuses pr2339_gpu_smoke.py unchanged from Gist snapshot c881bce7, targets e9be221f, and records your reported 4a8d76be run separately. I haven't run the smoke on a GPU.

Could you sanity-check just the acceptance table and environment assumptions before I seek a one-GPU reproduction? In particular, would running the unchanged smoke in an already-provisioned, compatible LMCache 0.4.5/ROCm environment preserve its intended coverage without the machine-specific outer runner?

The recipe retains the pack-thread/live-stream wait, producer_fenced=1, both save completions, and exact host-backed K/V + scale restoration. It explicitly excludes normal registration, scheduler-driven eviction, TP quorum, failed-save retries, model accuracy, and ordering-race or performance claims.

This is a source-review request only—no additional run requested from you and no changes to the fence implementation. Thanks for providing the concrete starting point.

AI-assisted source review and drafting; no GPU results of my own claimed.

@valarLip
valarLip merged commit 68e0df5 into ROCm:main Sep 25, 2026
60 of 63 checks passed
@NidhoggD1

Copy link
Copy Markdown
Contributor Author

@kvnloo reviewed the recipe at 4feea25b. Answers to your two questions, plus two things worth adding:

Acceptance table: correct. Each row matches an assertion in the script: one wait on a real torch.cuda.Event, taken inside batched_from_gpu on a thread other than the dispatcher with a live pack stream; exactly one store-side stats record with producer_fenced == 1 (loads go through batched_to_gpu and are not counted); two source-safe groups for the fixed 16-token / chunk-8 geometry; one succeeded dense.page.store plus one dense.page.source_quiescent for the same SaveOperationId; a 16/16 LocalCPUBackend lookup; and torch.equal on all eight tensors (2 layers × K/V/k_scale/v_scale) after zeroing.

Running without the outer runner: yes, coverage is unchanged. The runner only isolates the run and collects records. The smoke needs none of --ipc host, --network none or the 4 GiB shm: it is a single process with no lookup server or collectives. It reads only PR_COMMIT and SMOKE_RESULT from the environment, and your recipe exports both. Keep your import-origin check: the image we used also ships /app/ATOM, so without PR_SRC first on PYTHONPATH the smoke would silently test the image's copy.

Checksum provenance: now matched. The file executed for the 4a8d76be run has SHA-256 fc699dbc3465729fb8fa17593ee90c9830a28c5a7d41343a302edc43b968d342, recorded at launch time in that run's RUN.json, and the runner has SHA-256 daede71dd36bd6577c2b8a3e0c016151f324ab61cd8e9c9e077f7f106e57d627. Both match your captured snapshot c881bce7 byte for byte, so you can record that the historical run executed exactly the gist bytes.

Expected stderr line, not a failure: every run we have, on both machines, logs LMCache ERROR: Error closing backend LocalCPUBackend: tuple index out of range during LMCacheEngineBuilder.destroy(). This comes from LMCache 0.4.5's own teardown, after result.json is written. It is logged, not raised, and the process still exits 0. Your recipe says cleanup errors are not a pass, so it would help to list this line as a known exception rather than have a volunteer flag it.

Your point about container_removed=True is fair: the runner sets it without checking the result of docker rm. For our run, cleanup was confirmed separately and twice (no container left, GPU VRAM 0%, KFD process count back to baseline). The field on its own is not evidence.

Also, this PR merged as 68e0df5e, and its PR head is e9be221f, your target, so the recipe can pin either.

kvnloo commented Sep 25, 2026

Copy link
Copy Markdown
Contributor

Thanks for checking both the assertions and environment. I've updated the recipe with the historical checksum match and the exact LMCache 0.4.5 teardown message.

The write-up now separates the bounded transfer result, diagnostic review, and independently observed process cleanup. The known line alone does not invalidate otherwise valid transfer observations, but it does not certify cleanup or excuse another error. Your historical cleanup checks remain attributed to your run.

I'll keep e9be221f and the import-origin check intact, record the merge separately, and use the unchanged smoke without requiring the outer runner. No further review or GPU run requested from you here.

AI-assisted recipe maintenance; no GPU run of my own claimed.

yhl-amd pushed a commit that referenced this pull request Sep 27, 2026
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>
zejunchen-zejun added a commit that referenced this pull request Sep 28, 2026
A state load's lifecycle held three facts in three owners across two
processes: "this request owes a report on hash H" (StateOffloadIndex,
engine), "its state slot stays off the free list"
(BlockManager._orphan_load_slots, engine) and "have both legs landed"
(_JointPark, worker). No object could state

    dispatched == settled + outstanding

so no test could assert it.

The state leg now rides the request as `LMCacheReqMeta.state_load_spec`
and runs inside the KV leg's own task, so one dispatch emits exactly one
completion on every path, including a raise. Dense's `_do_load_req` is
split into `_load_kv_bytes` + `_finish_load` to give it that seam. A
state-only load travels the ordinary load path on a no-op KV spec
(hbm == lmc). `StateOffloadIndex` becomes the sole engine-side owner: it
takes the destination slot, absorbs orphan parking, and audits its own
invariant (surfaced as `state_offload_invariant_violations`).

Deleted: `_JointPark`, `metadata.state_loads`, the state-load disposition
channel, the state tier's load executor and staging-lane semaphore, the
engine-side load transport on BlockManager, `_publish_state_loads` /
`_settle_state_load` / `_abandon_state_load`, and the state-only park
branch. PP > 1 now refuses at startup instead of warning.

Rebased onto main (ba51495). The dense save producer fence this branch
carried is dropped in favour of main's #2339, which orders the pack stream
on the device rather than host-synchronizing; `producer_event` on
LMCacheReqMeta and its tests go with it, and staging.py's fence docs now
point at #2339's path. Conflicts with #2238 and #2305 resolved keeping both
sides.

Squashed from 20 commits; the pre-rebase history is at ea3f48d.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
valarLip pushed a commit that referenced this pull request Sep 28, 2026
* feat(offload): support native DSv4 checkpoints with LMCache MP

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>

* docs(offload): pass --eviction-policy to the LMCache MP server example

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>

* fix(lmcache_mp): skip native saves shorter than OFFLOAD_MIN_SAVE_TOKENS

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>

* fix(offload): keep early release safe without a BlockManager; one save 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>

* fix(lmcache_mp): bound unprovable transfers; fence native restores

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>

* fix(lmcache_mp): end sessions after the last save; unwind failed admissions

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>

* fix(lmcache_mp): keep draft PAGE regions out of the native image

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>

* fix(attention): take the descriptor slot in DSv4.1 paged-state copies

Review follow-up (#2250 §5).

descriptor_slot reached the base class, DSv4 and GDN but not DSv4.1, whose
execute_paged_state_copies(stores, restores) would raise TypeError inside the
native restore's exception handler and be reported as a failed restore.
StateCopies now keeps an independent pinned staging buffer and upload fence
per descriptor slot, like the other backends. A new AST sweep test fails if
any attention backend's implementation lacks the parameter.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix(offload): seed native and late-save hash chains from cache_seed

Review follow-up (#2250 §2).

The native boundary-hash fallback and BlockManager.acquire_offload_prefix both
started their chains at -1, while every BlockManager chain starts at
seq.cache_seed, which multimodal requests set. For those requests every
native save was skipped, every restored checkpoint was published under a key
no loader looks up, and late saves found nothing resident. The native
fallback now extends its chain through the new BlockManager.prefix_hash_chain
(the manager's own _chain_to), so both _boundary_hash branches share one
seed, algorithm and token slicing. acquire_offload_prefix seeds from
seq.cache_seed.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* fix(offload): fail closed on disabled paged-state checkpoints and composite state release

Review follow-ups (#2250 §6).

- The MP scheduler shell picked the native-state scheduler whenever the
  checkpoint coordinator existed, but the coordinator exists without prefix
  caching and then can never hold a READY image, so offload silently did
  nothing. A PAGE-only fallback would restore KV under stale recurrent state.
  Binding now fails with a message naming --enable-prefix-caching.
- MultiConnectorScheduler.can_partially_deallocate_state now requires every
  still-deferring sub to guarantee state safety. Before, one declaring sub was
  enough, and the state slot could be recycled while a non-declaring sub (e.g.
  a dense PAGE sub) still needed the request alive. The test now separates
  all-of from any-of.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

* perf(lmcache_mp): cache the native save frontier; reserve restore slots 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>

* docs(offload): document every offload env var; guard it with a test

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>

* test(lmcache_mp): make the raising-restore test independent of CUDA

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>

* fix(lmcache_mp): release transfer memory only on a report, fail stop 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>

* fix(offload): lease the final save source when prefix caching is off

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>

* fix(offload): reacquire a final save source only after a partial release

`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>

* fix(lmcache_mp): keep the native STATE image pinned until the STORE terminal

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>

* fix(envs): read an empty offload knob as its default everywhere

`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>

* fix(lmcache_mp): keep failed adoptions releasable and surface broken 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>

* test(lmcache_mp): run a native restore end to end and fix drifted copy 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>

* perf(lmcache_mp): invalidate a request's save frontier only on its own 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>

* refactor(attentions): one fenced descriptor staging for every checkpoint 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>

* refactor(lmcache_mp): validate PAGE views once for both registrations

`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>

* refactor(lmcache_mp): split the native worker's completion poll

`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>

* refactor(lmcache_mp): split mp/backend.py by responsibility

`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>

* fix(offload): keep clock reclaim off MP requests; warn on truncated late 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>

* refactor(kv_transfer): publish each PAGE region and its byte view together

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>

* refactor(kv_transfer): make a PAGE unit one value, its region and view 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>

---------

Co-authored-by: Honglie Yi <hyi@crsuse2-m2m-v2-020.us-east2-a.compute.internal>
Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
zejunchen-zejun added a commit that referenced this pull request Sep 28, 2026
A state load's lifecycle held three facts in three owners across two
processes: "this request owes a report on hash H" (StateOffloadIndex,
engine), "its state slot stays off the free list"
(BlockManager._orphan_load_slots, engine) and "have both legs landed"
(_JointPark, worker). No object could state

    dispatched == settled + outstanding

so no test could assert it.

The state leg now rides the request as `LMCacheReqMeta.state_load_spec`
and runs inside the KV leg's own task, so one dispatch emits exactly one
completion on every path, including a raise. Dense's `_do_load_req` is
split into `_load_kv_bytes` + `_finish_load` to give it that seam. A
state-only load travels the ordinary load path on a no-op KV spec
(hbm == lmc). `StateOffloadIndex` becomes the sole engine-side owner: it
takes the destination slot, absorbs orphan parking, and audits its own
invariant (surfaced as `state_offload_invariant_violations`).

Deleted: `_JointPark`, `metadata.state_loads`, the state-load disposition
channel, the state tier's load executor and staging-lane semaphore, the
engine-side load transport on BlockManager, `_publish_state_loads` /
`_settle_state_load` / `_abandon_state_load`, and the state-only park
branch. PP > 1 now refuses at startup instead of warning.

Rebased onto main (ba51495). The dense save producer fence this branch
carried is dropped in favour of main's #2339, which orders the pack stream
on the device rather than host-synchronizing; `producer_event` on
LMCacheReqMeta and its tests go with it, and staging.py's fence docs now
point at #2339's path. Conflicts with #2238 and #2305 resolved keeping both
sides.

Squashed from 20 commits; the pre-rebase history is at ea3f48d.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

6 participants