Skip to content

Continuous batching scheduler with priority and preemption

ModuleL10.2 · build · Rust · Pass 7 · 10 to 14 h
You buildrust/crates/tl-engine/src/sched.rs: the BlockSpace seam and its simplest implementation FreeListBlocks, the PrefillPolicy seam with WholePrompt, and Scheduler: admission by free blocks, a waiting queue ordered by priority with aging, a running batch re-formed every step, preemption by recompute, abort
Contractno Rust trait contract file yet: section 4 pins the API; X-TL-Priority in openai-subset.v1.yaml and PrefillRequest.priority in engine.proto feed Request.priority
Testscourse/tests/rust/l10_2.rs, 11 tests (what they check: section 4), most with a fake model whose outputs are known in advance
NeedsL10.1 the runner, SchedPolicy, and the KV pool (chapter) · ds.06 the lazy-deletion heap (chapter) · reading: lang.04 Rust, L8.3 paged KV in Python, M05.1 KV memory accounting · or --ref-deps
Used byL10.3 (chunked prefill is a PrefillPolicy), L10.4 (the block manager is a BlockSpace), L10.5 (the engine loop calls schedule and on_step every step)
MilestoneMS-L10 (64 concurrent requests complete, KV back at baseline)
Optional depthYu et al. 2022, Orca: iteration-level scheduling (free); Kwon et al. 2023, vLLM and PagedAttention (free), section 4.5 on preemption; vLLM scheduler (free)
  • The batch is re-formed at every step: finished requests leave and waiting ones join, so no sequence waits for the longest one in its batch (hand_example_schedule).
  • Admission is by free blocks: a request joins only when the pool can hold its whole sequence so far; the first request that does not fit stops admission, so big requests are not starved by small ones behind them (admission_refuses_what_can_never_run).
  • Priority with aging needs no re-sorting: ranking by priority + age / A orders requests exactly like the fixed key priority * A - arrival (aging_bounds_the_wait).
  • When a running request needs a block and none is free, the lowest-ranked running request is preempted: its blocks are freed and it is recomputed later from its tokens, with the same output (preemption_recomputes_and_keeps_outputs).
  • Batched greedy output equals single-request output, because the kernels are batch-invariant (batched_greedy_equals_single_request).
Terminal window
ol start L10.2 # stubs tl-engine/src/sched.rs
ol tests L10.2 # rung R0: read the catalog; the fake model is at the top of the test file
ol check L10.2
ol check L10.2 --ref-deps # if L10.1 or ds.06 is not passing yet

Add tl-ds (a path dependency) to your tl-engine/Cargo.toml if your L8.4 did not already: the waiting queue is your LazyHeap.


L10.1 runs one step over any batch you give it. Something has to decide what goes into each step: which requests run, which wait, and what happens when memory runs out. Serving requests one after another wastes the hardware: a decode step for one sequence reads every weight to produce a single token, and reading the weights is the expensive part, so 16 sequences in one step cost little more than 1. Static batching (collect 16 requests, run them to the end together) wastes the other way: short requests finish and their slots sit empty until the longest one ends, and a new request waits for the whole batch. Continuous batching decides again at every step (Orca, 2022). This module is that decision.

SymbolMeaningType
β\betatoken budget per step (max_batch_tokens)integer
Smax⁡S_{\max}running sequences at most (max_seqs)integer
BBtoken positions per KV block (block_size)integer
nrn_rtokens of request rr so far: prompt plus generatedinteger
crc_rpositions of rr whose K and V are computedinteger
PrP_r, ara_rpriority and arrival step of rrintegers
AAsteps per level of aging (aging_steps)integer
ttthe current stepinteger

A request is decoding when nr−cr=1n_r - c_r = 1 (its newest token has not been fed yet) and prefilling when nr−cr>1n_r - c_r > 1. A decode costs one token of the budget; a prefill chunk costs its length. Each step processes at most β\beta tokens and runs at most Smax⁡S_{\max} sequences.

  1. Decodes, highest rank first. Each needs room for one more position: append_slot(r, n_r) grows its block table to ⌈nr/B⌉\lceil n_r / B \rceil blocks. When the pool has no block, the lowest-ranked running request not yet scheduled is preempted, and the slot is tried again; when nothing ranks lower, the request preempts itself.
  2. Prefills in progress (only with chunked prefill, L10.3): each gets the chunk its PrefillPolicy grants from the remaining budget.
  3. Admissions, from the waiting queue in rank order, while fewer than Smax⁡S_{\max} run and budget is left. A request is admitted only when allocate finds blocks for all nrn_r positions (all or nothing); the first one that does not fit stops admission (no skipping ahead), so a large request at the head is not starved by a stream of small ones.

