Skip to content

perf(scheduler): coalesce prefills on TP - #2238

Merged
valarLip merged 8 commits into
mainfrom
hexwang/opt_schedule_agentic
Sep 22, 2026
Merged

valarLip merged 8 commits into
mainfrom
hexwang/opt_schedule_agentic

Conversation

@whx-sjtu

@whx-sjtu whx-sjtu commented Sep 15, 2026 •

Copy link
Copy Markdown
Contributor

Small agentic prefills interrupt decode.

Enable existing PrefillDelayer on TP/DCP (DP=1, PP=1) via nonzero ATOM_PREFILL_DECODE_INTERVAL. Add bounded HBM probes for batching, skip probing during decode protection, and wait for reusable in-flight prompt-end checkpoints. Interval 0 preserves default TP scheduling.

Kimi-K3 AgentX C64, TP8/DCP8; identical clients, full warmup + 3600s. Candidates use ATOM_PREFILL_DELAYER_MAX_QUEUE_MS=5000 and --state-checkpoint-interval-tokens 32768.

Metric Baseline Interval 4 Interval 8
ITL P50, ms 87.73 83.88 (-4.39%) 82.51 (-5.95%)
ITL P90, ms 113.80 106.58 (-6.34%) 99.98 (-12.14%)
Total token/s/chip (cached input included) 11086.34 11471.45 (+3.47%) 11620.80 (+4.82%)
Output token/s, node 631.96 663.52 (+5.00%) 675.48 (+6.89%)
TTFT P50, s 1.48 2.32 (+56.24%) 2.87 (+93.55%)
TTFT P90, s 4.02 5.94 (+47.79%) 10.47 (+160.49%)

Percentages are changes from baseline; TTFT increases are regressions. Recommend interval 4 for smaller TTFT cost. C64 combines code/configuration changes; no ablation was run for that comparison.

Kimi-K3 AgentX C16

Metric Interval 0 control Interval 4 Change
ITL P50, ms 17.798 16.608 -6.69%
ITL P90, ms 31.603 24.595 -22.18%
ITL P99, ms 57.477 41.182 -28.35%
Total token/s/chip (cached input included) 6701.83 7328.70 +9.35%
Output token/s, node 423.45 452.20 +6.79%
TTFT P50, ms 1182.455 1765.534 +49.31%
TTFT P90, ms 2994.960 3307.263 +10.43%
Successful requests 1462 1648 +12.72%

Percentages are changes from the C16 interval 0 control. Interval 4 lowers ITL while increasing total throughput, at a TTFT P50/P90 cost of 583/312 ms. Both runs have zero measurement errors and two end-of-run cancellations. All 1,462 common requests have identical output lengths; the table uses the full official AIPerf results. Total throughput is AIPerf total_token_throughput.avg / 8.

C16 scheduling diagnostics over the 3600s window: prefill batches decrease from 2,854 to 2,508 (-12.12%) while actual prefill tokens increase by 7.09%; the share of batches containing multiple requests rises from 10.62% to 31.06%.

C16's larger percentage ITL gain partly reflects its lower starting ITL: P90 drops by 7.01 ms (-22.18%), versus 7.22 ms (-6.34%) for C64 interval 4. The recipes also differ: C16 uses DSpark while C64 has no speculation, and the interval counts scheduling passes rather than emitted tokens. These results do not isolate concurrency as the cause of the different gains.

Validation: 355 passed, 1 skipped; GPU shared-prefix checks passed; full GSM8K C64, 5-shot: 1257/1319 and 1256/1319. Black passes; Ruff has existing findings, none added. Full-suite collection is blocked by missing SGLang/vLLM dependencies.

@github-actions

Copy link
Copy Markdown
Contributor

🏷️ CI Guide

Runs automatically on every eligible PR before approval:

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

Heavy model tests:

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

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

