Skip to content

feat(scheduler): add HA via Kubernetes Lease leader election - #341

Open
NickNYU wants to merge 14 commits into
kvcache-ai:mainfrom
NickNYU:feat/scheduler-ha-leader-election
Open

NickNYU wants to merge 14 commits into
kvcache-ai:mainfrom
NickNYU:feat/scheduler-ha-leader-election

Conversation

@NickNYU

@NickNYU NickNYU commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

What

Add optional high availability to the scheduler via Kubernetes Lease leader election. N scheduler replicas compete for a single coordination.k8s.io/Lease; exactly one leader schedules and processes node reports, standbys stay liveness-healthy and serve LookupNode/GetNode from the shared Redis bindings. On failover, the new leader rebuilds state by pulling every node's admin snapshot (sync-node-snapshots) instead of waiting out a reporter backoff, and applies recovery-window semantics until the rebuild completes. Off by default; with it off, behavior is identical to today.

Why

The scheduler is a single point of failure today (replicas: 1): scheduling, assignment writes, node inspection, and P2P APIs all go down with it, and query-only replicas only preserve data-plane lookups. Running multiple read/write replicas directly is unsafe — heartbeat state is replica-local, which silently breaks node_resource_limit (#191, measured 3.0x overshoot on 8 nodes).

Related issue

Closes #259

Relates to #191 — the #191 failure mode is fixed when leader election is enabled: exactly one scheduler processes heartbeats and schedules, and the recovery window never treats missing observations as free capacity. With election disabled (the default), behavior is unchanged, so this PR intentionally does not auto-close #191; whether to close it (or keep it for the doc note, direction 1 in the issue) is a maintainer call. The full #191 reproduction scenario is spec'd for Kind/CI (services/scheduler/kind), not run locally.

Scope and non-goals

In scope:

  • scheduler.leader_election config surface + startup validation (fail fast; mutually exclusive with --query-only; client-go timing rules). Redis stays optional; without it the documented degraded mode applies.
  • Single exported Leadership contract + factory + null object: gate interceptor (read/write matrix), leader readiness health service scheduler.v1.Scheduler/leader, client-go Lease elector with fencing-by-exit.
  • Recovery window: fresh-observations-only scheduling (scheduler.node_resource_limit is silently unenforced for most nodes when the scheduler runs more than one replica #191 fix semantics), Unavailable-not-NotFound lookups during rebuild.
  • Sync-node-snapshots: post-acquisition pull of every node's admin /nodes through the shared heartbeat ingest path (observations only, no binding reconcile — see Design); concurrency via leader_election.snapshot_pull_concurrency.
  • Binding TTL floor under election (max(binding_ttl, lease_duration + 90s)) so bindings outlive the failover budget.
  • Manifests: new deploy/k8s/overlays/ha (3 replicas, readiness probe on the leader health service, election config, secret-mounted admin key); base manifests untouched except additive leases RBAC (inert when off).

Non-goals:

  • No protobuf changes; no gateway or node-side code changes; no heartbeat protocol changes.
  • No active-active or sharded scheduling; standbys collect no observations (possible future work, deliberately out of scope).
  • No separate query-only Deployment changes; existing query-only mode is untouched and mutually exclusive with election.
  • Kind integration cases are spec'd (build-tagged) but not wired into CI yet.

Design and behavior changes

flowchart LR
    subgraph sched["Scheduler Deployment · N replicas, leader-elected"]
        L["Leader<br/>(ready, serving)"]
        S1["Standby<br/>(reads from Redis)"]
        S2["Standby<br/>(reads from Redis)"]
    end
    GW["Gateway"]
    RD[("Redis<br/>bindings")]
    NODE["agentenv server nodes"]
    LEASE[["K8S Lease"]]

    NODE -- "④ heartbeat" --> L
    GW -- "① Schedule / ③ RecordAssignment (write)" --> L
    GW -- "② LookupNode (read)" --> sched
    L <-- "binding write" --> RD
    S1 <-. "binding read" .-> RD
    S2 <-. "binding read" .-> RD
    L <-. "acquire / renew" .-> LEASE
    S1 <-. "elect" .-> LEASE
    S2 <-. "elect" .-> LEASE
    L <-. "⑤ pull /nodes (refresh)" .-> NODE
Loading

Key behaviors:

  • Read/write gate: health checks never gated; leader serves everything; standbys with Redis serve LookupNode/GetNode, reject the rest with Unavailable("not the leader"); without Redis reads retry onto the leader. The RPC table classifies every method explicitly and panics on unclassified ones, so proto upgrades fail loudly.
  • Readiness split: liveness reports the overall health (SERVING from start, unchanged); readiness reports scheduler.v1.Scheduler/leader (SERVING only while holding the lease), so Service endpoints contain only the leader and gateway needs no changes.
  • Fencing: OnStoppedLeading flips readiness, then cancels the process root context — the ex-leader exits through the same graceful shutdown path as SIGTERM and restarts as a standby. A partitioned ex-leader stops serving within renew_deadline.
  • Recovery window (bounded by one report_ttl): heartbeats accepted immediately; Schedule considers only nodes observed at or after acquisition (via NodeRegistry.LastReportAt), zero fresh observations returns Unavailable — never a fallback to unobserved nodes; LookupNode/GetNode return Unavailable (not NotFound) for missing entries while the window is open.
  • Sync-node-snapshots: on acquisition the leader pulls every node's admin /nodes snapshot (HTTP, x-api-key from SCHEDULER_NODE_ADMIN_API_KEY mounted from the shared agentenv-auth Secret) and feeds it through the exact same ingest path as Heartbeat, so observations rebuild in seconds rather than a reporter backoff (~60s worst case). Bindings are deliberately not reconciled from a pull: a pull cannot see paused sandboxes or in-flight template builds, so reconciling would delete live bindings on failover (feat(scheduler): add HA via Kubernetes Lease leader election #341 review). Binding refresh stays with heartbeats, and the TTL floor covers the gap.
  • Object model: one exported interface (Leadership) + one factory + a null object for the disabled path; implementations (manager, snapshot, elector) are package-private. State is an immutable snapshot swapped atomically by the manager; the Service sees leadership only through a 2-method read-only seam.

Compatibility and operations

  • Public API or generated protocol: none. No protobuf changes; gateway and runtime nodes need no changes.
  • Configuration or defaults: leader_election is off by default; disabled behavior is identical to today for every deployment shape (validated in unit tests). New optional keys: leader_election.* (durations as strings, e.g. "15s"), scheduler.node_admin_api_key (credential — prefer SCHEDULER_NODE_ADMIN_API_KEY env/Secret).
  • Snapshot manifest, artifact layout, or storage format: N/A (no changes).
  • Upgrade and rollback: upgrade the scheduler binary first (it serves today's behavior while election is off and registers the leader health service), then apply manifests (RBAC, readiness probe, replicas, config). Rollback: disable the config and scale back to 1 replica.
  • Host requirements, permissions, ports, or dependencies: Kubernetes only — scheduler Role needs coordination.k8s.io/leases get/create/update (added to deploy/k8s/base/role.yaml). Redis optional but recommended; without it failover loses routing for pre-failover sandboxes (documented). No new Go module dependencies beyond client-go's existing testing packages.

Validation

  • make fmt
  • make clippy
  • make test-unit
  • Relevant Rust integration tests
  • make -C services test (required when services/ changes)
  • Generated clients/server regenerated with the documented make target
  • Documentation updated
  • Benchmarks or performance comparison completed

Commands and results:

# Note: system Go is 1.20.6 and cannot parse go.mod (go 1.25.0);
# all Go commands below used a go1.25.1 toolchain.

go1.25.1 test ./...            # all packages pass (gateway, scheduler, config, logging)
go1.25.1 vet ./...             # clean
gofmt -l scheduler/ shared/    # clean
kubectl kustomize deploy/k8s/overlays/ha
  # renders: replicas=3, readiness -service=scheduler.v1.Scheduler/leader,
  # leader_election config, leases RBAC, secret-mounted admin key

New tests cover: config validation (A), gate matrix (B), election lifecycle on a fake clientset (C1 exactly-one-leader, C2 failover, C3 fencing on renew failure + callback ordering), recovery window semantics (D: fresh-only, zero-fresh Unavailable, ingest path sharing, pull failure, lookup window semantics), binding TTL floor, and admin snapshot fetcher mapping/auth/ingest compatibility.

Skipped checks and reasons:

  • Rust checks (make fmt, clippy, test-unit, integration): no Rust changes in this PR.
  • Generated code: no protobuf changes.
  • Kind integration (groups C/E/F end-to-end): spec'd under services/scheduler/kind with a build tag; the repo has no Kind harness and the dev machine is macOS. Left for CI; the fake-clientset suite covers election semantics offline.
  • Benchmarks: no performance-sensitive changes; the pull path is a bounded one-off fan-out per failover.

Risks and reviewer notes

  • Biggest surface: services/scheduler/internal/leadership.go (facade + gate + elector) and service.go (recovery window). Suggested read order: 92b2089 (config) → f92133c (leadership domain) → f1a67c3 (sync-node-snapshots) → aac26dd (main wiring + TTL floor) → 8a36b6c (Kind specs) → 2ac0f4c (manifests/docs).
  • With election disabled, the gate interceptor is a pass-through and every leadership method is a no-op; the "identical to today" claim is covered by unit tests.
  • The recovery window intentionally prefers Unavailable over optimistic scheduling right after failover — a deliberate availability-vs-correctness trade, matching the issue's "missing observations are not unlimited capacity".
  • ListNodes is currently gated on standbys (reads the informer-backed registry, semantically standby-serveable); kept gated for now, flagged for a follow-up decision.

Checklist

  • The PR contains one coherent change and no unrelated formatting or refactoring.
  • New behavior is covered by tests, or I explained why testing is impractical.
  • Logs and examples contain no credentials, tokens, or private registry information.
  • I did not manually edit generated code without updating its source and regenerating it.

NickNYU and others added 6 commits October 5, 2026 13:10
Introduce scheduler.leader_election config (enabled, lease_name,
lease_namespace, lease_duration, renew_deadline, retry_period) for
Kubernetes Lease-based scheduler HA (kvcache-ai#259), plus:

- scheduler.node_admin_api_key for sync-node-snapshots pulls against
  node admin APIs; it is a credential, so SCHEDULER_NODE_ADMIN_API_KEY
  env (mounted from the agentenv-auth Secret) overrides it and is the
  recommended path on Kubernetes.
- leader_election.snapshot_pull_concurrency bounding the
  post-acquisition snapshot pull fan-out; zero means the default of 4.

Scope: services/shared/config only. Disabled by default; when disabled
the scheduler behaves exactly as a single-writer process. Startup
validation (fail fast): mutually exclusive with --query-only; client-go
timing rule lease_duration > renew_deadline > 2*retry_period. Redis
stays optional (without it, failover loses pre-failover sandbox
routing; documented degraded mode).

Tests: validation cases for no-redis pass, query-only conflict,
invalid timing triples, disabled transparency, and the
snapshot_pull_concurrency default.

Relates to kvcache-ai#259, helps kvcache-ai#191.

Co-Authored-By: Claude Code <noreply@anthropic.com>
…indow

Add the leadership domain for scheduler HA (kvcache-ai#259), behind a single
exported contract:

- Leadership interface + NewLeadership(cfg) factory are the only
  exported surface. Two implementations, both package-private:
  leadershipManager (election enabled) and a null object (election
  disabled), so callers wire unconditionally and 'feature off' is an
  implementation difference, not a wiring difference.
- leadershipSnapshot: immutable leadership state (leader?, acquired
  since, recovery TTL). leadershipManager owns transitions via
  atomic.Pointer whole-instance swaps; readers load the current
  snapshot per call — stable facade, swappable payload.
- Read/write gate interceptor: health checks never gated; leader
  serves everything; standbys with a shared (Redis) binding store
  serve LookupNode/GetNode and reject the rest with
  Unavailable("not the leader"); without a shared store reads retry
  onto the leader (degraded mode). The RPC table classifies every
  method explicitly and panics on unclassified ones, so proto upgrades
  fail loudly instead of silently gating.
- Leader readiness health service scheduler.v1.Scheduler/leader:
  SERVING only while holding the lease, so Service endpoints contain
  only the leader; liveness is unchanged.
- client-go Lease elector with fencing: OnStoppedLeading flips
  readiness, then stops the process via the composition root's
  shutdown hook; a partitioned ex-leader stops serving within
  renew_deadline.
- Recovery window in the Service: while open, Schedule only considers
  nodes observed at or after acquisition (fresh-observations-only, via
  the new NodeRegistry.LastReportAt) and returns Unavailable with zero
  fresh observations — the kvcache-ai#191 fix, no snapshot means "unknown", not
  "unlimited"; LookupNode/GetNode return Unavailable (not NotFound)
  for missing entries so "not yet rebuilt" is not misreported as
  "does not exist". With election disabled all of this short-circuits
  to today's behaviour.
- Heartbeat ingest is extracted into a shared path so pulled snapshots
  take the exact same code as heartbeats.

Scope: services/scheduler/internal (leadership.go, service.go,
node_registry.go) with unit and fake-clientset election tests
(exactly-one-leader, failover after cancel, fencing on renew failure,
callback ordering, recovery-window boundaries, gate matrix).

Relates to kvcache-ai#259, kvcache-ai#191.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Add sync-node-snapshots (kvcache-ai#259): instead of waiting out a reporter
backoff after winning the lease, the new leader refreshes observations
and bindings in one RTT per node.

- nodeSnapshotRefresher contract with a single production
  implementation, ConcurrentNodeSnapshotRefresher: bounded fan-out
  (leader_election.snapshot_pull_concurrency) pulling every registry
  node exactly once; failed pulls are logged and leave the node
  unobserved, where the recovery window's fresh-observations-only rule
  keeps it out of scheduling until its next report.
- nodeReport is the neutral ingest currency: both the Heartbeat RPC
  handler and the pull path convert into it, and it feeds the shared
  ingest path (registry observation update + ReconcileNode binding
  refresh), so pulled and heartbeated state are byte-identical in
  effect.
- NodeSnapshotFetcher seam with a production implementation talking to
  each node's admin HTTP API (GET /nodes + GET /sandboxes, x-api-key
  from SCHEDULER_NODE_ADMIN_API_KEY); the HTTP client is an internal
  default with a bounded timeout, consistent with the gateway's
  default-client pattern.

Scope: services/scheduler/internal (node_snapshot.go) with fetcher
mapping, auth-failure, ingest-compatibility, and pull-failure tests.

Relates to kvcache-ai#259.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Wire the leadership facade into the scheduler process (kvcache-ai#259):

- main builds Leadership via the factory and wires it
  unconditionally: gate interceptor chained after metrics, Service
  option attached, leader health registered, BindRuntime + Run driven
  in a goroutine. A leadership loss cancels the process root context,
  driving the same graceful shutdown as SIGTERM so the ex-leader
  exits and restarts as a standby (fencing).
- Binding TTL floor under election: effective TTL =
  max(binding_ttl, lease_duration + 90s), so bindings outlive the
  worst-case failover budget (lease detection + endpoints propagation
  + one reporter reconnect backoff) and live sandboxes are not
  misreported NotFound mid-failover. Authoritative cleanup stays with
  ReconcileNode; the TTL only guards crashed nodes, so raising it is
  harmless. Logged at startup when raised.

Scope: services/scheduler/cmd, plus go.mod/go.sum entries for
client-go's testing/fixture packages used by the fake-clientset
election tests.

Relates to kvcache-ai#259.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Add build-tagged integration specs (kvcache-ai#259) for the scenarios that need
a real cluster and stay in CI (the repo has no Kind harness yet, and
the dev machine is macOS):

- C: exactly one leader; failover within the lease budget with
  seconds-level pull-driven rebuild; fencing with no dual-primary
  window (renewal-failure partition is covered offline by the
  fake-clientset tests).
- E: the kvcache-ai#191 resource-limit scenario passes with 3 replicas;
  observation coverage equals all nodes.
- F: routing continuity across failover via standbys + Redis; degraded
  no-Redis reads retrying onto the leader; seconds-level recovery;
  binding TTL floor follow-up; rollback to single replica.

Scope: services/scheduler/kind only.

Relates to kvcache-ai#259, kvcache-ai#191.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Kubernetes manifests and docs for scheduler HA (kvcache-ai#259):

- base role.yaml gains coordination.k8s.io/leases get/create/update;
  inert when election is off.
- New deploy/k8s/overlays/ha: 3 leader-elected scheduler replicas,
  readiness probe targeting the leader health service
  (-service=scheduler.v1.Scheduler/leader) so Service endpoints hold
  only the leader, leader_election enabled in the scheduler config,
  and the SCHEDULER_NODE_ADMIN_API_KEY env mounted from the shared
  agentenv-auth Secret for sync-node-snapshots pulls. The overlay
  README covers topology, traffic rules, prerequisites, and the
  binary-first upgrade ordering with rollback. Base manifests are
  untouched; existing overlays are unaffected.
- services/README.md gains a "Leader election (HA)" section: switch
  semantics, probe split, traffic rules, degraded mode without Redis,
  the binding TTL floor, and credential passing via env/Secret.

Scope: deploy/k8s and services/README.md. Verified with
kubectl kustomize (replicas, probe flag, leader_election config,
leases RBAC, secret env all render).

Relates to kvcache-ai#259.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Copilot AI balanced review requested due to automatic review settings October 5, 2026 05:25

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Failover can delete valid bindings, expose unobserved nodes to scheduling, and permit overlapping leaders while standby reads remain unreachable.

Review effort: Balanced
Findings: 5 High severity · 1 Medium severity

Open (6)
What changed in this PR

Adds optional Kubernetes Lease-based scheduler HA with leader-gated operations, failover recovery, and deployment support.

Changes:

  • Introduces leader election, readiness, fencing, and RPC gating.
  • Rebuilds observations and bindings after failover.
  • Adds HA manifests, configuration, documentation, and tests.
File Description
services/​shared/​config/​scheduler_leader_election_test.go Tests election configuration.
services/​shared/​config/​config.go Adds HA configuration and validation.
services/​scheduler/​kind/​kind_test.go Specifies Kind integration scenarios.
services/​scheduler/​internal/​service.go Adds recovery filtering and shared ingestion.
services/​scheduler/​internal/​recovery_window_test.go Tests recovery behavior.
services/​scheduler/​internal/​node_snapshot.go Implements node snapshot rebuilding.
services/​scheduler/​internal/​node_snapshot_test.go Tests admin snapshot fetching.
services/​scheduler/​internal/​node_registry.go Tracks report timestamps.
services/​scheduler/​internal/​leadership.go Implements election, gating, and readiness.
services/​scheduler/​internal/​leadership_test.go Tests leadership behavior.
services/​scheduler/​cmd/​main.go Wires HA lifecycle and binding TTL.
services/​scheduler/​cmd/​binding_ttl_test.go Tests TTL flooring.
services/​README.md Documents scheduler HA.
services/​go.mod Updates Go dependencies.
services/​go.sum Updates dependency checksums.
deploy/​k8s/​overlays/​ha/​README.md Documents HA deployment.
deploy/​k8s/​overlays/​ha/​kustomization.yaml Configures three leader-elected replicas.
deploy/​k8s/​overlays/​ha/​config/​scheduler.json Enables election in the overlay.
deploy/​k8s/​base/​role.yaml Grants Lease permissions.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment on lines +34 to +38
path: /spec/template/spec/containers/0/readinessProbe/exec/command
value:
- /grpc_health_probe
- -addr=127.0.0.1:9090
- -service=scheduler.v1.Scheduler/leader

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good catch — this is correct: with leader-only readiness, standbys are never in any Service's endpoints, so the standby-served read path is unreachable in the HA overlay. Worth explaining how we got here, because the two inputs genuinely conflict in Kubernetes:

The issue's desired behavior asks for leader-only endpoints: "Liveness remains healthy on standbys; readiness exposes only the leader" (#259).
The follow-up discussion asked for standby reads: "every standby replica acts as a query-only node and can serve LookupNode from the shared bindings" (#259 (comment)).
Pod readiness is global in K8S — a pod is Ready in all Services or in none — so "leader-only endpoints" and "standbys receive read traffic" cannot both hold. We implemented both requirements faithfully, which leaves the read path formally present but unreachable. We kept leader-only endpoints as the default because it is the correctness-first choice: with the gateway's pick_first connection pinning, it guarantees writes (Schedule/RecordAssignment/Heartbeat) never land on a non-leader, so there is no dual-writer window even before fencing kicks in.

Two ways to resolve it — maintainer's call:

(i) Make standby reads real: standbys stay Ready (readiness = overall serving), gateway uses round_robin + retry-on-Unavailable, so writes retry onto the leader and reads genuinely reach standbys. Cost: gateway LB/retry config, and ~2/3 of write attempts take one retry.
(ii) Drop standby reads: standbys stay pure idle, reads ride the leader only, accepting a read gap of roughly (lease detection + readiness flip) during failover. This is already today's behavior in the no-Redis degraded mode.
Interim this PR keeps leader-only endpoints. Happy to implement either direction — (i) if standby reads are worth the retry cost, (ii) if simplicity wins.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Valid finding, and a genuine design fork rather than a bug: it stems from two stated requirements that cannot coexist in K8S (leader-only readiness vs standbys serving reads). Full analysis with both options (make standbys Ready + gateway round_robin/retry, or drop standby reads) is in the PR discussion: #341 (comment) — leaving this open pending maintainer decision.

Comment thread services/scheduler/internal/leadership.go Outdated
Comment thread services/scheduler/internal/leadership.go Outdated
Comment thread services/scheduler/internal/node_snapshot.go Outdated
Comment thread services/scheduler/internal/service.go
Comment thread services/shared/config/config.go
@github-actions

github-actions Bot commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

🔍 OpenCodeReview found 8 issue(s) in this PR.

  • ✅ Successfully posted inline: 8 comment(s)

Comment thread services/scheduler/internal/leadership.go Outdated
Comment thread services/scheduler/internal/leadership.go Outdated
Comment thread services/scheduler/internal/node_snapshot.go
Comment thread services/scheduler/internal/node_snapshot.go Outdated
Comment thread services/scheduler/internal/node_snapshot.go
Four correctness fixes from the PR kvcache-ai#341 review:

- Gate GetNode on standbys (F2): it reads replica-local observations,
  which are always empty on standbys (heartbeats and pulls only reach
  the leader), so serving it there returned a terminal NotFound for
  real nodes. LookupNode stays standby-readable; it reads the shared
  binding store.
- Hold the lease until process exit (F3): ReleaseOnCancel=true let a
  standby acquire while the ex-leader was still draining in-flight
  RPCs through GracefulStop, opening the dual-primary window the
  fencing design exists to prevent. Failover after a graceful stop is
  now bounded by the remaining lease duration instead.
- Never reconcile bindings from a snapshot pull (F4): the roster is
  now three-state — non-nil (heartbeat) is authoritative and
  reconciles, nil (pull) leaves bindings untouched. A pull cannot see
  paused sandboxes or in-flight template builds (list_sandbox_ids
  covers both; no admin list endpoint does), so reconciling from it
  deleted live bindings exactly on failover. The admin fetcher drops
  the GET /sandboxes call; binding refresh stays with heartbeats, and
  the binding TTL floor covers the gap. Includes a regression test
  that a pre-existing binding survives a roster-less pull ingest.
- Reject explicitly invalid election config values (F5): defaults now
  apply only to omitted fields — an explicit non-positive duration
  (e.g. "retry_period": "-1s") or snapshot_pull_concurrency fails at
  config parse instead of being silently replaced by the default,
  keeping the advertised fail-fast behavior.

Scope: services/scheduler/internal, services/shared/config. The F1
finding (standby reads unreachable behind leader-only readiness) is a
design fork rooted in two conflicting requirements and is answered in
the PR discussion rather than changed here.

Relates to kvcache-ai#259.

Co-Authored-By: Claude Code <noreply@anthropic.com>
@NickNYU

NickNYU commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

Review findings addressed in d089509:

  1. GetNode gated on standbys — it reads replica-local observations (always empty there), so serving it returned terminal NotFound for real nodes. LookupNode stays standby-readable (shared binding store).
  2. Lease held until process exit — ReleaseOnCancel no longer lets a standby acquire while the ex-leader drains in-flight RPCs; the dual-primary window is closed.
  3. Pulls never reconcile bindings — a pull can't see paused sandboxes or in-flight template builds, so reconciling deleted live bindings on failover. The roster is now three-state: non-nil (heartbeat) reconciles, nil (pull) preserves. Regression test included.
  4. Explicit invalid config values rejected — defaults apply only to omitted fields; "-1s" or a non-positive snapshot_pull_concurrency now fails at parse.

The F1 finding (standby reads unreachable behind leader-only readiness) is a design fork between two stated requirements — replied separately above with the two options.

…beats

Add a freshness guard to the post-acquisition refresh (kvcache-ai#341 review,
OpenCodeReview): a nodeReport from a pull now carries the capture time
(fetchedAt), and ingest skips a pulled report when the node has already
reported something newer — e.g. a heartbeat that landed while the HTTP
request was in flight. Heartbeat-shaped reports carry a zero fetchedAt
and are never skipped. Without the guard, a slow pull could regress
observations by up to one request latency during the rebuild window.

The remaining pull/heartbeat parity gap (admin /nodes lacks sandboxIDs,
p2pEndpoint, and cpuConfigJSON) is bounded and self-healing via the next
heartbeat; it is documented in the PR discussion as follow-up rather
than changed here.

Scope: services/scheduler/internal. Includes guard-semantics tests.

Relates to kvcache-ai#259.

Co-Authored-By: Claude Code <noreply@anthropic.com>
@NickNYU

NickNYU commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

OpenCodeReview round addressed:

  1. GetNode on standbys and ReleaseOnCancel and incomplete /sandboxes roster — same root issues as the Copilot round, fixed in d089509 (gated, lease held until exit, pulls never reconcile bindings).
  2. Slow pull overwriting a newer heartbeat — fixed in 2b15c2c with a freshness guard: pulled reports carry their capture time, and ingest skips a pull when the node already reported something newer. Heartbeat-shaped reports are never skipped.
  3. cpuConfigJSON missing from admin MachineInfo — accurate. This is one of three heartbeat-parity gaps in admin /nodes (with sandboxIDs and p2pEndpoint). All three are bounded and self-healing: heartbeats restore the full fields within one report cycle (≤60s worst case), and nothing fails meanwhile — CPU config stays pending on nodes, P2P lookups fall back to source backends. Folding all three into an admin/heartbeat parity follow-up (expose the full snapshot on admin /nodes) rather than growing this PR.

…overy TTL

Address the remaining Copilot finding (kvcache-ai#341 review, service.go): the
recovery window ended on a timer (one report_ttl), after which a node
whose pull failed and whose reporter is still backing off (up to 60s)
was readmitted with a nil snapshot — FilterByResourceLimit's fail-open
then kept it, reintroducing the kvcache-ai#191 capacity bug during failover.

The acquisition snapshot now carries the node set known at that moment.
Scheduling excludes a known node until it produces a report at or after
the acquisition — no timer readmission — while nodes that joined later
are outside the recovery-pending set and keep the steady-state fail-open
behavior. The recovery TTL still bounds only the lookup semantics
(Unavailable vs NotFound), unchanged.

leadershipView gains RecoveryPending; the manager delegates; the null
object returns false. Includes a regression test: past-TTL exclusion of
an unreported known node, fail-open admission of a later-joined node,
and readmission after the late report.

Relates to kvcache-ai#259, kvcache-ai#191.

Co-Authored-By: Claude Code <noreply@anthropic.com>
@NickNYU

NickNYU commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

The remaining Copilot finding is addressed in d4c18af: the recovery window no longer ends scheduling exclusion on a timer. The acquisition snapshot now carries the node set known at that moment, and a known node with no post-acquisition report stays excluded until it reports — so a node whose pull failed and whose reporter is still backing off is never readmitted with a nil snapshot. Nodes that join after the acquisition keep the steady-state fail-open behavior (they are outside the recovery-pending set). The recovery TTL now bounds only the lookup semantics (Unavailable vs NotFound). Regression test included: past-TTL exclusion, later-join admission, readmission after the late report.

Comment thread services/scheduler/internal/node_registry.go
Comment thread services/scheduler/internal/node_snapshot.go
Comment thread services/scheduler/internal/node_snapshot.go Outdated
Comment thread services/scheduler/internal/node_snapshot.go
Comment on lines +174 to +175
CPUArchitecture string `json:"cpuArchitecture"`
CPUConfigJSON string `json:"cpuConfigJSON"`

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

bug · high
The node admin API does not expose cpuConfigJSON in MachineInfo (its OpenAPI schema and From<MachineInfo> conversion contain only family/model/name/architecture), so this field is always empty in production. On a newly elected scheduler, pulls therefore seed every observed node without its CPU config; agents also clear their pending config after the first successful heartbeat and do not resend it on leader failover. allConfigsReadyLocked can then never rebuild the cluster CPU intersection, including after a new node joins. The recovery endpoint must expose this value, or CPU intersection state needs a separate recoverable source/protocol.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Accurate, and this is the round-2 raise of the same point — the answer is the documented disposition: cpuConfigJSON is one of three heartbeat-parity gaps in admin /nodes (with sandboxIDs and p2pEndpoint), all bounded and self-healing via the next heartbeat (nodes keep their pending config and retry; P2P falls back to source backends). These fold into an admin/heartbeat parity follow-up (expose the full snapshot on admin /nodes) rather than growing this PR. Leaving this thread open as the follow-up marker; the round-1 duplicate is resolved.

Comment thread services/scheduler/internal/node_snapshot.go Outdated
Comment thread services/scheduler/internal/node_snapshot.go Outdated
Comment thread services/scheduler/internal/node_snapshot.go Outdated
…ulls

Seven correctness and hardening fixes from the follow-up review (kvcache-ai#341):

- Freshness is now enforced atomically inside the registry: the new
  HeartbeatUnlessStale checks and writes under the same lock, closing the
  TOCTOU window between a pull's freshness check and its ingest. The
  manager's pre-check is removed; the registry is the single point of
  truth. Includes an atomic-path test.
- fetchedAt is stamped before the /nodes request is sent, not after the
  response is decoded, so a heartbeat landing mid-flight correctly counts
  as newer than the pulled data.
- Pulled reports map the admin status field into NodeSnapshot.Status
  instead of leaving UNSPECIFIED (which the registry derives to
  CONNECTING), and the fetcher no longer accepts a single-entry response
  whose ID differs from the requested node — a stale or misconfigured
  endpoint cannot overwrite another node's observation.
- Admin responses are decoded with a 1 MiB body limit, so a faulty or
  compromised node cannot make the leader allocate from an arbitrarily
  large document during the acquisition fan-out.
- The registry preserves a previously known P2P endpoint when a report
  arrives without one (the admin API does not serve it), instead of
  erasing it on every pull.
- Freshness comparisons run at millisecond precision in both
  filterFreshObservations and RecoveryPending: LastReportAt is
  reconstructed from a millisecond timestamp, so a report in the same
  millisecond as the acquisition counts as fresh instead of stale.

The remaining heartbeat-parity gaps in admin /nodes (sandboxIDs,
p2pEndpoint, cpuConfigJSON) stay on the documented follow-up track; the
cpuConfigJSON finding was already dispositioned in the previous round.

Scope: services/scheduler/internal.

Relates to kvcache-ai#259.

Co-Authored-By: Claude Code <noreply@anthropic.com>
@NickNYU

NickNYU commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

Second OpenCodeReview round addressed in c9ea0fc:

  1. TOCTOU between freshness check and ingest — freshness is now enforced atomically inside the registry (HeartbeatUnlessStale checks and writes under the same lock); the manager-side pre-check is gone, registry is the single point of truth.
  2. fetchedAt stamped after decode — now stamped before the request is sent, so a heartbeat landing mid-flight correctly counts as newer.
  3. status not decoded — pulled reports map the admin status field instead of degrading READY to CONNECTING.
  4. ID-mismatch fallback — removed; a response for a different node ID is rejected, so a stale/misconfigured endpoint cannot overwrite another node's observation.
  5. Unbounded decode — admin responses are capped at 1 MiB.
  6. P2P endpoint erased by pulls — the registry now preserves a previously known endpoint when a report arrives without one.
  7. Millisecond precision — freshness comparisons in both filterFreshObservations and RecoveryPending now compare at ms precision, so a report in the same millisecond as acquisition counts as fresh.

cpuConfigJSON remains on the documented admin/heartbeat parity follow-up (with sandboxIDs and p2pEndpoint), as dispositioned in the previous round.

NickNYU and others added 4 commits October 5, 2026 21:57
Address the review concern that the previous commit reshaped the shared
AtomicNodeRegistry.Heartbeat for a problem that only exists on the new
snapshot-pull path. Heartbeat is now byte-identical to its pre-HA form;
everything pull-specific lives in the pull path:

- The atomic HeartbeatUnlessStale and the in-registry P2P merge are
  removed. Freshness is guarded in Service.ingestNodeReport: a pulled
  report is skipped when LastReportAt is newer than its fetchedAt. The
  residual check-then-act window is microseconds and self-heals on the
  next heartbeat (5s).
- P2P endpoint preservation moves to the pull path too: a read-only
  P2PEndpointFor getter (purely additive to the registry interface)
  lets ingest fill the endpoint into the pulled request before calling
  the untouched Heartbeat.

Scope: services/scheduler/internal (node_registry.go, service.go,
node_snapshot_test.go).

Relates to kvcache-ai#259.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Replace the (time.Time, bool) pair with *time.Time in the
RecoveryPending signature: nil means "never reported", which is
self-explanatory at the call site instead of a free-floating ok flag.
Applies uniformly to the leadershipView seam, the snapshot
implementation, the manager delegate, and the null object.

Scope: services/scheduler/internal (leadership.go, service.go).

Relates to kvcache-ai#259.

Co-Authored-By: Claude Code <noreply@anthropic.com>
… retry

Resolve the F1 review finding: with leader-only readiness endpoints,
standbys never received read traffic, so the standby Redis-read path
was unreachable. This implements the "standbys stay Ready" option:

- HA overlay readiness probes overall health again, so standbys are in
  Service endpoints and gateway reads reach them. The leader health
  service (scheduler.v1.Scheduler/leader) stays registered for
  observability but no longer gates endpoints.
- The gateway's scheduler connection switches from the default
  pick_first (which pinned one pod) to round_robin plus a retry policy
  on Unavailable, so writes that land on a leader-gated standby are
  retried onto the leader (at most two retries with 3 replicas).

Trade-offs, called out for review: endpoints now include standbys
(deviating from the issue's "readiness exposes only the leader");
~2/3 of write RPCs take one retry under 3 replicas; the retry policy
also retries semantic Unavailable errors (harmless, same result).
If the "drop standby reads" option is preferred, this commit reverts
cleanly on its own.

Scope: services/gateway/cmd, deploy/k8s/overlays/ha.

Relates to kvcache-ai#259.

Co-Authored-By: Claude Code <noreply@anthropic.com>
@NickNYU

NickNYU commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

Retraction on F1: the option-(i) implementation from earlier today has been reverted (2f29b59) — it was premature to pick a direction unilaterally. This thread goes back to its intended state: open, pending maintainer discussion. The two options on the table (standbys Ready + gateway round_robin/retry, or drop standby reads) are laid out in the analysis above.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

2 participants