Decodes go first because a running request’s next token is what its user is waiting for; the time per output token (TPOT) stays steady while new requests queue.

A sequence of nn positions occupies ⌈n/B⌉\lceil n / B \rceil blocks. One block holds BB positions of K and V for every layer and KV head, 2⋅L⋅Hkv⋅B⋅D⋅22 \cdot L \cdot H_{kv} \cdot B \cdot D \cdot 2 bytes in f16 (M05.1): for SmolLM2-135M (L=30L = 30, Hkv=3H_{kv} = 3, D=64D = 64) and B=16B = 16, that is 368,640 bytes, so a pool of 2048 blocks is 755 MB and holds 32,768 positions. The scheduler never touches bytes: a BlockSpace hands out block ids (FreeListBlocks here, the block manager of L10.4 later), and the runner writes into them.

add refuses at once what could never run: an empty prompt, a prompt that fills the context, and a request whose longest possible sequence, min⁡(prompt+max_tokens,context)\min(\text{prompt} + \text{max\_tokens}, \text{context}), needs more blocks than the pool holds (it would be preempted forever). A full waiting queue is QueueFull, which the server answers with 429.

Strict priority starves: a steady stream of priority-10 requests keeps a priority-0 request waiting forever. Aging raises a waiting request’s effective priority with its age: πr(t)=Pr+(t−ar)/A\pi_r(t) = P_r + (t - a_r) / A. Re-sorting the queue every step would cost O(Wlog⁡W)O(W \log W); instead notice that for two waiting requests

πr(t)>πs(t)  ⟺  Pr+t−arA>Ps+t−asA  ⟺  PrA−ar>PsA−as,\pi_r(t) > \pi_s(t) \iff P_r + \tfrac{t - a_r}{A} > P_s + \tfrac{t - a_s}{A} \iff P_r A - a_r > P_s A - a_s,

because tt cancels. So the waiting queue is a heap on the fixed key kr=PrA−ark_r = P_r A - a_r (higher first; equal keys by arrival), built once at push. It is your ds.06 LazyHeap, whose handles also make abort of a waiting request O(1). With Fcfs every key is 0 and arrival decides. A preempted request goes back with its original arrival, so it keeps its age.

When memory runs out mid-decode, something must give. vLLM’s two options: swap the victim’s blocks to host memory and back, or recompute: free its blocks and later prefill its whole sequence again (prompt plus every token generated so far). This engine recomputes: no second memory tier, and the recomputed K and V are bit-identical to the originals because the runner is chunk-invariant (L10.1), so the victim’s output does not change. Its sampler state is untouched (its tokens and its generator stay with the request), so even seeded sampling continues where it stopped.

The victim is the lowest-ranked running request, at equal rank the newest: it has had the least service, so evicting it wastes the least work.

BlockSpace is the scheduler’s only view of memory (block size, totals, allocate, append_slot, release, block_table); PrefillPolicy is its only view of how prompts are cut (chunk_len). Both are traits so this module can be finished and tested before L10.3 and L10.4 exist, and so tests can drive the scheduler with a fake model and a plain free list.

Blocks of B=4B = 4 (six of them, ids 0 to 5), budget β=8\beta = 8, Smax⁡=2S_{\max} = 2, Fcfs. Requests: A has a 5-token prompt and max_tokens 3; B a 3-token prompt and 2; C a 2-token prompt and 1.

Step 1. Nothing runs, so phase 3 admits. A needs ⌈5/4⌉=2\lceil 5/4 \rceil = 2 blocks: ids 0 and 1; its whole prompt (5 tokens) fits the budget, which drops to 3. B needs 1 block: id 2; 3 tokens, budget 0. Two sequences run, so C waits. The plan: prefill A 0..5 and B 0..3, both last chunks (sampled). Each samples one token: A has 6 tokens, B has 4.