@whx-sjtu whx-sjtu changed the title perf(scheduler): coalesce agentic prefills on TP perf(scheduler): coalesce prefills on TP Sep 15, 2026
@zufayu
zufayu requested a review from yitingw1 September 16, 2026 01:51
@whx-sjtu
whx-sjtu force-pushed the hexwang/opt_schedule_agentic branch from c084d85 to 41da984 Compare September 16, 2026 05:17
@whx-sjtu
whx-sjtu force-pushed the hexwang/opt_schedule_agentic branch 2 times, most recently from e768a0e to d75dad4 Compare September 17, 2026 07:40
@whx-sjtu
whx-sjtu force-pushed the hexwang/opt_schedule_agentic branch from d75dad4 to 4346669 Compare September 21, 2026 08:59
@valarLip

Copy link
Copy Markdown
Collaborator

Read this against the tree at head in an isolated worktree. The goal — let a
waiting request ride an in-flight prefill's prefix instead of recomputing it, and
extend the delayer's decode protection to plain TP — is worth having. Most of what
follows is about one thing: on the models this repo runs most, the feature does not
currently engage, while the machinery added to decide whether it should is paid for
on every tick.

The anchor is a state-cache position, not a prefix-cache boundary

_wait_for_inflight_prefix (scheduler.py:692) gates on
block_manager.enable_prefix_caching and then anchors on
producer.checkpoint_end_pos. Those are different preconditions.
BlockManager._record_checkpoint_end (block_manager.py:1879) opens with

seq.checkpoint_end_pos = 0
if not self.state.enabled:
    return

and returns again if state_checkpoint_interval_tokens == 0. So on any pure-attention
model — DeepSeek-R1, Qwen3, Llama — every producer carries anchor == 0, the filter
anchor - cached_tokens < self.max_num_batched_tokens is satisfied for every one of
them, each is continued, and the function returns False unconditionally. The wait
never happens; the for producer in self.running loop runs per waiting seq per tick
regardless.

Where the state cache is enabled, the anchor is still not the position the waiter
wants. _record_checkpoint_end picks it for state-resume keepability — floored to the
hash grid, stepped back by successor_room, capped at n_hash_blocks - 1 — and
checkpointers_at can decline it afterwards, so the waiter can park on a checkpoint
that never lands. What this feature needs is the position at which the producer's KV
prefix becomes reusable, which is a different quantity that happens to be available
nearby.

I would resolve this before weighing anything below, because it decides whether the
rest is a cost with a matching benefit or a cost on its own.

The skip is in the wrong place in the Phase-2 loop

Four separate consequences follow from where line 1531 sits, so they are worth taking
together.

It runs after can_allocate(seq) at 1519, which uses the default record=True. The
docstring for that parameter says record=False exists so "a probe cannot inflate the
operator-visible funnel", and the two other probe sites (_can_admit_head_prefill:787,
_local_prefill_pending_work:738) pass it. A seq that waits N ticks therefore commits
the joint boundary and bumps joint_boundaries / state_tier_boundaries N times for an
admission that never happens, then continues past the disown path at 1636-1646 that
would have cleared it. Moving the check above 1519, or probing with record=False, fixes
this on its own.

It is upstream of both _park_for_remote_load call sites (1612 and 1648) but downstream
of _query_connector_prefill_match (1431). A PD-consumer or LMCache request whose KV is
already materialized remotely, and which happens to share a prompt prefix with a local
in-flight prefill, has needs_remote_load=True and seq.offload_joint.kv_prefix_tokens
already written — and is then diverted into the prefix wait and never parked. The remote
transfer is not initiated, the connector lookup is re-issued from scratch every tick, and
the request waits on a producer whose output it does not need.

