Repository navigation
[Router] k8s e2e for cache-aware peer bootstrap (12/13) - #40698
Merged
Merged
Conversation
This was referenced Sep 22, 2026
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-11-component-test
branch
from
September 27, 2026 02:28
3650334 to
bb5bf77
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-12-k8s-e2e
branch
from
September 27, 2026 02:28
8bda96f to
5e5c384
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-11-component-test
branch
from
September 27, 2026 06:15
bb5bf77 to
13e8986
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-12-k8s-e2e
branch
from
September 27, 2026 06:15
5e5c384 to
f6d90bb
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-11-component-test
branch
from
September 27, 2026 07:04
13e8986 to
9334544
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-12-k8s-e2e
branch
from
September 27, 2026 07:04
f6d90bb to
c51ab53
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-11-component-test
branch
from
September 27, 2026 09:07
9334544 to
207c0ae
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-12-k8s-e2e
branch
from
September 27, 2026 09:07
c51ab53 to
bd6a6b5
Compare
ShangmingCai
force-pushed
the
router-peer-bootstrap-11-component-test
branch
3 times, most recently
from
September 28, 2026 10:49
0442da1 to
239ea27
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-11-component-test
branch
from
September 28, 2026 21:54
239ea27 to
db42de0
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-12-k8s-e2e
branch
from
September 28, 2026 21:54
bd6a6b5 to
232520a
Compare
…the pump The producer half now has a consumer. A booting replica asks its siblings for a cache-aware tree, vets what comes back, and hands it to the pump — which has been able to graft one since the previous change but had nothing to graft. The sweep retries rather than asking once. Worker discovery regularly completes before the peer watch has delivered its first EndpointSlice list, so a single pass sees zero candidates and abandons, and the joining replica boots cold beside warm siblings. The bootstrap deadline bounds the whole search, not each request, so a fleet of slow peers cannot outlast the readiness gate. The per-fetch timeout is a strict fraction of that deadline, not the deadline itself. `reqwest`'s total timeout would otherwise let the first unresponsive peer starve every other candidate — the outer deadline then cancels the sweep mid-fetch and the replica boots cold having tallied no peer outcome at all. `--kv-bootstrap-fetch-timeout-cap-ms` bounds the derivation from above, so a deliberately generous readiness budget still cannot park on one hung peer; the router warns at startup when the derived value lands below the floor a multi-megabyte body needs, because that failure otherwise looks like every peer being unreachable with the configured cap looking blameless. Accepting a snapshot is not the same as accepting a peer. A body can vet cleanly and still know nothing about the ranks being bootstrapped, so the sweep keeps looking rather than ending on it; a peer whose snapshot is permanently incompatible is never re-fetched; and one that simply has nothing right now sits out a few passes, because retrying every 250ms means re-downloading a multi-megabyte tree from a replica that is itself serving traffic. The sit-out doubles on each repeat miss by the same peer, up to 30s, so the sweep keeps its 250ms pickup of new siblings without hammering one that keeps failing. Every exit path — found, no peers, fleet cold, timed out — sends exactly one `PumpControl`, which is what releases the ranks from `Pending`. The cold-fleet verdict is discarded if the candidate set changed mid-pass: a peer that was never consulted holds state that, unlike post-subscription events, cannot be recovered later. One sweep per discovered worker for now. A fleet's worth of workers therefore fetches the same fleet-wide body once per worker; the coordinator that folds a discovery burst into a single fetch is the next change. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JBP3reyKKk4TPeZuppMPmW
…for it A new engine's ranks sit Pending on every router at once, and while a rank is Pending its deltas stay off the tree. No sibling's tree ever holds nodes for it, `covers_any` never succeeds, and each router waits out the deadline on siblings that are waiting too; the whole-tree `fleet_is_cold` verdict never fires because every sibling is warm. - `PeerSnapshot` gains `empty_ranks`: live ranks the export holds no node for, including ranks the producer is itself still bootstrapping. Additive on the wire: it defaults to empty and is omitted when empty, so an older producer reads as "no evidence". - A rank settles cold mid-sweep once every candidate names it there, or is hopeless as a whole. A silent, unanswered or unreachable peer vetoes. The sweep goes on for its other ranks. - A re-fetch demands an export newer than that peer's last answer. A floor fixed at sweep start was met by the producer's cached export for the whole sweep, so a retry was served the same non-covering answer until the deadline. The producer already pins the demanded instant on arrival, so a herd still shares one build. - Each pass sweeps only obligations still Pending on their own incarnation, and a sweep whose every rank left Pending (resolved from its origin, settled, forgotten) ends `ranks_resolved` and sends nothing. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0169aRN6guH335zF1pG456FC
A peer snapshot is fleet-wide: one body carries the whole worker table, every cursor, and the whole tree, so a single fetch already contains everything every pending rank needs. `add_worker` fires once per discovered engine, so sweeping from there re-downloads that same document once per engine — and the document grows with the fleet while the fetch count grows with it too. On a 168-engine fleet that is 168 fetches of a multi-megabyte body per booting replica, most of which time out and leave their ranks cold. `add_worker` now hands its obligations to a coordinator that serialises bootstrap into one sweep at a time and lets every rank pending when that sweep lands share its snapshot. The fetch is not delayed to collect a batch first because it does not need to be: the sweep IS the collection window — a multi-megabyte transfer takes far longer than an EndpointSlice watch event takes to deliver a fleet, so ranks discovered while it is in flight merge into the same delivery at no latency cost. Sharing is safe because the invariant is a SEQUENCE condition, not a wall-clock one. Vetting runs against the live-worker set at the moment the body arrives, so a worker discovered during the fetch is already covered; and the watermark check adjudicates every rank independently, so a rank whose publisher did advance in between is discarded to `Gap` rather than spliced over a hole. Sharing one snapshot only helps if it is fresh enough for every rank sharing it, and ranks in a batch begin holding at different moments — so the sweep asks with the strictest floor among them. That guarantee covers the ranks it STARTED with; a batch merged in afterwards was not represented in the request. Discovery accepts that trade, a gap retry does not, which is what `LateJoin` distinguishes. That retry is the other half of this change. A gap is the costliest failure — a snapshot fetched, grafted, then thrown away because the live stream did not join its watermark — and a fresher snapshot usually splices. It needs somewhere to put the rank back, which is the queue this change introduces, so until now a gapped rank was terminal. The tracker caps it at one retry per rank, and the retry is stamped NOW and refuses to ride an in-flight sweep, or it would be adjudicated against a snapshot taken before the gap and re-gap by construction. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01JBP3reyKKk4TPeZuppMPmW
…olved The coordinator folds ranks queued while a sweep runs into that sweep's delivery. A sweep whose own ranks all left Pending without it (resolved from their origin, settled cold per rank) ends `ranks_resolved` and sends no control message, so a late joiner folded into it would stay Pending with nothing left to release it, holding its batches and /readyz. Hand whatever the sweep never spoke for back for a sweep of its own. A rank the sweep settled cold mid-sweep is owed nothing either, but it can still read Pending until the pump drains that release, so the sweep now reports the obligations it settled and the coordinator drops them by obligation rather than by state. Otherwise such a rank would buy a second sweep here, or a coverage retry after a `found` verdict. Also pins that a rank back in Pending for a gap retry catches an in-place engine restart against its held queue: the retry keeps the first attempt's batches and no cursor. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0169aRN6guH335zF1pG456FC
A graft whose held queue was empty leaves its watermark unchecked until a live batch arrives, and nothing guarantees one ever does. An idle rank would serve grafted state forever with its continuity unproven, and a `BlockRemoved` lost in the subscribe window would survive as a permanent false cache hit that no later event corrects. Silence is not evidence, so the bounded wait ends in a QUESTION, not a verdict. Reaching the deferred path means nothing arrived between subscribing and grafting, and the subscriber is live before the fetch — so silence is far more often "this rank published nothing" than "we lost a delta". Discarding on a timer would cost every quiet fleet its warm tree, which is the regression this feature exists to prevent. So the pump asks the fleet whether the publisher moved past the watermark. Any peer's cursor is admissible: sequence numbers are the publisher's, so a peer reporting one above ours proves a batch was emitted that we never saw. A peer too cold to bootstrap from is still a valid witness, which is why the probe reads the wire cursor directly instead of vetting. It asks for the cursor table alone, not a snapshot — the question is answered completely by one integer per rank, and fetching a tree to read it made the proof cost scale with the tree, so a fleet large enough to need bootstrap was also the fleet that could not afford to prove it. Only a witness ABOVE the watermark discards. A graft no batch and no peer can speak to is KEPT, tallied `warm_unwitnessed` rather than `warm`, after `MAX_UNKNOWN_PROBES` unanswerable rounds — without a stop, a single-replica deployment would re-probe forever and its verdict would never resolve in the metrics. The probe runs off-pump because it does network I/O; the verdict comes back over the pump's own control channel so the tree write stays on the single writer. An outstanding probe blocks relaunch inside its timeout and stops blocking past it: a verdict that never lands (a panicked probe task, say) would otherwise freeze the rank in `Recovered` with no outcome ever recorded. Also re-attaches `demote_unproven_rank`'s doc comment, which sat above `requeue_gapped_rank` describing the wrong function. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01YZorgAox1CpLNHSjzpcxdb
…robe A splice probe answers for the graft it was launched about. Once a batch 0 has replaced that graft with a restarted engine's stream, a verdict still in flight, whose witnesses count in the dead numbering, must not demote the new stream's state. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_0169aRN6guH335zF1pG456FC
…t copied Component-level proof over the real transport: an axum server serving `/internal/kv_snapshot`, a reqwest client fetching it, JSON over a loopback socket. The interesting failure modes live in the wire format and the vetting step, not in the tree algebra — that is already covered by the unit tests on `export_snapshot` / `restore_snapshot`. The equivalence assertion is deliberately behavioural: for a large query set, `match_prefix` must return the same matched length and the same carrier set on both replicas. That is the property routing actually depends on; comparing node counts alone would pass while routing diverged. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01YZorgAox1CpLNHSjzpcxdb
A kind-cluster proof of the two questions the feature exists to answer: a replica joining a warm fleet ends up with the same cache-aware view as the replicas it bootstrapped from, and it keeps up with new engine events afterwards — the snapshot is spliced UNDER the live stream, not substituted for it. Plus the rolling-update hazard: new replicas must not bootstrap from each other and inherit an empty tree. Ships the KV-publishing fake worker the test drives, and the manifest that wires `POD_NAME`/`POD_NAMESPACE`/`POD_IP` from the downward API and the EndpointSlice RBAC the peer watch needs — deferred from the discovery commit to here, where there is something for it to exercise. Three anti-flake decisions, because the naive version of this test is flaky: events are driven by a POST rather than a timer, so nothing is ever in flight the test did not ask for; views are compared only after every replica's cursor has provably caught up to the worker's `last_seq`, never on a fixed sleep; and nothing asserts on a transient mid-rollout state — the scale-up case is asserted directly and the rollout case only after `rollout status` completes. Views are compared canonically, since snapshot node order follows per-shard hash-map iteration. Not yet run: this needs a kind cluster, so it is verified by construction only. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01YZorgAox1CpLNHSjzpcxdb
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-11-component-test
branch
from
September 29, 2026 21:48
db42de0 to
d411630
Compare
Kangyan-Zhou
force-pushed
the
router-peer-bootstrap-12-k8s-e2e
branch
from
September 29, 2026 21:48
232520a to
b8ab8e6
Compare
ShangmingCai
force-pushed
the
router-peer-bootstrap-11-component-test
branch
3 times, most recently
from
September 30, 2026 10:59
69964bc to
70c47a6
Compare
Base automatically changed from
router-peer-bootstrap-11-component-test
to
main
September 30, 2026 11:19
ShangmingCai
marked this pull request as ready for review
September 30, 2026 11:19
Collaborator
|
/tag-and-rerun-ci |
ShangmingCai
approved these changes
Oct 1, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Stack 12 of 13. Base:
router-peer-bootstrap-11-component-test. 6 files changed, 941 insertions(+), 2 deletions(-)The stack
export_snapshot/restore_snapshot; nothing calls them yetGET /internal/kv_snapshot— the producer half; nothing consumes it yet--kv-peer-selector, the peer registry, the_peersgauges; no behaviour changeBootstrapTracker+ the fetch client; the consumer's state and transportVettedSnapshot::from_wire— the only bridge from wire bytes to the treePendingrank's batches and grafts a snapshot handed to itmatch_prefixanswers as the source/readyzholds until bootstrap settles;--kv-bootstrap-seed-required; the metricsEach PR's base is the branch below it, so every diff shown here is that
PR's own change. Review bottom-up; GitHub retargets each child to
mainasits parent merges.
What the series does
A router replica subscribes to each worker's KV topic mid-stream, so every
block already resident in the engine's radix cache is invisible to it — and
engines publish
BlockStoredonly as they insert, so a prefix cached hoursago is never re-announced. A cache-blind replica then scatters the prefixes the
warm replicas were keeping consolidated, degrading the engines' locality for
the whole fleet; a rolling update does that to every replica in turn. This
series makes a booting replica pull a tree snapshot from a warm sibling over
HTTP and graft it beneath its live delta stream. Off unless
--kv-peer-selectoris set.
Supersedes #39750, which carried the same work as one branch on a stale base.
What this change does
A kind-cluster proof of the two questions the feature exists to answer: a
replica joining a warm fleet ends up with the same cache-aware view as the
replicas it bootstrapped from, and it keeps up with new engine events afterwards
— the snapshot is spliced UNDER the live stream, not substituted for it. Plus
the rolling-update hazard: new replicas must not bootstrap from each other and
inherit an empty tree.
Ships the KV-publishing fake worker the test drives, and the manifest that wires
POD_NAME/POD_NAMESPACE/POD_IPfrom the downward API and the EndpointSliceRBAC the peer watch needs — deferred from the discovery commit to here, where
there is something for it to exercise. The fake worker declares mirrors of the
engine's
EventBatch/BlockStoredstructs with the same msgspec options, somsgspec itself produces the engine's wire encoding — each event a tagged MAP
keyed by
type— rather than a hand-built shape; the router rejects the oldertagged-array event encoding outright, so a hand-built one would test nothing
but a decode failure.
Three anti-flake decisions, because the naive version of this test is flaky:
events are driven by a POST rather than a timer, so nothing is ever in flight
the test did not ask for; views are compared only after every replica's cursor
has provably caught up to the worker's
last_seq, never on a fixed sleep; andnothing asserts on a transient mid-rollout state — the scale-up case is asserted
directly and the rollout case only after
rollout statuscompletes. Views arecompared canonically, since snapshot node order follows per-shard hash-map
iteration.
Not yet run: this needs a kind cluster, so it is verified by construction only.
Fresh-engine fix
fake_kv_worker.pynow starts at seq 0 like the engine'sitertools.count(); its docstring claimed to match the engine while starting at 1.last_seqreads -1 before the first publish, the test's existing "none" value. The e2e flow is unchanged: the added replicas join after events are published, so their first held batch is past the origin and they still take a peer snapshot.Tests
cargo fmt --check,cargo clippy --all-targets -- -D warnings, and the lib +component + proxy suites all pass on this branch on its own, not only on the
tip of the stack.
Not yet run: the k8s e2e needs a kind cluster, so it is verified by
construction only.
Review pass
Reviewed with
/code-review --fixand/simplify; fixes were folded into this PR's own commit and the stack was re-verified tier by tier (cargo fmt --check,cargo clippy --all-targets -D warnings, lib + component + proxy tests on every branch).tests/e2e/k8s_integration/; it touches none of the splitsrc/files.🤖 Generated with Claude Code
https://claude.ai/code/session_016HmJvHV7QDPk3qAjQYzthd
CI States
Latest PR Test (Base): ✅ Run #36707894684
Latest PR Test (Extra): ❌ Run #36707894302
Latest PR Test (AMD ROCm 10): ➖ No AMD PR run found for this commit.