Step 2. Phase 1: A decodes (it needs room for 6 positions, which 2 blocks already give), then B (4 positions, 1 block). Budget 6 left, but Smax⁡=2S_{\max} = 2 blocks C again. B’s second token reaches its max_tokens: B finishes with Length, and block 2 goes back.

Step 3. A decodes (7 positions); C is admitted into block 2 and prefilled (2 tokens). A’s third token and C’s first end both: everything is free again. This is hand_example_schedule.

Aging. With A=4A = 4: a priority-0 request that arrived at step 0 has key 0⋅4−0=00 \cdot 4 - 0 = 0; a priority-10 request arriving at step tt has key 40−t40 - t. From step t=41t = 41 on, new priority-10 arrivals have keys below 0, so the old request is admitted first: aging bounded its wait at about P⋅A=40P \cdot A = 40 steps (aging_bounds_the_wait).

rust/crates/tl-engine/src/sched.rs
pub type RequestId = u64;
pub struct NoCapacity;
pub struct Allocation { pub cached_tokens: usize } // 0 here; prefix hits in L10.4
pub trait BlockSpace {
fn block_tokens(&self) -> usize; fn total_blocks(&self) -> usize; fn free_blocks(&self) -> usize;
fn allocate(&mut self, id: RequestId, tokens: &[u32]) -> Result<Allocation, NoCapacity>; // all or nothing
fn append_slot(&mut self, id: RequestId, total: usize) -> Result<(), NoCapacity>;
fn release(&mut self, id: RequestId, computed: &[u32]);
fn block_table(&self, id: RequestId) -> &[u32];
}
pub struct FreeListBlocks { /* ... */ } impl FreeListBlocks { pub fn new(block_tokens: usize, ids: Vec<u32>) -> Self; }
pub trait PrefillPolicy { fn chunk_len(&self, remaining: usize, budget: usize, alone: bool) -> usize; }
pub struct WholePrompt; // remaining if it fits the budget or the step is empty, else 0
pub struct SchedulerConfig { pub max_seqs: usize, pub max_batch_tokens: usize, pub max_waiting: usize,
pub policy: SchedPolicy, pub aging_steps: u64, pub max_model_len: usize }
pub struct Request { pub prompt: Vec<u32>, pub max_tokens: usize, pub priority: i32 }
pub enum AdmitError { QueueFull, Empty, TooLong { tokens: usize, limit: usize }, NeverFits { blocks: usize, total: usize } }
pub struct Chunk { pub id: RequestId, pub start: usize, pub len: usize, pub last: bool }
pub struct ScheduleOutput { pub prefill: Vec<Chunk>, pub decode: Vec<RequestId>, pub preempted: Vec<RequestId> }
pub struct StepOutput { pub id: RequestId, pub token: u32, pub stop: bool }
pub enum FinishReason { Stop, Length }
pub enum RequestEvent { Token { id: RequestId, token: u32 }, Finished { id: RequestId, reason: FinishReason } }
impl<B: BlockSpace> Scheduler<B> {
pub fn new(cfg: SchedulerConfig, blocks: B) -> Self;
pub fn with_prefill_policy(self, p: Box<dyn PrefillPolicy + Send>) -> Self;
pub fn add(&mut self, req: Request) -> Result<RequestId, AdmitError>;
pub fn abort(&mut self, id: RequestId);
pub fn schedule(&mut self) -> ScheduleOutput;
pub fn on_step(&mut self, outputs: &[StepOutput]) -> Vec<RequestEvent>;
pub fn tokens(&self, id: RequestId) -> Option<&[u32]>; pub fn prompt_len(&self, id: RequestId) -> Option<usize>;
pub fn unfinished(&self) -> Vec<RequestId>; pub fn has_unfinished(&self) -> bool; pub fn stats(&self) -> SchedStats;
pub fn blocks(&self) -> &B; pub fn blocks_mut(&mut self) -> &mut B;
}