The deferred seq goes to skipped_waiting_requests (1535), which is drained with
self.waiting.extend(...) at 1695 — the tail of a deque the loop pops from the head. The
two pre-existing users of that deque park on an event, so tail position is harmless for
them; this seq is immediately schedulable. It is demoted behind every request that has
arrived since, every tick, for as long as the producer prefills, and the wait carries no
age or tick deadline. ATOM_PREFILL_DELAYER_MAX_QUEUE_MS cannot rescue it:
_oldest_waiting_prefill_age_ms still reports it, so the delayer takes the queue_hot
must-fire exit every tick believing it released the request while Phase 2 skips it again —
the TTFT guard is defeated and coalescing is disabled for the duration. Every other requeue
in this loop uses appendleft(seq); break for exactly this reason.

And because the skip consumes neither a seq slot nor token budget before continueing, it
is the only branch that can drain the whole queue in one tick. With 128 requests against one
system prompt, each tick pops all 127 non-producers, pays the record=True hash walk and a
full-prefix compare for each, and re-appends them behind any later-arriving unrelated
request — which is admitted in their place. Arrival order is inverted for precisely the
traffic this feature targets. If every candidate is skipped, num_seqs_prefill == 0, so
notify_prefill_executed (engine_core.py:451) never fires, the decode interval is never
armed, and _fire()'s _reset() restarts the hold/stall episode from zero each tick — the
protection this PR is built around is inert in the case it is built for.

The fill signal was capped at four entries

_local_prefill_pending_work (scheduler.py:724) breaks at i >= 4, while
budget = n_prefillable * max_num_batched_tokens is unchanged from the
_waiting_new_token_count() it replaces — which scanned the whole queue and exited early
only at max_num_batched_tokens. With 256 queued 200-token prompts and a 16384 budget,
pending goes from 16384 (fill 1.0, immediate fire) to ~800 (fill 0.05), so the fill exit
becomes unreachable and release only ever comes from stall_ticks or ttft_max_ticks. The
coalescer degenerates from "batch when worthwhile" into a fixed N-tick delay — the latency
without the batching win, and visible in the stats as fire_fill ≈ 0.

The bound is also spent by entries that contribute nothing: the aborted /
WAITING_FOR_REMOTE_KVS / unschedulable branch continues, but i is the enumerate
index, so four such heads report pending=0 over an arbitrarily deep admittable queue.

Two smaller things in the same function: prefillable = True is set before
pending += max(0, remaining), so a wholly-tier-resident seq (where a tier hit is not
bounded by the prompt) reports (True, 0) and holds for a request the main loop would have
scheduled immediately; and the offload-resume branch skips the can_allocate fit probe
entirely, so a KV-pressured resume inflates the fill signal with work Phase 2 will refuse.

What it costs per tick

_local_prefill_pending_work runs up to four can_allocate(seq, record=False) prompt-hash
walks on every tick including pure-decode ticks, where _can_admit_head_prefill did at most
one — it returned on the first schedulable seq. At least one of those walks is then
recomputed from scratch at 1519 for the same sequence in the same tick, with nothing
carrying the hit count or the block hashes across. A fully-cached 32k prompt measures about
0.72 ms, so four is ~2.9 ms — roughly 15% of a 20 ms decode step, worse for a small model.
The early exit at pending >= budget only fires when the coalescer is not accumulating,
so the expensive path is the common one.

The prefix compare is

np.array_equal(
    np.frombuffer(seq.token_ids, dtype=np.int32, count=anchor),
    np.frombuffer(producer.token_ids, dtype=np.int32, count=anchor),
)

which evaluates (a1 == a2).all(): full width, no short-circuit, and an anchor-sized bool
temporary. Measured on 32k int32, equal costs 4.10 µs and first-element-differs costs
3.76 µs — mismatches pay nearly full price, and mismatches are the common case. Comparing
memoryview(...)[:anchor] gets C memcmp, which short-circuits and allocates nothing;
comparing the chained block hashes can_allocate has already computed would be better
still. Separately, the producer-side predicates (status, num_cached_tokens,
multimodal_data, and anchor itself) do not depend on the candidate at all, yet are
re-evaluated len(waiting) times per tick over a set that changes only when a prefill
completes.

"Local mode" has three definitions

