Skip to content

[Kimi-K3][LMCache] Fuse the state load leg, and fence the KV save against the forward - #2132

Open
zejunchen-zejun wants to merge 2 commits into
mainfrom
zejun/k3_lmcache_refactor
Open

zejunchen-zejun wants to merge 2 commits into
mainfrom
zejun/k3_lmcache_refactor

Conversation

@zejunchen-zejun

@zejunchen-zejun zejunchen-zejun commented Sep 3, 2026 •

Copy link
Copy Markdown
Collaborator

This PR restructures the Kimi-K3 LMCache offload tier that #2053 landed: the recurrent-state load leg now rides the KV load task, so one dispatch emits exactly one completion and the invariant dispatched == settled + outstanding is something one object can state and one test can assert.

The dense save fence this branch originally carried is dropped: main's #2339 fixed the same race better (the pack stream waits on the event on the device, no host sync). This PR is now the restructure only, plus the fixes listed under Also fixed.

The defect

A state load's lifecycle held three facts in three owners across two processes:

fact owner process
"this request owes a report on hash H" StateOffloadIndex engine
"its state slot stays off the free list" BlockManager._orphan_load_slots engine
"have both legs landed" _JointPark worker

No object could state dispatched == settled + outstanding, so no test could assert it. That is the shape behind the load-path findings this subsystem kept producing. _JointPark grew to 213 lines / 12 methods, every one of them added by a bug fix.

The change

The state leg rides the request as LMCacheReqMeta.state_load_spec, the same shape DSV4's slot_load_spec uses, and runs inside the KV leg's own task:

def _do_load_req(self, req):
    try:
        ok = self._load_kv_bytes(req)
        if ok and req.state_load_spec is not None:
            ok = self._load_state_bytes(req)      # same task, same thread
    except Exception:
        ok = False
    self._finish_load(req, ok)                    # ONE completion, every path

Dense's _do_load_req is split into _load_kv_bytes + _finish_load to give it that seam (behaviour-preserving for dense). A state-only load (KV resident, state in the tier) 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 (state_offload_invariant_violations in checkpoint_funnel).

Deleted (each greps to zero): _JointPark · metadata.state_loads · STATE_LOAD_DISPOSITION_CHANNEL · take_state_load_survived · _publish_state_loads · _settle_state_load · _abandon_state_load · _orphan_load_slots · reconcile_orphan_load_slots · take_state_loads · _tier_can_serve · staging_lanes. The state tier loses its own load executor.

Rules the fusion creates, each tested:

  • A fused ok=False may mean the KV leg, so it does not retract the state hash; only a real state get miss does. Every rank reports a verdict (a neutral one if it never ran the leg), or the TP quorum is never reached.
  • An abandon is not a miss.
  • A live request's slot is never yanked; reclaim frees only orphans.

Unchanged: _joint_kv_boundary, _chain_to, _gated_hit, _commit_joint_boundary, _no_joint, allocate, _state_leg_secured, disown_claimed_prefix. Joint resume at any checkpoint rung, and claiming HBM-resident KV so only the delta transfers, both stand. The store leg is out of scope.

Also fixed

  • K3 dropped the dense PAGE completion channels. KimiK3OffloadScheduler.connector_completion ended in a bare return False, so with early block release on, dense.page.source_safe / dense.page.store were logged as unhandled and their deferred blocks were only freed by the stall timeout. Pre-existing on main; it now delegates to the base.
  • The two-offload-sub refusal had holes. _offload_subconfig compared raw strings while sub-connectors resolve through the registry (strip, casefold, aliases), and lmcache_mp was missing. Both sets are now compared against the canonical name.
  • PP > 1 refuses at startup instead of warning (the offload is fully inert under PP).
  • Store failure legs (settle_state_store(ok=False)) had no test anywhere in the repo; now they do.