The caller’s loop is: schedule, build a ForwardBatch (one sequence per decode at position nr−1n_r - 1, one per chunk), ModelRunner::forward, sample every decode row and every chunk with last, then on_step with one StepOutput per sampled row. on_step advances crc_r, appends each token, and finishes requests that stopped (stop, set by the engine for EOS and stop strings) or reached max_tokens or the context.

TestKINDChecksWhy it matters downstream
hand_example_scheduleunitsection 3, step by step: chunks, decodes, block ids, the finishesthe worked example
free_list_is_all_or_nothingunitlowest ids first; no partial allocationadmission never half-allocates
every_request_finishes_and_no_block_leakspropertyrandom arrivals with forced preemptions: exact outputs, caps respected, every block free at the endthe engine runs for weeks
priority_runs_firstunitPriority admits by priority, Fcfs by arrivalX-TL-Priority from the gateway’s tenant tiers
aging_bounds_the_waitpropertywith aging a low-priority request waits at most about PAP A steps; without, it starvesfairness under load
preemption_recomputes_and_keeps_outputsfault, differentiala 14-block pool forces preemption; greedy outputs equal the roomy runmemory pressure never changes answers
batched_greedy_equals_single_requestdifferentialfour requests batched give each one’s alone tokensbatching is invisible to users
abort_frees_blocks_and_drops_waitingunitabort of running and waiting requests returns every blockclient disconnects (L10.5)
decodes_count_against_the_budgetboundarythree decodes leave room for 5 tokens: a 6-token prompt waitsthe budget bounds every step
stop_finishes_with_stopunitan end condition finishes with Stop at onceEOS and stop strings
admission_refuses_what_can_never_runboundaryEmpty, TooLong, NeverFits, QueueFull400 and 429 answers instead of a stuck queue
PitfallSymptomCaught by
Ordering the waiting heap lowest rank firstlow priority runs firstpriority_runs_first (mutant s01)
Ranking by priority alone under Prioritya steady high-priority stream starves the restaging_bounds_the_wait (mutant s02)
Preempting without releasing the victim’s blocksthe pool drains under pressure; requests never finishevery_request_finishes_and_no_block_leaks (mutant s03)
Recomputing from the prompt onlya preempted request repeats or loses its answerpreemption_recomputes_and_keeps_outputs (mutant s04)
Admitting past max_seqssteps grow without bound; memory per step explodeshand_example_schedule (mutant s05)
Aborting a running request without releasing its blocksevery disconnect leaks KVabort_frees_blocks_and_drops_waiting (mutant s06)
Queueing a request larger than the poolit is preempted forever and blocks the queueadmission_refuses_what_can_never_run (mutant s07)
Not counting decodes against the budgetsteps exceed max_batch_tokens; latency spikesdecodes_count_against_the_budget (mutant s08)
Counting generated tokens off by oneevery answer has one token too manyevery_request_finishes_and_no_block_leaks (mutant s09)
A free list that hands out arbitrary idsblock tables differ from run to run; harder to debugfree_list_is_all_or_nothing (mutant s10)
Skipping ahead in the queue when the head does not fitbig requests starve behind small oneshand_example_schedule
DirectionModuleHow it uses this
BackL10.1the runner executes each plan; SchedPolicy comes from EngineConfig
Backds.06LazyHeap is the waiting queue; its handles make abort O(1)
BackL8.3the paged KV semantics in Python, block by block
ForwardL10.3Chunked is a PrefillPolicy; plan turns a ScheduleOutput into a mixed batch
ForwardL10.4BlockManager is a BlockSpace with a prefix cache
ForwardL10.5the engine loop: schedule, forward, sample, on_step; AdmitError::QueueFull becomes 429
Your pieceProduction equivalentWhat it addsWhere to look
recompute onlyvLLM swap preemptionvictims’ blocks moved to host memory and backvllm v0 scheduler.py, _preempt_by_swap
one queue with agingmulti-level feedback queues, SLO-aware schedulersseparate queues per tenant tier, deadlinesSarathi-Serve, Llumnix
head-of-line admissionbest-fit admissionsmaller requests fill gaps when the head waitsthe vLLM max_num_batched_tokens discussion
a fixed budgetadaptive budgets from measured step timetargets a TPOT instead of a token countL10.7 metrics, load.01