engine_core.py:169 uses dp_size == 1 and pipeline_parallel_size == 1 and ATOM_PREFILL_DECODE_INTERVAL > 0; prefill_delayer.py:256 uses self.cpu_group is None
alone; scheduler.py:688 uses delayer.dp_size == 1 and delayer.cpu_group is None. None
derives from the others and none is asserted. A delayer with dp_size > 1 and
cpu_group is None — reachable if stateless_init_dp_group defers or fails, and
_init_prefill_delayer(config, self.dp_group) at engine_core.py:742 never asserts the
group is non-None — takes the protects_decode fast path while _local_prefill_coalescing
stays False, so that rank reports pending_tokens=0 and a coarse prefillable into a
decision whose n_prefillable < dp_size alignment gate is live. The failure mode is a
collective hang, not an error. Either protects_decode should require dp_size == 1 too,
or all three should read one PrefillDelayer.is_local.

Relatedly, in that branch prefillable = bool(self.waiting) or self._partial_prefill_count > 0
is the coarse signal _can_admit_head_prefill's own docstring (751-763) was written to
reject — it counts ABORTED and WAITING_FOR_REMOTE_KVS as prefill work, where
_local_prefill_pending_work filters them. A queue holding only aborted requests therefore
reports prefillable, should_allow_prefill takes the decode-interval return False at
prefill_delayer.py:377 instead of the vacuous-fire exit, and _reject_aborted_waiting
(1400) does not run — so the block tables and hybrid state groups the abort path exists to
reclaim immediately are held for the whole protection window, under KV pressure.

That decode-interval HOLD also returns before every must-fire bound the module docstring
lists — kv_high, kv_low, queue_hot/max_queue_ms, ttft_max_ticks,
partial_max_ticks — which matters more now that the path is reachable on plain TP rather
than DP-only. Meanwhile _kv_usage() and _oldest_waiting_prefill_age_ms() (an O(n) scan
calling _unschedulable_reason per seq) are still computed at scheduler.py:1336-1338 and
then discarded during protection, so the claimed saving is partial.

Docs and tests

docs/environment_variables.md:28 is still headed "Prefill delayer (DP attention)" and
still says "Active only when data_parallel_size > 1"; the ATOM_PREFILL_DECODE_INTERVAL
row at 41 still describes only the decode interval. The new coupling — that this variable
now also switches on TP coalescing, double-gated on ATOM_ENABLE_PREFILL_DELAYER with a
silent early return, with no way to get one without the other — exists only as a comment in
envs.py. A TP=8/DP=1 operator reading those docs concludes the coalescer cannot be active.
In the same change, prefill_delayer.py's decision pseudocode and its "Single-rank / TP-only
mode" section never mention protects_decode, scheduler.py:1312's comment describes one of
the three new branches, and block_manager.py:870's "the sole probe caller reads only the
>= 0 return" invariant is now false twice over — there is a second probe caller, and it
reads the hit count.

There are no tests for any of it. protects_decode appears in tests/ only in a stub;
_local_prefill_pending_work, _wait_for_inflight_prefix, _local_prefill_coalescing and
_init_prefill_delayer appear nowhere. The one test touched
(tests/test_scheduler_partial_prefill_tail.py:48) sets _VetoDelayer.dp_size = 2 and
protects_decode → False, which routes it down the old _can_admit_head_prefill branch.
The 624 scheduler/delayer/block_manager tests pass against this head precisely because
nothing exercises the new paths. That stub also has no cpu_group and survives
delayer.dp_size == 1 and delayer.cpu_group is None only by and short-circuit, so the
natural next stub — dp_size=1, for the single-rank scheduler this PR is about — raises
AttributeError inside set_prefill_delayer.

What I would want to see

The anchor question first, since a corrected anchor may change what the rest should look
like. Then the skip moved above the record=True probe and below the remote-load park, with
appendleft + break like its neighbours and a tick deadline. Then one definition of local
mode that the other two read. The per-tick cost is worth attacking only after that, and most
of it goes away by reusing the hashes can_allocate already computed rather than walking the
prompt a second time.

