Repository navigation
feat(scheduler): admit ragged cohorts — the admission layer #103 named, and a general throughput bug - #104
feat(scheduler): admit ragged cohorts — the admission layer #103 named, and a general throughput bug#104heydryft wants to merge 3 commits into
Conversation
… an MTP one The batch admission layer assumed one batch has one cache length, in two places. `scheduler/default_scheduler.rs` partitioned the running set by `Sequence::cache_bucket_len` and ran exactly ONE bucket per step, so sequences at different lengths could not batch together at all; and `engine/mod.rs` issued `CacheInstruction::In` only when the completion-id list changed, so a stable cohort's front-alignment was never recomputed. Both now read one `RaggedAdmission`, decided once in `Engine::new` from the pipeline's own declaration (`CacheManagerMixin::ragged_batch_admission`, default a refusal carrying a reason). With admission refused the bucket key and the `pre_op` predicate are the pre-change expressions exactly. The cost of the old rule is a law, measured by a CPU scheduler simulation: with B sequences over D distinct cache lengths far enough apart that the coalescence override is refused, B/D of them run per step. At B=128 over 128 lengths 64 tokens apart that is 1.00, in steady state. `KvAdvance::PerSequence` is reachable for a V4 target for the first time. For every other architecture the declaration still refuses, from `model_masks_ragged_batches` — no other model in this tree threads a ragged-batch mask into its forward, and that is where the refusal now is. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Code Metrics Report━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ Language Files Lines Code Comments Blanks ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ C Header 5 305 210 52 43 CSS 2 1181 1036 34 111 CUDA 72 24328 17592 4018 2718 Dockerfile 1 39 22 8 9 JavaScript 16 3546 2676 482 388 Jinja2 7 694 656 5 33 JSON 74 4600 4597 0 3 Makefile 1 6 5 0 1 Metal Shading Lan| 33 12224 9431 1142 1651 PowerShell 1 300 227 30 43 Python 143 14830 12217 797 1816 Shell 23 5756 3964 1412 380 Plain Text 4 3801 0 2479 1322 TOML 33 1485 1292 43 150 YAML 3 25 23 2 0 ───────────────────────────────────────────────────────────────────────────────── HTML 4 2687 2604 43 40 |- CSS 2 543 479 37 27 |- JavaScript 1 1233 1215 12 6 (Total) 4463 4298 92 73 ───────────────────────────────────────────────────────────────────────────────── Jupyter Notebooks 4 122 83 23 16 |- Markdown 1 60 30 22 8 |- Python 1 122 113 1 8 (Total) 304 226 46 32 ───────────────────────────────────────────────────────────────────────────────── Markdown 189 37703 0 28931 8772 |- BASH 70 1618 1189 313 116 |- C 2 12 12 0 0 |- CUDA 2 84 56 16 12 |- JSON 18 708 708 0 0 |- PowerShell 1 1 1 0 0 |- Python 23 1008 787 113 108 |- Rust 65 2048 1713 77 258 |- TOML 6 207 164 0 43 |- YAML 4 38 33 5 0 (Total) 43427 4663 29455 9309 ───────────────────────────────────────────────────────────────────────────────── Rust 662 307804 266560 13893 27351 |- Markdown 477 22853 471 19669 2713 (Total) 330657 267031 33562 30064 ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ Total 1277 451971 330166 73659 48146 ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ |
…ort a broken probe as a zero The first run of this experiment produced no usable ratio, and the cause was the harness. Counting `usage.completion_tokens` off COMPLETED responses means that at ~47 tok/s with max_tokens=128 and B>=32 no request finishes inside a 45 s steady window, so every row reads 0 tokens / 0 completions / 0 errors — indistinguishable from a dead server. Proven by another agent with a fake SSE endpoint that never completes, which reproduced the signature exactly. The driver now streams and counts chunks on arrival, so the measurement does not depend on any request completing. Every row carries a `verdict`, and the summary returns nan rather than a confident ratio for any cell marked HARNESS_PRODUCED_NOTHING (D18: a broken probe is not a zero). Mixed prompt lengths capped at 512 words. The 91,933 'errors' in arm B were a retry storm off a pre-existing long-prompt fault at ~1,055 words, present on older builds and unrelated to this chain; a 32..512 ladder is still a 16x spread, which is what the experiment actually needs. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…ort it PR #93 landed `XsRollingCache::trim_tail_to(new_base: usize)` and `reconcile_xs_bases` on master while this branch was making `base`/`tokens` per-row `Vec<usize>`. The 7 resulting errors are semantic, not textual: `base` is the same identifier naming two different quantities. - On master `tail` is `tokens - base` wide, so `base` IS the physical left edge of the buffer. Divergent `base` at equal `tokens` gives physical widths 4 vs 132 and the `slice_set` mismatch #93 exists to fix, reachable on plain decode through the prefix cacher. - Here `tail` is `[B, W, hidden]` and end-anchored, with `base[i] >= tokens[i] - W`. `base` is a logical resume point decoupled from the buffer; `BatchSrc::of` sets `v_slack_dim: Some(1)` with `v_slack_at_front: true`, so the widths already agree. So the reconciliation is required on one path and a regression on the other. Trimming every row to `base_max` raises the shortest-reach row's rollback floor to the batch's -- a per-row quantity replaced by a batch-wide scalar, the same shape as the cohort min-rollback #92 removed and the padding traps #102/#103/#104 each caught. This is the last layer still holding one. - `reconcile_xs_bases` takes `xs_per_seq` and returns early when set. Flag off it runs verbatim; #93's two mutation-proven tests pass unchanged (they run flag-off), with only a private-field read swapped for the `resumable_from()` accessor, which is `base[0]` for a single-row cache. - `trim_tail_to` REFUSES a multi-row cache rather than reading `base[0]`. A scalar `new_base` cannot describe a trim of rows at different resume points, and proceeding returns a right-shaped tensor -- exactly what the caller checks -- so it names the reason instead (D18). `per_row_xs_bases_survive_a_batch_round_trip_and_are_not_flattened` is the test that discriminates: flag on, two rows at divergent `base` and equal `tokens` must batch AND each keep its own `base` through `split_row`. Both mutations were run rather than assumed -- trimming the tensors is caught by the width assertion, but reconciling only the logical `base` passes every shape check and fails solely on the per-row read. 8/8 mutations caught. One survivor on the first pass (`trim_tail_to` narrowing from the wrong end) was real: no test read tail CONTENT, because every fixture fed zeros. `trimming_the_retained_window_drops_the_oldest_rows_not_the_newest` feeds a ramp and closes it. 430 lib tests pass, 14 synthetic, workspace check clean, scoped clippy exits 0, zero new clippy diagnostics in the touched files, and rustfmt drift is 6 before and 6 after -- no upstream reformat churn. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
The 512-word cap was set on the assumption that the ~1,055-word long-prompt fault was a hard blocker on the mixed arms. It did not reproduce on a provenance-verified build (1,100 words in 6.6 s), and a five-cohort differential to 256x spread, paged on and off, found zero faults. The spread IS the independent variable — D, the number of distinct cache lengths in flight — so capping it weakens exactly the contrast the experiment exists to draw. Back to 32..2048 (64x), with MIXED_WORDS_MAX to narrow it if a particular box turns out to have the fault. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Holding this PR pending a review of its interaction with #122 — not a CI problem, and CI will not find it#122 (ragged dense decode) has landed on a master that already contained #100. Its files and this PR's overlap in two places, and unlike the rest of the stack the overlap is semantic, not incidental:
Both are mechanisms for "how much of row Git auto-merged #122 and #100 cleanly with 466/466 green, which establishes their texts do not conflict and nothing more. What's needed before this lands: somebody reads both and answers whether they compose, duplicate, or contend — specifically whether the thread-local That's a review question, not a CI question, and I'm the merge chain rather than a reviewer — so this holds rather than proceeding on a green gate. Knock-on: #114 and #116 are stacked above this one, so both are held by ordering. On their own merits they are clear of #122 — #114 touches |
The held question, answered — read on master with #122 landed. Compose, duplicate, AND contend — in three different places.The hold asked: do these compose, duplicate, or contend — specifically, is the thread-local 1.
|
| #122 (on master) | this PR | |
|---|---|---|
| layout | front-align all rows to a shared width once | keep per-row lengths |
| dead prefix | masked by CausalMasker via RAGGED_LEAD_PAD |
stripped by drop_dead_prefix |
clone_in |
membership change only | every step |
| cost | one alignment per cohort | one re-assembly per token |
What this needs — and it is now a measurement, not a review. Master already ships design A. The question is whether design B buys anything design A does not, priced against B's per-step rebuild. Until someone runs that A/B, merging this PR is a throughput regression with no test to catch it. I am leaving the hold in place, but the question is no longer open-ended: it is one number.
🟠 Separate finding, ON MASTER, not this PR's fault — two strip implementations with opposite failure policies
drop_dead_prefix refuses a non-Normal slot with a named error ("Nothing front-pads it, so nothing should be stripping it either"). #122's strip_pending_lead_pad reimplements the same narrow inline and silently skips any non-Normal slot — if let KvCache::Normal { k, v } = slot with no else. Same operation, one refuses by name and one says nothing. Filing separately; it is a defect on master, not a reason to hold this PR.
Nothing deleted. This branch, and #113 / #114 / #116 / #121 above it, are untouched.
|
Closing as superseded, not abandoned. The admission layer this PR proposed is already on master by another route: What is given up: nothing functional. Reopen if the ragged-cohort admission needs to become a named policy object rather than a boolean gate. |
Stacked on #103 (
feat/mtp-draft-chain-per-seq) → #102 → #100 → #95 → #92. Review those first.The bug, in one sentence
#103 named this as
batch_admission_admits_ragged_cohorts()and framed it as the sixth MTP blocker. It is not primarily an MTP blocker. It is the reason a fleet's admitted batch is not the batch that runs.📊 Bucket shattering is real, and it is a law
bucket_shattering_law(default_scheduler.rs) is a CPU simulation of the scheduler's own logic — not a hardware measurement (D14). It admitsBsequences atDdistinct cache lengths 64 tokens apart (far enough thatselect_running_bucket's coalescence override is refused) and counts how many the scheduler puts in each forward:effective_batch = B / D. At B=128 with every sequence at its own length, one sequence runs; the other 127 are admitted, resident in KV, counted asrunningby the logger, and idle.admitted_batch_size_stops_buying_anything_once_lengths_divergestates the same thing as a curve:[1, 1, 1]at B = 8/32/128 refused,[8, 32, 128]granted.The
grantedcolumn is the negative control in the direction that matters — same sequences, same steps, same arithmetic, only the bucket key differs.B=32 → 47.45, B=128 → 43.73curve (FACTS.mdsession 8). That run had--mtp-depth 3, prefix cache off, and--max-seqs 128pinning the B=128 row. What the simulation establishes is that the scheduler has a mechanism that would produce exactly that shape, and how large it is. Whether it fired on that box is a hardware question — see the GPU ask.The fix — two sites, one decision
scheduler/default_scheduler.rsBucketKey = (cache_bucket_len, imgs, offset), one bucket per stepragged; every completion sequence takes one shared key (ragged_bucket_len), so the batch stays whole however far apart its caches areengine/mod.rspre_op = Iniff the completion-id list changedissues_cache_in(no_kv_cache, ids_changed, granted)— a granted run re-assembles every step, because the front-alignment is a function of lengths that diverge, andInis the only place a per-rowlead_padis reported at allBoth read one
RaggedAdmission, decided once inEngine::newfromCacheManagerMixin::ragged_batch_admission(). Deciding it in one place is the point: a scheduler that formed ragged cohorts for a pipeline the engine never re-assembled for would front-align once and then decode against a stale alignment.🔑 Why flag-off is byte-identical, structurally
Not "a test observed the same numbers" — the two predicates are the pre-change expressions:
ragged_bucket_lenreturnsNonewhen refused, sobucket_key_ofis(cache_bucket_len, imgs && is_prompt, token_offset, false)— the historical three-part key with a constant appended. No sequence can move between buckets, at any length, in any state (a_refused_run_builds_the_pre_change_bucket_keyasserts it component-for-component).issues_cache_in(no_kv, changed, false) == !no_kv && changedover every input (a_refused_run_issues_exactly_the_pre_change_cache_instructions).b1_is_untouched_by_either_admission_rule).No seventh flag
The declaration rides
ARC_MTP_PER_SEQ_KV, the flag #100/#102/#103 already coordinate on, plusARC_V4_XS_PER_SEQthroughcache_supports_per_sequence_advance. A run that did not ask for per-sequence advance declares a refusal by name: its rows never diverge, so ragged admission would cost aclone_in_cacheper step and buy nothing.🔴 The default declaration is a refusal, not
falseCacheManagerMixin::ragged_batch_admissiondefaults toErr(reason)naming the pipeline. Serving a ragged cohort means left-aligning it, and the zero-filled dead prefix that leaves is not a masked row — it scores logit 0 and takes real softmax weight. Opting in is the safety property;the_default_declaration_is_a_refusal_that_names_itselfpins it, and the engine logs the reason so a run that asked and did not get it is told which layer said no (D18).🔴 The trap — the barrier re-forming as padding, at the admission key
#102 found it in the target cache, #103 one cache down. The admission layer's copy is the nastier one:
Sequence::cache_bucket_len— the single functionbucket_key_ofreads — isnormal_cache()[0].current_seq_len(), and afterfront_align_batchthat is the batch maximum for every row. A caller that records the padded length makes every sequence report the same value, and then:the_scheduler_reads_each_row_s_own_length_not_the_padded_oneruns the pad → forward → commit → strip loop over real tensors and asserts on the growth of the spread:3 + 2nwith the strip, a constant 2 without it, at n = 1/6/12. It readscurrent_seq_len()back off the slot rather than recomputing it — which is exactly how #103's first draft of the equivalent test passed with the strip stubbed out.Mutation runs — 26 mutations, 26 caught
Four survived the first run and the fixture was the problem in three of them, as four PRs in this chain have now found:
raggedfrom the coalescing filter survived because nothing tested that the override must be refused across the ragged flag (two buckets that can never merge — one a[B,1]decode, one a[B,t]prefill — must not trigger a merge-justified override).the_coalescence_override_never_targets_a_bucket_it_cannot_merge_withis set up as exactly the case where the override would fire if the keys looked mergeable;ARC_V4_XS_PER_SEQon, where the cache check and the mask check agree. With the flag off a V4 cache still masks a ragged batch but itsXsRollingslots carry one length, and the two checks separate;🔴
PerSequenceis reachable — and where the refusal movedper_sequence_advance_is_reachable_once_the_engine_grants_admissionruns the whole chain through the real trait methods on a V4-shapedMtpSpeculativePipeline: declare →RaggedAdmission::decide→set_ragged_batch_admission→kv_advance() == KvAdvance::PerSequence. Its control is the identical pipeline with the grant withheld, which must stayCohort. Five waves refused this mode; it is now on.I am not claiming this is the last blocker, and for the general throughput problem it plainly is not. The refusal moved to
model_masks_ragged_batches: no architecture in this tree except DeepSeek V4 threads a ragged-batch mask into its forward, so every other model still declares a refusal and keeps length-bucketed admission. The admission layer is fixed for everyone; the masks are not. That is now a per-architecture job, not a layer —CausalMasker::make_causal_mask_matrixalready routes onPastKvLenCache::per_seq_kv_lens, so what each model needs is to hand itRaggedKvLensinstead of its cache, plusclone_in_cache/clone_out_cachelearningfront_align_batch+drop_dead_prefixfor the plain (non-MTP) decode path.the_declaration_still_refuses_every_model_that_cannot_mask_a_ragged_batchis the test that will fail the day someone relaxes that without doing the work.📊 Measured — CPU tests and simulation, not hardware (D14)
cargo test -p mistralrs-core --lib→ 454 passed, 0 failed (was 442 at #103's tip) ·--test synthetic_load_smoke→ 13 passed ·cargo check --workspace --testsgreen · scoped clippy lane exits 0 ·mistralrs-coreclippy diagnostics diffed per touched file before/after: zero new (engine/mod.rs0→0,kv_cache/mod.rs4→4,pipeline/mod.rs9→9,mtp_pipeline.rs7→6,default_scheduler.rs1→1,scheduler/mod.rs0→0,ragged_admission.rs0,sequence.rs0→0; workspace total 238→237) · zero rustfmt drift added:rustfmtfollowsmoddeclarations and reformatted six upstream files it should not have, so that was reverted and the remaining noise stripped withgit merge-fileagainst a formatted baseline — the diff is+54/-4inengine/mod.rsand+40/-0inpipeline/mod.rs, all of it new.🔴 The GPU ask — the experiment that proves or kills the hypothesis
arc-tools/ragged_admission_curve.sh. I did not run it (D15 — I never callruncrate). Self-contained,setsid nohup … < /dev/null, status file at/root/logs/ragged_admission_curve.status, gated on/root/arc-tools/gpu_box_preflight.sh(refuses to measure if it fails or is missing).Four arms × B ∈ {1, 8, 32, 64, 128},
--max-seqs 256so the cap never pins the B=128 row (it did in session 8),--mtp-depth 3held constant across all arms so the MTP chain is not a confound:ARC_V4_XS_PER_SEQ=1 ARC_MTP_PER_SEQ_KV=1🔑 A vs B is the whole test and it needs none of this PR's code — same binary, same flags, same model, only the prompt-length distribution differs.
The one number:
tok_s(B) / tok_s(A)as a function of B. Shattering predicts it falls as B grows (refused effective batch isB/D, uniform isB). If it is flat at ~1.0 across the sweep, the hypothesis is dead for this workload — which is a real answer and should be recorded as one. Thentok_s(D)/tok_s(B)at B=128 is what this PR is worth on the mixed workload, andtok_s(*, B=1)must agree across all four arms.IntervalLoggerprints"<N> running, <M> waiting", andNis literally how many sequences the scheduler put in the forward — no throughput reasoning needed. Butthroughput_logging_enabledis only reachable from the Rust SDK builder (with_throughput_logging);mistralrs serveexposes no flag for it, so the line is never emitted by a server run.scrape.pylooks for it and reportsNAwith the reason when it finds nothing — never0.00, because a zero would read as "no sequences ran" when it means "nobody was logging" (D18).