From review of this PR (3ea7ab115):

  • State-only resumes under kv_connector: multi. Fusing moved the park decision to the KV load owner, and a state-only load has none. That is the production agentic shape ([mooncake producer, lmcache_offload]), which main parked unconditionally. The composite now asks the tier sub when the engine secured a state leg and no sub owns KV.
  • The build-time pin check (fix(offload): memoize the external-tier lookup per HBM frontier (native + plugin) #2305) dropped state-only loads after the request parked when HBM held more of the prefix than the KV tier. A spec that reads nothing from the tier now skips the check.
  • A stale offload_load_cancelled could suppress a state leg on re-arbitration; each arbitration now starts clean.
  • State verdicts are keyed by load generation, not request id (a re-admitted request's miss was being tombstoned).
  • orphan() matches the slot, not just the request id; reclaim counts as an abandon.

Behaviour changes to sign off on

  1. PP > 1 raises at startup.
  2. A request torn down mid-load is orphaned, not abandoned: its index entry stays outstanding until the report lands or reclaim gives up.
  3. A new seam in dense (_load_kv_bytes / _finish_load). The only observable reordering is [OFFLOAD-LOAD-PROF] logging before the unpin.

Known limitations

  • After a process restart the engine-side index is empty, so K3 cannot read back KV LMCache still holds (the index is not persisted, because LocalDiskBackend never scans its directory).
  • Under multi, if a P/D sub wins the KV leg while a state leg is also pending, the unpark follows that sub's report alone. This gap is also on main; it is noted at the call site in scheduler.py.

Testing

CPU suite (--ignore=tests/plugin, two modules with missing deps excluded): 6780+ passed; the failure set is identical to origin/main (28 pre-existing, test_pd_pp / test_mooncake_rail_address). tests/plugin: failure set identical to the parent commit. Every new regression test was checked to fail against the code before its fix.

Needs GPU validation before merge (owned by the tester): Kimi-K3 TP8/DCP8 on the AtoMesh multi config, comparing against main:

  • state-tier load count and TTFT (the multi state-only path must be live, not recomputing);
  • two-pass GSM8K delta_acc inside the ±0.024 temp=0 band, with state_offload_invariant_violations == 0.

🤖 Generated with Claude Code

@zejunchen-zejun zejunchen-zejun changed the title [NOT READY][NO REVIEW NEEDED] refactor K3 LMCache [NOT READY] refactor K3 LMCache Sep 3, 2026
@github-actions

github-actions Bot commented Sep 3, 2026

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 2132 --add-label <label>

zejunchen-zejun added a commit that referenced this pull request Sep 3, 2026
…arration

Responding to review of #2132.

**F1 is rejected on the facts, and its remedy adopted anyway.** The claim was
that `stats()` has zero callers repo-wide, so the invariant audit never runs.
It does run: `PagedStateCheckpointCoordinator.checkpoint_fates` reaches it as
`getattr(self.offload, "stats", None)` (page_unit_checkpoint.py:1025), which
`BlockManager.state_checkpoint_fates` folds into `checkpoint_funnel`. Verified
by construction, not by reading: corrupting the accounting makes
`state_offload_invariant_violations` go to 1 in the funnel dict.

But a competent reviewer grepped and concluded the opposite, which is itself the
finding -- a chain whose middle link is a `getattr` cannot be seen, so it now
has a test that asserts a violation reaches the funnel, and the docstring names
the prefixed key rather than the bare one.

**F2 was asked as a question and the answer is that the path is covered.** A
state load armed and then refused after the arm settles through
`Scheduler._drop_state_load` -> `abandon_load` before `build_connector_meta`
clears its bookkeeping. Now tested: the dispatch is settled, the invariant
holds, and -- the part worth pinning -- an abandon does NOT retract the hash,
because nothing was attempted and the bytes are still there.

**F3 accepted, narrowly.** Two blocks in `state_offload.py` argued against the
previous design rather than describing this one, which after a squash-merge a
reader cannot check. Both had one load-bearing sentence buried in the history:
why `stores_refused` is separate from `stores_failed`, and why capability is
derived from config rather than the connector's name. Kept those, dropped the
narration.

4259 passed.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@zejunchen-zejun zejunchen-zejun changed the title [NOT READY] refactor K3 LMCache [Kimi-K3][LMCache] Fuse the state load leg, and fence the KV save against the forward Sep 11, 2026
zejunchen-zejun added a commit that referenced this pull request Sep 14, 2026
…arration

Responding to review of #2132.

**F1 is rejected on the facts, and its remedy adopted anyway.** The claim was
that `stats()` has zero callers repo-wide, so the invariant audit never runs.
It does run: `PagedStateCheckpointCoordinator.checkpoint_fates` reaches it as
`getattr(self.offload, "stats", None)` (page_unit_checkpoint.py:1025), which
`BlockManager.state_checkpoint_fates` folds into `checkpoint_funnel`. Verified
by construction, not by reading: corrupting the accounting makes
`state_offload_invariant_violations` go to 1 in the funnel dict.

But a competent reviewer grepped and concluded the opposite, which is itself the
finding -- a chain whose middle link is a `getattr` cannot be seen, so it now
has a test that asserts a violation reaches the funnel, and the docstring names
the prefixed key rather than the bare one.

**F2 was asked as a question and the answer is that the path is covered.** A
state load armed and then refused after the arm settles through
`Scheduler._drop_state_load` -> `abandon_load` before `build_connector_meta`
clears its bookkeeping. Now tested: the dispatch is settled, the invariant
holds, and -- the part worth pinning -- an abandon does NOT retract the hash,
because nothing was attempted and the bytes are still there.

**F3 accepted, narrowly.** Two blocks in `state_offload.py` argued against the
previous design rather than describing this one, which after a squash-merge a
reader cannot check. Both had one load-bearing sentence buried in the history:
why `stores_refused` is separate from `stores_failed`, and why capability is
derived from config rather than the connector's name. Kept those, dropped the
narration.

4259 passed.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@zejunchen-zejun
zejunchen-zejun force-pushed the zejun/k3_lmcache_refactor branch from 0b54083 to 94caab3 Compare September 14, 2026 01:42
@zejunchen-zejun
zejunchen-zejun marked this pull request as ready for review September 14, 2026 03:20
Copilot AI lite review requested due to automatic review settings September 14, 2026 03:20

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🟡 Changes recommended

Unresolved critical routing and completion issues, plus a state-index invariant issue, remain.

Get a fresh assessment by requesting another Copilot review.

Pull request overview

This PR fences dense KV saves against forward writes and fuses Kimi-K3 recurrent-state loads into the KV load lifecycle.

Changes:

  • Adds producer-event synchronization for dense KV gathers.
  • Centralizes state-load ownership, slot reclamation, and completion accounting.
  • Expands connector, scheduler, lifecycle, and regression-test coverage.
File summaries
File Reviewed changes and final findings
tests/test_state_offload_index.py Tests lifecycle invariants and reclamation.
tests/test_state_checkpoint.py Updates state checkpoint coverage.
tests/test_page_unit_checkpoint.py Tests page-unit checkpoint behavior.
tests/test_multi_connector.py Tests composite connector routing.
tests/test_lmcache_offload_connector.py.names Updates test-name metadata.
tests/test_lmcache_offload_connector.py Tests fused load and state-tier behavior.
tests/test_kv_drain_liveness.py Tests idle transfer draining and completion liveness.
tests/test_block_manager.py Tests state-slot ownership changes.
atom/model_engine/state_offload.py Centralizes state-load lifecycle accounting. Moderate (2 votes): the invariant audit compares cardinalities rather than key sets.
atom/model_engine/scheduler.py Integrates fused-load completion and reclamation handling.
atom/model_engine/pp_engine_core.py Propagates connector completion events.
atom/model_engine/engine_core.py Updates idle KV-work draining.
atom/model_engine/block_manager.py Transfers orphaned state-slot ownership.
atom/kv_transfer/offload/metadata.py Adds state-load and producer-event metadata.
atom/kv_transfer/offload/hybrid/kimi_k3/state_tier.py Implements synchronous state loads and reports.
atom/kv_transfer/offload/hybrid/kimi_k3/state_object.py Updates state transfer integration.
atom/kv_transfer/offload/hybrid/kimi_k3/connector.py Fuses K3 KV/state loads. Critical (2 votes): hash-only verdict keys can be tombstoned after the first quorum, dropping later repeated-load verdicts. Critical (2 votes): load verdicts returning True can be misclassified as terminal saves.
atom/kv_transfer/offload/dense/connector.py Adds KV producer fencing and load refactoring seams.
atom/kv_transfer/offload/connector.py Updates scheduler-shell forwarding.
atom/kv_transfer/offload/_offload_common.py Revises shared completion handling.
atom/kv_transfer/disaggregation/types.py Adjusts connector metadata types.
atom/kv_transfer/disaggregation/multi/multi_connector.py Updates composite connector routing. Critical (1 vote): the dense/K3 routing can resume the forward before K3’s state-only task completes.
Review details
  • Files reviewed: 22/22 changed files
  • Comments generated: 4
  • Review effort level: Lite

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread atom/kv_transfer/disaggregation/multi/multi_connector.py
Comment thread atom/kv_transfer/offload/hybrid/kimi_k3/connector.py Outdated
Comment thread atom/kv_transfer/offload/hybrid/kimi_k3/connector.py Outdated
Comment thread atom/model_engine/state_offload.py
@zejunchen-zejun
zejunchen-zejun marked this pull request as draft September 14, 2026 13:44
@valarLip

Copy link
Copy Markdown
Collaborator

Review — Kimi-K3 state-offload refactor

Reviewed at 8ed12fcbf against merge base 09cabf1d2. Findings marked [verified] I read out or measured myself; [reported] ones match the code shape but I did not independently confirm. Measurement scripts noted inline.

The direction of this PR is right, and worth saying before the findings: it is a consolidation, not a bolt-on.

state_tier.py            491 ->  249   (-242)
block_manager.py        2669 -> 2498   (-171)
kimi_k3/connector.py     912 ->  800   (-112)
scheduler.py            3804 -> 3737   (-67)
state_offload.py         419 ->  612   (+193)   <- the one owner that gains

The state-load lifecycle used to be spread across block_manager, scheduler and state_tier; it now has a single owner. I went through StateOffloadIndex's 14 public methods and they are cohesive — index (note_stored/forget/could_serve), lifecycle (request_load/orphan/complete_load/fail_load/abandon_load/reclaim), introspection. Deleting _JointPark and _orphan_load_slots is the right call. Everything below is "the consolidation didn't finish", not "this is the wrong shape".


Blockers

1. The state-load verdict is emitted only on ranks whose KV leg succeeded, so the quorum for that key is never reached — permanently [verified]

# hybrid/kimi_k3/connector.py:310
ok = self._load_kv_bytes(req)
if ok and req.state_load_spec is not None:
    ok = self._load_state_bytes(req)      # the only thing that writes _hash_verdicts

_load_kv_bytes's verdict is rank-local: dense/connector.py:287 is loaded = bool(ret_mask[hbm:lmc].all().item()), out of self._engine.retrieve(...) — this rank's own LMCache instance and its own LRU. Rank 3 missing a chunk short-circuits and records no (h, rid); ranks 0-2, 4-7 do.

Then:

# disaggregation/aggregator.py:79-83
for key, reports in list(self._reports.items()):
    if len(reports) < self._world_size:
        continue                           # not deleted, no TTL, no eviction

reset() is never called on the serving path. So the key sits at world_size - 1 reports for the process lifetime:

  • one permanently leaked _reports entry per occurrence, and
  • the hash is never retracted, so could_serve keeps voting for bytes LMCache has dropped — every later request over that prefix parks, misses, and recomputes. Forever.

This needs no unusual configuration; it happens in a fully homogeneous TP build the first time one rank's LRU differs from another's, which is the normal state of independent per-rank caches.

A rank that does not run the state leg must still emit a neutral report on the key (or the key must be self-clearing).

Related hardening, latent rather than reachable today: the tier-None load path at connector.py:391 returns before the verdict loop with no equivalent of the store leg's _store_failed_no_tier synthesis at :349-357. The two legs are asymmetric.

2. reclaim ages entries from dispatch time, not orphan time — so the safety window it exists to enforce can already be zero [verified]

# state_offload.py:254   (inside request_load — the only assignment to `at`)
prefix_hash=int(h), slot=int(slot), at=monotonic()

# state_offload.py:259-271  orphan() sets orphaned=True and does NOT refresh `at`

# state_offload.py:347-351
deadline = monotonic() - timeout_s
... if pending.orphaned and pending.at <= deadline

I grepped the file: at is stamped once, at dispatch, and never refreshed. So the guard window is timeout_s minus the load's in-flight age. A load that has already been outstanding longer than timeout_s — which is precisely the hung/slow-worker case reclaim exists for — becomes eligible the instant it is orphaned, and the next _maybe_reconcile tick frees the slot while the worker may still be scattering into it. The next request gets a buffer someone else is filling, with has_initial_state already true over it: silent wrong output.

The method's own docstring states the requirement it violates:

state_offload.py:341-343 — "timeout_s <= 0 disables reclamation … and must not be tighter than that window or this becomes the hazard it exists to prevent."

The deleted _orphan_load_slots stamped at orphan time. One-line fix: refresh at in orphan().


Structure

3. One fact now has two owners with opposite dispositions [verified]

state_tier_capability() (engine, state_offload.py:504) decides tier availability from config alone, at BlockManager.__init__, and its docstring enumerates the six conditions the worker refuses on — pp_size > 1 among them. Its stated policy:

"A false negative (tier off) is much cheaper than a false positive"

_build_state_tier() (worker, connector.py:148-166) decides the same condition and now raise ValueErrors. Its own comment admits the duplication:

# connector.py:152
# The engine agrees independently (`state_tier_capability`), so both legs
# of a K3 request are then declined and the offload does nothing at all

So the engine has already turned the tier off, gracefully, before the worker exists — and the worker then raises fatally over a case that is already handled.

The justification for making it fatal is contradicted one screen away:

# connector.py:667-668
if not getattr(seq, "has_per_req_cache", False):
    return super()._decide_load_after_alloc(seq, ls)

Sequences without a per-request cache delegate to the dense path, which is unaffected by the state tier being off — so "the offload would be inert" is a property of hybrid sequences, not of the connector. (I did not check what fraction of a K3 deployment that is; I am only reporting the delegation.)

And raising out of register_kv_caches during a worker RPC is the documented server-wedge shape: the server hangs between load model runner success and ready, and wait_server_ready.sh does not classify it — so this carefully written message is exactly what the operator never sees. Under multi, MultiConnector.register_kv_caches fans to subs with no per-sub guard, so a co-configured P/D producer goes down with it. [reported]

Let the engine-side owner decide this alone.

4. state_offload.py holds two unrelated jobs [verified]

:36-416    StateOffloadIndex            runtime lifecycle      <- what this module is for
:418-612   StateTierCapability
           _offload_subconfig           connector-config       <- ~195 lines
           _SubConnectorView              introspection
           state_tier_capability
           state_tier_chunk_tokens

The second half parses kv_transfer_config dicts, walks multi's sub-connector list, and builds a _SubConnectorView.__getattr__ config proxy — that is kv_transfer's business, living in a model_engine/ module whose stated job is the state-load index.

And finding 3 grows on exactly this seam: the two halves of this file answer questions belonging to two different layers, and one of those questions is also answered in the other layer. Moving :418-612 to something like atom/kv_transfer/offload/capability.py surfaces the duplicate ownership on its own.

5. _build_state_tier is 134 lines with six refusals, five graceful and one fatal [verified]

:149  pp_size > 1                 ->  raise ValueError      <- the only one
:169  backend is None             ->  warning + return
:180  not layout_id               ->  warning + return
:193  ...                         ->  warning + return
:204  no page_unit_views          ->  warning + return
:218  image_bytes != entry_bytes  ->  warning + return

Six judgments of the same kind in one function, five degrading and one killing the process, with no principle separating them. The function also does the aiter import, geometry probing and codec construction. The six refusals want to be one predicate — and that predicate already exists, as state_tier_capability.


Performance

6. check_invariant costs 31 ms near the hash cap, and stats() runs it on every metrics read [verified — measured]

Chain: stats() (:378) → audit_invariant() (:171) → check_invariant() (:148) → set(self.hashes) != set(self._hash_lru), with _hash_cap = 1 << 20 (:89).

n=    1024  set-compare    0.017 ms   len-compare 0.000085 ms
n=   16384  set-compare    0.302 ms   len-compare 0.000054 ms
n=  262144  set-compare    6.334 ms   len-compare 0.000134 ms
n= 1048576  set-compare   31.394 ms   len-compare 0.000201 ms

31.4 ms on the engine loop thread, on the periodic push clock and on every client cache-stats call — a stall that grows with uptime and is invisible in any GPU trace.

Meanwhile stats()'s docstring says:

:375 — "it costs two integer comparisons per metrics read"

On the same path, check_invariant's own comment at :143-147 correctly argues why the key-set compare is necessary and why a length compare is blind to the divergence it catches. So the reasoning about why it must be expensive is right, and the cost was then described as if the cheap option had been taken.

The fix is not to drop the set compare — the argument at :143-147 holds. It is to stop running it on every metrics read: keep the cheap arm (dispatched != settled + outstanding) every time, and decimate the set compare or run it only when the cheap arm looks wrong.

7. pending_loads rebuilds a dict, and is called inside the failure loop [verified — measured]

:127-129 is a property that rebuilds {req_id: prefix_hash} on every read, and scheduler.py:3163 calls it inside the failed_loading loop:

2048 outstanding, per rebuild:                                54.1 us
x 2048 failed requests in one step (LMCache backend down):     0.11 s   on the scheduler thread

Quadratic exactly when the engine most needs to drain. A hash_of(req_id) accessor keeps _outstanding private and removes it.

8. metadata.py grew 25 lines for two write-only fields, justified by a check that does not exist [verified]

StateLoadSpec.boundary_tokens and chunk_tokens have no reader — grep -rn state_load_spec atom/ returns four hits, and the only consumer, _load_state_bytes (connector.py:319-333), uses boundary_hash and destination_slot and validates nothing. The docstring justifies chunk_tokens with "the worker validates the transfer against it."

So the hazard it names — engine and worker disagreeing about the boundary, "silent wrong output rather than an error anyone would see" — is exactly as unguarded after this PR as before, while the comment tells the next reviewer it is covered. Either write the check in _load_state_bytes, or drop both fields and the paragraph.


Lint and tests

I ran these; they were not covered in the review that preceded this comment.

  • black --check — clean (75 files).
  • ruff check on the 14 changed atom/ files — 13 findings, all in engine_core.py at lines 261 / 323 / 1104-1119 / 1179 / 1367-1377. This PR's only hunk in that file is @@ -553,16 +553,15 @@, so none land in the diff context and none are introduced here. Lint is clean for this PR.
  • Test suite: 9 failed, 6170 passed on the branch; all 8 touched files pass standalone (663 passed). The 9 failures are in test_deepseek_v4_wo_a_dequant.py and test_heavy_ci_gate.py, untouched here — they read as pre-existing but were not confirmed against the merge base. [reported]

Two test issues worth fixing regardless [reported]:

  • tests/test_lmcache_offload_connector.py.names is a committed 194-line scratch artifact — git ls-files tests/ | grep -v '\.py$' returns exactly this one file. Not collected, not linted, not kept in sync.
  • test_a_single_stage_still_builds_the_tier_path sets c._state_tier = None and then asserts it is None — replacing _build_state_tier's body with return keeps it green. test_a_load_lands_while_a_store_is_still_stuck asserts codec.loaded.is_set() after a synchronous tier.load_state(...) that sets it, so the regression its name describes would still pass.

Also worth a look (not independently verified)

  • scheduler.py:3151 — take_missed_state_hashes() is drained once per step and matched only against that step's failed_loading; anything unmatched is discarded. Per rank the verdict is written before the load completion but drained after it, and the window between the two writes contains _lookup_unpin, so a get_finished() landing there publishes the verdict on step N and the failure on step N+1. Step N drops the miss; step N+1 sees an empty missed, so forget() never runs. Because verdict_step <= failure_step always, the skew is one-directional and never self-corrects. Note this is purely an intra-rank write-before/drain-after inversion across two threads — both channels use the same drain(), the same world_size gate, and one aggregate() call, so there is no quorum asymmetry to go looking for. missed needs to be retained across steps and matched per (hash, req_id).
  • block_manager.py:2417 — a preempt/re-admit inside the orphan window leaks a state slot permanently: the index withholds the slot recorded at request_load time while deallocate withholds seq.state_slots[0], and after re-admission those are two different slots. The deleted _orphan_load_slots was a list per req_id precisely because seq.id is per-request, not per-admission.
  • connector.py:536 — a state load already counted as dispatched is dropped when the build-time KV refusal fires (chunked_scheduler.py:333-336 clears the pending load, so nothing reaches meta.requests, and the post-pass clears _state_load_seqs). The index entry is live and not orphaned, so reclaim deliberately skips it, and the only abandon_load call site is on the admission-time path. The request hangs in WAITING_FOR_REMOTE_KVS with no recovery.
  • connector.py:520 — the post-pass matches _state_load_seqs against meta.requests by req_id only, never checking req.load_spec is not None, so the StateLoadSpec can be attached to a save-only request that start_load_kv never dispatches a load for — same unsettleable hang by a second route.
  • multi_connector.py:567 — under kv_connector: multi, a K3 state-only load can never succeed: every sub returns (0, False), winner is None, should_park_for_load_after_alloc returns False, and _drop_state_load disowns it into a full prompt recompute. The state tier's headline capability is silently dead under the composite shell while it works under the single-connector shell. The K3 sub's own hook would have returned True, so this is purely multi's routing.
  • state_tier.py:114 — deleting the lmc-state-load lane did not remove the queue, it moved the state leg onto _load_executor, a max_workers=1 pool shared with every KV retrieve. On a burst of N hybrid admissions, request N's state H2D waits behind N-1 full KV retrieves and each joint load's TTFT goes from max(kv, state) to kv + state. _LOAD_WAIT_WARN_MS — the only metric that would have shown it — was deleted in the same hunk. Fusing the completion does not require fusing the execution.
  • dense/connector.py:195 — the new producer fence's torch.cuda.Event() / record() is unguarded inside start_load_kv (one event per saving request per step), and a raise there escapes through ModelRunner.process_kvconnector_output, which has no handler — dropping every KV load and save dispatched that step. Both other in-tree fences wrap it and hoist one event per step. Secondary: the event defaults to blocking=False, so producer_event.synchronize() busy-spins a core on the save worker for the whole fenced forward.
  • connector.py:531 — destination_slot=int(seq.state_slot) is snapshotted at metadata-build time and never validated non-negative; seq.state_slot returns -1 once state_slots is empty, and codec.get(h, -1) negative-indexes into the state pool. Today this is covered only incidentally, by deallocate clearing offload_joint.load_hash one line later. An assertion at the codec boundary would make it non-incidental.
  • Doc drift: staging.py:43-59 — cited by the new fence comment as its authority — still names save_ready_event, a symbol that does not exist on this leg; STATE_LOAD_VERDICT_TAG's comment says ("state_load", hash) while the code emits the 3-tuple (tag, hash, rid); K3 close()'s docstring still says "executors" plural though only the store executor remains; offload/README.md still describes state_tier's deleted load executor.

Findings 1 and 2 are the two I would hold the merge on: the first is unrecoverable and needs no unusual configuration, the second violates a constraint its own docstring writes down. Both are small fixes. Finding 3 is the one I would fix next, and moving state_offload.py:418-612 out (finding 4) makes it fall out naturally rather than having to be argued.

zejunchen-zejun added a commit that referenced this pull request Sep 15, 2026
…arration

Responding to review of #2132.

**F1 is rejected on the facts, and its remedy adopted anyway.** The claim was
that `stats()` has zero callers repo-wide, so the invariant audit never runs.
It does run: `PagedStateCheckpointCoordinator.checkpoint_fates` reaches it as
`getattr(self.offload, "stats", None)` (page_unit_checkpoint.py:1025), which
`BlockManager.state_checkpoint_fates` folds into `checkpoint_funnel`. Verified
by construction, not by reading: corrupting the accounting makes
`state_offload_invariant_violations` go to 1 in the funnel dict.

But a competent reviewer grepped and concluded the opposite, which is itself the
finding -- a chain whose middle link is a `getattr` cannot be seen, so it now
has a test that asserts a violation reaches the funnel, and the docstring names
the prefixed key rather than the bare one.

**F2 was asked as a question and the answer is that the path is covered.** A
state load armed and then refused after the arm settles through
`Scheduler._drop_state_load` -> `abandon_load` before `build_connector_meta`
clears its bookkeeping. Now tested: the dispatch is settled, the invariant
holds, and -- the part worth pinning -- an abandon does NOT retract the hash,
because nothing was attempted and the bytes are still there.

**F3 accepted, narrowly.** Two blocks in `state_offload.py` argued against the
previous design rather than describing this one, which after a squash-merge a
reader cannot check. Both had one load-bearing sentence buried in the history:
why `stores_refused` is separate from `stores_failed`, and why capability is
derived from config rather than the connector's name. Kept those, dropped the
narration.

4259 passed.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@zejunchen-zejun
zejunchen-zejun force-pushed the zejun/k3_lmcache_refactor branch from 44ca569 to 62ae93c Compare September 15, 2026 10:12
@zejunchen-zejun

Copy link
Copy Markdown
Collaborator Author

Reply to the review on PR #2132

Thanks — this is the most useful review this branch has had. Every finding I
checked reproduced, including the measurements, and one of them ("also worth a
look") turned out to be load-bearing enough that chasing it caught a crash I had
introduced two commits earlier. Details below.

Four commits on top of the reviewed 8ed12fcbf:

commit covers
09af0e092 findings 1, 2, 6, 7, 8, test hygiene
cfe9ff20c dense/connector.py:195 fence; all four doc-drift items
6288d3853 self-reported: backs out a guard 09af0e092 added
44ca56923 scheduler.py:3151; connector.py:520 hardening

Blockers

Finding 1 — verdict quorum unreachable

Confirmed exactly as described. _load_kv_bytes returns
ret_mask[hbm:lmc].all() over this rank's own LMCache instance, so a rank whose
KV leg missed short-circuited before _load_state_bytes and recorded no
(hash, req_id) while the others did. _TPCompletionGroup.drain skips an
incomplete key with a bare continue — no TTL, no eviction, reset unused on
the serving path — so the key stayed at world_size - 1 for the process
lifetime.

Fix (09af0e092): a rank that skips the state leg now reports neutrally
(StateOffloadTier.note_load_unrun) instead of not reporting at all. Neutral
and not a miss, on your own framing: only an empty get is evidence LMCache
dropped the bytes, and a rank that never asked holds none. The quorum stays
failure-dominant, so a genuine miss on any rank still retracts. Two regression
tests, both verified to fail against the previous behaviour.

On the latent asymmetry you flagged alongside it — the tier-None load path
returning before the verdict loop with no analogue of _store_failed_no_tier:
_note_state_leg_unrun is deliberately silent with no tier, because a rank with
no tier reports on no verdict key at all, so there is no partial quorum to
complete. That leaves it reachable only if ranks disagree about whether a tier
built, which the engine-side gate already decides globally. Left, and said so in
the code, rather than synthesising a report for a key that does not exist.

Finding 2 — reclaim window measured from the wrong instant

Confirmed. at is assigned only in request_load; orphan() set orphaned
and left at alone. So the window was timeout_s minus the load's in-flight
age, and a load already outstanding longer than timeout_s — the hung-worker
case — was eligible the moment it was orphaned.

Fix (09af0e092): orphan() restamps, as the deleted _orphan_load_slots did.
One line, one regression test.


Performance

Finding 6 — check_invariant on every metrics read

Confirmed and reproduced — I measure 59 ms at n = 1 << 20, worse than your 31.

Provenance, since it matters: this one is mine and was one commit old. The
key-set compare landed in 8ed12fcbf answering an earlier review comment, and I
did not notice stats() is both the periodic metrics path and every client
cache-stats call. The docstring still advertising "two integer comparisons" was
true before that commit and false after it.

Fix (09af0e092) is the one you proposed. The argument at check_invariant
holds, so the set compare is kept and decimated: the counter identity — what
actually catches a parked request or a leaked slot — runs every read; the set
compare runs one read in 256. Divergence is a bug, not a race: it does not heal,
so decimation costs only report latency. 59.0 ms → 0.23 ms amortised.
check_invariant keeps deep=True by default so tests still check the whole
contract. Docstring corrected.

Finding 7 — pending_loads rebuilt per failure

Confirmed. Added hash_of(req_id) and switched the call site. pending_loads
stays for introspection, with a docstring saying which is which.

Finding 8 — two write-only fields, and a docstring covering for them

Confirmed: four hits for state_load_spec, and the only consumer uses
boundary_hash and destination_slot.

Took the "write the check" branch rather than dropping the fields — the hazard
is real and having the engine's numbers travel with the spec is the right shape;
what was missing was the check. _load_state_bytes now refuses a boundary off
the chunk grid, and refuses rather than clamps, because there is no safe
reinterpretation of a boundary the two sides disagree about. Folded in your
separate note about destination_slot never being validated non-negative.


"Also worth a look" — triaged, not waved past

You marked these unverified. Seven of the nine are in code this branch
introduces (producer_event, state_load_spec, take_missed_state_hashes are
all absent at the merge base; _orphan_load_slots went from 15 references to
0), so I owed them a verdict rather than a shrug. Reachability was checked
before deciding.

item verdict action
dense/connector.py:195 fence unguarded reachable fixed, cfe9ff20c
scheduler.py:3151 missed-hash skew reachable, reproduced fixed, 44ca56923
connector.py:520 save-only spec not reachable hardened anyway, 44ca56923
connector.py:536 dropped state load not reachable as described see below — it caught something else
connector.py:531 negative slot reachable fixed, 09af0e092
doc drift (4 items) n/a fixed, cfe9ff20c
block_manager.py:2417 re-admit slot leak accepted not fixed — see Deferred
state_tier.py:114 shared load executor accepted not fixed — see Deferred
multi_connector.py:567 state-only dead under multi accepted not fixed — see Deferred

dense/connector.py:195 had three problems in one hunk. A raise escaped —
ModelRunner.process_kvconnector_output has no handler, so a failure building
the event dropped every KV load and save dispatched that step, a far larger
blast radius than the corruption the fence prevents. It now falls back to
torch.cuda.synchronize(): strictly stronger ordering, so degrading to it
cannot produce the torn gather. One event per request bought no ordering over
one per step, so it is hoisted, as both other in-tree fences do. And the event
defaulted to blocking=False while its consumer synchronize()s it for the
whole fenced forward — a spun core for that window, per save; now blocking=True.
Two regression tests, the second verified to fail by the raise escaping.

scheduler.py:3151 is real and I reproduced it end-to-end. The worker writes
the miss verdict in load_state, then _finish_load runs _lookup_unpin — a
real LMCache call, under a different lock — before recording the failure, while
get_finished drains the two in the opposite order. A step landing in that
window publishes the verdict on step N and the failure on step N+1; step N
discarded the unmatched miss, step N+1 saw an empty set. verdict_step <= failure_step always, so it never self-corrects. Same end state as finding 1.
Fix: a miss is evidence about the hash and forget needs no request id, so
the retraction no longer goes through the per-request report at all. missing=
stays on fail_load, which uses it to decide slot reuse — that one is genuinely
per-request.

connector.py:520 is not reachable: the shape needs the sid popped from
_reqs_need_recv during the build, whose only in-build pop is the
_decide_load_after_alloc-False branch — which already ran one call earlier in
should_park_for_load_after_alloc, whose False routes to _drop_state_load,
which clears load_hash before the build. Hardened anyway, because relying on
that made the loop's correctness depend on a field a different owner writes; it
now tests req.load_spec is not None itself.


Something I have to report against myself

09af0e092 — the same commit that fixed findings 1, 2, 6, 7 and 8 — also
restored a guard from the merge base that I had deleted earlier on this branch
(ba5df6304). I restored it verbatim without checking that it still composes.
It does not, in two independent ways, and 6288d3853 backs it out.

It called BlockManager.cancel_state_load — a method this branch deleted.
hasattr(BlockManager, "cancel_state_load") is False, and the only remaining
reference in the tree was the restored call itself. So it raised AttributeError
instead of guarding, and that escapes _process_engine_step_inner,
_process_engine_step (only a finally) and busy_loop (no except): a dead
engine core.

And its predicate, load_hash != -1 and not boundary_tokens, is exactly the
fused state-only shape. At the merge base that shape returned
(False, "per_req_cache_state_boundary"), so needs_remote_load was False and
the guard was never reached for it. The fused path returns
(True, "state_only_load") — so the guard matches every state-only load, and
repairing the AttributeError would have silently dropped the capability this
branch exists to add.

Why it is gone rather than repaired: what the guard was for — two transfers
reporting once — cannot arise while one connector owns both legs. A state-only
load moves no KV (_load_kv_bytes treats lmc <= hbm as the no-op success it
is) and reports once. It arises under kv_connector: multi, where a PD sub wins
the KV leg and _update_waiting_for_remote_kv unparks on that sub's report
alone. That unpark takes no state-leg predicate at the merge base either, so it
is a pre-existing scheduler-wide gap, not something to patch K3-shaped here. The
reasoning is left at the call site so the next reader does not restore it again.

Nothing dynamic could have caught this: test_scheduler.py does not import in
this environment. So 6288d3853 adds
tests/test_scheduler_block_manager_contract.py, a static check that every
BlockManager method the scheduler calls exists — verified to fail against the
deleted-method call.


Deferred, with reasons

Findings 3, 4, 5. I agree with all three, including the part I expected to
argue with: making pp_size > 1 fatal in the worker is redundant against an
engine-side gate that has already declined gracefully, and raising out of
register_kv_caches wedges the server between "load model runner success" and
"ready", so the carefully written message is exactly the one the operator never
sees. Not defending it.

Not doing it here. Moving state_offload.py:418-612 into kv_transfer,
collapsing _build_state_tier's six refusals onto the predicate that already
exists, and deleting the duplicate owner is a structural change with its own
blast radius, and this PR's whole argument is that it is narrower than the one
it replaces. Happy to open it immediately if you would rather review them
together.

The multi unpark gap. _update_waiting_for_remote_kv unparks on one
connector's report with no state-leg predicate, and MultiConnector passes
finished_recving through ungated. I diffed both against the merge base:
identical. Your multi_connector.py:567 note is the same root cause by a
different branch (no sub wins, the state-only load is disowned into a full
recompute — silently dead rather than racing); the other branch is a PD sub
winning the KV leg and unparking while the state H2D is still writing. One fix
in the scheduler, not two K3-shaped patches.

I would fold one more thing into that change: _offload_subconfig compares
kv_connector as a raw string while sub-connectors are built through
KVConnectorFactory.canonical_name, which strips, casefolds and resolves
aliases. So the "two offload sub-connectors are unrepresentable" guarantee that
multi_connector.py leans on in three comments does not hold:

['lmcache_offload', 'lmcache_offload']     -> raises
['lmcache_offload', 'LMCacheConnectorV1']  -> passes    # alias
['lmcache_mp',      'lmcache_offload']     -> passes    # not in _STATE_TIER_BACKENDS

state_tier.py:114. Accepted as stated, and worth naming plainly: fusing
the completion did not require fusing the execution, and putting the state
leg on the shared max_workers=1 load pool turns a joint load's TTFT from
max(kv, state) into kv + state under a burst. _LOAD_WAIT_WARN_MS went in
the same hunk, so nothing reports it. I have not measured it and I am not going
to claim it is small. It wants its own change with a number attached.

block_manager.py:2417. Accepted. Restoring per-admission identity means
the index holds a list per req_id again rather than a single pending — a data
structure change, not a patch.


Verification

  • black clean; ruff clean on every changed file.
  • Full suite 143 failed / 5167 passed, against 143 / 5157 at the merge
    base 09cabf1d2 — identical failure set, +10 new tests. That also settles the
    question you left open: the 9 failures in test_deepseek_v4_wo_a_dequant.py
    and test_heavy_ci_gate.py are present at the merge base, so they are
    pre-existing.
  • Every regression test added here was run against the previous behaviour and
    confirmed to fail first.

Heavy accuracy CI is still skipped (draft, unapproved, no ci:* label). Given
that the accuracy fix is the headline, I would rather run ci:atom before you
approve than after — say the word and I will add the label.

zejunchen-zejun added a commit that referenced this pull request Sep 18, 2026
…arration

Responding to review of #2132.

**F1 is rejected on the facts, and its remedy adopted anyway.** The claim was
that `stats()` has zero callers repo-wide, so the invariant audit never runs.
It does run: `PagedStateCheckpointCoordinator.checkpoint_fates` reaches it as
`getattr(self.offload, "stats", None)` (page_unit_checkpoint.py:1025), which
`BlockManager.state_checkpoint_fates` folds into `checkpoint_funnel`. Verified
by construction, not by reading: corrupting the accounting makes
`state_offload_invariant_violations` go to 1 in the funnel dict.

But a competent reviewer grepped and concluded the opposite, which is itself the
finding -- a chain whose middle link is a `getattr` cannot be seen, so it now
has a test that asserts a violation reaches the funnel, and the docstring names
the prefixed key rather than the bare one.

**F2 was asked as a question and the answer is that the path is covered.** A
state load armed and then refused after the arm settles through
`Scheduler._drop_state_load` -> `abandon_load` before `build_connector_meta`
clears its bookkeeping. Now tested: the dispatch is settled, the invariant
holds, and -- the part worth pinning -- an abandon does NOT retract the hash,
because nothing was attempted and the bytes are still there.

**F3 accepted, narrowly.** Two blocks in `state_offload.py` argued against the
previous design rather than describing this one, which after a squash-merge a
reader cannot check. Both had one load-bearing sentence buried in the history:
why `stores_refused` is separate from `stores_failed`, and why capability is
derived from config rather than the connector's name. Kept those, dropped the
narration.

4259 passed.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@zejunchen-zejun
zejunchen-zejun force-pushed the zejun/k3_lmcache_refactor branch from 44b86db to bbf5531 Compare September 18, 2026 09:09
@zejunchen-zejun
zejunchen-zejun marked this pull request as ready for review September 21, 2026 02:41
Copilot AI review requested due to automatic review settings September 21, 2026 02:41

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot review overview

🔵 Needs a closer look

It changes correctness-critical offload synchronization and multi-component load lifecycle semantics, which warrants careful human review plus targeted GPU validation.

Review effort: Lite
Findings: 1 High severity · 1 Medium severity

Open (2)
Resolved since last review (2)

zejunchen-zejun added a commit that referenced this pull request Sep 21, 2026
…arration

Responding to review of #2132.

**F1 is rejected on the facts, and its remedy adopted anyway.** The claim was
that `stats()` has zero callers repo-wide, so the invariant audit never runs.
It does run: `PagedStateCheckpointCoordinator.checkpoint_fates` reaches it as
`getattr(self.offload, "stats", None)` (page_unit_checkpoint.py:1025), which
`BlockManager.state_checkpoint_fates` folds into `checkpoint_funnel`. Verified
by construction, not by reading: corrupting the accounting makes
`state_offload_invariant_violations` go to 1 in the funnel dict.

But a competent reviewer grepped and concluded the opposite, which is itself the
finding -- a chain whose middle link is a `getattr` cannot be seen, so it now
has a test that asserts a violation reaches the funnel, and the docstring names
the prefixed key rather than the bare one.

**F2 was asked as a question and the answer is that the path is covered.** A
state load armed and then refused after the arm settles through
`Scheduler._drop_state_load` -> `abandon_load` before `build_connector_meta`
clears its bookkeeping. Now tested: the dispatch is settled, the invariant
holds, and -- the part worth pinning -- an abandon does NOT retract the hash,
because nothing was attempted and the bytes are still there.

**F3 accepted, narrowly.** Two blocks in `state_offload.py` argued against the
previous design rather than describing this one, which after a squash-merge a
reader cannot check. Both had one load-bearing sentence buried in the history:
why `stores_refused` is separate from `stores_failed`, and why capability is
derived from config rather than the connector's name. Kept those, dropped the
narration.

4259 passed.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Copilot AI review requested due to automatic review settings September 21, 2026 02:50
@zejunchen-zejun
zejunchen-zejun force-pushed the zejun/k3_lmcache_refactor branch from bbf5531 to a82902b Compare September 21, 2026 02:50

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot review overview

🔵 Needs a closer look

It makes correctness-critical changes across connector/scheduler/engine boundaries (including GPU stream ordering), and should receive final human review despite strong test coverage.

Review effort: Lite
Findings: 1 High severity · 1 Medium severity · 1 Low severity

Open (3)

Comment thread atom/model_engine/scheduler.py Outdated
zejunchen-zejun added a commit that referenced this pull request Sep 22, 2026
…arration

Responding to review of #2132.

**F1 is rejected on the facts, and its remedy adopted anyway.** The claim was
that `stats()` has zero callers repo-wide, so the invariant audit never runs.
It does run: `PagedStateCheckpointCoordinator.checkpoint_fates` reaches it as
`getattr(self.offload, "stats", None)` (page_unit_checkpoint.py:1025), which
`BlockManager.state_checkpoint_fates` folds into `checkpoint_funnel`. Verified
by construction, not by reading: corrupting the accounting makes
`state_offload_invariant_violations` go to 1 in the funnel dict.

But a competent reviewer grepped and concluded the opposite, which is itself the
finding -- a chain whose middle link is a `getattr` cannot be seen, so it now
has a test that asserts a violation reaches the funnel, and the docstring names
the prefixed key rather than the bare one.

**F2 was asked as a question and the answer is that the path is covered.** A
state load armed and then refused after the arm settles through
`Scheduler._drop_state_load` -> `abandon_load` before `build_connector_meta`
clears its bookkeeping. Now tested: the dispatch is settled, the invariant
holds, and -- the part worth pinning -- an abandon does NOT retract the hash,
because nothing was attempted and the bytes are still there.

**F3 accepted, narrowly.** Two blocks in `state_offload.py` argued against the
previous design rather than describing this one, which after a squash-merge a
reader cannot check. Both had one load-bearing sentence buried in the history:
why `stores_refused` is separate from `stores_failed`, and why capability is
derived from config rather than the connector's name. Kept those, dropped the
narration.

4259 passed.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
@zejunchen-zejun
zejunchen-zejun force-pushed the zejun/k3_lmcache_refactor branch from a82902b to f70842a Compare September 22, 2026 12:03
Copilot AI review requested due to automatic review settings September 22, 2026 12:03

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot review overview

🟡 Changes recommended

A fused load where KV succeeds but the state leg fails can currently fail to force a recompute in vLLM because no load-error blocks are recorded in the state-only/no-op KV shape.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 2 High severity · 2 Medium severity

Open (4)
Resolved since last review (1)

Comment thread atom/kv_transfer/offload/hybrid/kimi_k3/connector.py
Comment thread atom/model_engine/state_offload.py
Copilot AI review requested due to automatic review settings September 22, 2026 12:22

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot review overview

🔵 Needs a closer look

It changes correctness-critical ordering and cross-thread/cross-process offload load/save lifecycles across many components, so it warrants final human review plus targeted GPU validation.

Review effort: Lite
Findings: 2 High severity · 2 Medium severity

Open (4)

@zejunchen-zejun
zejunchen-zejun force-pushed the zejun/k3_lmcache_refactor branch from ea3f48d to 08dc08e Compare September 28, 2026 04:08
Copilot AI review requested due to automatic review settings September 28, 2026 04:08

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Comment thread atom/kv_transfer/offload/hybrid/kimi_k3/connector.py
Comment thread atom/kv_transfer/offload/hybrid/kimi_k3/state_tier.py Outdated
Comment thread atom/model_engine/block_manager.py Outdated
Comment thread atom/model_engine/state_offload.py
Comment thread atom/kv_transfer/offload/chunked_scheduler.py
Copilot AI review requested due to automatic review settings September 28, 2026 06:06

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Comment on lines +553 to +555
# Hashes whose state `get` missed on some rank. Drained by the engine,
# which is the only owner of the index that advertises them.
self._state_load_missed: set[int] = set()
Comment on lines +650 to +652
# Anything left never reached the metadata (its KV leg was refused after
# the arm), so drop it rather than attach it to a later step's request.
self._state_load_seqs.clear()
Comment on lines +344 to +348
def fail_load(self, req_id, *, missing: bool = False) -> None:
"""No usable load came back.

`missing` is what separates the two failures the fused load can report.
The verdict on `failed_loading` covers BOTH legs, so it may mean the KV
@valarLip

Copy link
Copy Markdown
Collaborator

Review — fuse the Kimi-K3 state load leg

Reviewed 3ea7ab115 against origin/main (ba5149567) in a detached worktree: 26 files, +2453/−1674.

Fusing the two legs into one task and one completion is the right move, and the deletions are most of why: one report instead of two, one quorum, one place that decides a load is done. The verdict key carrying the load generation (_state_verdict_key) is exactly the kind of detail that is usually missed.

Most of what follows collapses into two roots rather than fifteen independent defects, so I have organised it that way.

Root 1 — the dispatch decision and the arm live in different owners. BlockManager._start_state_load calls request_load() and sets load_hash during allocate, while the connector decides at update_state_after_alloc / build_connector_meta whether the leg actually travels. Four findings below are all "the connector declined and nothing told the index." Having the connector call request_load() when it attaches the StateLoadSpec — or having every non-dispatch path call abandon_load + _drop_state_load — collapses all four.

Root 2 — the state-only load is expressed as a sentinel KV spec. lmcache_cached_tokens == hbm_cached_tokens has already grown five special cases across shared infrastructure. The connector's own comment at line 631 names the only obstacle: start_load_kv dispatches on load_spec is not None. Widening that to load_spec is not None or state_load_spec is not None deletes all five special cases and two of the findings below.

Separate from both, one ordering hazard sits at the heart of the new verdict protocol, so it goes first.

Provenance: [verified] = I read the deciding lines at this head. [reported] = the shape matches the code but I did not trace it end to end.


1. The neutral verdict is written before the KV leg runs, and can tombstone the key the real verdict needs [verified]

_do_load_req orders the work:

self._note_state_leg_unrun(req)      # :332  -> verdict True for key K
ok = self._load_kv_bytes(req)        # :333  -> a multi-MB LMCache retrieve
if ok and req.state_load_spec is not None:
    ok = self._load_state_bytes(req) # :335  -> the real verdict for key K

Both reports use the same _state_verdict_key(req). The comment above them says the merge in note_load_unrun is failure-dominant "so when the leg DOES run below, its real verdict overwrites this one" — which holds only while both are still in the dict. take_hash_verdicts() (state_tier.py:157-169) drains and clears it, and the comment at connector.py:518-520 states what happens next:

"The aggregator tombstones every connector-completion key it takes quorum on, and drops a report whose key is already tombstoned, so a bare hash is reportable exactly once per process, and a bare request id once per request."

So a get_finished tick landing during _load_kv_bytes publishes the neutral True on every rank, reaches quorum, and tombstones key K. The real False from _load_state_bytes is then dropped at report()'s if key in self._tombstones: return. take_missed_state_hashes stays empty, StateOffloadIndex.forget(H) never runs, and the index keeps advertising a hash whose bytes LMCache dropped — every later request over that prefix parks, misses and recomputes, permanently.

The "reportable exactly once per key" property is not incidental here; it is what makes the second report unreportable. The window is a full LMCache retrieve against a ~10 ms engine step. Whether the tick actually interleaves is yours to rule out — but the safe shape is to write the neutral verdict only on the paths that are about to skip the leg, not unconditionally ahead of it.

2. Root 1: four places where the connector declined and the index was not told

A multi sub that loses the KV arbitration leaves the state leg armed and lets the winner's completion count as a state restore [reported]. With [moriio, lmcache_offload(kimi_k3)], moriio wins get_num_new_matched_tokens, the composite fans cancel_pending_load to K3, chunked_scheduler.py:1064 sets offload_load_cancelled, and update_state_after_alloc returns at :599 without attaching a spec — but _start_state_load already did request_load() and set load_hash. should_park_for_load_after_alloc routes to _load_owner (moriio, no hook) and returns True, so the request parks; _drop_state_load is only reached on the non-parking branch. moriio's finished_loading then drives offload.complete_load(req_id): loads_completed incremented, load_hash cleared, offload_loaded=True. The forward resumes with has_initial_state true over the previous occupant's bytes. No error, no counter.

An armed leg whose request already parked is dropped with nothing left to report it [reported]. _decide_load_after_alloc returns joint_state_and_kv, so _park_for_remote_load puts the seq in WAITING_FOR_REMOTE_KVS. Later in the same step _ensure_lookup_pin (chunked_scheduler.py:439) finds the tier evicted the prefix, returns False, _clear_pending_load(sid), continue — no LMCacheReqMeta. The K3 post-pass pops nothing and _state_load_seqs.clear() discards the arm. No load task runs, so no finished_loading / failed_loading ever arrives; _update_waiting_for_remote_kv only moves seqs present in the finished/failed sets and there is no park timeout, so the request holds its blocks and its park slot forever, and reclaim() refuses to touch it because it filters on orphaned. The deleted _publish_state_loads no-carrier branch (abandon_load + failed_recving_kv_req_ids.append) was this shape's only watchdog.

deallocate orphans without abandoning, so re-admission is refused a state load [reported]. Old deallocate called abandon_load(seq.id), popping the entry. New code calls only orphan(), which sets orphaned=True and keeps it. On re-admission request_load hits if req_id in self._outstanding (state_offload.py:284), logs "already has a load in flight", returns False, and the boundary is disowned into a full recompute even though the hash is present in LMCache — for the whole abandon window, or forever if the first load was torn down before build_connector_meta attached its spec (the allocate()-returned-False requeue at scheduler.py:1795). Related: the orphan call is nested inside if seq.has_per_req_cache and seq.state_slots: (block_manager.py:2437), so a teardown with slots already empty leaks the _outstanding entry permanently — reclaim skips it because orphaned was never set.

offload_load_cancelled is a permanent latch outside multi [reported]. It is set by the generic ChunkedOffloadSchedulerBase.cancel_pending_load and cleared in exactly one place: MultiConnectorScheduler.get_num_new_matched_tokens. In the vLLM plugin, AtomOffloadConnector.handle_preemptions (plugin/vllm/kv_transfer/connector.py:1043) calls cancel_pending_load(seq) on every preemption and then reset_for_preemption(), which resets nine placement fields but not this one — and there is no MultiConnector in that chain. From the first preemption onward the request takes the :599 early return for the rest of its life: permanent degradation to full recompute, while _start_state_load keeps dispatching index entries for it, which is the first item in this section. The in-code comment "The flag lives exactly as long as the request does, so nothing has to clean it up" is the bug — the request outlives the cancellation.

3. Root 2: the sentinel KV spec, and the two costs it is already imposing

A failed state-only load records load-error blocks over a fully-resident, possibly shared KV range [reported]. _decide_load_after_alloc's state-only branch sets hbm_cached_tokens = lmcache_cached_tokens = hbm, where hbm is a hash-block multiple (16) and need not be a multiple of virtual_block_size (64/128). A state miss then routes through _finish_load(req, False) → _record_load_error_blocks, which computes start = hbm // bs, end = ceil(hbm / bs) = start + 1, and adds block_ids[start] to _failed_load_blocks. That block is valid HBM KV, shared with other running sequences after a prefix-cache hit, and the consumer truncates num_computed_tokens at it. Only the state leg failed; the KV leg moved nothing and cannot have corrupted anything. On main this path returned early before any error-block recording.

A state-only load ships the full prompt and the full block table for a leg defined to move zero bytes [reported]. With lmcache_cached_tokens = hbm and transfer_end_tokens = None, transfer_end == hbm, so the builder runs token_ids=list(seq.token_ids[:hbm]) + block_ids=list(seq.block_table). A state-only load exists because the KV prefix is fully resident, so hbm is large — at K3's 1M context that is a several-hundred-thousand-element Python list copy on the engine loop thread, pickled through ZMQ to every TP rank, unpickled, sliced a second time by _load_kv_bytes (toks = req.token_ids[:lmc], above the lmc <= hbm -> return True early exit), then discarded. At TP8 that is eight copies of a multi-MB list per state-only resume, all on the TTFT path. The deleted metadata.state_loads comment said it outright: "a state load shares no shape with a KV transfer — no token ids, no block ids, no chunking." Short of the one-condition fix above: emit empty token_ids/block_ids for state-only, and hoist the lmc <= hbm check above the slice.

4. A default-argument change that was not swept to its callers [verified]

fail_load gained *, missing: bool = False, and only missing=True retracts the hash. grep -rn 'missing=True' atom/ returns only docstring prose — the two production callers, kda_state.py:778 (plugin) and scheduler.py:3760 (engine), both omit it, and the sole caller that passes it is tests/test_state_offload_index.py:99.

On the ATOM path retraction survives through take_missed_state_hashes (scheduler.py:3729). On the plugin path it does not exist at all — take_missed_state_hashes has no implementation anywhere under atom/plugin/. On main, fail_load unconditionally called forget(h), so the index self-healed on the first miss; now could_serve(h) keeps voting a dead hash and every later request over that prefix parks, misses and recomputes, permanently.

The reasoning in the new docstring is right — a joint failed_loading does not prove the state bytes are gone. It just needs a caller that can tell the difference on both paths.

5. PP>1 went from graceful degradation to a hard raise inside the worker [verified]

Main logged "the state tier is unsupported under pipeline parallelism (pipeline_parallel_size=%d); paged KV is unaffected" and returned, leaving the dense paged-KV leg working. This head raises ValueError from _build_state_tier, which is called by register_kv_caches — on the worker, after weights are loaded, with no handler on that path.

This repo's recorded failure mode for that shape is a worker dying between "load model runner success" and the ready signal: the RPC deadlocks and wait_server_ready.sh does not report an error. The operator gets a hang after minutes of weight loading rather than a config error.

The raise was also moved above the soft return refusals (state_backend is None, layout, views), so a PP>1 config with no state backend at all — previously a soft no-tier — now fails to start too. The justification ("with the tier off a K3 request's KV leg is declined too") is asserted rather than demonstrated, and it contradicts the line it replaces.

Refusing loudly is the right instinct; the right depth is the engine's config-validation path, which already computes state_tier_capability, before workers are spawned.

6. Slot lifetime: two ways a destination slot is released while it may still be written [reported]

_outstanding is keyed by bare req_id while an orphaned entry outlives its request. Request R's state load is in flight when R is preempted → orphan(R, slotA); the entry stays with orphaned=True and slotA is held. R is re-admitted, request_load is refused (§2), so R has no state leg — but its ordinary KV load is armed, and that load's finished_loading arrives with req_id R → scheduler.py:3719 offload.complete_load(R) → _settle pops the first generation's orphaned entry and releases slotA while the first load's H2D may still be writing it. The next admission gets slotA with has_initial_state true over half-written bytes, loads_completed is credited for a load nobody observed, and the first load's real report later finds nothing to settle. This PR recognises exactly this generation hazard for the verdict key — _state_verdict_key carries a LoadOperationId — and leaves the index on a bare request id.

reclaim() releases an orphaned destination slot on a wall-clock timeout. orphan() holds the slot precisely because "the worker may still be scattering into it", and a timeout cannot distinguish a dead worker from a slow one (backed-up CPU tier, long D2H queue, stalled disk backend). On expiry the pool hands the slot to a new admission with has_initial_state true, the new request's forward reads it, and then the old worker's unpack lands on top — a live sequence's recurrent state overwritten with another request's history, silently, with the late complete_load finding nothing to settle. The docstring borrows the store-side twin's safety argument ("a leaked pin breaks no BlockPool invariant"); that argument does not transfer to a write destination.

And the reclaimer itself may never run. offload.reclaim() — the only orphan-slot reclaimer left after BlockManager.reconcile_orphan_load_slots was deleted — is nested inside if offload is not None and take is not None: at scheduler.py:3836. The deleted reclaimer had no dependency on the store-report channel. For any deployment where BlockManager builds a StateOffloadIndex (decided from config alone) but the connector exposes no take_state_reports — a kv_role: kv_consumer shell, a duck-typed plugin connector, a multi composite whose _state_tier_sub() is None — reclaim is never called on any step, and a lost report strands the slot forever. Enough of those and can_allocate's state gate refuses every per-request-cache admission and the engine wedges with no error.

7. A rank without a tier votes on nothing, and the key never dies [reported]

_build_state_tier has four soft return refusals below the PP raise, leaving _state_tier=None while state_tier_capability (config-only, in the engine process) still grants loads. Then _note_state_leg_unrun returns silently on tier is None, _load_state_bytes returns False, and get_finished's if self._state_tier is None: return out sits above the take_hash_verdicts loop — so that rank emits no verdict at all — while the scheduler's fail_load(req_id) no longer forgets (§4). Steady state: every request over those prefixes parks, fails and recomputes, forever.

If tier presence is asymmetric across TP ranks it is worse: the tier-bearing ranks report the key, the tier-less rank does not, and drain skips an incomplete key with a bare continue — no TTL, no eviction, reset never called on the serving path — so the key sits at world_size - 1 reports forever. That is precisely the immortal-key hazard the comment at connector.py:320-334 documents for the KV-short-circuit case.

8. Deleting the tier's own executor serialises the legs and unbounds the staging buffers [reported]

Two costs from dropping the dedicated 1-worker load executor and its _staging_budget semaphore.

Latency. TTFT for a joint load goes from max(KV_retrieve, state_get) to KV_retrieve + state_get, and a state leg — which reaches LMCache's disk backend and ends in a blocking producer.synchronize() — now head-of-line-blocks every other request's KV retrieve on a pool that _offload_common.py:171 defaults to one thread. A state-only load, which moves zero KV, previously had a lane KV could never block.

Memory. StagedTransfer._tls is threading.local() and release_after_transfer defaults off, so a staging buffer is allocated per thread that ever runs a transfer. The deleted comment quantified the old budget: "both executors are max_workers=1 ... standing HBM is exactly two buffers." With the documented OFFLOAD_LOAD_WORKERS=4 that becomes five buffers — roughly 165 MiB of extra standing HBM per rank, taken out of the KV pool, with no knob and no accounting — and the removed semaphore lets n_load loads plus one store gather concurrently where two was the cap. The first state leg on each load thread also pays a cold torch.empty(entry_bytes) inline on its own TTFT.

A cheaper fusion keeps the state leg on its own executor and fuses only the report: submit inside _do_load_req, fut.result() after _load_kv_bytes. One completion, two concurrent transfers.

9. _canonical_connector swallows everything and returns "" [reported]

The catch is bare except Exception, not the ValueError that KVConnectorFactory.canonical_name raises for an unknown name (factory.py:113). Any other failure — the registry not yet populated when state_tier_capability runs at BlockManager.__init__, a non-str config value, a future rename of the classmethod — returns "" for every name: state_tier_capability reports "connector hosts no state tier", _offload_subconfig reports "multi lists no offload connector", no StateOffloadIndex is built, and offload is inert while the operator believes it is on. That is the exact scenario StateTierCapability's own docstring says it exists to prevent, and the same "" makes two offload subs both miss _OFFLOAD_BACKENDS, skipping the loud refusal.

Fix-then-sweep: state_offload.py:601 and :662 and multi_connector.py:105 still compare the raw string against "multi", which is what the helper was added to eliminate. An alias spelling takes the non-multi branch, never unwraps the offload sub, reads a zero chunk grid, and silently disables the joint KV load.

10. Smaller

  • tests/test_scheduler_block_manager_contract.py reads scheduler.py as text and AST-walks it to assert method existence. A test that greps the code it tests passes for reasons unrelated to behaviour. Its own docstring says it exists because test_scheduler.py does not import in this environment — that import failure is the defect worth fixing, since a collection error aborts the entire CI run rather than one module.
  • Both commits carry a Co-Authored-By: Claude Opus 5.5 (1M context) trailer. Worth stripping before merge; it is permanent in upstream history afterwards.
  • StateLoadSpec.chunk_tokens is written and never read, and its docstring claims "the worker validates the transfer against it" — a validation this same PR deleted (connector.py:373 explains why).
  • Stale docs: offload/connector.py:152 still names the deleted lmc-state-load thread; README.md:102 still says "store/finished/failed hash sets"; _offload_common.py:428 still says the face routes "state stores and loads".
  • _state_verdict_key (connector.py:419) re-derives an identity _load_completion_id already owns, with two divergences (is not None vs or, and a str() the base does not apply) — two spellings of a key whose whole point is being byte-identical across ranks and correlatable with the KV completion.
  • _confirm_remote_load_after_alloc (scheduler.py:2607) lost its if not needs_remote_load: return False short-circuit, and should_park_for_load_after_alloc is mutating (rewrites ls.hbm_cached_tokens, can start a handoff, sets offload_loaded_tokens, clears the spec). Dense / DSV4 / lmcache_mp sequences now take those mutations on a path they never reached before; the change is needed for K3 state-only but was not scoped to it.
  • kimi_k3/connector.py is 927 lines (main: 912) against the 800-line soft ceiling with no stated exception; _build_state_tier is 134 lines and _load_state_bytes 60 with roughly a 3:1 comment-to-code ratio.

zejunchen-zejun and others added 2 commits September 28, 2026 19:18
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>
…pin check

Review round on 08dc08e.

A state-only load could be parked and then never reported. Main's #2305
re-takes the lookup pin at build time and drops a load whose tier hit is
below `lmcache_cached_tokens`. The state-only no-op spec aims that at the
HBM length, and HBM can hold more of the prefix than the KV tier does
(1024 resident, 768 in the tier, boundary 768), so the load was dropped
after the request had parked. `_ensure_lookup_pin` now skips a spec that
reads nothing from the tier, and the state-only branch clears a stale
`transfer_end_tokens` so the spec really reads nothing.

Under `kv_connector: multi` every state-only resume was declined. Fusing
the legs moved the park decision to the KV load owner, and a state-only
load has none: every sub's KV answer is 0. That is the production agentic
shape ([mooncake producer, lmcache_offload]); main parked it
unconditionally. The composite now asks the tier sub when the engine
secured a state leg (`offload_joint.load_hash`) and no sub owns KV.

`seq.offload_load_cancelled` was never cleared, so a request whose tier
sub lost one arbitration kept its state leg suppressed after winning a
later one: the KV load went out alone and the forward would resume over a
slot nothing filled. Each arbitration now starts clean.

State verdicts are keyed by the load generation (`LoadOperationId`), not
the request id: a re-admitted request loading the same hash again reused
its first load's tombstoned key and its miss was dropped.

`orphan()` claims a slot only if it is the one the load writes into. The
request id names the request, not the admission, so a later admission's
teardown was told the index had its slot, and leaked it.

`reclaim` settles as an abandon, so the outcome counters add up to
`settled` again.

Each new test fails against the previous code.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Copilot AI review requested due to automatic review settings September 28, 2026 12:11
@zejunchen-zejun
zejunchen-zejun force-pushed the zejun/k3_lmcache_refactor branch from 3ea7ab1 to 3a61b3d Compare September 28, 2026 12:11

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@zejunchen-zejun

Copy link
Copy Markdown
Collaborator Author

@valarLip
Thanks for the review. Grouping by root cause was the right call: fixing the two roots closed most of the list, and the remaining items did not need separate machinery.

All of it is in b41397a9d on top of the two existing commits. The Co-Authored-By trailers are gone from all three. Every new regression test was run against the previous head (3a61b3dce); all but one fail there, and the one that passes pins behaviour this round did not change. The full CPU suite (--ignore=tests/plugin) and tests/plugin have the same failure sets as origin/main (1423fceb0): 28 environment-dependent failures in test_pd_pp / test_mooncake_rail_address, none of them in files this PR touches.

Where I verified a [reported] item differently from how it was written, I say so below.


§1 — neutral verdict written before the KV leg — agreed, fixed (together with §8)

Confirmed. get_finished drains on the RPC thread while the load runs on the load pool, so a drain during _load_kv_bytes could publish the neutral True, reach quorum, tombstone the key, and drop the real False.

The neutral verdict is now written only on the paths where the leg does not run: the leg was refused (coverage, destination_slot < 0), or submitting it raised. A leg that runs reports its own verdict and nothing is written ahead of it. See _submit_state_leg, and the docstring on StateOffloadTier.note_load_unrun.

Tests: test_a_failed_kv_leg_fails_the_load_the_state_leg_ran_beside (asserts no neutral verdict when the leg runs) and test_a_refused_state_leg_still_reports_on_the_verdict_key.

§2 / Root 1 — the connector declined and nothing told the index

Parked load dropped at build time — agreed, fixed generally. This was not only a K3 problem. #2305's tier_lost_prefix drop leaves a parked dense (or DSV4) request with no report in exactly the same way. The dense and DSV4 builders now record such a load (_drop_undispatched_load), and the scheduler drains take_undispatched_loads() at the top of every schedule() and _update_from_kv_xfer_finished. Each one is failed like a failed load: it goes to failed_recving_kv_req_ids, and the state leg gets abandon_load (nothing was attempted, so it is not a miss). The shell and multi forward the new method. Test: test_a_parked_load_the_tier_lost_is_reported_rather_than_stranded.

multi sub loses the KV arbitration — agreed on the outcome, one detail differs. A P/D winner reports on finished_recving, not finished_loading, so complete_load is never called. Instead the index entry stays outstanding for good, and the request resumes on the winner's report over a slot nothing restored. The outcome is silent wrong output either way. The fix restores what main did with cancel_state_load. The scheduler asks carries_state_load(seq) before parking. MultiConnector answers True only when the load owner is the tier sub, or when there is no KV owner (a state-only load, which the tier sub carries). Otherwise _release_uncarried_state_load abandons the state load, disowns the boundary and sets num_cached_tokens = 0, and the request then waits only for the winner's KV. Tests: test_the_state_leg_is_carried_only_by_the_tier_subs_load and TestAStateLegNoLoadCarriesIsGivenBack.

Re-admission refused while an orphan is outstanding — not reachable on the native path today. A torn-down admission's slot is still protected (see §6). In the native scheduler, a request with a state load in flight is never deallocated and then re-admitted: allocate returns False only on the no-load disown path, a parked request is not in running so it is never preempted, and abort defers the free until the report lands (_awaiting_aborted_load_cleanup). The allocate()-returned-False requeue at scheduler.py:1795 therefore never has an outstanding entry. The empty-state_slots leak you noted is unreachable for the same reason: request_load records seq.state_slot at allocate, and the slots are only cleared by deallocate itself.

offload_load_cancelled latch outside multi — the latch is not reachable, but the placement was wrong, so it moved. The plugin never uses KimiK3OffloadScheduler (its K3 path is kda_state.py), so nothing on the plugin path reads the flag. On the native path the generic cancel is called only from _reject_aborted_waiting, whose request is finished. Still, a generic base setting a flag that only one composite scenario needs is wrong. ChunkedOffloadSchedulerBase.cancel_pending_load no longer sets it. MultiConnector.get_num_new_matched_tokens sets it, per arbitration, only when a tier sub exists and loses, and clears it at the start of the next arbitration. The SeqView slot this had forced is removed.

§3 / Root 2 — the sentinel KV spec — agreed, done the way you suggested

A state-only load is now a request of its own:

  • K3's should_park_for_load_after_alloc parks it and drops any KV spec the lookup left behind.
  • The next build emits LMCacheReqMeta(load_spec=None, token_ids=[], block_ids=[], load_operation=…, state_load_spec=…), registered in _active_load_operations like any load.
  • The worker dispatches on load_spec is not None or state_load_spec is not None (DenseOffloadConnector._loads).
  • _load_kv_bytes returns early on load_spec is None, and now also on lmc <= hbm before slicing token_ids.

This deletes the _ensure_lookup_pin exemption, the transfer_end_tokens reset and the state-only branch of _decide_load_after_alloc. Both costs you listed go with it: a failed state-only load records no load-error block (_record_load_error_blocks already returns on load_spec is None), and nothing ships the resident prompt.

Tests: test_a_state_only_load_travels_as_a_request_of_its_own (it also asserts the tier is never looked up), test_the_worker_runs_a_request_with_only_a_state_leg, and test_a_failed_state_only_load_marks_no_kv_block_as_errored.

A correction to something I said earlier: in the previous round I told Copilot that the [hbm, lmc) error range is empty for a state-only load. That was wrong whenever hbm is not aligned to virtual_block_size. It is moot now.

§4 — fail_load(missing=) not swept to callers — agreed, fixed

Confirmed: the plugin's state_failed is the KDA leg's own verdict, and it lost the retraction that main had. KdaBoundaryPlanner.on_load_result now passes missing=True. On the native path, retraction keeps going through take_missed_state_hashes, where it is correct because a fused failed_loading covers both legs. Test: test_a_failed_kda_load_retracts_the_hash.

§5 — PP > 1 raised inside the worker — agreed, moved

The worker is back to main's warning and return (now in tier_builder.py). The refusal is refuse_state_tier_under_pp(config), called at the top of EngineCore.__init__ before AsyncIOProcManager spawns any worker. It reuses the config-only capability check with the PP gate off, so a PP config that would not host a tier anyway (for example a dense layout) still starts. Tests: test_the_engine_refuses_a_state_tier_under_pp_before_workers_exist and test_a_worker_turns_the_tier_off_under_pp_rather_than_raising.

§6 — slot lifetime

Bare req_id keying — agreed on the mechanism, fixed. Not reachable today (see §2). When a request is admitted again, BlockManager.allocate calls StateOffloadIndex.detach_orphan(seq.id). That re-keys an earlier admission's orphan under a key no request id equals, so no completion report can settle it; only reclaim can. request_load refuses the request a second state load until then. I did not key the index by LoadOperationId: request_load runs at allocate, before the operation exists, and the scheduler receives bare ids after process_completions. Detaching at the one event that makes the id ambiguous is exact without either. Test: test_a_readmitted_requests_reports_cannot_settle_its_old_orphan.

reclaim() on a wall-clock timeout — kept deliberately. You are right that a timeout cannot tell a dead worker from a slow one. The alternative is a permanent leak per lost report, which drains the pool until the state gate refuses every hybrid admission. That is the same trade main's reconcile_orphan_load_slots made, with the same window (pin timeout plus margin, minutes, against transfers that take milliseconds), and the window starts at orphan time, not dispatch. I would rather keep reclamation and document it than trade it for a guaranteed wedge. If you think that is the wrong call, I will make it opt-in.

Reclaim nested under take is not None — agreed, fixed. It is hoisted out, so it runs whenever an index exists.

§7 — a rank without a tier votes on nothing — not reachable

A tier-less rank cannot be asked for a state load. A load needs could_serve(h), which needs h indexed, and a store is indexed only on a TP quorum that is failure-dominant. A rank without a tier reports every store as failed on STATE_INDEX_CHANNEL, and as source-released on STATE_SOURCE_CHANNEL so the units come back (_store_failed_no_tier). So if any rank lacks a tier, no hash is ever indexed and no load is ever granted, whether tiers are symmetric or not. The comment in _submit_state_leg states this where the no-tier branch is.

§8 — deleting the tier's executor — agreed, restored

StateOffloadTier has its own single-worker lmc-state-load lane again. _do_load_req submits the state leg first, runs the KV leg, then joins. The join runs on every path, a raise included, because reporting before the H2D lands would let the request resume or recompute over a slot still being written. The result is one completion, a joint load costs the longer leg rather than the sum, a state get no longer holds up other requests' KV retrieves, and standing staging is back to two buffers (one per lane).

The one semantic change: the legs no longer gate each other, so the state leg can land while the KV leg fails. That is safe, because a failed load recomputes from 0 and never reads the restored state. The docstring says so.

Tests: test_the_state_leg_lands_before_the_load_reports and test_a_submitted_load_runs_on_its_own_lane_past_a_stuck_store.

§9 — _canonical_connector — agreed, fixed

It now catches only ValueError (an unknown or empty name), so anything else surfaces. The three raw "multi" compares (state_offload.py ×2, multi_connector.py) compare canonical names. Tests: test_only_an_unknown_name_reads_as_no_connector and test_a_multi_spelled_as_an_alias_still_unwraps_the_offload_sub.

§10 — smaller

  • Contract test — removed. test_scheduler.py does import here (180 passed), so the reason for the static check is gone.
  • Co-Authored-By — stripped from all three commits.
  • StateLoadSpec.chunk_tokens — removed, with the docstring claim.
  • Stale docs — fixed in the README (both lines) and in _offload_common.py. The lmc-state-load mention in offload/connector.py is accurate again now that the lane is back.
  • _state_verdict_key — gone; the key is _load_completion_id(req).
  • _confirm_remote_load_after_alloc — short-circuits again unless the sequence has a lookup hit or a state leg (load_hash != -1), so dense, DSV4 and lmcache_mp sequences no longer reach the mutating hook. Test: test_only_a_state_leg_can_need_a_park_the_lookup_did_not_ask_for.
  • File size — _build_state_tier moved to hybrid/kimi_k3/tier_builder.py, which also documents the rule that it never raises. kimi_k3/connector.py is now 786 lines.

GPU validation is still owed before merge, and more so now that the worker's load path changed. On the AtoMesh multi config against main, it needs to cover state-tier load count and TTFT (the state-only path must be live), and two-pass GSM8K delta_acc with state_offload_invariant_violations == 0.

This branch has not been deployed

No deployments
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.

3 participants