Happy to look at a revision.

@whx-sjtu
whx-sjtu force-pushed the hexwang/opt_schedule_agentic branch from 4346669 to 5c220c6 Compare September 21, 2026 14:06
@valarLip

Copy link
Copy Markdown
Collaborator

Review — schedule optimization for agentic workloads

Reviewed 5c220c6e7 in a detached worktree. Full unit suite on the PR head: 2926 passed, 67 skipped, one pre-existing unrelated failure in tests/models/deepseek_v41/test_compilation.py.

The two ideas here are good ones — hold a TP-local prefill batch until it is worth firing, and let a consumer wait for the producer that is about to hand it a checkpoint. But both are implemented by predicting what Phase 2 will do rather than by asking it, and that has a single clean statement:

_local_prefill_pending_work is a second implementation of Phase 2's admission logic. It disagrees with the real one in five places, and every disagreement produces a fill signal that can never be realized — so the coalescer fires every tick while _stat_fire_fill reports that it is working.

The wait has the mirror-image problem: its bound is a property of the slot, not of the request.

Provenance: [verified] means the deciding code was read in the PR-head worktree or the behavior was measured. [reported] means the shape matches the code but was not traced end to end. Line numbers are PR-head.


1. The estimator and Phase 2 disagree in five places [verified]

The estimator's only job is to predict how much Phase 2 will admit. Where it diverges:

_local_prefill_pending_work Phase 2
probe budget tripped (:794) returns the full budget the real work may be ~1k tokens
sequence that does not fit (:800) continue, keeps summing break
long_prefill_token_threshold not applied applied at :1638
offload-resume one-block clamp (:785) unconditional min(...) only when num_new_tokens <= 0 (:1546-1562)
free slots (:779) max_num_seqs - len(running) also counts _num_parked_remote_kv (:1598-1603)

Each row is independently sufficient to produce fill = 1.0 forever. The first is the worst, because it fires on the workload this PR targets: with max_num_batched_tokens=8192 and a queue of 32k prompts each carrying a ~31k prefix-cache hit, iteration 1 probes seq1, probed_tokens = 32768 >= 8192 — note probed_tokens += seq.num_tokens charges the whole prompt including cached tokens, so one sequence trips the cap — and iteration 2 takes return prefillable, budget if prefillable else pending. fill = 8192/8192 = 1.0 >= target_fill → _fire("fill"), every tick. TP coalescing never coalesces for the highest-hit-rate queue there is.

The continue/break row has a second edge: it does not decrement slots, so the loop can walk the entire waiting deque. And with chunked prefill off, a preempted head sequence (num_tokens > num_prompt_tokens, so remaining > budget, but _unschedulable_reason only tests num_prompt_tokens) is skipped by the estimator, summed past, reported as a full batch — and then Phase 2 reaches it, _prefill_chunk_for_budget returns None, waiting.appendleft(A); break, zero tokens admitted.

The offload row is a plain arithmetic overstatement: hash_block_size=256, num_tokens=1000, num_cached_tokens=900 → Phase 2 schedules 100, the estimator reports 232.

Nothing asserts the two agree. If the estimator is to stay, the honest shape is one walk that both the signal and the admission read, or at minimum a debug-mode assertion that the predicted admission matches the realized one.

2. ttft_max_ticks is a property of the slot, not of the request [verified]

_inflight_prefix_wait is a single (seq_id, deadline) pair, re-armed from the current tick whenever a different sequence reaches the head (:752):

Seq A waits behind producer P; the slot is (A, tick+200). KV pressure makes preempt() do self.waiting.appendleft(B) (:2712) — or _park_ready_offload_partial_prefills does extendleft (:3543). B also matches P, and because self._inflight_prefix_wait[0] != B.id the slot becomes (B, tick+200), destroying A's deadline. When A returns to the head it gets a fresh 200 ticks. Neither ever reaches self._schedule_tick >= deadline. The slot is also never cleared on admit, abort or finish — only in set_prefill_delayer.

