From 42a3a26e100fde601ed38eb43689513564a66eaf Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Thu, 24 Sep 2026 22:03:21 +0800 Subject: [PATCH 01/14] fix(vllm): warm finite-window benchmark caches with native prefills Signed-off-by: Yiming Liu --- .../src/dynamo/vllm/instrumented_scheduler.py | 527 ++++++++--- .../tests/test_vllm_instrumented_scheduler.py | 847 +++++++++++++++++- .../observability/environment-variables.mdx | 12 +- 3 files changed, 1238 insertions(+), 148 deletions(-) diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index b6c44213a04e..fbe344102582 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -106,7 +106,7 @@ from vllm.v1.core.sched.async_scheduler import AsyncScheduler from vllm.v1.core.sched.output import CachedRequestData, NewRequestData, SchedulerOutput from vllm.v1.core.single_type_kv_cache_manager import CrossAttentionManager -from vllm.v1.kv_cache_interface import MambaSpec +from vllm.v1.kv_cache_interface import FullAttentionSpec, MambaSpec, SlidingWindowSpec from vllm.v1.request import Request, RequestStatus from dynamo.common.forward_pass_metrics import ( @@ -2747,7 +2747,12 @@ def _bench_build_grid(self) -> None: len(warmup_points), ) self._bench_expected_points = len(self._bench_grid) - len(warmup_points) - # Finalize execution order BEFORE numbering: IDs then follow execution + # Native plans need a distinct key for every execution, including eager + # replicas and duplicate coordinates. These temporary IDs are translated + # to the public IDs after warm-up planning finalizes execution order. + for execution_id, point in enumerate(self._bench_grid, start=1): + point.benchmark_id = execution_id + # Finalize execution order BEFORE public numbering: IDs follow execution # order, so a soft-timeout artifact holds the contiguous prefix 1..k # the native-artifact contract requires, and the grid digest covers # the order every rank will actually run. @@ -2759,13 +2764,25 @@ def _bench_build_grid(self) -> None: # ID order are independent. real_id = 0 warmup_id = self._bench_expected_points + public_ids = {} for point in self._bench_grid: + execution_id = point.benchmark_id if EAGER_WARMUP_REASON in point.sample_reasons: warmup_id += 1 point.benchmark_id = warmup_id else: real_id += 1 point.benchmark_id = real_id + public_ids[execution_id] = point.benchmark_id + if self._kvwarm_native and getattr(self, "_kvwarm_plan", None): + # Rebase from the old dictionary in one pass: updating it in place + # could overwrite another execution's key during the permutation. + self._kvwarm_plan = { + (public_ids[execution_id], batch, kv_tokens): depth + for (execution_id, batch, kv_tokens), depth in self._kvwarm_plan.items() + } + for fallback in self._kvwarm_meta.get("capacity_fallbacks", []): + fallback["benchmark_id"] = public_ids[fallback["benchmark_id"]] grid_payload = json.dumps( [asdict(point) for point in self._bench_grid], sort_keys=True, @@ -4125,6 +4142,8 @@ def _bench_cleanup_requests(self) -> None: [rid for rid in self._bench_active_req_ids if rid in self.requests] ) self._bench_active_req_ids.clear() + self._kvwarm_native_active = False + self._kvwarm_native_resume_ids = set() self._schedule_times.clear() self._bench_extra_steps_left = 0 @@ -4780,6 +4799,11 @@ def _bench_step_prefill(self) -> SchedulerOutput | None: # Socket-level timeout: a stalled endpoint must fail the download (and # with it the warm-up) instead of blocking the scheduler indefinitely. _KVWARM_DOWNLOAD_TIMEOUT_S = 60 + _kvwarm_native: bool = False + # Only the current fleet resumed from a ready native stage has this guarantee. + _kvwarm_native_active: bool = False + _kvwarm_stage_point: BenchmarkPoint | None = None + _kvwarm_native_resume_ids: set[str] | None = None _kvwarm_stage_t0: float | None _kvwarm_stage_batch: int | None # Local outcome ``(batch, ok, detail)`` of the active stage once this @@ -4850,6 +4874,29 @@ def _kvwarm_state_layer_groups(self) -> list[str]: names.append(spec_name) return names + def _kvwarm_native_layout(self) -> bool: + """Finite windows must continue their exact native prefill requests. + + A deep parked chain may have evicted the history a shallower point + needs. Inkling's convolution cache is a SlidingWindowSpec too, so + native allocation and forward execution initialize all four streams + without assuming a tensor layout or borrowing writable state. + """ + groups = self.kv_cache_manager.kv_cache_config.kv_cache_groups + specs = [group.kv_cache_spec for group in groups] + return ( + not self._bench_random_kda + and any( + isinstance(spec, SlidingWindowSpec) and spec.sliding_window > 0 + for spec in specs + ) + and all( + isinstance(spec, (FullAttentionSpec, SlidingWindowSpec)) + and (not isinstance(spec, SlidingWindowSpec) or spec.sliding_window > 0) + for spec in specs + ) + ) + def _kvwarm_seed_regime(self, point) -> str: """Row-level KV seed provenance for artifact consumers. @@ -4906,10 +4953,12 @@ def _kvwarm_release_heavy_state(self) -> None: setattr(self, attr, None) def _kvwarm_warm_eligible(self) -> bool: - """Select real attention-KV warm-up for EP MoE or explicit random-state mode. + """Select native finite-window or shared-prefix attention-KV warm-up. Random-state mode also admits hybrid MoE without EP: its attention prefixes are real, while recurrent states remain private and synthetic. + Finite-window layouts instead prefill and continue each point's own + requests; this needs neither expert parallelism nor prefix caching. The verdict travels in the capacity envelope (see ``_bench_make_local_capacity``), so every host-local input the stage @@ -4951,7 +5000,12 @@ def _kvwarm_warm_eligible(self) -> bool: False, ) ) - if not has_experts: + self._kvwarm_native = self._kvwarm_native_layout() + if self._kvwarm_native: + reason = self._kvwarm_probe_content() + eligible = reason is None + meta["initialization_strategy"] = "native_exact_context" + elif not has_experts: reason = "dense_model_content_insensitive" elif not ep_enabled and not self._bench_random_kda: reason = "moe_tp_balanced_by_construction" @@ -5029,7 +5083,9 @@ def _kvwarm_depth_cap(self) -> int: token), so the cap keeps drift headroom below the model length. The negotiated length applies once it exists; before negotiation the local length stands in, an upper bound of the group's.""" - return self._bench_capacity_limit("max_model_len") - 4 + return self._bench_capacity_limit("max_model_len") - ( + 3 if self._kvwarm_native else 4 + ) # ------- Dataset: three-tier resolution + even-half pool + lazy tokenize ------- @@ -5223,114 +5279,142 @@ def _kvwarm_prepare(self, mode: str) -> None: # ``_kvwarm_point_need`` and the block check in # ``_kvwarm_register_shadow``, otherwise a covered point's shadow can # need one block more than its chain holds. - repeats = self._kvwarm_giant_repeats() - margin = 1 + repeats - plan: dict = {} - rung_ctxs: dict = {} - for p in decode_pts: - ctxs = self._bench_decode_context_lengths( - p.total_kv_read_tokens, p.batch_size - ) - want = min(max(ctxs) + margin, self._kvwarm_depth_cap()) - plan[p.batch_size] = max(plan.get(p.batch_size, 0), want) - rung_ctxs.setdefault(p.batch_size, set()).update(int(c) for c in ctxs) - # Shadows own private tail blocks (the admission write plus the steady - # headroom) on top of the shared chain prefix, drawn from the same pool - # while the chains are parked. Reserve them per request and per KV - # group; otherwise a rung whose chains fill the pool dies at injection - # ("Cannot get N free blocks from the pool"). - # The pool figure is the group's negotiated one (the smallest rank's; - # local before negotiation), like the depth cap: the plan decides - # which rung every rank builds and which points it warms, and the - # stage round (``_kvwarm_stage_round``) relies on every rank - # agreeing on both. - worst_case_tail = self._kvwarm_shadow_tail_blocks(repeats) - for batch, depth in list(plan.items()): - # Reserve the tails the rung's OWN points take (exact per-group arithmetic - # at each measured context), bounded by the worst case; the worst case alone - # (two blocks per group) over-reserves on hybrids. - shadow_tail_blocks = min( - worst_case_tail, - max( + if self._kvwarm_native: + plan = {} + for point in decode_pts: + contexts = self._bench_decode_context_lengths( + point.total_kv_read_tokens, point.batch_size + ) + injected = [max(1, context - 1) for context in contexts] + required = self._kvwarm_native_required_blocks(injected) + usable = self._bench_grid_usable_blocks( + point.batch_size, reserve_watermark=True + ) + depth = max(injected) + plan[self._kvwarm_plan_key(point)] = ( + depth + if required <= usable and depth <= self._kvwarm_depth_cap() + else 0 + ) + if required > usable: + meta.setdefault("capacity_fallbacks", []).append( + { + "benchmark_id": point.benchmark_id, + "batch": point.batch_size, + "depth": depth, + "required_blocks": required, + "usable_blocks": usable, + } + ) + else: + repeats = self._kvwarm_giant_repeats() + margin = 1 + repeats + plan: dict = {} + rung_ctxs: dict = {} + for p in decode_pts: + ctxs = self._bench_decode_context_lengths( + p.total_kv_read_tokens, p.batch_size + ) + want = min(max(ctxs) + margin, self._kvwarm_depth_cap()) + plan[p.batch_size] = max(plan.get(p.batch_size, 0), want) + rung_ctxs.setdefault(p.batch_size, set()).update(int(c) for c in ctxs) + # Shadows own private tail blocks (the admission write plus the steady + # headroom) on top of the shared chain prefix, drawn from the same pool + # while the chains are parked. Reserve them per request and per KV + # group; otherwise a rung whose chains fill the pool dies at injection + # ("Cannot get N free blocks from the pool"). + # The pool figure is the group's negotiated one (the smallest rank's; + # local before negotiation), like the depth cap: the plan decides + # which rung every rank builds and which points it warms, and the + # stage round (``_kvwarm_stage_round``) relies on every rank + # agreeing on both. + worst_case_tail = self._kvwarm_shadow_tail_blocks(repeats) + for batch, depth in list(plan.items()): + # Reserve the tails the rung's OWN points take (exact per-group arithmetic + # at each measured context), bounded by the worst case; the worst case alone + # (two blocks per group) over-reserves on hybrids. + shadow_tail_blocks = min( + worst_case_tail, + max( + ( + # the shadow is admitted at ctx - 1 with ``repeats`` steady + # steps (``_bench_step_decode`` / ``_kvwarm_inject_borrowed``); + # the recurrent read slot can cross a block boundary between + # ctx and ctx - 1, so reserve at the admission geometry + self._kvwarm_shadow_tail_blocks_for( + max(1, ctx - 1), max(1, repeats) + ) + for ctx in rung_ctxs.get(batch, ()) + ), + default=worst_case_tail, + ), + ) + usable = self._bench_grid_usable_blocks(batch, reserve_watermark=True) + # Plan against a margin of the pool: the per-request footprint estimate is a + # lower bound (block-boundary rounding, transient Mamba boundary blocks), + # and a stage whose chains do not ALL fit loses its real coverage, so a + # small margin buys full stages. Under attention-DP every rank must build + # the same stages: leave more headroom (a per-rank stall would desynchronize + # the ranks) -- 15% for DP>1, 5% otherwise. + pool_margin = 0.85 if getattr(self, "_bench_dp_size", 1) > 1 else 0.95 + pool = int(usable * pool_margin) + while depth > 8 and ( ( - # the shadow is admitted at ctx - 1 with ``repeats`` steady - # steps (``_bench_step_decode`` / ``_kvwarm_inject_borrowed``); - # the recurrent read slot can cross a block boundary between - # ctx and ctx - 1, so reserve at the admission geometry - self._kvwarm_shadow_tail_blocks_for( - max(1, ctx - 1), max(1, repeats) + self._bench_blocks_per_req( + depth, apply_admission_cap=True, resident_chain=True ) - for ctx in rung_ctxs.get(batch, ()) - ), - default=worst_case_tail, - ), - ) - usable = self._bench_grid_usable_blocks(batch, reserve_watermark=True) - # Plan against a margin of the pool: the per-request footprint estimate is a - # lower bound (block-boundary rounding, transient Mamba boundary blocks), - # and a stage whose chains do not ALL fit loses its real coverage, so a - # small margin buys full stages. Under attention-DP every rank must build - # the same stages: leave more headroom (a per-rank stall would desynchronize - # the ranks) -- 15% for DP>1, 5% otherwise. - pool_margin = 0.85 if getattr(self, "_bench_dp_size", 1) > 1 else 0.95 - pool = int(usable * pool_margin) - while depth > 8 and ( - ( + + shadow_tail_blocks + ) + * batch + > pool + ): + depth -= 1 + required = ( self._bench_blocks_per_req( depth, apply_admission_cap=True, resident_chain=True ) + shadow_tail_blocks - ) - * batch - > pool - ): - depth -= 1 - required = ( - self._bench_blocks_per_req( - depth, apply_admission_cap=True, resident_chain=True - ) - + shadow_tail_blocks - ) * batch - # The margin only steers the depth trim above; whether a stage is built at - # all is decided against the full pool, as upstream does (a rung already at - # the depth floor is not demoted by the margin). - if required > usable: - # Reaching the depth floor does not prove the fleet fits. - # This also covers an initially short chain below the floor. - # Preserve the points with explicit fake-KV provenance, but - # do not build a stage that violates the warmup pool budget. - meta.setdefault("capacity_fallbacks", []).append( - { - "batch": batch, - "depth": depth, - "required_blocks": required, - "usable_blocks": usable, - } - ) - logger.warning( - "KVWARM: batch=%d depth=%d needs %d blocks including " - "shadow reserves, pool has %d; using fake-KV fallback", - batch, - depth, - required, - usable, - ) - depth = 0 - plan[batch] = depth - # Slot budget: a stage parks ``batch`` chains on the worker and measures each of - # its points by injecting ``batch`` shadow requests on top, so ``2 * batch`` - # request slots must exist (worker asserts "No free indices" otherwise). Rungs - # above that fall back to fake injection (measured after the chains are - # released). On models whose max_num_seqs is memory-capped (Mamba/KDA state - # blocks) this bites at batch > max_num_seqs / 2. - try: - slots = int(self._bench_capacity_limit("max_num_running_reqs")) - except (AttributeError, TypeError, ValueError): - # capacity unknown (e.g. partially constructed scheduler): leave the plan alone - slots = 0 - if slots > 0: - for batch in [b for b in plan if 2 * b > slots]: - plan.pop(batch) + ) * batch + # The margin only steers the depth trim above; whether a stage is built at + # all is decided against the full pool, as upstream does (a rung already at + # the depth floor is not demoted by the margin). + if required > usable: + # Reaching the depth floor does not prove the fleet fits. + # This also covers an initially short chain below the floor. + # Preserve the points with explicit fake-KV provenance, but + # do not build a stage that violates the warmup pool budget. + meta.setdefault("capacity_fallbacks", []).append( + { + "batch": batch, + "depth": depth, + "required_blocks": required, + "usable_blocks": usable, + } + ) + logger.warning( + "KVWARM: batch=%d depth=%d needs %d blocks including " + "shadow reserves, pool has %d; using fake-KV fallback", + batch, + depth, + required, + usable, + ) + depth = 0 + plan[batch] = depth + # Slot budget: a stage parks ``batch`` chains on the worker and measures each of + # its points by injecting ``batch`` shadow requests on top, so ``2 * batch`` + # request slots must exist (worker asserts "No free indices" otherwise). Rungs + # above that fall back to fake injection (measured after the chains are + # released). On models whose max_num_seqs is memory-capped (Mamba/KDA state + # blocks) this bites at batch > max_num_seqs / 2. + try: + slots = int(self._bench_capacity_limit("max_num_running_reqs")) + except (AttributeError, TypeError, ValueError): + # capacity unknown (e.g. partially constructed scheduler): leave the plan alone + slots = 0 + if slots > 0: + for batch in [b for b in plan if 2 * b > slots]: + plan.pop(batch) self._kvwarm_plan = plan # Second reordering: all warmed points first, fake fallbacks last -- # fake injection fills the whole pool and evicts the chains' cached @@ -5393,6 +5477,8 @@ def _kvwarm_prepare(self, mode: str) -> None: self._kvwarm_chain_prompts: dict = {} self._kvwarm_borrowed_ids: set = set() self._kvwarm_stage_batch = None + self._kvwarm_stage_point = None + self._kvwarm_native_resume_ids = set() self._kvwarm_building = False self._kvwarm_stage_local = None self._kvwarm_round_seq = 0 @@ -5513,15 +5599,60 @@ def _kvwarm_shadow_tail_blocks(self, repeats: int) -> int: for manager in managers ) + def _kvwarm_plan_key(self, point) -> int | tuple[int, int, int]: + if self._kvwarm_native: + return ( + point.benchmark_id, + point.batch_size, + self._bench_decode_steady_kv_tokens( + point.batch_size, point.total_kv_read_tokens + ), + ) + return point.batch_size + + def _kvwarm_native_required_blocks(self, context_lengths: list[int]) -> int: + """Native fleet peak, including admission and repeated steady writes.""" + repeats = min( + self._kvwarm_giant_repeats(), + max( + 1, + self._bench_capacity_limit("max_model_len") - 2 - max(context_lengths), + ), + ) + return sum( + self._bench_blocks_per_req( + context + 1 + repeats + self.num_lookahead_tokens, + apply_admission_cap=True, + ) + for context in context_lengths + ) + + def _kvwarm_native_pool_shortfall(self) -> int: + """Check the remaining native allocation headroom before ranks agree.""" + manager = self.kv_cache_manager + contexts = [ + len(self._kvwarm_chain_prompts[req_id]) for req_id in self._kvwarm_chain_ids + ] + resident = sum( + not block.is_null + for group in manager.coordinator.single_type_managers + for req_id in self._kvwarm_chain_ids + for block in group.req_to_blocks[req_id] + ) + required = self._kvwarm_native_required_blocks(contexts) + return max(0, required - resident - manager.block_pool.get_num_free_blocks()) + def _kvwarm_plan_covers(self, point) -> bool: """Plan-level coverage decision (independent of live chains): the shared source of truth for chain building and injection dispatch.""" plan = getattr(self, "_kvwarm_plan", None) if not plan: return False - depth = plan.get(point.batch_size, 0) + depth = plan.get(self._kvwarm_plan_key(point), 0) if not depth: return False + if self._kvwarm_native: + return True ctxs = self._bench_decode_context_lengths( point.total_kv_read_tokens, point.batch_size ) @@ -5590,19 +5721,47 @@ def _kvwarm_step_busy(self) -> bool: if self._kvwarm_chain_ids: self._kvwarm_shed_chains() return self._bench_frees_pending() - if self._kvwarm_stage_batch != nxt.batch_size: + new_native_point = self._kvwarm_native and ( + not self._kvwarm_chain_ids + or self._kvwarm_stage_point is None + or self._kvwarm_plan_key(self._kvwarm_stage_point) + != self._kvwarm_plan_key(nxt) + ) + if self._kvwarm_stage_batch != nxt.batch_size or new_native_point: self._kvwarm_shed_chains() if self._bench_frees_pending(): return True # the next fleet would draw from blocks still fenced - self._kvwarm_start_stage(nxt.batch_size, self._kvwarm_plan[nxt.batch_size]) + if self._kvwarm_native: + self._kvwarm_start_stage( + nxt.batch_size, + self._kvwarm_plan[self._kvwarm_plan_key(nxt)], + point=nxt, + ) + else: + self._kvwarm_start_stage( + nxt.batch_size, self._kvwarm_plan[nxt.batch_size] + ) return True return False - def _kvwarm_start_stage(self, batch: int, depth: int) -> None: - """Launch chain prefills for one ``(batch, depth)`` warm-up stage.""" + def _kvwarm_start_stage( + self, batch: int, depth: int, *, point: BenchmarkPoint | None = None + ) -> None: + """Launch real prefills, at each admission context for finite windows.""" + depths = [depth] * batch + if self._kvwarm_native: + if point is None: + raise ValueError("native KV warm-up requires an exact benchmark point") + depths = [ + max(1, context - 1) + for context in self._bench_decode_context_lengths( + point.total_kv_read_tokens, batch + ) + ] + self._kvwarm_stage_point = point t0 = time.monotonic() - for i in range(batch): - tokens = self._kvwarm_chain_token_ids(i, depth) + for i, request_depth in enumerate(depths): + tokens = self._kvwarm_chain_token_ids(i, request_depth) req_id = f"__kvwarm_chain_{self._kvwarm_seq}" self._kvwarm_seq += 1 req = Request( @@ -5611,9 +5770,13 @@ def _kvwarm_start_stage(self, batch: int, depth: int) -> None: sampling_params=SamplingParams(max_tokens=100_000, ignore_eos=True), pooling_params=None, block_hasher=self._bench_block_hasher, - # Salts are stable per (rank, chain index): a new generation's chain - # hits the old chain's cached blocks and computes only the extension - cache_salt=f"__kvwarm_{self._fpm_dp_rank}_{i}", + # Only shared-prefix stages reuse salts across generations. + # Native points own all writable blocks, including partial tails. + cache_salt=( + f"__kvwarm_native_{self._fpm_dp_rank}_{req_id}" + if self._kvwarm_native + else f"__kvwarm_{self._fpm_dp_rank}_{i}" + ), ) self.add_request(req) self._kvwarm_chain_ids.append(req_id) @@ -5637,7 +5800,10 @@ def _kvwarm_monitor_build(self) -> bool: ) vanished.append(req_id) continue - if req.num_computed_tokens >= len(self._kvwarm_chain_prompts[req_id]): + target = len(self._kvwarm_chain_prompts[req_id]) + if self._kvwarm_native and req.num_computed_tokens > target: + return self._kvwarm_stage_outcome(False, {"context_overshoot": req_id}) + if req.num_computed_tokens >= target: running = self.running # type: ignore[has-type] if any(r.request_id == req_id for r in running): # Park: leave the scheduler's view; blocks and requests @@ -5665,7 +5831,11 @@ def _kvwarm_monitor_build(self) -> bool: return self._kvwarm_stage_outcome(False, {"vanished": len(vanished)}) if pending: return True - if self._bench_synchronizer is not None: + if self._kvwarm_native: + shortfall = self._kvwarm_native_pool_shortfall() + if shortfall: + return self._kvwarm_stage_outcome(False, {"pool_shortfall": shortfall}) + elif self._bench_synchronizer is not None: # Under attention-DP the rung's verdict is shared, so the pool # check injection repeats per point runs here for the whole rung # first: a rank skipping one point on its own would fork the @@ -5701,6 +5871,12 @@ def _kvwarm_stage_outcome(self, ok: bool, detail: dict) -> bool: (``_kvwarm_stage_round``). True either way: the step is idle. """ batch = self._kvwarm_stage_batch + if self._kvwarm_native and self._kvwarm_stage_point is not None: + detail = { + "benchmark_id": self._kvwarm_stage_point.benchmark_id, + "total_kv_read_tokens": self._kvwarm_stage_point.total_kv_read_tokens, + **detail, + } self._kvwarm_building = False if not ok: self._kvwarm_shed_chains() @@ -5774,7 +5950,9 @@ def _kvwarm_stage_settle(self, batch: int | None, ok: bool, detail: dict) -> Non detail.get("build_seconds", 0.0), ) return - if batch is not None: + if self._kvwarm_native and self._kvwarm_stage_point is not None: + self._kvwarm_plan[self._kvwarm_plan_key(self._kvwarm_stage_point)] = 0 + elif batch is not None: self._kvwarm_plan[batch] = 0 meta["stages"].append({"batch": batch, "failed": True, **detail}) logger.warning( @@ -5820,6 +5998,8 @@ def _kvwarm_shed_chains(self) -> None: self._kvwarm_chain_prompts = {} self._kvwarm_stage_batch = None self._kvwarm_building = False + # Preserve the exact point until its local/group outcome settles: + # failed native stages retire only that point's plan entry. # ------- Shadow injection: borrow chain blocks, original two-step flow ------- @@ -5842,6 +6022,19 @@ def _kvwarm_covers(self, point, injected_lengths) -> bool: chains = self._kvwarm_chain_ids if len(chains) < point.batch_size: return False + if self._kvwarm_native: + return ( + self._kvwarm_stage_point is not None + and self._kvwarm_plan_key(self._kvwarm_stage_point) + == self._kvwarm_plan_key(point) + and all( + (request := self.requests.get(chains[i])) is not None + and request.num_computed_tokens == injected + and request.num_output_placeholders == 0 + and len(self._kvwarm_chain_prompts[chains[i]]) == injected + for i, injected in enumerate(injected_lengths) + ) + ) need = self._kvwarm_point_need() return all( injected + need <= len(self._kvwarm_chain_prompts[chains[i]]) @@ -6127,6 +6320,17 @@ def _kvwarm_inject_borrowed(self, context_lengths) -> "SchedulerOutput": ) return output + def _kvwarm_resume_native(self) -> SchedulerOutput | None: + """Continue the prefills themselves; their private window state is exact.""" + req_ids = self._kvwarm_chain_ids + self._kvwarm_native_active = True + self._kvwarm_native_resume_ids = set(req_ids) + self._bench_active_req_ids.update(req_ids) + self.running.extend(self.requests[req_id] for req_id in req_ids) + self._kvwarm_chain_ids = [] + self._kvwarm_chain_prompts = {} + return self._bench_make_steady_step() + def _bench_make_steady_step(self) -> SchedulerOutput | None: """One production-shaped decode step for the injected requests. @@ -6154,29 +6358,56 @@ def _bench_make_steady_step(self) -> SchedulerOutput | None: return None kvwarm_borrowed: set[str] = getattr(self, "_kvwarm_borrowed_ids", set()) new_blocks: dict[str, Any] = {} - for request in reqs: - if request.request_id in kvwarm_borrowed: - # KVWARM shadow: blocks borrowed from a parked chain; depth headroom - # already covers the steady write -- zero allocation. - new_blocks[request.request_id] = None - continue - blocks = self.kv_cache_manager.allocate_slots( - request, - 1, - num_lookahead_tokens=getattr(self, "num_lookahead_tokens", 0), - delay_cache_blocks=True, - ) - if blocks is None: - return None - new_blocks[request.request_id] = blocks + try: + for request in reqs: + if request.request_id in kvwarm_borrowed: + # KVWARM shadow: blocks borrowed from a parked chain; depth headroom + # already covers the steady write -- zero allocation. + new_blocks[request.request_id] = None + continue + blocks = self.kv_cache_manager.allocate_slots( + request, + 1, + num_lookahead_tokens=getattr(self, "num_lookahead_tokens", 0), + delay_cache_blocks=True, + ) + if blocks is None: + if self._kvwarm_native_active: + # The stage already proved fleet headroom on every rank. + # A partial allocation cannot be retried safely or skipped + # on one rank while peers execute a collective. + raise RuntimeError( + "native KV warm-up allocation failed after stage readiness" + ) + if self._kvwarm_native: + self._kvwarm_take_cow_copies() + return None + new_blocks[request.request_id] = blocks + except Exception: + if self._kvwarm_native: + self._kvwarm_take_cow_copies() + raise + native_resumes = self._kvwarm_native_resume_ids or set() + # V1 dropped parked requests from its persistent batch. Restore their + # full native tables and sampled token history. V2 retains the table + # until finish and only accepts appended block deltas. + resumed = ( + native_resumes if native_resumes and not self.use_v2_model_runner else set() + ) cached = CachedRequestData( req_ids=[request.request_id for request in reqs], - resumed_req_ids=set(), + resumed_req_ids=resumed, new_token_ids=[], - all_token_ids={}, + all_token_ids={ + request.request_id: request.all_token_ids.copy() + for request in reqs + if request.request_id in native_resumes + }, new_block_ids=[ ( - new_blocks[request.request_id].get_block_ids(allow_none=True) + self.kv_cache_manager.get_block_ids(request.request_id) + if request.request_id in resumed + else new_blocks[request.request_id].get_block_ids(allow_none=True) if new_blocks[request.request_id] is not None else None ) @@ -6204,6 +6435,12 @@ def _bench_make_steady_step(self) -> SchedulerOutput | None: else None ), ) + if self._kvwarm_native: + copies, retained = self.kv_cache_manager.take_kv_cache_block_copies() + if copies: + self._free_cow_retained_blocks(retained, self.sched_step_seq + 1) + output.kv_cache_block_copies = copies + self._kvwarm_native_resume_ids = set() if self.connector is not None: output.kv_connector_metadata = self.connector.build_connector_meta(output) if self.ec_connector is not None: @@ -6329,8 +6566,12 @@ def _bench_step_decode(self) -> SchedulerOutput | None: # Model-length cap: after step k, total = ctx+k <= max_model_len, # and runner bookkeeping writes through +1 (points at the cap # fall back to the legacy single steady step automatically). + # Use the negotiated limit so every DP rank executes the same steps. max_ctx = max(injected_lengths) + 1 - repeats = min(repeats, max(1, self.max_model_len - 1 - max_ctx)) + repeats = min( + repeats, + max(1, self._bench_capacity_limit("max_model_len") - 1 - max_ctx), + ) if not kvwarm_real: multi = sum( self._bench_blocks_per_req( @@ -6350,15 +6591,17 @@ def _bench_step_decode(self) -> SchedulerOutput | None: point.batch_size, ) output = ( - self._kvwarm_inject_borrowed(injected_lengths) + self._kvwarm_resume_native() + if kvwarm_real and self._kvwarm_native + else self._kvwarm_inject_borrowed(injected_lengths) if kvwarm_real else self._bench_inject_fake_decode(injected_lengths) ) - if output.total_num_scheduled_tokens != point.batch_size: + if output is None or output.total_num_scheduled_tokens != point.batch_size: logger.warning( "Skipping benchmark decode point after request injection produced " "%d of %d requests: %s", - output.total_num_scheduled_tokens, + output.total_num_scheduled_tokens if output is not None else 0, point.batch_size, point, ) diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index e56e2825585d..abf107e469e2 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -18,13 +18,22 @@ import time import uuid from collections import deque -from dataclasses import replace +from dataclasses import asdict, replace from types import SimpleNamespace from unittest.mock import MagicMock, call import pytest +import torch from vllm.config import CUDAGraphMode # noqa: E402 -from vllm.v1.request import RequestStatus # noqa: E402 +from vllm.sampling_params import SamplingParams +from vllm.v1.core import kv_cache_manager +from vllm.v1.kv_cache_interface import ( + FullAttentionSpec, + KVCacheConfig, + KVCacheGroupSpec, + SlidingWindowSpec, +) +from vllm.v1.request import Request, RequestStatus # noqa: E402 @pytest.fixture(autouse=True) @@ -98,6 +107,9 @@ def _install_test_capacity_preflight(stub, capacity=None): capacity = capacity or _benchmark_capacity() stub._bench_make_local_capacity = lambda: capacity stub._bench_synchronizer = None + stub.kv_cache_manager = SimpleNamespace( + kv_cache_config=SimpleNamespace(kv_cache_groups=[]) + ) # ``_bench_build_grid`` re-filters the decode capture list against the # negotiated request limit before generating the grid; stubs that don't # model captures still need the attribute to exist. @@ -5041,6 +5053,836 @@ def _kvwarm_gate_stub(*, state_groups=(), experts=8, ep=True, prefix=True): return stub +def _kvwarm_native_gate_stub(*, experts=8, ep=False, prefix=False): + stub = _kvwarm_gate_stub(experts=experts, ep=ep, prefix=prefix) + specs = [ + FullAttentionSpec( + block_size=16, num_kv_heads=1, head_size=8, dtype=torch.bfloat16 + ), + SlidingWindowSpec( + block_size=16, + num_kv_heads=1, + head_size=8, + dtype=torch.bfloat16, + sliding_window=512, + ), + SlidingWindowSpec( + block_size=4, + num_kv_heads=1, + head_size=8, + dtype=torch.bfloat16, + sliding_window=4, + ), + ] + stub.kv_cache_manager.kv_cache_config.kv_cache_groups = [ + SimpleNamespace(kv_cache_spec=spec) for spec in specs + ] + return stub + + +@pytest.mark.core +@pytest.mark.parametrize( + "experts,ep,prefix", [(8, False, False), (8, True, True), (0, False, False)] +) +def test_kvwarm_sliding_window_uses_native_prefills_without_ep_or_prefix_cache( + experts, ep, prefix +): + stub = _kvwarm_native_gate_stub(experts=experts, ep=ep, prefix=prefix) + + assert stub._kvwarm_warm_eligible() + assert stub._kvwarm_native + assert stub._kvwarm_meta["initialization_strategy"] == "native_exact_context" + + +@pytest.mark.core +def test_kvwarm_sliding_window_does_not_admit_mamba_state(): + stub = _kvwarm_native_gate_stub(ep=True, prefix=True) + stub.kv_cache_manager.kv_cache_config.kv_cache_groups.append( + SimpleNamespace(kv_cache_spec=type("MambaSpec", (), {})()) + ) + + assert not stub._kvwarm_warm_eligible() + assert not stub._kvwarm_native + assert stub._kvwarm_meta["skip_reason"] == "hybrid_state_layers_unsupported" + + +@pytest.mark.core +@pytest.mark.parametrize("contexts", [(4, 3), (5, 4), (16, 15), (512, 511), (513, 512)]) +def test_kvwarm_native_stage_prefills_exact_heterogeneous_contexts(contexts): + stub = InstrumentedScheduler.__new__(InstrumentedScheduler) + stub._kvwarm_native = True + stub._kvwarm_seq = 0 + stub._kvwarm_chain_ids = [] + stub._kvwarm_chain_prompts = {} + stub._fpm_dp_rank = 0 + stub._bench_block_hasher = None + stub._kvwarm_chain_token_ids = lambda index, depth: [index + 1] * depth + stub.add_request = MagicMock() + point = BenchmarkPoint( + point_type="decode", + benchmark_id=5, + batch_size=2, + total_kv_read_tokens=sum(contexts) + 2, + ) + + stub._kvwarm_start_stage(2, max(contexts), point=point) + + requests = [call_.args[0] for call_ in stub.add_request.call_args_list] + assert [len(request.prompt_token_ids) for request in requests] == list(contexts) + assert [request.num_computed_tokens for request in requests] == [0, 0] + assert len({request.cache_salt for request in requests}) == 2 + first_salts = {request.cache_salt for request in requests} + stub._kvwarm_chain_ids = [] + stub._kvwarm_chain_prompts = {} + stub.add_request.reset_mock() + + stub._kvwarm_start_stage(2, max(contexts), point=replace(point, benchmark_id=6)) + + assert first_salts.isdisjoint( + call_.args[0].cache_salt for call_ in stub.add_request.call_args_list + ) + + +@pytest.mark.core +def test_kvwarm_native_plan_prices_one_fleet_and_rebuilds_each_point(): + stub = _kvwarm_planner_stub(usable_blocks=4) + stub._kvwarm_native = True + points = [ + BenchmarkPoint( + point_type="decode", + benchmark_id=index, + batch_size=2, + total_kv_read_tokens=2 * context, + ) + for index, context in enumerate((31, 17, 2), start=1) + ] + stub._bench_grid = deque(points) + + stub._kvwarm_prepare("decode") + + assert not stub._kvwarm_plan_covers(points[0]) # Six blocks including repeats. + assert stub._kvwarm_plan_covers(points[1]) # Four native blocks, no shadows. + assert stub._kvwarm_plan_covers(points[2]) + stub._bench_active_req_ids = set() + stub._bench_current_point = None + stub._bench_soft_timeout_elapsed = lambda: False + stub._bench_frees_pending = lambda: False + stub._kvwarm_shed_chains = MagicMock() + stub._kvwarm_start_stage = MagicMock() + stub._kvwarm_stage_batch = 2 + stub._kvwarm_stage_point = points[1] + stub._kvwarm_chain_ids = [] # Previous point was promoted and retired. + stub._bench_grid = deque([points[2]]) + + assert stub._kvwarm_step_busy() + stub._kvwarm_start_stage.assert_called_once_with(2, 1, point=points[2]) + + +@pytest.mark.core +@pytest.mark.parametrize("fail_first_warmup", [False, True]) +def test_kvwarm_native_grid_numbering_preserves_each_execution( + monkeypatch, fail_first_warmup +): + monkeypatch.setenv("DYN_BENCH_KV_WARMUP", "on") + points = { + "schema_version": 1, + "prefill": [], + "decode": [ + {"batch_size": batch, "total_kv_read_tokens": context} + for batch, context in ((1, 1), (2, 9), (1, 7), (20, 40), (2, 33), (1, 7)) + ], + } + stubs = [] + for rank, num_blocks in enumerate((64, 96)): + stub = _explicit_grid_stub("decode", points) + gate = _kvwarm_native_gate_stub() + stub.vllm_config = gate.vllm_config + stub._kvwarm_load_texts = gate._kvwarm_load_texts + stub._kvwarm_tokenizer = gate._kvwarm_tokenizer + stub._kvwarm_resolve_dataset = _kvwarm_no_dataset_resolution + stub.max_model_len = 64 + stub.block_size = stub._bench_hash_block_size = 16 + stub.cache_config = SimpleNamespace(enable_prefix_caching=False, block_size=16) + stub._bench_decode_capture_sizes = [] + stub._bench_decode_cudagraph_mode = "NONE" + stub._bench_cudagraph_capture_sizes = [] + config = KVCacheConfig( + num_blocks=num_blocks, + kv_cache_tensors=[], + kv_cache_groups=[ + KVCacheGroupSpec( + layer_names=[str(index)], kv_cache_spec=group.kv_cache_spec + ) + for index, group in enumerate( + gate.kv_cache_manager.kv_cache_config.kv_cache_groups + ) + ], + ) + stub.kv_cache_manager = kv_cache_manager.KVCacheManager( + config, + max_model_len=stub.max_model_len, + scheduler_block_size=16, + hash_block_size=16, + max_in_flight_tokens=16, + enable_caching=False, + ) + # Exercise the real capacity/eligibility probe, preparation, numbering, + # and dispatch. Only dataset/tokenizer I/O and engine execution are doubles. + del stub._bench_make_local_capacity + stub._fpm_dp_rank = rank + stub._bench_active_req_ids = set() + stub._bench_current_point = None + stub._bench_drain_pending = False + stub._bench_phase = _BenchPhase.DECODE_SWEEP + stub._bench_start_timing = lambda: None + stub._bench_soft_timeout_elapsed = lambda: False + stub._bench_stop_at_timeout_boundary = lambda phase: False + stub._bench_inject_fake_decode = MagicMock( + side_effect=AssertionError("numbered native point fell back to fake decode") + ) + stub._bench_block_hasher = None + stub._schedule_times = deque() + stub.deferred_frees = deque() + stub.requests = {} + stub.running = [] + stub.finished_req_ids = set() + stub.use_v2_model_runner = False + stub.needs_kv_cache_zeroing = False + stub.connector = stub.ec_connector = None + stub.sched_step_seq = 0 + stubs.append(stub) + + common = instrumented_scheduler_module._BenchmarkCapacityEnvelope.common( + [stub._bench_make_local_capacity() for stub in stubs] + ) + histories = [] + for stub in stubs: + sync = MagicMock() + sync.negotiate_capacity.return_value = common + sync.stage_round.side_effect = lambda seq, batch, phase, ok: SimpleNamespace( + all_done=phase == "done", ok=ok + ) + stub._bench_synchronizer = sync + + def add_request(request, stub=stub): + stub.requests[request.request_id] = request + stub.running.append(request) + + def finish_requests(req_ids, stub=stub): + for req_id in req_ids: + request = stub.requests.pop(req_id) + stub.kv_cache_manager.free(request) + if request in stub.running: + stub.running.remove(request) + + stub.add_request = add_request + stub._bench_finish_requests = finish_requests + stub._bench_build_grid() + grid = list(stub._bench_grid) + assert stub._kvwarm_meta["warm_eligible"] + assert stub._kvwarm_meta["initialization_strategy"] == "native_exact_context" + assert [(p.batch_size, p.total_kv_read_tokens) for p in grid] == [ + (1, 1), + (2, 9), # Eager replicas retain their original order. + (2, 33), + (2, 9), + (1, 7), + (1, 7), + (1, 1), + (20, 40), + (20, 40), # Warmup and real point cannot fit repeated native writes. + ] + assert [p.benchmark_id for p in grid] == [7, 8, 1, 2, 3, 4, 5, 9, 6] + payload = json.dumps( + [asdict(point) for point in grid], sort_keys=True, separators=(",", ":") + ).encode() + assert stub._bench_grid_digest == hashlib.sha256(payload).hexdigest() + sync.synchronize_grid.assert_called_once_with( + grid_digest=stub._bench_grid_digest, expected_points=6, missing_phases=[] + ) + seen_requests = set() + seen_salts = set() + history = [] + for index, point in enumerate(grid): + if point.batch_size == 20: + # Planned fallbacks stay synthetic; this test only dispatches + # the eligible points (fallback execution has separate coverage). + assert not stub._kvwarm_plan_covers(point) + assert not stub._kvwarm_step_busy() + assert stub._bench_grid.popleft() is point + continue + # Start through the same state-machine entry as schedule(), without + # preassigning IDs or mocking coverage/the native stage start. + assert stub._bench_step() is None + assert stub._kvwarm_building, "numbered point must start a native stage" + assert stub._kvwarm_plan_covers(point) + assert len(stub._kvwarm_plan) == len(grid) + injected = [ + max(1, length - 1) + for length in stub._bench_decode_context_lengths( + point.total_kv_read_tokens, point.batch_size + ) + ] + requests = [stub.requests[req_id] for req_id in stub._kvwarm_chain_ids] + assert [len(request.prompt_token_ids) for request in requests] == injected + assert all(request.num_computed_tokens == 0 for request in requests) + req_ids = {request.request_id for request in requests} + salts = {request.cache_salt for request in requests} + assert seen_requests.isdisjoint(req_ids) + assert seen_salts.isdisjoint(salts) + seen_requests.update(req_ids) + seen_salts.update(salts) + for request, length in zip(requests, injected): + assert stub.kv_cache_manager.allocate_slots(request, length) is not None + request.num_computed_tokens = length + request.append_output_token_ids(7) + request.status = RequestStatus.RUNNING + if fail_first_warmup and index == 0: + # A failed replica must not demote the later measured execution + # at identical coordinates, even though both began with ID zero. + requests[0].num_computed_tokens += 1 + assert stub._bench_step() is None + assert not stub._kvwarm_plan_covers(point) + assert stub._kvwarm_plan_covers(grid[6]) + assert stub._bench_grid.popleft() is point + continue + assert stub._bench_step() is None # Park and settle the real prefill. + assert not stub._kvwarm_building + output = stub._bench_step() + assert output is not None + assert output.scheduled_new_reqs == [] + assert output.scheduled_cached_reqs.num_computed_tokens == injected + assert set(output.scheduled_cached_reqs.req_ids) == req_ids + assert stub._kvwarm_seed_regime(stub._bench_current_point) == "real_kv" + history.append((point.benchmark_id, injected)) + # Advance native state before retiring it: an independent repetition + # must receive fresh prefills, never these already-advanced requests. + for request in requests: + request.num_computed_tokens += 1 + request.append_output_token_ids(7) + stub._bench_cleanup_requests() + stub._bench_current_point = None + assert not stub.requests + assert stub.kv_cache_manager.block_pool.get_num_free_blocks() == ( + stub.kv_cache_manager.kv_cache_config.num_blocks - 1 + ) + assert stub._kvwarm_meta["points_fake_fallback"] == 0 + assert stub._kvwarm_meta["capacity_fallbacks"] == [ + { + "benchmark_id": benchmark_id, + "batch": 20, + "depth": 1, + "required_blocks": 80, + "usable_blocks": 63, + } + for benchmark_id in (9, 6) + ] + histories.append((history, sync.stage_round.call_args_list)) + assert stubs[0]._bench_grid_digest == stubs[1]._bench_grid_digest + assert stubs[0]._kvwarm_plan == stubs[1]._kvwarm_plan + assert histories[0] == histories[1] + + +def _kvwarm_native_resume_stub(use_v2): + stub = InstrumentedScheduler.__new__(InstrumentedScheduler) + stub._kvwarm_native = True + stub._kvwarm_chain_ids = ["chain-a", "chain-b"] + stub._kvwarm_chain_prompts = {"chain-a": [1] * 3, "chain-b": [2] * 4} + requests = [ + SimpleNamespace( + request_id=req_id, + num_computed_tokens=length, + num_output_tokens=1, + num_output_placeholders=0, + all_token_ids=list(range(length + 1)), + is_finished=lambda: False, + ) + for req_id, length in (("chain-a", 3), ("chain-b", 4)) + ] + stub.requests = {request.request_id: request for request in requests} + stub.running = [] + stub._bench_active_req_ids = set() + stub.use_v2_model_runner = use_v2 + stub._bench_random_kda = False + stub._kvwarm_borrowed_ids = set() + blocks = MagicMock() + blocks.get_block_ids.return_value = ([11], [21]) + kv = MagicMock() + kv.allocate_slots.return_value = blocks + kv.get_block_ids.side_effect = lambda req_id: ( + ([1, 11], [2, 21]) if req_id == "chain-a" else ([3, 12], [4, 22]) + ) + kv.take_kv_cache_block_copies.return_value = ([(8, 9)], ["retained"]) + kv.num_kv_cache_groups = 2 + stub.kv_cache_manager = kv + stub.sched_step_seq = 10 + stub._free_cow_retained_blocks = MagicMock() + stub.num_lookahead_tokens = 0 + stub.needs_kv_cache_zeroing = False + stub.finished_req_ids = set() + stub.connector = None + stub.ec_connector = None + return stub, requests + + +@pytest.mark.core +@pytest.mark.parametrize("use_v2", [False, True]) +def test_kvwarm_native_resume_continues_private_requests_and_preserves_runner_state( + use_v2, +): + stub, requests = _kvwarm_native_resume_stub(use_v2) + + output = stub._kvwarm_resume_native() + + assert output.total_num_scheduled_tokens == 2 + assert stub.running == requests + assert all(stub.requests[request.request_id] is request for request in requests) + assert stub._bench_active_req_ids == {"chain-a", "chain-b"} + assert stub._kvwarm_chain_ids == [] + assert stub._kvwarm_borrowed_ids == set() + cached = output.scheduled_cached_reqs + assert cached.num_computed_tokens == [3, 4] + assert cached.all_token_ids == { + request.request_id: request.all_token_ids for request in requests + } + if use_v2: + assert cached.resumed_req_ids == set() + assert cached.new_block_ids == [([11], [21]), ([11], [21])] + else: + assert cached.resumed_req_ids == {"chain-a", "chain-b"} + assert cached.new_block_ids == [([1, 11], [2, 21]), ([3, 12], [4, 22])] + assert output.kv_cache_block_copies == [(8, 9)] + stub._free_cow_retained_blocks.assert_called_once_with(["retained"], 11) + assert stub._kvwarm_native_resume_ids == set() + + +def _kvwarm_native_build_stub(): + stub = _kvwarm_planner_stub(usable_blocks=100) + stub._kvwarm_native = True + point = BenchmarkPoint( + point_type="decode", benchmark_id=1, batch_size=2, total_kv_read_tokens=9 + ) + next_point = replace(point, benchmark_id=2, total_kv_read_tokens=5) + stub._bench_grid = deque([point, next_point]) + stub._kvwarm_prepare("decode") + stub._kvwarm_meta_init()["stages"] = [] + stub._kvwarm_stage_point = point + stub._kvwarm_stage_batch = 2 + stub._kvwarm_building = True + stub._kvwarm_stage_t0 = time.monotonic() + stub._kvwarm_chain_ids = ["chain-a", "chain-b"] + stub._kvwarm_chain_prompts = {"chain-a": [1] * 4, "chain-b": [2] * 3} + stub.requests = { + req_id: SimpleNamespace( + request_id=req_id, + num_computed_tokens=len(tokens), + num_output_placeholders=1, + ) + for req_id, tokens in stub._kvwarm_chain_prompts.items() + } + stub.running = list(stub.requests.values()) + stub._kvwarm_native_pool_shortfall = MagicMock(return_value=0) + stub._bench_active_req_ids = set() + stub._bench_current_point = None + stub._bench_soft_timeout_elapsed = lambda: False + stub._bench_frees_pending = lambda: False + stub._bench_finish_requests = MagicMock() + return stub, point, next_point + + +@pytest.mark.core +def test_kvwarm_native_parks_exact_contexts_until_all_async_outputs_drain(): + stub, point, _ = _kvwarm_native_build_stub() + assert stub._kvwarm_step_busy() + assert stub.running == [] + assert stub._kvwarm_building + assert not stub._kvwarm_covers(point, [4, 3]) + stub._kvwarm_native_pool_shortfall.assert_not_called() + + stub.requests["chain-a"].num_output_placeholders = 0 + assert stub._kvwarm_step_busy() + assert stub._kvwarm_building + stub.requests["chain-b"].num_output_placeholders = 0 + assert stub._kvwarm_step_busy() + assert not stub._kvwarm_building + assert stub._kvwarm_covers(point, [4, 3]) + assert not stub._kvwarm_covers(point, [3, 3]) + assert not stub._kvwarm_step_busy() + + +@pytest.mark.core +@pytest.mark.parametrize("failure", ["overshoot", "vanished", "capacity", "timeout"]) +def test_kvwarm_native_failure_releases_fleet_and_preserves_other_contexts(failure): + stub, point, next_point = _kvwarm_native_build_stub() + for request in stub.requests.values(): + request.num_output_placeholders = 0 + if failure == "overshoot": + stub.requests["chain-a"].num_computed_tokens += 1 + elif failure == "vanished": + del stub.requests["chain-a"] + elif failure == "capacity": + stub._kvwarm_native_pool_shortfall.return_value = 1 + else: + stub._bench_soft_timeout_elapsed = lambda: True + + assert stub._kvwarm_step_busy() + + stub._bench_finish_requests.assert_called_once_with(["chain-a", "chain-b"]) + assert stub._kvwarm_chain_ids == [] + assert not stub._kvwarm_plan_covers(point) + assert stub._kvwarm_plan_covers(next_point) + assert stub._kvwarm_meta_init()["stages"][0]["failed"] + assert stub._kvwarm_meta_init()["stages"][0]["benchmark_id"] == point.benchmark_id + assert not stub._kvwarm_step_busy() # The failed context is not rebuilt. + + +@pytest.mark.core +def test_kvwarm_native_dp_waits_for_peers_and_demotes_only_the_failed_context(): + stub, point, next_point = _kvwarm_native_build_stub() + for request in stub.requests.values(): + request.num_output_placeholders = 0 + sync = MagicMock() + stub._bench_synchronizer = sync + sync.stage_round.side_effect = [ + SimpleNamespace(all_done=False), + SimpleNamespace(all_done=True, ok=False), + SimpleNamespace(all_done=False), + ] + + assert stub._kvwarm_step_busy() # Locally ready, still waiting for peer. + assert stub._kvwarm_plan_covers(point) + assert stub._kvwarm_stage_local[:2] == (2, True) + assert stub._kvwarm_step_busy() # Peer failed; group demotes this point. + assert not stub._kvwarm_plan_covers(point) + assert stub._kvwarm_plan_covers(next_point) + assert stub._kvwarm_chain_ids == [] + assert not stub._kvwarm_step_busy() + stub._bench_finish_requests.assert_called_once_with(["chain-a", "chain-b"]) + + +@pytest.mark.core +def test_kvwarm_native_capacity_counts_resident_blocks_not_evicted_placeholders(): + stub = InstrumentedScheduler.__new__(InstrumentedScheduler) + stub._kvwarm_chain_ids = ["chain"] + stub._kvwarm_chain_prompts = {"chain": [1] * 8} + stub._kvwarm_native_required_blocks = lambda contexts: 4 + group = SimpleNamespace( + req_to_blocks={ + "chain": [SimpleNamespace(is_null=True)] * 20 + + [SimpleNamespace(is_null=False)] * 2 + } + ) + stub.kv_cache_manager = SimpleNamespace( + coordinator=SimpleNamespace(single_type_managers=[group]), + block_pool=SimpleNamespace(get_num_free_blocks=lambda: 1), + ) + + assert stub._kvwarm_native_pool_shortfall() == 1 + + +@pytest.mark.core +def test_kvwarm_native_capacity_uses_sliding_window_admission_caps(): + stub = InstrumentedScheduler.__new__(InstrumentedScheduler) + stub.max_model_len = 1024 + stub.num_lookahead_tokens = 0 + stub.kv_cache_manager = SimpleNamespace( + coordinator=SimpleNamespace( + single_type_managers=[ + SimpleNamespace(block_size=16), + SimpleNamespace(block_size=4, _max_admission_blocks_per_request=4), + ] + ) + ) + + # 511 prefilled tokens + admission + three steady writes; the full + # group holds 33 blocks, while the native finite window needs only four. + assert stub._kvwarm_native_required_blocks([511]) == 37 + + +@pytest.mark.core +@pytest.mark.parametrize("after_admission", [False, True]) +@pytest.mark.parametrize("allocation_raises", [False, True]) +def test_kvwarm_native_failed_resume_remains_owned_for_cleanup( + allocation_raises, after_admission +): + stub, _ = _kvwarm_native_resume_stub(False) + if after_admission: + stub._kvwarm_resume_native() + stub.kv_cache_manager.allocate_slots.return_value = None + if allocation_raises: + stub.kv_cache_manager.allocate_slots.side_effect = RuntimeError( + "allocation failed" + ) + stub._bench_finish_requests = MagicMock() + stub._schedule_times = deque() + with pytest.raises(RuntimeError, match="allocation failed"): + if after_admission: + stub._bench_make_steady_step() + else: + stub._kvwarm_resume_native() + assert stub._bench_active_req_ids == {"chain-a", "chain-b"} + stub.kv_cache_manager.block_pool.free_blocks.assert_called_once_with(["retained"]) + + stub._bench_cleanup_requests() + + assert stub._bench_active_req_ids == set() + assert stub._kvwarm_native_resume_ids == set() + assert set(stub._bench_finish_requests.call_args.args[0]) == {"chain-a", "chain-b"} + + +@pytest.mark.core +@pytest.mark.parametrize("after_native", [False, True]) +def test_kvwarm_synthetic_fallback_releases_copies_after_failed_steady_allocation( + after_native, +): + stub, requests = _kvwarm_native_resume_stub(False) + if after_native: + stub._kvwarm_resume_native() + stub._bench_finish_requests = MagicMock() + stub._schedule_times = deque() + stub._bench_cleanup_requests() + # The next fleet uses synthetic fallback, even on a native-capable layout. + stub.running = requests + stub._bench_active_req_ids = {request.request_id for request in requests} + stub._free_cow_retained_blocks.reset_mock() + blocks = stub.kv_cache_manager.allocate_slots.return_value + stub.kv_cache_manager.allocate_slots.side_effect = [blocks, None] + stub.kv_cache_manager.take_kv_cache_block_copies.return_value = ( + [(31, 32)], + ["fallback-retained"], + ) + + assert stub._bench_make_steady_step() is None + stub.kv_cache_manager.block_pool.free_blocks.assert_called_once_with( + ["fallback-retained"] + ) + stub._free_cow_retained_blocks.assert_not_called() + assert stub._bench_active_req_ids == {request.request_id for request in requests} + + +@pytest.mark.core +def test_kvwarm_native_decode_dispatch_uses_real_continuation_and_provenance(): + stub, requests = _kvwarm_native_resume_stub(False) + point = BenchmarkPoint( + point_type="decode", benchmark_id=7, batch_size=2, total_kv_read_tokens=9 + ) + for request, context in zip(requests, (4, 3), strict=True): + request.num_computed_tokens = context + request.all_token_ids = list(range(context + 1)) + stub._kvwarm_chain_prompts[request.request_id] = list(range(context)) + stub._kvwarm_stage_point = point + stub._kvwarm_stage_batch = 2 + stub._kvwarm_building = False + stub._kvwarm_plan = {stub._kvwarm_plan_key(point): 4} + stub._bench_grid = deque([point]) + stub._bench_current_point = None + stub._bench_drain_pending = False + stub._bench_frees_pending = lambda: False + stub._bench_stop_at_timeout_boundary = lambda phase: False + stub.max_model_len = 8192 + + output = stub._bench_step_decode() + + assert output.scheduled_new_reqs == [] + assert output.scheduled_cached_reqs.num_computed_tokens == [4, 3] + assert stub._bench_admission_kv_tokens == 7 + assert stub._bench_expected_fpms == 4 # Discard admission, measure three steps. + assert stub._kvwarm_seed_regime(stub._bench_current_point) == "real_kv" + assert stub._kvwarm_meta["points_real_kv"] == 1 + + +@pytest.mark.core +@pytest.mark.parametrize("context,repeats", [(124, 3), (125, 2), (126, 1)]) +def test_kvwarm_native_dp_dispatch_caps_repeats_by_negotiated_model_length( + context, repeats, monkeypatch +): + monkeypatch.setenv("DYN_BENCH_KV_WARMUP", "on") + capacities = [_benchmark_capacity(max_model_len=limit) for limit in (128, 256)] + negotiated = instrumented_scheduler_module._BenchmarkCapacityEnvelope.common( + capacities + ) + point = BenchmarkPoint( + point_type="decode", + benchmark_id=1, + batch_size=2, + total_kv_read_tokens=2 * context, + ) + dispatches = [] + for capacity in capacities: + stub, _ = _kvwarm_native_resume_stub(False) + stub.max_model_len = capacity.max_model_len + stub._bench_negotiated_capacity = negotiated + stub.requests = {} + for req_id in stub._kvwarm_chain_ids: + tokens = [1] * (context - 1) + request = Request( + request_id=req_id, + prompt_token_ids=tokens, + sampling_params=SamplingParams(max_tokens=8, ignore_eos=True), + pooling_params=None, + ) + request.num_computed_tokens = len(tokens) + request.append_output_token_ids(17) + stub.requests[req_id] = request + stub._kvwarm_chain_prompts[req_id] = tokens + stub._kvwarm_stage_point = point + stub._kvwarm_stage_batch = point.batch_size + stub._kvwarm_building = False + stub._kvwarm_plan = {stub._kvwarm_plan_key(point): context - 1} + stub._bench_grid = deque([point]) + stub._bench_current_point = None + stub._bench_drain_pending = False + stub._bench_point_deadline = 0.0 + stub._bench_frees_pending = lambda: False + stub._bench_stop_at_timeout_boundary = lambda phase: False + stub._bench_save_current_point = MagicMock() + stub._bench_finish_requests = MagicMock() + stub._bench_transition_to_timeout_done = lambda: False + stub._schedule_times = deque() + + forwards = [] + # Drive admission, steady forwards, and entry into result collection + # with real vLLM Request and SchedulerOutput types, without GPU execution. + for _ in range(5): + output = stub._bench_step_decode() + if output is None: + stub._bench_save_current_point.assert_called_once_with() + assert stub._bench_drain_pending + break + stub._bench_save_current_point.assert_not_called() + assert output.total_num_scheduled_tokens == point.batch_size + forwards.append(output.scheduled_cached_reqs.num_computed_tokens) + for request in stub.running: + request.num_computed_tokens += 1 + request.append_output_token_ids(17) + stub._bench_current_fpms.append({}) + else: + pytest.fail("decode point did not reach result collection") + dispatches.append(forwards) + + expected = [[context - 1 + step] * point.batch_size for step in range(1 + repeats)] + assert dispatches == [expected, expected] + + +@pytest.mark.core +def test_kvwarm_native_capacity_fallback_saves_available_steady_samples(monkeypatch): + monkeypatch.setenv("DYN_BENCH_KV_WARMUP", "on") + monkeypatch.setenv("DYN_BENCH_GIANT_KV_THRESHOLD", "0") + monkeypatch.setenv("DYN_BENCH_GIANT_KV_REPEATS", "3") + point = BenchmarkPoint( + point_type="decode", benchmark_id=1, batch_size=1, total_kv_read_tokens=14 + ) + stub = _benchmark_save_stub(point, []) + specs = [ + FullAttentionSpec( + block_size=16, num_kv_heads=1, head_size=8, dtype=torch.bfloat16 + ), + SlidingWindowSpec( + block_size=16, + num_kv_heads=1, + head_size=8, + dtype=torch.bfloat16, + sliding_window=512, + ), + ] + config = KVCacheConfig( + num_blocks=3, # The null block leaves two usable blocks. + kv_cache_tensors=[], + kv_cache_groups=[ + KVCacheGroupSpec(layer_names=[str(index)], kv_cache_spec=spec) + for index, spec in enumerate(specs) + ], + ) + stub.kv_cache_manager = kv_cache_manager.KVCacheManager( + config, + max_model_len=128, + scheduler_block_size=16, + hash_block_size=16, + max_in_flight_tokens=16, + enable_caching=False, + ) + stub.max_model_len = 128 + stub.num_lookahead_tokens = 0 + stub.cache_config = SimpleNamespace(enable_prefix_caching=False, block_size=16) + stub.vllm_config = SimpleNamespace( + model_config=SimpleNamespace(hf_config=SimpleNamespace(num_experts=8)), + parallel_config=SimpleNamespace(enable_expert_parallel=False), + ) + stub._kvwarm_probe_content = lambda: None + stub._bench_random_kda = False + stub._bench_grid = deque([point]) + stub._bench_active_req_ids = set() + stub._bench_current_point = None + stub._bench_drain_pending = False + stub.deferred_frees = deque() + stub._bench_stop_at_timeout_boundary = lambda phase: False + stub._bench_seq = 0 + stub._bench_block_hasher = None + stub._bench_synthetic_token_ids = lambda salt, length: [1] * length + stub.requests = {} + stub.running = [] + stub.finished_req_ids = set() + stub.connector = None + stub.ec_connector = None + stub.needs_kv_cache_zeroing = False + stub.use_v2_model_runner = False + stub.sched_step_seq = 0 + stub._bench_point_deadline = 0.0 + stub._schedule_times = deque() + stub._bench_transition_to_timeout_done = lambda: False + + def finish_requests(req_ids): + for req_id in req_ids: + stub.kv_cache_manager.free(stub.requests.pop(req_id)) + stub.running = [] + + stub._bench_finish_requests = finish_requests + stub._kvwarm_prepare("decode") + assert stub._kvwarm_native + assert not stub._kvwarm_plan_covers(point) + + # Real CPU allocator and request bookkeeping; only worker outputs are + # supplied here. The final repeat crosses the 16-token block boundary. + for context in (13, 14, 15): + output = stub._bench_step_decode() + assert output is not None + contexts = ( + [request.num_computed_tokens for request in output.scheduled_new_reqs] + if output.scheduled_new_reqs + else output.scheduled_cached_reqs.num_computed_tokens + ) + assert contexts == [context] + for request in stub.running: + request.num_computed_tokens += 1 + request.append_output_token_ids(7) + stub._bench_current_fpms.append( + { + "counter_id": point.benchmark_id, + "dp_rank": 0, + "wall_time": 1.0, + "scheduled_requests": { + "num_decode_requests": 1, + "sum_decode_kv_tokens": context, + }, + } + ) + + assert stub._bench_expected_fpms == 4 + assert stub.kv_cache_manager.block_pool.get_num_free_blocks() == 0 + assert stub._bench_step_decode() is None + assert stub._bench_results == [] # Wait for the normal result deadline. + stub._bench_point_deadline = time.monotonic() - 1 + assert stub._bench_step_decode() is None + assert len(stub._bench_results) == 1 + result = stub._bench_results[0] + assert result.fpms[0]["kvwarm_giant_median_of"] == 2 + assert stub._kvwarm_seed_regime(result.point) == "fake_fallback" + assert stub._bench_skipped_points == [] + assert stub._bench_active_req_ids == set() + assert stub.requests == {} + assert stub.kv_cache_manager.block_pool.get_num_free_blocks() == 2 + + def _kvwarm_sharegpt_file(tmp_path, bodies): """Write a ShareGPT-shaped dump whose conversations all land in the collection half of the hash split (``_kvwarm_load_texts`` keeps only @@ -5273,6 +6115,7 @@ def _kvwarm_planner_stub(usable_blocks, groups=1, block_size=16): stub._bench_synchronizer = None stub._bench_negotiated_capacity = None stub.max_model_len = 8192 + stub.num_lookahead_tokens = 0 stub.cache_config = SimpleNamespace(block_size=block_size) stub.kv_cache_manager = SimpleNamespace( coordinator=SimpleNamespace( diff --git a/docs/fern/pages/reference/observability/environment-variables.mdx b/docs/fern/pages/reference/observability/environment-variables.mdx index a46d90995c03..2e8c43852670 100644 --- a/docs/fern/pages/reference/observability/environment-variables.mdx +++ b/docs/fern/pages/reference/observability/environment-variables.mdx @@ -291,9 +291,13 @@ See [Forward Pass Metrics Trace Reference](forward-pass-metrics-traces.mdx) for These variables tune the vLLM worker's self-benchmark collector (`--benchmark-mode`), which measures whole-model forward-pass latencies for the FPM performance database. They have no effect on serving. - KV warm-up master switch. When on, eligible decode benchmark points read attention KV produced by real prefill chains. By default, eligibility requires a MoE model with expert parallelism and no recurrent-state layers. With `DYN_BENCHMARK_RANDOMIZE_KDA_STATE=true`, hybrid MoE models are also eligible without expert parallelism; their attention KV is real and their recurrent states remain synthetic. Set `off` to force the legacy synthetic-KV path. + KV warm-up master switch. When on, eligible decode benchmark points read attention KV produced by real prefill. Set `off` to force the legacy synthetic-KV path. - All eligible configurations require prefix caching, a loadable tokenizer, and enough seeding dataset content for the model length. Configurations that do not meet these requirements skip warm-up and record the reason. Under attention data parallelism the decision is group-wide: if any rank is ineligible, every rank uses synthetic KV; if any rank's warm-up stage fails to build or cannot reserve its shadow blocks, every rank falls back to synthetic KV for that stage. + When the runtime cache layout contains at least one finite `SlidingWindowSpec` and only `FullAttentionSpec` or `SlidingWindowSpec` groups, warm-up uses native prefill at each point's exact admission context. The same live requests then continue into decode, preserving their private convolution state instead of borrowing state from a deeper context. This covers Inkling's convolution cache layout and applies to dense models and MoE models with or without expert parallelism. It does not require prefix caching. Each point rebuilds its state, adding untimed collection work. The artifact records `kvwarm.initialization_strategy = "native_exact_context"`. + + Other eligible configurations use shared prefill chains and require prefix caching. By default, chain eligibility requires a MoE model with expert parallelism and no recurrent-state layers. With `DYN_BENCHMARK_RANDOMIZE_KDA_STATE=true`, hybrid MoE models are also eligible without expert parallelism; their attention KV is real and their recurrent states remain synthetic. This option disables the native exact-context strategy. + + Both strategies require a loadable tokenizer, enough seeding dataset content for the model length, and sufficient cache capacity. Configurations that do not meet these requirements skip warm-up and record the reason. Under attention data parallelism the decision is group-wide: if any rank is ineligible, every rank uses synthetic KV; if any rank's warm-up stage fails, every rank falls back for that stage. Failed initialization never marks synthetic state as `real_kv`. @@ -319,11 +323,11 @@ These variables tune the vLLM worker's self-benchmark collector (`--benchmark-mo - Total-KV-read threshold (tokens) above which a fake-injected decode point is measured with repeated steady steps. Points served from the real-KV warm-up chains always repeat. Set `0` to apply repeat protection to every point. + Total-KV-read threshold (tokens) at or above which fake-injected decode points use `DYN_BENCH_GIANT_KV_REPEATS`. Real-KV warm-up points use that count regardless of this threshold, subject to model-length headroom. Set `0` to apply it to every point. - Steady-step repeat count for real-KV points and for fake-injected points above the threshold; the recorded latency is the median. The warm-up chains reserve this many steady writes per request. + Requested steady-step count for real-KV points and for fake-injected points at or above the threshold; the recorded latency is the median. The count is capped by model-length headroom, using the negotiated limit under attention data parallelism, so points may have a single sample. Native warm-up reserves the capped count; shared prefill chains reserve the configured count. Fake-injected points also fall back to one sample when cache capacity cannot fit the repeats. From 59bdde05c3b053326b6c1ddbd24e844bb7f4c049 Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Fri, 2 Oct 2026 07:47:58 +0800 Subject: [PATCH 02/14] fix(vllm): declare the warm-up plan type for mypy The native plan is keyed per execution and the shared-chain plan per batch rung. mypy inferred the native key type for both branches of `_kvwarm_prepare` and a tuple-keyed type for `self._kvwarm_plan`, then rejected the shared-chain annotation of `plan` as a redefinition (11 errors after the rebase onto main). Declare `_kvwarm_plan: dict` on the class and annotate the first `plan` binding instead. No runtime change. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- components/src/dynamo/vllm/instrumented_scheduler.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index fbe344102582..24d9260c200c 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -4806,6 +4806,8 @@ def _bench_step_prefill(self) -> SchedulerOutput | None: _kvwarm_native_resume_ids: set[str] | None = None _kvwarm_stage_t0: float | None _kvwarm_stage_batch: int | None + # Stage depth by ``_kvwarm_plan_key``: a batch rung, or one native execution. + _kvwarm_plan: dict # Local outcome ``(batch, ok, detail)`` of the active stage once this # rank's build has closed, held until the group's round says every rank # is done (``_kvwarm_stage_round``). @@ -5280,7 +5282,7 @@ def _kvwarm_prepare(self, mode: str) -> None: # ``_kvwarm_register_shadow``, otherwise a covered point's shadow can # need one block more than its chain holds. if self._kvwarm_native: - plan = {} + plan: dict = {} for point in decode_pts: contexts = self._bench_decode_context_lengths( point.total_kv_read_tokens, point.batch_size @@ -5309,7 +5311,7 @@ def _kvwarm_prepare(self, mode: str) -> None: else: repeats = self._kvwarm_giant_repeats() margin = 1 + repeats - plan: dict = {} + plan = {} rung_ctxs: dict = {} for p in decode_pts: ctxs = self._bench_decode_context_lengths( From 609c223e983776fa7ed9c6aa6b0dfca616bb6463 Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Fri, 2 Oct 2026 07:52:50 +0800 Subject: [PATCH 03/14] fix(vllm): keep native plan numbering valid when attention-DP drops points Main's attention-DP filter in `_kvwarm_prepare` removes decode points the warm-up plan cannot cover. `_bench_build_grid` then rebased every native plan key and capacity fallback through `public_ids`, so a native layout under attention-DP with one uncovered point failed grid building with a KeyError. Rebase only the executions left in the grid and record a dropped point's capacity fallback with a null `benchmark_id`. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../src/dynamo/vllm/instrumented_scheduler.py | 6 ++- .../tests/test_vllm_instrumented_scheduler.py | 52 +++++++++++++++++++ 2 files changed, 57 insertions(+), 1 deletion(-) diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index 24d9260c200c..e12185294429 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -2777,12 +2777,16 @@ def _bench_build_grid(self) -> None: if self._kvwarm_native and getattr(self, "_kvwarm_plan", None): # Rebase from the old dictionary in one pass: updating it in place # could overwrite another execution's key during the permutation. + # Attention-DP removes uncovered points from the grid + # (``_kvwarm_prepare``): drop their plan entries and keep their + # capacity fallbacks without a public ID. self._kvwarm_plan = { (public_ids[execution_id], batch, kv_tokens): depth for (execution_id, batch, kv_tokens), depth in self._kvwarm_plan.items() + if execution_id in public_ids } for fallback in self._kvwarm_meta.get("capacity_fallbacks", []): - fallback["benchmark_id"] = public_ids[fallback["benchmark_id"]] + fallback["benchmark_id"] = public_ids.get(fallback["benchmark_id"]) grid_payload = json.dumps( [asdict(point) for point in self._bench_grid], sort_keys=True, diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index abf107e469e2..0c743123c22b 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -6790,6 +6790,58 @@ def test_kvwarm_dp_filter_marks_decode_missing_when_nothing_is_covered(monkeypat assert stub._bench_missing_phases == ["decode"] +@pytest.mark.core +def test_kvwarm_dp_filter_rebases_native_plan_without_dropped_points(monkeypatch): + """Attention-DP removes the native points the plan cannot cover before + ``_bench_build_grid`` assigns public IDs, so only the remaining executions + keep plan entries, and a dropped point's capacity fallback has no ID.""" + monkeypatch.setenv("DYN_BENCH_KV_WARMUP", "on") + monkeypatch.setenv("DYN_BENCH_GIANT_KV_REPEATS", "3") + stub = InstrumentedScheduler.__new__(InstrumentedScheduler) + _install_test_capacity_preflight( + stub, _benchmark_capacity(usable_blocks_without_watermark=4) + ) + stub._bench_config = BenchmarkConfig(mode="decode") + stub._bench_explicit_points = None + stub._bench_grid = deque() + stub._bench_grid_built = False + stub._bench_missing_phases = [] + stub._bench_dp_size = 2 + stub.num_lookahead_tokens = 0 + stub._bench_blocks_per_req = lambda tokens, **_: -(-tokens // 16) + stub._kvwarm_warm_eligible = lambda: True + stub._kvwarm_native = True + points = [ + BenchmarkPoint( + point_type="decode", + benchmark_id=0, + batch_size=2, + total_kv_read_tokens=2 * context, + ) + for context in (31, 17, 2) + ] + stub._bench_generate_decode_grid = lambda: stub._bench_grid.extend(points) + stub._bench_eager_warmup_points = lambda: [] + + stub._bench_build_grid() + + grid = list(stub._bench_grid) + # Context 31 needs six blocks with its repeated writes; the pool has four. + assert [(p.benchmark_id, p.total_kv_read_tokens) for p in grid] == [(1, 34), (2, 4)] + assert stub._bench_expected_points == 2 + assert sorted(stub._kvwarm_plan) == [(1, 2, 34), (2, 2, 4)] + assert all(stub._kvwarm_plan_covers(point) for point in grid) + assert stub._kvwarm_meta["capacity_fallbacks"] == [ + { + "benchmark_id": None, + "batch": 2, + "depth": 30, + "required_blocks": 6, + "usable_blocks": 4, + } + ] + + def test_giant_fake_off_by_batch_correction_requires_a_steady_sample(): """The admission step also measures ``declared - batch``; only a recorded steady sample (``kvwarm_steady_sample``, set by both save paths) may be accepted at the From b5d25b83d9adf832206d02dcd2fc8016525e19e4 Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Fri, 2 Oct 2026 08:47:56 +0800 Subject: [PATCH 04/14] fix(vllm): keep dense attention-only models on synthetic KV Native exact-context warm-up ran before the dense-model check, so dense sliding-window models (Gemma-3) paid a full prefill per point although their decode timing does not depend on KV content. Use native only when the model has experts or recurrent-state layers; dense attention-only models keep the dense_model_content_insensitive skip. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../src/dynamo/vllm/instrumented_scheduler.py | 11 ++++++++++- .../tests/test_vllm_instrumented_scheduler.py | 16 +++++++++++++--- 2 files changed, 23 insertions(+), 4 deletions(-) diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index e12185294429..8ccdb51b85b0 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -4887,6 +4887,10 @@ def _kvwarm_native_layout(self) -> bool: needs. Inkling's convolution cache is a SlidingWindowSpec too, so native allocation and forward execution initialize all four streams without assuming a tensor layout or borrowing writable state. + + The layout check is necessary but not sufficient: the gate in + ``_kvwarm_warm_eligible`` also requires experts or recurrent-state + layers, because dense attention-only models are content-insensitive. """ groups = self.kv_cache_manager.kv_cache_config.kv_cache_groups specs = [group.kv_cache_spec for group in groups] @@ -5006,7 +5010,12 @@ def _kvwarm_warm_eligible(self) -> bool: False, ) ) - self._kvwarm_native = self._kvwarm_native_layout() + # Native prefill only pays off when decode timing depends on KV + # content: expert routing, or recurrent state that synthetic KV + # cannot reproduce. Dense attention-only models stay synthetic. + self._kvwarm_native = self._kvwarm_native_layout() and ( + has_experts or bool(self._kvwarm_state_layer_groups()) + ) if self._kvwarm_native: reason = self._kvwarm_probe_content() eligible = reason is None diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index 0c743123c22b..16d419614ede 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -5081,9 +5081,7 @@ def _kvwarm_native_gate_stub(*, experts=8, ep=False, prefix=False): @pytest.mark.core -@pytest.mark.parametrize( - "experts,ep,prefix", [(8, False, False), (8, True, True), (0, False, False)] -) +@pytest.mark.parametrize("experts,ep,prefix", [(8, False, False), (8, True, True)]) def test_kvwarm_sliding_window_uses_native_prefills_without_ep_or_prefix_cache( experts, ep, prefix ): @@ -5094,6 +5092,18 @@ def test_kvwarm_sliding_window_uses_native_prefills_without_ep_or_prefix_cache( assert stub._kvwarm_meta["initialization_strategy"] == "native_exact_context" +@pytest.mark.core +def test_kvwarm_native_skips_dense_attention_only_layouts(): + # Dense attention-only models are content-insensitive: synthetic KV is + # correct by construction, so native prefill would add cost and no fidelity. + stub = _kvwarm_native_gate_stub(experts=0, ep=False, prefix=False) + + assert not stub._kvwarm_warm_eligible() + assert not stub._kvwarm_native + assert stub._kvwarm_meta["skip_reason"] == "dense_model_content_insensitive" + assert "initialization_strategy" not in stub._kvwarm_meta + + @pytest.mark.core def test_kvwarm_sliding_window_does_not_admit_mamba_state(): stub = _kvwarm_native_gate_stub(ep=True, prefix=True) From 1338e23d9220fffa0b0cbfdde3b995e26f227a6c Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Fri, 2 Oct 2026 09:22:45 +0800 Subject: [PATCH 05/14] fix(vllm): warm Mamba and KDA state layouts with native prefills Hybrid layouts with recurrent-state groups fell back to random state or skipped the warm-up, and #14614's live-state mode borrows a deeper chain's state. Native exact-context prefill computes the state exactly and needs no shadows, so admit Mamba/KDA groups next to full attention and finite windows when random-KDA is off. Native takes precedence over --benchmark-hybrid-live-state, which is logged as unused. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../src/dynamo/vllm/instrumented_scheduler.py | 53 +++++--- .../tests/test_vllm_instrumented_scheduler.py | 125 +++++++++++++++++- .../observability/environment-variables.mdx | 4 +- 3 files changed, 154 insertions(+), 28 deletions(-) diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index 8ccdb51b85b0..abbc6b15704d 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -4880,29 +4880,38 @@ def _kvwarm_state_layer_groups(self) -> list[str]: names.append(spec_name) return names + @staticmethod + def _kvwarm_is_state_spec(spec) -> bool: + """Recurrent (Mamba/KDA) state group; same predicate as live-state mode.""" + return isinstance(spec, MambaSpec) or "Mamba" in type(spec).__name__ + def _kvwarm_native_layout(self) -> bool: - """Finite windows must continue their exact native prefill requests. + """Layouts whose decode state must come from each point's own prefill. A deep parked chain may have evicted the history a shallower point - needs. Inkling's convolution cache is a SlidingWindowSpec too, so - native allocation and forward execution initialize all four streams - without assuming a tensor layout or borrowing writable state. - - The layout check is necessary but not sufficient: the gate in - ``_kvwarm_warm_eligible`` also requires experts or recurrent-state - layers, because dense attention-only models are content-insensitive. + needs (finite windows, including Inkling's convolution cache), and a + recurrent state cannot be borrowed from a chain at another depth. + Native allocation and forward execution build both exactly. Necessary, + not sufficient: the gate also requires experts or state layers. """ - groups = self.kv_cache_manager.kv_cache_config.kv_cache_groups - specs = [group.kv_cache_spec for group in groups] + specs = [ + group.kv_cache_spec + for group in self.kv_cache_manager.kv_cache_config.kv_cache_groups + ] + + def finite_window(spec) -> bool: + return isinstance(spec, SlidingWindowSpec) and spec.sliding_window > 0 + return ( not self._bench_random_kda and any( - isinstance(spec, SlidingWindowSpec) and spec.sliding_window > 0 + finite_window(spec) or self._kvwarm_is_state_spec(spec) for spec in specs ) and all( - isinstance(spec, (FullAttentionSpec, SlidingWindowSpec)) - and (not isinstance(spec, SlidingWindowSpec) or spec.sliding_window > 0) + isinstance(spec, FullAttentionSpec) + or finite_window(spec) + or self._kvwarm_is_state_spec(spec) for spec in specs ) ) @@ -4963,12 +4972,14 @@ def _kvwarm_release_heavy_state(self) -> None: setattr(self, attr, None) def _kvwarm_warm_eligible(self) -> bool: - """Select native finite-window or shared-prefix attention-KV warm-up. + """Select native exact-context or shared-prefix attention-KV warm-up. Random-state mode also admits hybrid MoE without EP: its attention prefixes are real, while recurrent states remain private and synthetic. - Finite-window layouts instead prefill and continue each point's own - requests; this needs neither expert parallelism nor prefix caching. + Native warm-up instead prefills and continues each point's own requests + for models with experts or recurrent-state layers whose layout + qualifies (``_kvwarm_native_layout``); this needs neither expert + parallelism nor prefix caching. The verdict travels in the capacity envelope (see ``_bench_make_local_capacity``), so every host-local input the stage @@ -5010,13 +5021,17 @@ def _kvwarm_warm_eligible(self) -> bool: False, ) ) - # Native prefill only pays off when decode timing depends on KV - # content: expert routing, or recurrent state that synthetic KV - # cannot reproduce. Dense attention-only models stay synthetic. + # Native for MoE with or without EP (keeps Inkling's exact convolution + # state) or recurrent-state layers; dense attention-only stays synthetic. self._kvwarm_native = self._kvwarm_native_layout() and ( has_experts or bool(self._kvwarm_state_layer_groups()) ) if self._kvwarm_native: + if self._bench_hybrid_live_state: + logger.info( + "KVWARM: native exact-context warm-up computes the " + "recurrent state; --benchmark-hybrid-live-state is unused" + ) reason = self._kvwarm_probe_content() eligible = reason is None meta["initialization_strategy"] = "native_exact_context" diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index 16d419614ede..5f8f16388237 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -14,6 +14,7 @@ import hashlib import json +import logging import threading import time import uuid @@ -5105,15 +5106,87 @@ def test_kvwarm_native_skips_dense_attention_only_layouts(): @pytest.mark.core -def test_kvwarm_sliding_window_does_not_admit_mamba_state(): - stub = _kvwarm_native_gate_stub(ep=True, prefix=True) +def test_kvwarm_native_reads_experts_from_the_text_config(): + # Inkling's HF config: no top-level expert keys, text_config.n_routed_experts. + stub = _kvwarm_native_gate_stub(ep=False, prefix=False) + stub.vllm_config.model_config = SimpleNamespace( + hf_config=SimpleNamespace(), + hf_text_config=SimpleNamespace(n_routed_experts=256), + ) + + assert stub._kvwarm_warm_eligible() + assert stub._kvwarm_native + + +def _with_state_group(stub): stub.kv_cache_manager.kv_cache_config.kv_cache_groups.append( SimpleNamespace(kv_cache_spec=type("MambaSpec", (), {})()) ) + return stub + + +@pytest.mark.core +@pytest.mark.parametrize("experts", [8, 0]) +def test_kvwarm_native_admits_recurrent_state_layouts(experts): + stub = _with_state_group(_kvwarm_native_gate_stub(experts=experts, ep=False)) + + assert stub._kvwarm_warm_eligible() + assert stub._kvwarm_native + assert stub._kvwarm_meta["initialization_strategy"] == "native_exact_context" + + +@pytest.mark.core +def test_kvwarm_native_admits_full_attention_plus_state_without_windows(): + # GLM-5.3-Flash / Qwen3-Next shape: full (or MLA) attention + KDA/Mamba state. + stub = _kvwarm_gate_stub(experts=8, ep=False, prefix=False) + stub.kv_cache_manager.kv_cache_config.kv_cache_groups = [ + SimpleNamespace( + kv_cache_spec=FullAttentionSpec( + block_size=16, num_kv_heads=1, head_size=8, dtype=torch.bfloat16 + ) + ), + SimpleNamespace(kv_cache_spec=type("MambaSpec", (), {})()), + ] + + assert stub._kvwarm_warm_eligible() + assert stub._kvwarm_native + + +@pytest.mark.core +def test_kvwarm_random_kda_keeps_shared_chain_for_state_layouts(): + stub = _with_state_group(_kvwarm_native_gate_stub(experts=8, ep=True, prefix=True)) + stub._bench_random_kda = True + + stub._kvwarm_warm_eligible() + + assert not stub._kvwarm_native + assert "initialization_strategy" not in stub._kvwarm_meta + + +@pytest.mark.core +def test_kvwarm_native_wins_over_hybrid_live_state(caplog): + stub = _with_state_group(_kvwarm_native_gate_stub(experts=8, ep=True, prefix=True)) + stub._bench_hybrid_live_state = True + + with caplog.at_level( + logging.INFO, logger=instrumented_scheduler_module.logger.name + ): + assert stub._kvwarm_warm_eligible() + + assert stub._kvwarm_native + assert any("live-state" in record.getMessage() for record in caplog.records) + + +@pytest.mark.core +def test_kvwarm_native_rejects_unsupported_spec_with_state(): + stub = _with_state_group(_kvwarm_native_gate_stub(experts=8, ep=True, prefix=True)) + stub.kv_cache_manager.kv_cache_config.kv_cache_groups.append( + SimpleNamespace(kv_cache_spec=type("ChunkedLocalAttentionSpec", (), {})()) + ) + + stub._kvwarm_warm_eligible() - assert not stub._kvwarm_warm_eligible() assert not stub._kvwarm_native - assert stub._kvwarm_meta["skip_reason"] == "hybrid_state_layers_unsupported" @pytest.mark.core @@ -5436,15 +5509,28 @@ def _kvwarm_native_resume_stub(use_v2): @pytest.mark.core +@pytest.mark.parametrize("with_state", [False, True]) @pytest.mark.parametrize("use_v2", [False, True]) def test_kvwarm_native_resume_continues_private_requests_and_preserves_runner_state( - use_v2, + use_v2, with_state ): stub, requests = _kvwarm_native_resume_stub(use_v2) + tables = {"chain-a": ([1, 11], [2, 21]), "chain-b": ([3, 12], [4, 22])} + delta = ([11], [21]) + if with_state: + # Second group is a Mamba/KDA state group holding one live state block + # per request; this decode step reuses it, so its delta is empty. + tables = {"chain-a": ([1, 11], [31]), "chain-b": ([3, 12], [32])} + delta = ([11], []) + stub.kv_cache_manager.get_block_ids.side_effect = tables.__getitem__ + blocks = stub.kv_cache_manager.allocate_slots.return_value + blocks.get_block_ids.return_value = delta output = stub._kvwarm_resume_native() assert output.total_num_scheduled_tokens == 2 + # Cached path only: runners rebuild request state just for new requests. + assert output.scheduled_new_reqs == [] assert stub.running == requests assert all(stub.requests[request.request_id] is request for request in requests) assert stub._bench_active_req_ids == {"chain-a", "chain-b"} @@ -5456,11 +5542,13 @@ def test_kvwarm_native_resume_continues_private_requests_and_preserves_runner_st request.request_id: request.all_token_ids for request in requests } if use_v2: + # V2 keeps the parked tables and accepts only the appended delta. assert cached.resumed_req_ids == set() - assert cached.new_block_ids == [([11], [21]), ([11], [21])] + assert cached.new_block_ids == [delta, delta] else: + # V1 re-adds the parked requests: every group's full table, unchanged. assert cached.resumed_req_ids == {"chain-a", "chain-b"} - assert cached.new_block_ids == [([1, 11], [2, 21]), ([3, 12], [4, 22])] + assert cached.new_block_ids == [tables["chain-a"], tables["chain-b"]] assert output.kv_cache_block_copies == [(8, 9)] stub._free_cow_retained_blocks.assert_called_once_with(["retained"], 11) assert stub._kvwarm_native_resume_ids == set() @@ -5609,6 +5697,29 @@ def test_kvwarm_native_capacity_uses_sliding_window_admission_caps(): assert stub._kvwarm_native_required_blocks([511]) == 37 +@pytest.mark.core +def test_kvwarm_native_capacity_counts_mamba_align_state_blocks(): + stub = InstrumentedScheduler.__new__(InstrumentedScheduler) + full = SimpleNamespace(block_size=16) + state = SimpleNamespace( + block_size=16, + mamba_cache_mode="align", + num_speculative_blocks=0, + kv_cache_spec=SimpleNamespace(num_prefill_checkpoint_blocks=1), + ) + stub.kv_cache_manager = SimpleNamespace( + coordinator=SimpleNamespace(single_type_managers=[full, state]) + ) + stub.block_size = 16 + stub.num_lookahead_tokens = 0 + stub._kvwarm_giant_repeats = lambda: 3 + stub._bench_capacity_limit = lambda name: 4096 + + # Per request: ceil((ctx + 1 + 3) / 16) full-attention blocks + 3 align + # state blocks (2 + 0 speculative + 1 checkpoint). + assert stub._kvwarm_native_required_blocks([31, 12]) == (3 + 3) + (1 + 3) + + @pytest.mark.core @pytest.mark.parametrize("after_admission", [False, True]) @pytest.mark.parametrize("allocation_raises", [False, True]) diff --git a/docs/fern/pages/reference/observability/environment-variables.mdx b/docs/fern/pages/reference/observability/environment-variables.mdx index 2e8c43852670..b91f55a9f65b 100644 --- a/docs/fern/pages/reference/observability/environment-variables.mdx +++ b/docs/fern/pages/reference/observability/environment-variables.mdx @@ -293,9 +293,9 @@ These variables tune the vLLM worker's self-benchmark collector (`--benchmark-mo KV warm-up master switch. When on, eligible decode benchmark points read attention KV produced by real prefill. Set `off` to force the legacy synthetic-KV path. - When the runtime cache layout contains at least one finite `SlidingWindowSpec` and only `FullAttentionSpec` or `SlidingWindowSpec` groups, warm-up uses native prefill at each point's exact admission context. The same live requests then continue into decode, preserving their private convolution state instead of borrowing state from a deeper context. This covers Inkling's convolution cache layout and applies to dense models and MoE models with or without expert parallelism. It does not require prefix caching. Each point rebuilds its state, adding untimed collection work. The artifact records `kvwarm.initialization_strategy = "native_exact_context"`. + Models whose decode timing depends on KV content (MoE models, or models with recurrent-state layers) use native prefill when every cache group is `FullAttentionSpec`, a finite `SlidingWindowSpec`, or a Mamba/KDA state group, and at least one group is a finite window or a state group: each point prefills its own requests to the exact admission context, and the same requests continue into decode, so window history and recurrent state are computed by the model instead of borrowed from a chain at another depth. This covers Inkling's convolution cache and hybrid models such as GLM-5.3-Flash. It does not require expert parallelism or prefix caching, and `--benchmark-hybrid-live-state` has no effect when it applies. Dense attention-only models keep synthetic KV, which is correct by construction for them. The artifact records `kvwarm.initialization_strategy = "native_exact_context"`. - Other eligible configurations use shared prefill chains and require prefix caching. By default, chain eligibility requires a MoE model with expert parallelism and no recurrent-state layers. With `DYN_BENCHMARK_RANDOMIZE_KDA_STATE=true`, hybrid MoE models are also eligible without expert parallelism; their attention KV is real and their recurrent states remain synthetic. This option disables the native exact-context strategy. + Other eligible configurations use shared prefill chains and require prefix caching. By default, chain eligibility requires a MoE model with expert parallelism and no recurrent-state layers. With `DYN_BENCHMARK_RANDOMIZE_KDA_STATE=true`, hybrid MoE models keep the shared-chain random-state strategy, also without expert parallelism: their attention KV is real and their recurrent states remain synthetic. This option disables the native exact-context strategy, even for layouts that qualify for it. Both strategies require a loadable tokenizer, enough seeding dataset content for the model length, and sufficient cache capacity. Configurations that do not meet these requirements skip warm-up and record the reason. Under attention data parallelism the decision is group-wide: if any rank is ineligible, every rank uses synthetic KV; if any rank's warm-up stage fails, every rank falls back for that stage. Failed initialization never marks synthetic state as `real_kv`. From 5a0456e2941907309100dcaa2562f9f9ccfbbf82 Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Fri, 2 Oct 2026 17:03:00 +0800 Subject: [PATCH 06/14] fix(vllm): reserve align headroom for native stages, fix live-state help An align-mode Mamba/KDA prefill that ends on a hash boundary inside a Mamba block registers its own partial tail, and the first decode's admission check then asks for one block more than it allocates. With vLLM's default zero watermark, a stage planned at exactly the pool edge raised instead of falling back, so native stages reserve one block when any group runs in align mode. The --benchmark-hybrid-live-state help now says the flag applies only to layouts that do not qualify for native exact-context warm-up. Main's recurrent-state gate test uses a layout native cannot take, so it keeps pinning hybrid_state_layers_unsupported, and the gate comment and the state-group docstring say why TP-only MoE and state layers stay native. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- components/src/dynamo/vllm/backend_args.py | 5 ++- .../src/dynamo/vllm/instrumented_scheduler.py | 29 +++++++++---- .../tests/test_vllm_instrumented_scheduler.py | 41 ++++++++++++++----- 3 files changed, 55 insertions(+), 20 deletions(-) diff --git a/components/src/dynamo/vllm/backend_args.py b/components/src/dynamo/vllm/backend_args.py index 470aa33a82bc..422d230fd995 100644 --- a/components/src/dynamo/vllm/backend_args.py +++ b/components/src/dynamo/vllm/backend_args.py @@ -490,8 +490,9 @@ def add_arguments(self, parser) -> None: "recurrent-state groups forked from the parked chain's live state block " "instead of skipping the warm-up. Attention KV is the chain's real prefix; " "the recurrent state is a valid but deeper-context state (shallow points " - "read a few percent fast). Mutually exclusive with " - "--benchmark-randomize-kda-state." + "read a few percent fast). Applies only to layouts that do not qualify " + "for native exact-context warm-up, which takes precedence. Mutually " + "exclusive with --benchmark-randomize-kda-state." ), ) add_argument( diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index abbc6b15704d..1a3acd4beb90 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -4865,10 +4865,11 @@ def _kvwarm_meta_init(self) -> dict: return meta def _kvwarm_state_layer_groups(self) -> list[str]: - """Recurrent groups whose state must never be borrowed from a chain. + """Recurrent groups, whose exact state no shared chain can provide. - Random-state mode gives shadows private slots; without that mode, - these groups remain ineligible for the real-KV warm-up. + Native exact-context warm-up admits them: each point computes its own + state. The shared-chain path still needs random-state mode (private + synthetic slots) or live-state mode (the chain's deeper live state). """ manager = getattr(self, "kv_cache_manager", None) config = getattr(manager, "kv_cache_config", None) @@ -5021,8 +5022,10 @@ def _kvwarm_warm_eligible(self) -> bool: False, ) ) - # Native for MoE with or without EP (keeps Inkling's exact convolution - # state) or recurrent-state layers; dense attention-only stays synthetic. + # Native for MoE with or without EP, or recurrent-state layers; dense + # attention-only stays synthetic. TP-only MoE stays native, not balanced by + # construction: Inkling's convolution cache is a SlidingWindowSpec that the + # layout cannot tell apart from other finite windows, and must stay exact. self._kvwarm_native = self._kvwarm_native_layout() and ( has_experts or bool(self._kvwarm_state_layer_groups()) ) @@ -5641,7 +5644,13 @@ def _kvwarm_plan_key(self, point) -> int | tuple[int, int, int]: return point.batch_size def _kvwarm_native_required_blocks(self, context_lengths: list[int]) -> int: - """Native fleet peak, including admission and repeated steady writes.""" + """Native fleet peak, including admission and repeated steady writes. + + An align-mode state group adds one block per stage: a prefill that ends + on a hash boundary inside a Mamba block registers its own partial tail, + and the first decode's admission check then asks for one block more + than it allocates. + """ repeats = min( self._kvwarm_giant_repeats(), max( @@ -5649,7 +5658,13 @@ def _kvwarm_native_required_blocks(self, context_lengths: list[int]) -> int: self._bench_capacity_limit("max_model_len") - 2 - max(context_lengths), ), ) - return sum( + align_headroom = int( + any( + getattr(manager, "mamba_cache_mode", None) == "align" + for manager in self.kv_cache_manager.coordinator.single_type_managers + ) + ) + return align_headroom + sum( self._bench_blocks_per_req( context + 1 + repeats + self.num_lookahead_tokens, apply_admission_cap=True, diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index 5f8f16388237..55419299006d 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -5698,14 +5698,24 @@ def test_kvwarm_native_capacity_uses_sliding_window_admission_caps(): @pytest.mark.core -def test_kvwarm_native_capacity_counts_mamba_align_state_blocks(): +@pytest.mark.parametrize("align", [True, False]) +def test_kvwarm_native_capacity_counts_mamba_state_blocks(align): stub = InstrumentedScheduler.__new__(InstrumentedScheduler) full = SimpleNamespace(block_size=16) - state = SimpleNamespace( - block_size=16, - mamba_cache_mode="align", - num_speculative_blocks=0, - kv_cache_spec=SimpleNamespace(num_prefill_checkpoint_blocks=1), + # Three state blocks per request in both modes: align holds 2 + 0 + # speculative + 1 checkpoint; "none" (block size = max_model_len) holds + # 1 + 2 speculative. + state = ( + SimpleNamespace( + block_size=16, + mamba_cache_mode="align", + num_speculative_blocks=0, + kv_cache_spec=SimpleNamespace(num_prefill_checkpoint_blocks=1), + ) + if align + else SimpleNamespace( + block_size=4096, mamba_cache_mode="none", num_speculative_blocks=2 + ) ) stub.kv_cache_manager = SimpleNamespace( coordinator=SimpleNamespace(single_type_managers=[full, state]) @@ -5715,9 +5725,12 @@ def test_kvwarm_native_capacity_counts_mamba_align_state_blocks(): stub._kvwarm_giant_repeats = lambda: 3 stub._bench_capacity_limit = lambda name: 4096 - # Per request: ceil((ctx + 1 + 3) / 16) full-attention blocks + 3 align - # state blocks (2 + 0 speculative + 1 checkpoint). - assert stub._kvwarm_native_required_blocks([31, 12]) == (3 + 3) + (1 + 3) + # Per request: ceil((ctx + 1 + 3) / 16) full-attention blocks + 3 state + # blocks. Align adds one block per stage: a prefill ending on a hash + # boundary inside a Mamba block registers its own partial tail, and the + # first decode's admission check asks for one block more than it allocates. + expected = (3 + 3) + (1 + 3) + int(align) + assert stub._kvwarm_native_required_blocks([31, 12]) == expected @pytest.mark.core @@ -6028,10 +6041,14 @@ def _kvwarm_collection_bodies(count, length=80): return bodies -def test_kvwarm_gate_rejects_recurrent_state_layers(monkeypatch): +@pytest.mark.core +def test_kvwarm_gate_rejects_state_layers_outside_native_layouts(monkeypatch): + # A chunked-local group keeps this hybrid off the native path, and shared + # chains cannot supply its recurrent state without random or live-state mode. monkeypatch.setenv("DYN_BENCH_KV_WARMUP", "on") - stub = _kvwarm_gate_stub(state_groups=("FullAttentionSpec", "MambaSpec")) + stub = _kvwarm_gate_stub(state_groups=("ChunkedLocalAttentionSpec", "MambaSpec")) assert InstrumentedScheduler._kvwarm_warm_eligible(stub) is False + assert not stub._kvwarm_native assert stub._kvwarm_meta["skip_reason"] == "hybrid_state_layers_unsupported" @@ -6930,6 +6947,8 @@ def test_kvwarm_dp_filter_rebases_native_plan_without_dropped_points(monkeypatch stub._bench_dp_size = 2 stub.num_lookahead_tokens = 0 stub._bench_blocks_per_req = lambda tokens, **_: -(-tokens // 16) + # No align-mode state group, so native stages need no extra headroom. + stub.kv_cache_manager.coordinator = SimpleNamespace(single_type_managers=[]) stub._kvwarm_warm_eligible = lambda: True stub._kvwarm_native = True points = [ From 1f82903ac7c54ec025f85b4954532a9dbf207067 Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Fri, 2 Oct 2026 17:35:32 +0800 Subject: [PATCH 07/14] fix(vllm): label gate-skipped decode points skip: Decode points of a configuration the warm-up gate rejected were stamped kvwarm_fake_fallback, so a dense model's rows read "wanted real KV, fell back" although synthetic KV is the intended input, and the skip: regime was unreachable. Stamp only real KV and genuine fallbacks, and count gate-skipped points separately. The giant off-by-batch correction keyed on that stamp; it now applies to every warm-up point without real KV, so gate-skipped giant points keep it. Docs and log texts now say that under attention data parallelism failed-stage points are skipped and uncovered points dropped, not faked, and the explicit-point error names the native footprint as a cause. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../src/dynamo/vllm/instrumented_scheduler.py | 33 +++-- .../tests/test_vllm_instrumented_scheduler.py | 120 +++++++++++++++++- .../observability/environment-variables.mdx | 4 +- 3 files changed, 142 insertions(+), 15 deletions(-) diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index 1a3acd4beb90..364b0b584a0e 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -4858,6 +4858,7 @@ def _kvwarm_meta_init(self) -> dict: "stages": [], "points_real_kv": 0, "points_fake_fallback": 0, + "points_gate_skipped": 0, "giant_kv_threshold": self._kvwarm_giant_threshold(), "giant_kv_repeats": self._kvwarm_giant_repeats(), } @@ -5473,7 +5474,8 @@ def _kvwarm_prepare(self, mode: str) -> None: raise RuntimeError( "KVWARM: attention-DP runs measure real-KV points only, and the " f"warm-up plan cannot cover explicit decode point(s) {coords} " - "(rung depth trimmed by the per-rank pool / slot budget); lower " + "(rung depth trimmed by the per-rank pool / slot budget, or a " + "native point's footprint exceeds the per-rank pool); lower " "the requested batch or context, or run without attention-DP" ) logger.warning( @@ -5869,7 +5871,8 @@ def _kvwarm_monitor_build(self) -> bool: # fallback) and release the survivors. logger.warning( "KVWARM: %d chain(s) of stage batch=%s vanished during the " - "build; failing the stage, its points fall back to fake injection", + "build; failing the stage, its points fall back to fake injection " + "(skipped under attention-DP)", len(vanished), self._kvwarm_stage_batch, ) @@ -6002,7 +6005,7 @@ def _kvwarm_stage_settle(self, batch: int | None, ok: bool, detail: dict) -> Non meta["stages"].append({"batch": batch, "failed": True, **detail}) logger.warning( "KVWARM: stage batch=%s failed (%s); its points fall back to fake " - "injection", + "injection (skipped under attention-DP)", batch, detail, ) @@ -6581,17 +6584,21 @@ def _bench_step_decode(self) -> SchedulerOutput | None: return None if self._kvwarm_flag_on(): meta = self._kvwarm_meta_init() + stamp: str | None if kvwarm_real: meta["points_real_kv"] += 1 + stamp = "kvwarm_real_kv" + elif meta.get("warm_eligible") is False: + # The gate rejected the configuration: synthetic KV is the + # design, not a fallback, and the row regime reads + # skip: (``_kvwarm_seed_regime``). + meta["points_gate_skipped"] += 1 + stamp = None else: meta["points_fake_fallback"] += 1 - point = replace( - point, - sample_reasons=[ - *point.sample_reasons, - "kvwarm_real_kv" if kvwarm_real else "kvwarm_fake_fallback", - ], - ) + stamp = "kvwarm_fake_fallback" + if stamp is not None: + point = replace(point, sample_reasons=[*point.sample_reasons, stamp]) self._bench_current_point = point self._bench_current_fpms = [] self._bench_extra_steps_left = 1 @@ -6893,10 +6900,12 @@ def _bench_fpm_validation_failure( # all-or-nothing publish gate. The admission step has the same total, so # the correction requires a recorded steady sample (``kvwarm_steady_sample``, # set by both save paths): a point that hit its deadline with the admission - # FPM alone stays a validation skip. + # FPM alone stays a validation skip. Fake means any warm-up point without + # real KV: gate-skipped points take the same path but carry no stamp. measured = scheduled.get("sum_decode_kv_tokens") if ( - "kvwarm_fake_fallback" in (point.sample_reasons or ()) + self._kvwarm_flag_on() + and "kvwarm_real_kv" not in (point.sample_reasons or ()) and point.total_kv_read_tokens >= self._kvwarm_giant_threshold() and measured == point.total_kv_read_tokens - point.batch_size and bool(fpm.get("kvwarm_steady_sample")) diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index 55419299006d..e41eebeec53d 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -5825,6 +5825,91 @@ def test_kvwarm_native_decode_dispatch_uses_real_continuation_and_provenance(): assert stub._kvwarm_meta["points_real_kv"] == 1 +def _kvwarm_uncovered_dispatch_stub(warm_eligible, skip_reason=None): + """Native-layout decode dispatch of a point no warm-up stage covers.""" + stub, _ = _kvwarm_native_resume_stub(False) + point = BenchmarkPoint( + point_type="decode", benchmark_id=7, batch_size=2, total_kv_read_tokens=9 + ) + # The gate's verdict as ``_kvwarm_warm_eligible`` leaves it before dispatch. + stub._kvwarm_meta_init().update( + warm_eligible=warm_eligible, skip_reason=skip_reason + ) + stub._kvwarm_eligible_cache = warm_eligible + stub._kvwarm_covers = lambda point, lengths: False + stub._bench_inject_fake_decode = MagicMock( + return_value=SimpleNamespace(total_num_scheduled_tokens=point.batch_size) + ) + stub._bench_grid = deque([point]) + stub._bench_current_point = None + stub._bench_drain_pending = False + stub._bench_frees_pending = lambda: False + stub._bench_stop_at_timeout_boundary = lambda phase: False + stub.max_model_len = 8192 + return stub + + +@pytest.mark.core +def test_kvwarm_gate_skipped_decode_point_records_skip_regime(monkeypatch): + monkeypatch.setenv("DYN_BENCH_KV_WARMUP", "on") + stub = _kvwarm_uncovered_dispatch_stub(False, "dense_model_content_insensitive") + + assert stub._bench_step_decode() is not None + + point = stub._bench_current_point + stub._bench_inject_fake_decode.assert_called_once_with([4, 3]) + assert "kvwarm_fake_fallback" not in point.sample_reasons + assert "kvwarm_real_kv" not in point.sample_reasons + assert stub._kvwarm_seed_regime(point) == "skip:dense_model_content_insensitive" + assert stub._kvwarm_meta["points_gate_skipped"] == 1 + assert stub._kvwarm_meta["points_fake_fallback"] == 0 + + +@pytest.mark.core +def test_kvwarm_eligible_capacity_fallback_still_records_fake_fallback(monkeypatch): + monkeypatch.setenv("DYN_BENCH_KV_WARMUP", "on") + stub = _kvwarm_uncovered_dispatch_stub(True) + # Capacity fallback: the planner zeroed the depth of a point the gate accepted. + stub._kvwarm_plan = {stub._kvwarm_plan_key(stub._bench_grid[0]): 0} + + assert stub._bench_step_decode() is not None + + point = stub._bench_current_point + stub._bench_inject_fake_decode.assert_called_once_with([4, 3]) + assert "kvwarm_fake_fallback" in point.sample_reasons + assert stub._kvwarm_seed_regime(point) == "fake_fallback" + assert stub._kvwarm_meta["points_fake_fallback"] == 1 + assert stub._kvwarm_meta["points_gate_skipped"] == 0 + + +@pytest.mark.core +def test_kvwarm_group_skip_labels_decode_points_like_the_ineligible_rank( + monkeypatch, +): + monkeypatch.setenv("DYN_BENCH_KV_WARMUP", "on") + # Attention-DP: rank 0's own gate failed and rank 1's passed; the + # negotiated envelope carries the group verdict to both. + ranks = [ + _kvwarm_uncovered_dispatch_stub(False, "dataset_empty"), + _kvwarm_uncovered_dispatch_stub(True), + ] + points = [] + for stub in ranks: + stub._bench_dp_size = 2 + stub._bench_negotiated_capacity = SimpleNamespace(kvwarm_eligible=False) + stub._kvwarm_prepare("decode") + + assert stub._bench_step_decode() is not None + + points.append(stub._bench_current_point) + assert stub._kvwarm_meta["points_gate_skipped"] == 1 + assert stub._kvwarm_meta["points_fake_fallback"] == 0 + # Ranks compare whole points: the READY/RESULT digests and the rank merge. + assert points[0] == points[1] + assert ranks[0]._kvwarm_seed_regime(points[0]) == "skip:dataset_empty" + assert ranks[1]._kvwarm_seed_regime(points[1]) == "skip:peer_ineligible" + + @pytest.mark.core @pytest.mark.parametrize("context,repeats", [(124, 3), (125, 2), (126, 1)]) def test_kvwarm_native_dp_dispatch_caps_repeats_by_negotiated_model_length( @@ -6987,7 +7072,9 @@ def test_giant_fake_off_by_batch_correction_requires_a_steady_sample(): sample (``kvwarm_steady_sample``, set by both save paths) may be accepted at the measured coordinate. A giant fake point that reached its deadline with the admission FPM alone is a validation skip, not a decode measurement.""" - stub = SimpleNamespace(_kvwarm_giant_threshold=lambda: 1000) + stub = SimpleNamespace( + _kvwarm_giant_threshold=lambda: 1000, _kvwarm_flag_on=lambda: True + ) point = BenchmarkPoint( point_type="decode", benchmark_id=1, @@ -7068,6 +7155,37 @@ def test_two_fpm_save_path_skips_an_admission_only_giant_sample(monkeypatch): ] +@pytest.mark.core +@pytest.mark.parametrize( + "warmup,reasons,accepted", + [ + ("on", ["kvwarm_fake_fallback"], True), + ("on", [], True), + ("on", ["kvwarm_real_kv"], False), + ("off", [], False), + ], + ids=["fake_fallback", "gate_skipped", "real_kv", "legacy"], +) +def test_giant_off_by_batch_correction_accepts_every_synthetic_warmup_point( + monkeypatch, warmup, reasons, accepted +): + """Gate-skipped points run the same synthetic injection and giant repeats as + fake fallbacks but carry no stamp; the correction still accepts them, and + still rejects real-KV and legacy points.""" + monkeypatch.setenv("DYN_BENCH_KV_WARMUP", warmup) + monkeypatch.setenv("DYN_BENCH_GIANT_KV_THRESHOLD", "1000") + stub = InstrumentedScheduler.__new__(InstrumentedScheduler) + point = replace(_giant_fake_point(), sample_reasons=reasons) + steady = { + "scheduled_requests": {"num_decode_requests": 2, "sum_decode_kv_tokens": 1998}, + "kvwarm_steady_sample": True, + } + + assert stub._bench_fpm_validation_failure(point, steady) == ( + None if accepted else "measured_decode_context_mismatch" + ) + + def test_kvwarm_shadow_registration_rejects_too_shallow_chain(): stub, mgr, pool, chain = _shadow_stub() with pytest.raises(RuntimeError, match="too shallow"): diff --git a/docs/fern/pages/reference/observability/environment-variables.mdx b/docs/fern/pages/reference/observability/environment-variables.mdx index b91f55a9f65b..ae5dd62b2af7 100644 --- a/docs/fern/pages/reference/observability/environment-variables.mdx +++ b/docs/fern/pages/reference/observability/environment-variables.mdx @@ -297,7 +297,7 @@ These variables tune the vLLM worker's self-benchmark collector (`--benchmark-mo Other eligible configurations use shared prefill chains and require prefix caching. By default, chain eligibility requires a MoE model with expert parallelism and no recurrent-state layers. With `DYN_BENCHMARK_RANDOMIZE_KDA_STATE=true`, hybrid MoE models keep the shared-chain random-state strategy, also without expert parallelism: their attention KV is real and their recurrent states remain synthetic. This option disables the native exact-context strategy, even for layouts that qualify for it. - Both strategies require a loadable tokenizer, enough seeding dataset content for the model length, and sufficient cache capacity. Configurations that do not meet these requirements skip warm-up and record the reason. Under attention data parallelism the decision is group-wide: if any rank is ineligible, every rank uses synthetic KV; if any rank's warm-up stage fails, every rank falls back for that stage. Failed initialization never marks synthetic state as `real_kv`. + Both strategies require a loadable tokenizer, enough seeding dataset content for the model length, and sufficient cache capacity. Configurations that do not meet these requirements skip warm-up and record the reason. Under attention data parallelism the decision is group-wide: if any rank is ineligible, every rank uses synthetic KV. In an eligible group, points the plan cannot cover are dropped from the grid, and points whose warm-up stage fails on any rank are skipped (`stage_failed_under_attention_dp`); outside attention data parallelism, both fall back to synthetic KV (`fake_fallback`). Failed initialization never marks synthetic state as `real_kv`. @@ -342,7 +342,7 @@ These variables tune the vLLM worker's self-benchmark collector (`--benchmark-mo Optional tag mixed into the real-text pool draw so repeated collections sample different content windows under otherwise identical settings. -Every result and skipped-point entry in the artifact carries a `kv_seed_regime` field plus `kvwarm` and `synthetic_prompts` metadata blocks, so downstream consumers can filter by seeding provenance. Decode rows carry `real_kv`, `fake_fallback`, `legacy` (warm-up switched off), `skip:` (gate rejected the configuration), or `unstamped` (decode point that never reached injection). Prefill rows that read past KV carry `real_prefix` (real-KV seeding on) or `fake_prefix` (synthetic prefix blocks); other prefill rows are `not_applicable`. +Every result and skipped-point entry in the artifact carries a `kv_seed_regime` field plus `kvwarm` and `synthetic_prompts` metadata blocks, so downstream consumers can filter by seeding provenance. Decode rows carry `real_kv`, `fake_fallback`, `legacy` (warm-up switched off), `skip:`, or `unstamped` (decode point that never reached injection). `fake_fallback` occurs only outside attention data parallelism: the configuration was eligible, but this point fell back to synthetic KV because of cache capacity, a failed warm-up stage, or a warm-up plan that does not cover it. `skip:` means the warm-up gate rejected the configuration and synthetic KV is the intended input, for example `skip:dense_model_content_insensitive`; under attention data parallelism, ranks whose own check passed record `skip:peer_ineligible`. Prefill rows that read past KV carry `real_prefix` (real-KV seeding on) or `fake_prefix` (synthetic prefix blocks); other prefill rows are `not_applicable`. ## Request tracing From 0a6544eccddcb08a40f5e37c74e2ba4abbcca202 Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Fri, 2 Oct 2026 18:07:38 +0800 Subject: [PATCH 08/14] test(vllm): tidy native warm-up tests and document the native cost Document the untimed prefill a native sweep adds, keep a caller's KV cache manager in the capacity test helper, and drop a comment that only restated its assertion. Also reword the docs and comments that described only the non-DP shared-chain outcomes: skip: no longer reads as the intended input, gate-rejected random-KDA runs read skip:, explicit points under attention-DP raise instead of being dropped, a failed native stage retires one point, and TP-only MoE stays native only with a qualifying layout. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../src/dynamo/vllm/instrumented_scheduler.py | 36 ++++++++++++------- .../tests/test_vllm_instrumented_scheduler.py | 10 +++--- .../observability/environment-variables.mdx | 8 ++--- 3 files changed, 33 insertions(+), 21 deletions(-) diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index 364b0b584a0e..9f3db5274625 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -5024,9 +5024,10 @@ def _kvwarm_warm_eligible(self) -> bool: ) ) # Native for MoE with or without EP, or recurrent-state layers; dense - # attention-only stays synthetic. TP-only MoE stays native, not balanced by - # construction: Inkling's convolution cache is a SlidingWindowSpec that the - # layout cannot tell apart from other finite windows, and must stay exact. + # attention-only stays synthetic. TP-only MoE with a qualifying layout + # stays native, not balanced by construction: Inkling's convolution cache + # is a SlidingWindowSpec that the layout cannot tell apart from other + # finite windows, and must stay exact. self._kvwarm_native = self._kvwarm_native_layout() and ( has_experts or bool(self._kvwarm_state_layer_groups()) ) @@ -5865,10 +5866,14 @@ def _kvwarm_monitor_build(self) -> bool: else: pending = True if vanished: - # A partially built fleet cannot serve its rung: surviving chains - # would only pin KV that fake injection needs. Fail the whole - # stage (every point of this rung takes the fake-injection - # fallback) and release the survivors. + # A partially built fleet cannot serve its stage: surviving chains + # would only pin KV that the next stage or fake injection needs. Fail + # the whole stage and release the survivors. The points the stage + # served lose real-KV coverage: every point of the batch rung for + # shared chains, the one point for native prefills. Outside + # attention-DP they take the fake-injection fallback; under + # attention-DP they are skipped, because fake injection is not + # rank-consistent. logger.warning( "KVWARM: %d chain(s) of stage batch=%s vanished during the " "build; failing the stage, its points fall back to fake injection " @@ -5986,9 +5991,12 @@ def _kvwarm_stage_round(self) -> bool: return True def _kvwarm_stage_settle(self, batch: int | None, ok: bool, detail: dict) -> None: - """Record the final outcome of a stage. A failed rung has its plan - depth zeroed so every point of the rung takes the fake-injection - fallback (``_kvwarm_plan_covers`` reads the plan).""" + """Record the final outcome of a stage. A failed stage has its plan + depth zeroed, so the points it served lose coverage + (``_kvwarm_plan_covers`` reads the plan): the whole batch rung for + shared chains, the one point for native prefills. Outside attention-DP + they take the fake-injection fallback; under attention-DP they are + skipped, because fake injection is not rank-consistent.""" meta = self._kvwarm_meta_init() if ok: meta["stages"].append({"batch": batch, **detail}) @@ -6589,9 +6597,11 @@ def _bench_step_decode(self) -> SchedulerOutput | None: meta["points_real_kv"] += 1 stamp = "kvwarm_real_kv" elif meta.get("warm_eligible") is False: - # The gate rejected the configuration: synthetic KV is the - # design, not a fallback, and the row regime reads - # skip: (``_kvwarm_seed_regime``). + # The gate rejected the whole configuration, so every point uses + # synthetic KV and none fell back: no real-KV attempt existed. The + # row regime reads skip: (``_kvwarm_seed_regime``); the + # reason says why, for example dense_model_content_insensitive + # (synthetic KV is correct by construction) or dataset_empty. meta["points_gate_skipped"] += 1 stamp = None else: diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index e41eebeec53d..c3f85799a98f 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -108,9 +108,11 @@ def _install_test_capacity_preflight(stub, capacity=None): capacity = capacity or _benchmark_capacity() stub._bench_make_local_capacity = lambda: capacity stub._bench_synchronizer = None - stub.kv_cache_manager = SimpleNamespace( - kv_cache_config=SimpleNamespace(kv_cache_groups=[]) - ) + if getattr(stub, "kv_cache_manager", None) is None: + # Keep a caller's manager: _bench_blocks_per_req reads its groups. + stub.kv_cache_manager = SimpleNamespace( + kv_cache_config=SimpleNamespace(kv_cache_groups=[]) + ) # ``_bench_build_grid`` re-filters the decode capture list against the # negotiated request limit before generating the grid; stubs that don't # model captures still need the attribute to exist. @@ -6089,7 +6091,7 @@ def finish_requests(req_ids): assert stub._bench_expected_fpms == 4 assert stub.kv_cache_manager.block_pool.get_num_free_blocks() == 0 assert stub._bench_step_decode() is None - assert stub._bench_results == [] # Wait for the normal result deadline. + assert stub._bench_results == [] stub._bench_point_deadline = time.monotonic() - 1 assert stub._bench_step_decode() is None assert len(stub._bench_results) == 1 diff --git a/docs/fern/pages/reference/observability/environment-variables.mdx b/docs/fern/pages/reference/observability/environment-variables.mdx index ae5dd62b2af7..5c3b70cafecf 100644 --- a/docs/fern/pages/reference/observability/environment-variables.mdx +++ b/docs/fern/pages/reference/observability/environment-variables.mdx @@ -293,11 +293,11 @@ These variables tune the vLLM worker's self-benchmark collector (`--benchmark-mo KV warm-up master switch. When on, eligible decode benchmark points read attention KV produced by real prefill. Set `off` to force the legacy synthetic-KV path. - Models whose decode timing depends on KV content (MoE models, or models with recurrent-state layers) use native prefill when every cache group is `FullAttentionSpec`, a finite `SlidingWindowSpec`, or a Mamba/KDA state group, and at least one group is a finite window or a state group: each point prefills its own requests to the exact admission context, and the same requests continue into decode, so window history and recurrent state are computed by the model instead of borrowed from a chain at another depth. This covers Inkling's convolution cache and hybrid models such as GLM-5.3-Flash. It does not require expert parallelism or prefix caching, and `--benchmark-hybrid-live-state` has no effect when it applies. Dense attention-only models keep synthetic KV, which is correct by construction for them. The artifact records `kvwarm.initialization_strategy = "native_exact_context"`. + Models whose decode timing depends on KV content (MoE models, or models with recurrent-state layers) use native prefill when every cache group is `FullAttentionSpec`, a finite `SlidingWindowSpec`, or a Mamba/KDA state group, and at least one group is a finite window or a state group: each point prefills its own requests to the exact admission context, and the same requests continue into decode, so window history and recurrent state are computed by the model instead of borrowed from a chain at another depth. This covers Inkling's convolution cache and hybrid models such as GLM-5.3-Flash. It does not require expert parallelism or prefix caching, and `--benchmark-hybrid-live-state` has no effect when it applies. Dense attention-only models keep synthetic KV, which is correct by construction for them. The artifact records `kvwarm.initialization_strategy = "native_exact_context"`. Each decode execution, including eager warm-up replicas and repeated coordinates, prefills its own requests and shares no prefix with other executions, so the untimed prefill per sweep is about the sum of `total_kv_read_tokens` over the decode executions. Other eligible configurations use shared prefill chains and require prefix caching. By default, chain eligibility requires a MoE model with expert parallelism and no recurrent-state layers. With `DYN_BENCHMARK_RANDOMIZE_KDA_STATE=true`, hybrid MoE models keep the shared-chain random-state strategy, also without expert parallelism: their attention KV is real and their recurrent states remain synthetic. This option disables the native exact-context strategy, even for layouts that qualify for it. - Both strategies require a loadable tokenizer, enough seeding dataset content for the model length, and sufficient cache capacity. Configurations that do not meet these requirements skip warm-up and record the reason. Under attention data parallelism the decision is group-wide: if any rank is ineligible, every rank uses synthetic KV. In an eligible group, points the plan cannot cover are dropped from the grid, and points whose warm-up stage fails on any rank are skipped (`stage_failed_under_attention_dp`); outside attention data parallelism, both fall back to synthetic KV (`fake_fallback`). Failed initialization never marks synthetic state as `real_kv`. + Both strategies require a loadable tokenizer, enough seeding dataset content for the model length, and sufficient cache capacity. Configurations that do not meet these requirements skip warm-up and record the reason. Under attention data parallelism the decision is group-wide: if any rank is ineligible, every rank uses synthetic KV. In an eligible group, points the plan cannot cover are dropped from the grid, and points whose warm-up stage fails on any rank are skipped (`stage_failed_under_attention_dp`); outside attention data parallelism, both fall back to synthetic KV (`fake_fallback`). Under attention data parallelism, explicit benchmark points (`DYN_BENCHMARK_POINTS_FILE`) that the plan cannot cover, or whose warm-up stage fails, raise an error instead of being dropped or skipped. Failed initialization never marks synthetic state as `real_kv`. @@ -307,7 +307,7 @@ These variables tune the vLLM worker's self-benchmark collector (`--benchmark-mo With `DYN_BENCH_KV_WARMUP=on`, hybrid MoE models can use real attention KV from prefill chains even without expert parallelism. Their recurrent states remain synthetic; the source chain's states are never borrowed or modified. Existing dataset, prefix-cache, and capacity requirements still apply. Without a usable chain, attention KV follows the existing synthetic fallback. - Artifacts distinguish `real_attention_kv_random_kda` from `fake_attention_kv_random_kda` and record the distribution in `recurrent_state`. These values are performance-test inputs, not valid context history or a guarantee of real-workload performance parity. + Artifacts distinguish `real_attention_kv_random_kda` from `fake_attention_kv_random_kda` and record the distribution in `recurrent_state`; a run the warm-up gate rejects labels its decode rows `skip:`, not `fake_attention_kv_random_kda`. These values are performance-test inputs, not valid context history or a guarantee of real-workload performance parity. @@ -342,7 +342,7 @@ These variables tune the vLLM worker's self-benchmark collector (`--benchmark-mo Optional tag mixed into the real-text pool draw so repeated collections sample different content windows under otherwise identical settings. -Every result and skipped-point entry in the artifact carries a `kv_seed_regime` field plus `kvwarm` and `synthetic_prompts` metadata blocks, so downstream consumers can filter by seeding provenance. Decode rows carry `real_kv`, `fake_fallback`, `legacy` (warm-up switched off), `skip:`, or `unstamped` (decode point that never reached injection). `fake_fallback` occurs only outside attention data parallelism: the configuration was eligible, but this point fell back to synthetic KV because of cache capacity, a failed warm-up stage, or a warm-up plan that does not cover it. `skip:` means the warm-up gate rejected the configuration and synthetic KV is the intended input, for example `skip:dense_model_content_insensitive`; under attention data parallelism, ranks whose own check passed record `skip:peer_ineligible`. Prefill rows that read past KV carry `real_prefix` (real-KV seeding on) or `fake_prefix` (synthetic prefix blocks); other prefill rows are `not_applicable`. +Every result and skipped-point entry in the artifact carries a `kv_seed_regime` field plus `kvwarm` and `synthetic_prompts` metadata blocks, so downstream consumers can filter by seeding provenance. Decode rows carry `real_kv`, `fake_fallback`, `legacy` (warm-up switched off), `skip:`, or `unstamped` (decode point that never reached injection). `fake_fallback` occurs only outside attention data parallelism: the configuration was eligible, but this point fell back to synthetic KV because of cache capacity, a failed warm-up stage, or a warm-up plan that does not cover it. `skip:` means the warm-up gate rejected the whole configuration, so every decode point used synthetic KV; `` says why, for example `skip:dense_model_content_insensitive` (synthetic KV is correct by construction) or `skip:dataset_empty`. Under attention data parallelism, ranks whose own check passed record `skip:peer_ineligible`. Prefill rows that read past KV carry `real_prefix` (real-KV seeding on) or `fake_prefix` (synthetic prefix blocks); other prefill rows are `not_applicable`. ## Request tracing From 4d72a2a2d2cedd9ca05baf76c3eead08f69c237d Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Fri, 2 Oct 2026 19:00:59 +0800 Subject: [PATCH 09/14] fix(vllm): address the final review of the native warm-up update Six minor findings from the final whole-branch review: 1. _kvwarm_native_required_blocks reserves one align headroom block per align-mode Mamba manager instead of one flat block; the capacity test gains a two-align-manager case that expects +2. 2. A native capacity fallback entry carries total_kv_read_tokens, so a point dropped under attention-DP keeps its coordinate, and the fallback now logs a warning; both expected dicts are updated. 3. The env-var page no longer lists cache capacity as a requirement whose shortfall skips warm-up, documents the per-point capacity_fallbacks record, and merges the attention-DP sentences. 4. The env-var page states that the native layout check accepts subclasses (MLA, sliding-window MLA, k-pool tail) and rejects circular-buffer groups. 5. The test helper comment names the manager attribute that _bench_blocks_per_req actually reads. 6. The direct native stage test keeps one context pair, with its intra-stage and cross-call salt assertions. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../src/dynamo/vllm/instrumented_scheduler.py | 26 +++++--- .../tests/test_vllm_instrumented_scheduler.py | 60 +++++++++++-------- .../observability/environment-variables.mdx | 4 +- 3 files changed, 55 insertions(+), 35 deletions(-) diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index 9f3db5274625..66499613c4da 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -5336,11 +5336,20 @@ def _kvwarm_prepare(self, mode: str) -> None: { "benchmark_id": point.benchmark_id, "batch": point.batch_size, + "total_kv_read_tokens": point.total_kv_read_tokens, "depth": depth, "required_blocks": required, "usable_blocks": usable, } ) + logger.warning( + "KVWARM: native point batch=%d kv=%d needs %d blocks, " + "pool has %d; using fake-KV fallback", + point.batch_size, + point.total_kv_read_tokens, + required, + usable, + ) else: repeats = self._kvwarm_giant_repeats() margin = 1 + repeats @@ -5649,10 +5658,11 @@ def _kvwarm_plan_key(self, point) -> int | tuple[int, int, int]: def _kvwarm_native_required_blocks(self, context_lengths: list[int]) -> int: """Native fleet peak, including admission and repeated steady writes. - An align-mode state group adds one block per stage: a prefill that ends - on a hash boundary inside a Mamba block registers its own partial tail, - and the first decode's admission check then asks for one block more - than it allocates. + Each align-mode state group adds one block per stage: a prefill that + ends on a hash boundary inside a Mamba block registers its own partial + tail in that group's manager, and the first decode's admission check + then asks the group for one block more than it allocates. The + coordinator sums the asks over groups. """ repeats = min( self._kvwarm_giant_repeats(), @@ -5661,11 +5671,9 @@ def _kvwarm_native_required_blocks(self, context_lengths: list[int]) -> int: self._bench_capacity_limit("max_model_len") - 2 - max(context_lengths), ), ) - align_headroom = int( - any( - getattr(manager, "mamba_cache_mode", None) == "align" - for manager in self.kv_cache_manager.coordinator.single_type_managers - ) + align_headroom = sum( + getattr(manager, "mamba_cache_mode", None) == "align" + for manager in self.kv_cache_manager.coordinator.single_type_managers ) return align_headroom + sum( self._bench_blocks_per_req( diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index c3f85799a98f..adc813052c12 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -109,7 +109,8 @@ def _install_test_capacity_preflight(stub, capacity=None): stub._bench_make_local_capacity = lambda: capacity stub._bench_synchronizer = None if getattr(stub, "kv_cache_manager", None) is None: - # Keep a caller's manager: _bench_blocks_per_req reads its groups. + # Keep a caller's manager: _bench_blocks_per_req reads its coordinator's + # managers, not the groups installed here. stub.kv_cache_manager = SimpleNamespace( kv_cache_config=SimpleNamespace(kv_cache_groups=[]) ) @@ -5192,8 +5193,8 @@ def test_kvwarm_native_rejects_unsupported_spec_with_state(): @pytest.mark.core -@pytest.mark.parametrize("contexts", [(4, 3), (5, 4), (16, 15), (512, 511), (513, 512)]) -def test_kvwarm_native_stage_prefills_exact_heterogeneous_contexts(contexts): +def test_kvwarm_native_stage_prefills_exact_heterogeneous_contexts(): + contexts = (513, 512) # Admission lengths (ctx - 1) of contexts 514 and 513. stub = InstrumentedScheduler.__new__(InstrumentedScheduler) stub._kvwarm_native = True stub._kvwarm_seq = 0 @@ -5456,6 +5457,7 @@ def finish_requests(req_ids, stub=stub): { "benchmark_id": benchmark_id, "batch": 20, + "total_kv_read_tokens": 40, "depth": 1, "required_blocks": 80, "usable_blocks": 63, @@ -5700,27 +5702,34 @@ def test_kvwarm_native_capacity_uses_sliding_window_admission_caps(): @pytest.mark.core -@pytest.mark.parametrize("align", [True, False]) -def test_kvwarm_native_capacity_counts_mamba_state_blocks(align): +@pytest.mark.parametrize( + "modes", + [("none",), ("align",), ("align", "align")], + ids=["none", "align", "two_align"], +) +def test_kvwarm_native_capacity_counts_mamba_state_blocks(modes): stub = InstrumentedScheduler.__new__(InstrumentedScheduler) full = SimpleNamespace(block_size=16) - # Three state blocks per request in both modes: align holds 2 + 0 - # speculative + 1 checkpoint; "none" (block size = max_model_len) holds - # 1 + 2 speculative. - state = ( - SimpleNamespace( - block_size=16, - mamba_cache_mode="align", - num_speculative_blocks=0, - kv_cache_spec=SimpleNamespace(num_prefill_checkpoint_blocks=1), - ) - if align - else SimpleNamespace( - block_size=4096, mamba_cache_mode="none", num_speculative_blocks=2 + # Three state blocks per request and state group in every mode: align holds + # 2 + 0 speculative + 1 checkpoint; "none" (block size = max_model_len) + # holds 1 + 2 speculative. + states = [ + ( + SimpleNamespace( + block_size=16, + mamba_cache_mode="align", + num_speculative_blocks=0, + kv_cache_spec=SimpleNamespace(num_prefill_checkpoint_blocks=1), + ) + if mode == "align" + else SimpleNamespace( + block_size=4096, mamba_cache_mode="none", num_speculative_blocks=2 + ) ) - ) + for mode in modes + ] stub.kv_cache_manager = SimpleNamespace( - coordinator=SimpleNamespace(single_type_managers=[full, state]) + coordinator=SimpleNamespace(single_type_managers=[full, *states]) ) stub.block_size = 16 stub.num_lookahead_tokens = 0 @@ -5728,10 +5737,12 @@ def test_kvwarm_native_capacity_counts_mamba_state_blocks(align): stub._bench_capacity_limit = lambda name: 4096 # Per request: ceil((ctx + 1 + 3) / 16) full-attention blocks + 3 state - # blocks. Align adds one block per stage: a prefill ending on a hash - # boundary inside a Mamba block registers its own partial tail, and the - # first decode's admission check asks for one block more than it allocates. - expected = (3 + 3) + (1 + 3) + int(align) + # blocks per state group. Each align group adds one block per stage: a + # prefill ending on a hash boundary inside a Mamba block registers its own + # partial tail, and the first decode's admission check asks that group for + # one block more than it allocates; the coordinator sums the asks. + state_blocks = 3 * len(modes) + expected = (3 + state_blocks) + (1 + state_blocks) + modes.count("align") assert stub._kvwarm_native_required_blocks([31, 12]) == expected @@ -7062,6 +7073,7 @@ def test_kvwarm_dp_filter_rebases_native_plan_without_dropped_points(monkeypatch { "benchmark_id": None, "batch": 2, + "total_kv_read_tokens": 62, "depth": 30, "required_blocks": 6, "usable_blocks": 4, diff --git a/docs/fern/pages/reference/observability/environment-variables.mdx b/docs/fern/pages/reference/observability/environment-variables.mdx index 5c3b70cafecf..74292ef892f6 100644 --- a/docs/fern/pages/reference/observability/environment-variables.mdx +++ b/docs/fern/pages/reference/observability/environment-variables.mdx @@ -293,11 +293,11 @@ These variables tune the vLLM worker's self-benchmark collector (`--benchmark-mo KV warm-up master switch. When on, eligible decode benchmark points read attention KV produced by real prefill. Set `off` to force the legacy synthetic-KV path. - Models whose decode timing depends on KV content (MoE models, or models with recurrent-state layers) use native prefill when every cache group is `FullAttentionSpec`, a finite `SlidingWindowSpec`, or a Mamba/KDA state group, and at least one group is a finite window or a state group: each point prefills its own requests to the exact admission context, and the same requests continue into decode, so window history and recurrent state are computed by the model instead of borrowed from a chain at another depth. This covers Inkling's convolution cache and hybrid models such as GLM-5.3-Flash. It does not require expert parallelism or prefix caching, and `--benchmark-hybrid-live-state` has no effect when it applies. Dense attention-only models keep synthetic KV, which is correct by construction for them. The artifact records `kvwarm.initialization_strategy = "native_exact_context"`. Each decode execution, including eager warm-up replicas and repeated coordinates, prefills its own requests and shares no prefix with other executions, so the untimed prefill per sweep is about the sum of `total_kv_read_tokens` over the decode executions. + Models whose decode timing depends on KV content (MoE models, or models with recurrent-state layers) use native prefill when every cache group is `FullAttentionSpec`, a finite `SlidingWindowSpec`, or a Mamba/KDA state group, and at least one group is a finite window or a state group: each point prefills its own requests to the exact admission context, and the same requests continue into decode, so window history and recurrent state are computed by the model instead of borrowed from a chain at another depth. This covers Inkling's convolution cache and hybrid models such as GLM-5.3-Flash. It does not require expert parallelism or prefix caching, and `--benchmark-hybrid-live-state` has no effect when it applies. The cache-group check accepts subclasses, so MLA attention (`MLAAttentionSpec`), sliding-window MLA (`SlidingWindowMLASpec`), and k-pool tail (`KpoolTailSpec`) groups pass the per-group check; circular-buffer (`CircularBufferSpec`) groups do not. Dense attention-only models keep synthetic KV, which is correct by construction for them. The artifact records `kvwarm.initialization_strategy = "native_exact_context"`. Each decode execution, including eager warm-up replicas and repeated coordinates, prefills its own requests and shares no prefix with other executions, so the untimed prefill per sweep is about the sum of `total_kv_read_tokens` over the decode executions. Other eligible configurations use shared prefill chains and require prefix caching. By default, chain eligibility requires a MoE model with expert parallelism and no recurrent-state layers. With `DYN_BENCHMARK_RANDOMIZE_KDA_STATE=true`, hybrid MoE models keep the shared-chain random-state strategy, also without expert parallelism: their attention KV is real and their recurrent states remain synthetic. This option disables the native exact-context strategy, even for layouts that qualify for it. - Both strategies require a loadable tokenizer, enough seeding dataset content for the model length, and sufficient cache capacity. Configurations that do not meet these requirements skip warm-up and record the reason. Under attention data parallelism the decision is group-wide: if any rank is ineligible, every rank uses synthetic KV. In an eligible group, points the plan cannot cover are dropped from the grid, and points whose warm-up stage fails on any rank are skipped (`stage_failed_under_attention_dp`); outside attention data parallelism, both fall back to synthetic KV (`fake_fallback`). Under attention data parallelism, explicit benchmark points (`DYN_BENCHMARK_POINTS_FILE`) that the plan cannot cover, or whose warm-up stage fails, raise an error instead of being dropped or skipped. Failed initialization never marks synthetic state as `real_kv`. + Both strategies require a loadable tokenizer and enough seeding dataset content for the model length. Configurations that do not meet these requirements skip warm-up and record the reason. A point whose warm-up footprint exceeds the cache pool is recorded in `kvwarm.capacity_fallbacks` and measured on synthetic KV (`fake_fallback`), or dropped from the grid under attention data parallelism. Under attention data parallelism the decision is group-wide: if any rank is ineligible, every rank uses synthetic KV; in an eligible group, a generated point that the plan cannot cover is dropped from the grid and one whose warm-up stage fails on any rank is skipped (`stage_failed_under_attention_dp`), while an explicit benchmark point (`DYN_BENCHMARK_POINTS_FILE`) in either case raises an error instead. Outside attention data parallelism, points the plan cannot cover and points whose warm-up stage fails fall back to synthetic KV (`fake_fallback`). Failed initialization never marks synthetic state as `real_kv`. From 1fa6362129ad5f0995deebaa2cbc93d0c8de68a9 Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Sat, 10 Oct 2026 12:12:37 +0800 Subject: [PATCH 10/14] test(vllm): give the native grid-numbering stub the state #15110 reads test_kvwarm_native_grid_numbering_preserves_each_execution runs the real capacity probe and decode dispatch on a stub built with InstrumentedScheduler.__new__. After the merge of main, both paths read state that #15110 added, and the test failed with AttributeError: - the grid-invariants digest now includes _bench_measurement_protocol(), which reads _bench_vocab_size; - _bench_pop_next() opens an eager_shape warmup record for eager replicas, which reads _bench_results. Set both on the stub, as _digest_stub and _benchmark_save_stub already do. No production code changes. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index 6bb3e455bf1d..fd7607fb9eb2 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -5551,6 +5551,7 @@ def test_kvwarm_native_grid_numbering_preserves_each_execution( stub._bench_decode_capture_sizes = [] stub._bench_decode_cudagraph_mode = "NONE" stub._bench_cudagraph_capture_sizes = [] + stub._bench_vocab_size = 0 config = KVCacheConfig( num_blocks=num_blocks, kv_cache_tensors=[], @@ -5577,6 +5578,7 @@ def test_kvwarm_native_grid_numbering_preserves_each_execution( stub._fpm_dp_rank = rank stub._bench_active_req_ids = set() stub._bench_current_point = None + stub._bench_results = [] stub._bench_drain_pending = False stub._bench_phase = _BenchPhase.DECODE_SWEEP stub._bench_start_timing = lambda: None From af3c1195ef7e0ee11c7c0ee0fcc199487367e26a Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Sat, 10 Oct 2026 12:14:44 +0800 Subject: [PATCH 11/14] fix(vllm): record prompt evidence for native warm-up decode points #15110 hashes the injected prompts of every benchmark point into benchmark_measurement.prompts, in each admission path: _bench_inject_prefill, _bench_inject_fake_decode and _kvwarm_inject_borrowed. Native exact-context points are admitted by _kvwarm_resume_native, which continues the stage's own prefilled requests and recorded nothing. Every native real-KV row therefore reported prompts.status "unavailable", although its prompts are known and reproducible: the slot's chain text, seeded by the content seed and the DP rank, cut to the injected context. Hash the stage prompts when the requests resume. Each request's num_tokens is its prefill length, the injected context ctx - 1. The admission token is a sampled continuation, which the protocol lists as unobserved. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../src/dynamo/vllm/instrumented_scheduler.py | 5 +++ .../tests/test_vllm_instrumented_scheduler.py | 37 +++++++++++++++++++ 2 files changed, 42 insertions(+) diff --git a/components/src/dynamo/vllm/instrumented_scheduler.py b/components/src/dynamo/vllm/instrumented_scheduler.py index 9f385f429543..9ca9d13df7c5 100644 --- a/components/src/dynamo/vllm/instrumented_scheduler.py +++ b/components/src/dynamo/vllm/instrumented_scheduler.py @@ -7113,6 +7113,11 @@ def _kvwarm_inject_borrowed(self, context_lengths) -> "SchedulerOutput": def _kvwarm_resume_native(self) -> SchedulerOutput | None: """Continue the prefills themselves; their private window state is exact.""" req_ids = self._kvwarm_chain_ids + # The stage prompts are this point's injected input. The admission + # token is a sampled continuation, which prompt evidence never covers. + self._bench_record_prompt_evidence( + [self._kvwarm_chain_prompts[req_id] for req_id in req_ids] + ) self._kvwarm_native_active = True self._kvwarm_native_resume_ids = set(req_ids) self._bench_active_req_ids.update(req_ids) diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index fd7607fb9eb2..8b7c92d686c2 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -6099,6 +6099,43 @@ def test_kvwarm_native_decode_dispatch_uses_real_continuation_and_provenance(): assert stub._kvwarm_meta["points_real_kv"] == 1 +@pytest.mark.core +def test_kvwarm_native_decode_records_stage_prompt_evidence(): + stub, requests = _kvwarm_native_resume_stub(False) + point = BenchmarkPoint( + point_type="decode", benchmark_id=7, batch_size=2, total_kv_read_tokens=9 + ) + stage_prompts = [] + for request, context in zip(requests, (4, 3), strict=True): + request.num_computed_tokens = context + request.all_token_ids = list(range(context + 1)) + stage_prompts.append(list(range(10, 10 + context))) + stub._kvwarm_chain_prompts[request.request_id] = stage_prompts[-1] + stub._kvwarm_stage_point = point + stub._kvwarm_stage_batch = 2 + stub._kvwarm_building = False + stub._kvwarm_plan = {stub._kvwarm_plan_key(point): 4} + stub._bench_grid = deque([point]) + stub._bench_current_point = None + stub._bench_drain_pending = False + stub._bench_frees_pending = lambda: False + stub._bench_stop_at_timeout_boundary = lambda phase: False + stub.max_model_len = 8192 + + assert stub._bench_step_decode() is not None + + # The stage prompts are the injected input. The admission token is a + # sampled continuation, which prompt evidence does not cover. + expected = InstrumentedScheduler.__new__(InstrumentedScheduler) + expected._bench_record_prompt_evidence(stage_prompts) + assert stub._bench_prompt_evidence == expected._bench_prompt_evidence + assert stub._bench_prompt_evidence["status"] == "recorded" + assert [row["num_tokens"] for row in stub._bench_prompt_evidence["requests"]] == [ + 4, + 3, + ] + + def _kvwarm_uncovered_dispatch_stub(warm_eligible, skip_reason=None): """Native-layout decode dispatch of a point no warm-up stage covers.""" stub, _ = _kvwarm_native_resume_stub(False) From e55735bbd1938f5bd2675d37f2476327d0395be4 Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Sat, 10 Oct 2026 12:23:17 +0800 Subject: [PATCH 12/14] test(vllm): run the native admission-cap test with vLLM 0.31.0's name vLLM 0.31.0 renamed the single-type KV cache manager's _max_admission_blocks_per_request to max_admission_blocks_per_request. The scheduler reads both through _kvwarm_admission_cap, but test_kvwarm_native_capacity_uses_sliding_window_admission_caps stubbed only the old name, so it did not cover the native capacity path on the pinned vLLM. Parametrize it over both names, as main's 0.31.0 bump did for test_capacity_digest_tracks_admission_cap. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../dynamo/vllm/tests/test_vllm_instrumented_scheduler.py | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py index 8b7c92d686c2..856577f040c9 100644 --- a/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py +++ b/components/src/dynamo/vllm/tests/test_vllm_instrumented_scheduler.py @@ -5944,7 +5944,11 @@ def test_kvwarm_native_capacity_counts_resident_blocks_not_evicted_placeholders( @pytest.mark.core -def test_kvwarm_native_capacity_uses_sliding_window_admission_caps(): +@pytest.mark.parametrize( + "cap_attribute", + ["_max_admission_blocks_per_request", "max_admission_blocks_per_request"], +) +def test_kvwarm_native_capacity_uses_sliding_window_admission_caps(cap_attribute): stub = InstrumentedScheduler.__new__(InstrumentedScheduler) stub.max_model_len = 1024 stub.num_lookahead_tokens = 0 @@ -5952,7 +5956,7 @@ def test_kvwarm_native_capacity_uses_sliding_window_admission_caps(): coordinator=SimpleNamespace( single_type_managers=[ SimpleNamespace(block_size=16), - SimpleNamespace(block_size=4, _max_admission_blocks_per_request=4), + SimpleNamespace(block_size=4, **{cap_attribute: 4}), ] ) ) From 8d484b41c2b373250aac3e2db9c1a4cbf674be68 Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Sat, 10 Oct 2026 12:43:58 +0800 Subject: [PATCH 13/14] docs(vllm): describe native warm-up in the benchmark evidence docs The measurement-protocol section of the forward pass metrics trace reference described decode real-KV warmup as shared chains only. Say that the unobserved decode warmup history covers shared chains and native exact-context stages, that both record their stages in kvwarm.stages (kvwarm.initialization_strategy names the native strategy), and that rows measured after a native stage hash the stage's prefill prompts, ctx - 1 tokens per request, because the admission token is sampled. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../reference/observability/forward-pass-metrics-traces.mdx | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/docs/fern/pages/reference/observability/forward-pass-metrics-traces.mdx b/docs/fern/pages/reference/observability/forward-pass-metrics-traces.mdx index ecd5a1663c7c..e5cb6b2de5a2 100644 --- a/docs/fern/pages/reference/observability/forward-pass-metrics-traces.mdx +++ b/docs/fern/pages/reference/observability/forward-pass-metrics-traces.mdx @@ -170,7 +170,7 @@ The shared protocol contains: | `timing_metric` | `"scheduler_wall_time"`; the existing scheduler timing boundary is unchanged, not replaced with CUDA-event timing | | `execution_evidence` | Versioned contract for per-sample observations and warmup records; identifies their field names, process-local clock and forward indices, graph-statistics source, and validation scope | | `input_evidence_scope` | `"injected_prompt_token_ids"`; hashes cover the prompt ids available at request injection | -| `unobserved` | Lists inputs or context not established by this evidence: sampled continuation tokens, KV-cache tensors, recurrent-state tensors, decode real-KV chain warmup history, equivalent execution history, and equivalent graph, kernel, or cache warming | +| `unobserved` | Lists inputs or context not established by this evidence: sampled continuation tokens, KV-cache tensors, recurrent-state tensors, decode real-KV warmup history (shared chains or native exact-context stages), equivalent execution history, and equivalent graph, kernel, or cache warming | | `preparation` | Warm-up iteration count, prefill real-seed and decode real-KV settings, and giant-KV threshold and repeat count | Each retained FPM's `benchmark_measurement` contains: @@ -179,7 +179,7 @@ Each retained FPM's `benchmark_measurement` contains: |---|---| | `point_key`, `dp_rank` | Coordinate identity and the rank that produced this observation; the point key is independent of the grid digest and benchmark point id | | `prompts.status` | `"recorded"` when injected prompt ids were captured, otherwise `"unavailable"`; missing evidence is not a match | -| `prompts.sha256`, `prompts.requests` | Combined prompt digest and the request slots, lengths, and individual prompt digests | +| `prompts.sha256`, `prompts.requests` | Combined prompt digest and the request slots, lengths, and individual prompt digests; on rows measured after a native exact-context stage, the hashed prompts are the stage's prefill prompts, `ctx - 1` tokens for a request with context length `ctx`, because the admission token is sampled | | `preparation` | Actual grid digest, number of completed points before this point, KV-seeding regime, and `warmup_records_before`, the number of prior warmup records, including failed attempts | | `expected_internal_samples` | Number of scheduler steps expected for this point, including an admission step when the path uses one | | `raw_fpms` | Individual FPMs before reduction, in collection order, including discarded admission measurements and slow samples; each carries its own `benchmark_sample` execution evidence | @@ -225,7 +225,7 @@ Observed dispatch is separate from the point's `expected_cudagraph_mode` and `ex ### Warmup Completion Evidence -Each rank artifact's `warmup_evidence` contains `status` (`recorded` or `unavailable`) and an ordered `records` list. The list covers global warmup, discarded eager-shape replicas, real-prefix seeding, and real-prefix same-shape preparation. An empty recorded list means none of these attempts was recorded. Decode KV-chain preparation retains its existing `kvwarm` metadata; its individual warmup attempts are not recorded here. +Each rank artifact's `warmup_evidence` contains `status` (`recorded` or `unavailable`) and an ordered `records` list. The list covers global warmup, discarded eager-shape replicas, real-prefix seeding, and real-prefix same-shape preparation. An empty recorded list means none of these attempts was recorded. Decode real-KV warmup, by shared prefill chains or native exact-context stages, is not in this list: each stage is recorded in `kvwarm.stages`, and `kvwarm.initialization_strategy` is `native_exact_context` when the native strategy is selected. | Record field | Meaning | |---|---| From 5d654ec31930d3bf8b1a3e7998addae04b78870a Mon Sep 17 00:00:00 2001 From: Yiming Liu Date: Sat, 10 Oct 2026 13:20:40 +0800 Subject: [PATCH 14/14] docs(vllm): limit the native prompt-length note to real_kv rows A point whose native stage failed falls back to synthetic KV, and its fake_fallback row hashes the full ctx-token prompts of fake injection. The note on native prompt evidence covered every row measured after a native stage; limit it to real_kv rows prepared by one. Signed-off-by: Yiming Liu Co-Authored-By: Claude Opus 5.5 --- .../reference/observability/forward-pass-metrics-traces.mdx | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/fern/pages/reference/observability/forward-pass-metrics-traces.mdx b/docs/fern/pages/reference/observability/forward-pass-metrics-traces.mdx index e5cb6b2de5a2..c069b42773be 100644 --- a/docs/fern/pages/reference/observability/forward-pass-metrics-traces.mdx +++ b/docs/fern/pages/reference/observability/forward-pass-metrics-traces.mdx @@ -179,7 +179,7 @@ Each retained FPM's `benchmark_measurement` contains: |---|---| | `point_key`, `dp_rank` | Coordinate identity and the rank that produced this observation; the point key is independent of the grid digest and benchmark point id | | `prompts.status` | `"recorded"` when injected prompt ids were captured, otherwise `"unavailable"`; missing evidence is not a match | -| `prompts.sha256`, `prompts.requests` | Combined prompt digest and the request slots, lengths, and individual prompt digests; on rows measured after a native exact-context stage, the hashed prompts are the stage's prefill prompts, `ctx - 1` tokens for a request with context length `ctx`, because the admission token is sampled | +| `prompts.sha256`, `prompts.requests` | Combined prompt digest and the request slots, lengths, and individual prompt digests; on `real_kv` rows prepared by a native exact-context stage, the hashed prompts are the stage's prefill prompts, `ctx - 1` tokens for a request with context length `ctx`, because the admission token is sampled | | `preparation` | Actual grid digest, number of completed points before this point, KV-seeding regime, and `warmup_records_before`, the number of prior warmup records, including failed attempts | | `expected_internal_samples` | Number of scheduler steps expected for this point, including an admission step when the path uses one | | `raw_fpms` | Individual FPMs before reduction, in collection order, including discarded admission measurements and slow samples; each carries its own `benchmark_sample` execution evidence |