The only wall-clock escape is delayer.max_queue_ms (:716), which defaults to None, and even when set is tested against the head sequence's arrive_time only.

3. The wait is enforced with appendleft + break, so one sequence stalls the whole queue [verified]

At :1630, a wait that is a property of request A blocks everything behind it. Agentic traffic: A is a 128k prompt matching in-flight producer P; behind it sit 30 short unrelated requests with no producer, each admittable in one pass. All are deferred up to 200 scheduler passes (~4 s at 20 ms/pass) for no benefit, and the max_queue_ms escape never even looks at them.

ttft_max_ticks was sized as a coalescer hold bound — decode keeps running and the hold ends on a fill target. Reused here it has no fill-based early exit, so it is a hard 200-pass ceiling rather than a rarely-reached one. A per-sequence deferral wants a per-sequence mechanism (skip A, keep scanning), not a queue-head barrier.

4. checkpointers_at is passed the consumer's remainder as the producer's forward size [verified]

checkpointers_at's contract (block_manager.py:2226) is that next_forward_tokens is what the producer's next forward carries, compared against each class's successor_room. :732 passes min(producer.num_prompt_tokens - anchor, seq.num_tokens - anchor).

With a V4-style class (successor_room = 131), the producer term is already guaranteed to pass by _record_checkpoint_end; but a consumer whose prompt is anchor + 64 tokens — a short turn-2 continuation — makes the min yield 64 < 131, keepers comes back empty, and the wait is skipped. That is the shape _record_checkpoint_end's own docstring says dominates (93.5% of cc-trace resumes land on a previous prompt end). Silent no-op, no counter, indistinguishable from "no producer matched".

5. Three state/semantic issues around the new gate

not seq.offload_joint.kv_prefix_tokens does not mean "no remote prefix" (:1627) [verified]. _query_connector_prefill_match assigns kv_prefix_tokens = int(ext_tokens) + int(seq.num_cached_tokens) (:2232), and reset_joint explicitly leaves it alone (sequence.py:116-118). A preempted or partially-cached sequence with num_cached_tokens > 0 and ext_tokens == 0 — no remote match at all — is silently excluded from the wait, in exactly the multi-turn traffic this targets. Conversely, with kv_connector is None the field is never written, so the guard's meaning depends on whether a connector is configured. needs_remote_load (already in the same condition) or oj.load_hash != -1 is the predicate you want.

Deferring record_allocation removes the per-pass reset_joint() clean slate (:1648) [reported]. _commit_joint_boundary opens with reset_joint() and its docstring names this as "the same clean-slate the old code got from resetting at the top of every boundary walk"; it is the engine's only caller. Now: tick N, seq S commits boundary_tokens=B, allocate() returns False (:1649), deallocate(S) frees the block table but leaves boundary_tokens = B, requeue. Tick N+1, can_allocate(record=False) no longer resets, and S is refused at the new wait gate (:1628) or the chunk gate (:1645) — so it sits in waiting carrying a boundary/claim span describing blocks freed on tick N. _consume_failed_remote_kv (:2087) and the joint-disown block (:1698) both branch on boundary_tokens being truthy. I could not construct a live read before the next record_allocation, so this is latent — but the invariant "a joint span on a seq was computed against the current pool" is gone and nothing restores it.

The decode-protection prefillable drops the admission contract (:1398) [verified]. prefillable = self._partial_prefill_count > 0 or next(self._waiting_prefills(), None) is not None filters on status and _unschedulable_reason only — no KV fit, no token budget — while the delayer's documented contract (prefill_delayer.py:295) is "fresh head that can allocate". With the KV pool full and a non-empty queue, n_prefillable = 1 → should_allow_prefill returns False at :393 for the whole interval, instead of taking the n_prefillable == 0 vacuous-fire path that also calls _reset(). Since Phase 1 is also gated by delayer_allows (:1441), an in-flight chunked prefill holding KV blocks gets frozen for ATOM_PREFILL_DECODE_INTERVAL ticks per chunk — the opposite of the partial_max_ticks bound at :413, which the interval branch returns before ever reaching. A stale _hold_ticks carried across a KV-pressure gap can also trip the ttft must-fire bound on the first tick of the next burst.

6. A perf(scheduler) change that adds per-tick CPU on the serial path [verified, measured]

The replaced _can_admit_head_prefill probed at most one sequence and returned. _local_prefill_pending_work (:796) probes until cumulative seq.num_tokens >= budget, and the cap is tested before the probe against the previous total — so the first sequence is always probed at full width. Each probe is a full O(prompt) chained-xxhash can_allocate(record=False).

Standalone measurement: ~2.9 ms for a 128k prompt at hash_block_size=64 (script left at /app/logs_claude/pr2238_probe_cost.py). With one 128k prompt at the head that is ~2.9 ms of Python/xxhash per tick for the entire hold episode, and on the tick the delayer fires Phase 2 walks the same sequences again (~5.8 ms total).

block_manager.py:966 documents this exact incident as the reason the work was moved below the refusal; this reintroduces it one call site up. The results table in the PR body has no per-tick scheduler CPU number — given that this is the metric the change is most likely to regress, it should.

7. Under PD disaggregation the feature is inert, and the startup log says otherwise [verified]

disagg_is_decode defaults to False (config.py:1862), so on a prefill node the DP=1/PP=1/interval>0 guard passes, _init_prefill_delayer runs, and startup prints PrefillDelayer initialized: dp_size=1 .... One line later PrefillEngineCore.__init__ (:958) replaces self.scheduler with PrefillScheduler, which is not a Scheduler subclass (scheduler.py:3767) and has neither prefill_delayer nor set_prefill_delayer.

On the decode side the base sets self.scheduler = None and DecodeEngineCore builds its DecodeScheduler after super().__init__() without ever calling _init_prefill_delayer.

So setting ATOM_PREFILL_DECODE_INTERVAL=4 uniformly, as docs/environment_variables.md now invites ("For a single scheduler (DP=1, PP=1), including TP/DCP, it is opt-in"), does nothing on either PD role — and on the prefill node it logs that it did something. Related: DPEngineCoreProc.__init__ calls self._init_prefill_delayer(config, self.dp_group) unconditionally, which is an AttributeError on a None scheduler for DP + disagg-decode. Pre-existing in shape, but now funnelled through the new helper.

8. Fix-then-sweep miss, and the abort drain only reaches the head [verified]

The PR correctly widens status == WAITING_FOR_REMOTE_KVS to status in (ABORTED, WAITING_FOR_REMOTE_KVS) at :767, :838, :926 and :991. get_next_batch_info (:3721) keeps the old single-status filter and so reports an aborted corpse as prefill work.

Separately, the new drain while self.waiting and self.waiting[0].status == ABORTED pops only a contiguous head prefix. An aborted sequence sitting behind one live waiting request is invisible to it, and is also filtered out of _waiting_prefills — so during decode protection it holds a full block table plus a hybrid state group, which is precisely the leak the comment at :1483-1488 warns "wedges every hybrid request once enough disconnects accumulate". That comment claims the loop fixes this; it fixes the head case only.

9. Five places where a comment or doc now contradicts the code [verified]

  • envs.py:726-733 still calls ATOM_PREFILL_DELAYER_MAX_QUEUE_MS a "TTFT SLA guard ... Bounds worst-case TTFT" — the exact claim this PR deleted from prefill_delayer.py and docs/environment_variables.md ("it is not an end-to-end TTFT guarantee").
  • block_manager.py:873-877 is newly written and is false: it asserts record=False commits no "funnel counters", but _record_checkpoint_demand runs before the record gate and moves demands_recorded / demands_declined_no_room (:1857-1870). It now runs for every waiting sequence every tick, so /debug/cache_stats counts requests that never ran.
  • block_manager.py:961-973 still says "a real admission computes the boundary and commits it" (Phase 2 now passes record=False) and cites is_mixed_batch, which exists nowhere in the tree.
  • scheduler.py:912 and :980-985 still enumerate two statuses and still promise "a request never starves in the queue".
  • docs/scheduling_kv_cache_guide.md:72-73 and :407 still show the old can_allocate(seq) admission flow and signature.

10. The new local path has zero executed test lines [verified]

grep -rn '_local_prefill_pending_work|_inflight_prefix_wait|_local_prefill_coalescing|protects_decode' tests/
→ one stub, in tests/test_scheduler_partial_prefill_tail.py:48

and that stub pins is_local = False as a class attribute, so _local_prefill_coalescing stays False and the test exercises only the pre-existing cross-DP veto path. Nothing in the suite ever sets it True; protects_decode is never invoked anywhere in the repo's tests.

§1, §4 and the :800 and :779 rows are all the kind a single scheduler-level test with a local delayer would catch — a queue whose realized admission is compared against the reported fill would catch the entire §1 table at once.

set_prefill_delayer also silently became a three-member duck-typed protocol (is_local, protects_decode, max_queue_ms) on an untyped parameter with no Protocol and no base class, so any other double raises AttributeError at attach time.


Checked and not flagging, so they do not get re-litigated: seq.token_ids is array.array("i"), so memoryview(...)[:64] and np.frombuffer(..., count=anchor) both index in tokens — units are correct and no BufferError is reachable today; _extend_hash_chain copies rather than extending the caller's block_hashes, so the deferred commit sees the same chain; the not self._rejected early return can neither produce a mishandled empty batch (engine_core gates on len(req_ids) > 0) nor loop; the ABORTED drain cannot spin; there is no double delayer init on DP; the new ValueError is unreachable from production DP; config.pipeline_parallel_size exists; _schedule_tick cannot wrap. kv_usage=0.0 / age=0.0 under protects_decode are also unreachable today — but only because the interval branch returns first, which is the non-local invariant §5 is about.

Of the above, §1 and §2 are the ones I would hold on: they are not edge cases but the normal operating point of the workloads named in the PR title, and §10 means neither would show up as a failing test if it regressed further.

Update the local coalescing flag when installing or removing a delayer, and compare shared token prefixes through NumPy buffer views to avoid copying both arrays. Keep the existing cross-DP test stub explicit about its topology.
Keep per-request deadlines through retries and allocation rollbacks while
allowing bounded lookahead for independent requests. Wait for useful in-flight
checkpoints beyond smaller offload hits, without letting prior queue age
disable prefix reuse. Clarify coalescing and checkpoint wait limits.
Three abort tests bypass Scheduler.__init__ and missed its new deadline map.
Initialize the map so all six parameterized offload cancellation cases exercise
resource cleanup with a complete scheduler fixture.
@whx-sjtu
whx-sjtu force-pushed the hexwang/opt_schedule_agentic branch from 371071f to 66e6bd0 Compare September 22, 2026 08:42
@valarLip
valarLip merged commit ee49a46 into main Sep 22, 2026
57 of 69 checks passed
@valarLip
valarLip deleted the hexwang/opt_schedule_agentic branch September 22, 2026 14:40
zejunchen-zejun added a commit that referenced this pull request Sep 28, 2026
A state load's lifecycle held three facts in three owners across two
processes: "this request owes a report on hash H" (StateOffloadIndex,
engine), "its state slot stays off the free list"
(BlockManager._orphan_load_slots, engine) and "have both legs landed"
(_JointPark, worker). No object could state

    dispatched == settled + outstanding

so no test could assert it.

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

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

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

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

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

    dispatched == settled + outstanding

so no test could assert it.

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

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

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

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

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

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants