diff --git a/docs/specs/012-drift-roll-budget/spec.md b/docs/specs/012-drift-roll-budget/spec.md index 69a7866c..d7f20d2e 100644 --- a/docs/specs/012-drift-roll-budget/spec.md +++ b/docs/specs/012-drift-roll-budget/spec.md @@ -29,7 +29,7 @@ An operator bumps the cell sidecar image with the budget at 25%. In each namespa **Acceptance Scenarios**: 1. **Given** a namespace with N Running nodes and budget B, **When** a drift lands on all of them, **Then** at most `max(1, N*B/100)` nodes update at once, and the rest report `NodeUpdateInProgress=False/UpdateDeferred` naming the slot holders. -2. **Given** a waiting node, **When** a node holding a slot finishes, **Then** the waiting node starts within one status poll (30 seconds). +2. **Given** a waiting node, **When** a node holding a slot finishes, **Then** the waiting node re-checks its slot at once, not on its next status poll (30 seconds). ### User Story 2 - The budget changes nothing until a cell sets it (Priority: P1) @@ -45,6 +45,7 @@ An operator bumps the cell sidecar image with the budget at 25%. In each namespa 2. WHEN a drifted node holds no slot, THE controller SHALL build no plan for it and SHALL set `NodeUpdateInProgress=False` with reason `UpdateDeferred` and a message that names the slot holders. 3. THE controller SHALL list the namespace's SeiNodes with an uncached read. Because the node controller reconciles one node at a time, each slot decision then sees every earlier node's persisted status, and the slots never overfill. 4. THE budget SHALL NOT gate a data reset, a hold change, a config update, an init plan, or a resize. +5. WHEN a node leaves `NodeUpdateInProgress=True`, changes `spec.paused`, changes phase, or is deleted, THE controller SHALL enqueue every node in its namespace that reports `UpdateDeferred`, so the next node in slot order does not wait for its status poll. ### Requirement 2: A ConfigMap-configured node keeps its running images while it waits @@ -65,11 +66,13 @@ An operator bumps the cell sidecar image with the budget at 25%. In each namespa - **SC-001**: Unit tests cover the slot count and order, the exclusions, the plan gate, and the template gate; a reconciler test shows a waiting nodeConfig node keeps its StatefulSet template. *Verifier:* `make test`. - **SC-002**: A harbor sidecar bump with the budget at 25% never exceeds the slot count in any namespace. *Verifier:* judgement — the roll log in PLT-1399. +- **SC-003**: Unit tests cover the slot-release predicate and the peer mapping (Req 1.5). *Verifier:* `make test`. ## Known limits - While a `nodeConfig` node runs another plan, or has a reset or hold change pending, the controller applies its drifted template. The pod then rolls with that plan, outside the budget, and `observe-image` records the new images, so no second roll follows. - A drift update stuck in progress keeps its slot, so a namespace's further drift waits until an operator clears it. The `UpdateDeferred` message names the holder. +- The node controller reconciles one node at a time, so a woken node still waits behind the nodes queued before it. In the 2026-10-08 rollout, before Req 1.5, the median wait from a freed slot to the next start was 1 to 28 seconds per namespace, and the longest was 52 seconds in the 51-node prod cell. - Cells still roll through separate pull requests; the budget is per namespace, not per chain across cells. ## Out of scope diff --git a/internal/controller/node/controller.go b/internal/controller/node/controller.go index 545e2dea..9bb05526 100644 --- a/internal/controller/node/controller.go +++ b/internal/controller/node/controller.go @@ -514,6 +514,7 @@ func (r *SeiNodeReconciler) SetupWithManager(mgr ctrl.Manager) error { Owns(&corev1.PersistentVolumeClaim{}). Watches(&seiv1alpha1.SeiNodeTaskWorkflow{}, &workflowTargetHandler{}). Watches(&corev1.Pod{}, handler.EnqueueRequestsFromMapFunc(podToSeiNode), builder.WithPredicates(podReadyChanged)). + Watches(&seiv1alpha1.SeiNode{}, handler.EnqueueRequestsFromMapFunc(r.deferredPeers), builder.WithPredicates(slotReleased)). Named(seiNodeControllerName). Complete(r) } diff --git a/internal/controller/node/driftwake.go b/internal/controller/node/driftwake.go new file mode 100644 index 00000000..108384e6 --- /dev/null +++ b/internal/controller/node/driftwake.go @@ -0,0 +1,66 @@ +package node + +import ( + "context" + + "k8s.io/apimachinery/pkg/api/meta" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/event" + "sigs.k8s.io/controller-runtime/pkg/log" + "sigs.k8s.io/controller-runtime/pkg/predicate" + "sigs.k8s.io/controller-runtime/pkg/reconcile" + + seiv1alpha1 "github.com/sei-protocol/sei-k8s-controller/api/v1alpha1" + "github.com/sei-protocol/sei-k8s-controller/internal/planner" +) + +// A drifted node that waits for a roll slot re-checks on its status poll. The +// watch below wakes it as soon as a slot frees instead, so the next node in +// slot order starts without waiting out the poll interval. + +// slotReleased passes a SeiNode event that can free a drift-roll slot or add +// one: the node leaves NodeUpdateInProgress=True, its pause flips, its phase +// changes (the slot count counts Running, unpaused nodes), or it is deleted. +var slotReleased = predicate.Funcs{ + CreateFunc: func(event.CreateEvent) bool { return false }, + UpdateFunc: func(e event.UpdateEvent) bool { + oldNode, ok := e.ObjectOld.(*seiv1alpha1.SeiNode) + if !ok { + return false + } + newNode, ok := e.ObjectNew.(*seiv1alpha1.SeiNode) + if !ok { + return false + } + return (nodeUpdating(oldNode) && !nodeUpdating(newNode)) || + oldNode.Spec.Paused != newNode.Spec.Paused || + oldNode.Status.Phase != newNode.Status.Phase + }, + DeleteFunc: func(event.DeleteEvent) bool { return true }, + GenericFunc: func(event.GenericEvent) bool { return false }, +} + +func nodeUpdating(node *seiv1alpha1.SeiNode) bool { + return meta.IsStatusConditionTrue(node.Status.Conditions, seiv1alpha1.ConditionNodeUpdateInProgress) +} + +// deferredPeers maps a node that freed its slot to every node in its namespace +// that reports UpdateDeferred. Each woken node runs the slot decision again +// with an uncached read, so a stale cache here costs at most one extra +// reconcile. +func (r *SeiNodeReconciler) deferredPeers(ctx context.Context, obj client.Object) []reconcile.Request { + var nodes seiv1alpha1.SeiNodeList + if err := r.List(ctx, &nodes, client.InNamespace(obj.GetNamespace())); err != nil { + log.FromContext(ctx).Error(err, "listing nodes waiting for a drift-roll slot", "namespace", obj.GetNamespace()) + return nil + } + var reqs []reconcile.Request + for i := range nodes.Items { + n := &nodes.Items[i] + c := meta.FindStatusCondition(n.Status.Conditions, seiv1alpha1.ConditionNodeUpdateInProgress) + if c != nil && c.Reason == planner.ReasonUpdateDeferred { + reqs = append(reqs, reconcile.Request{NamespacedName: client.ObjectKeyFromObject(n)}) + } + } + return reqs +} diff --git a/internal/controller/node/driftwake_test.go b/internal/controller/node/driftwake_test.go new file mode 100644 index 00000000..8e5ec47d --- /dev/null +++ b/internal/controller/node/driftwake_test.go @@ -0,0 +1,85 @@ +package node + +import ( + "context" + "testing" + + . "github.com/onsi/gomega" + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" + "sigs.k8s.io/controller-runtime/pkg/event" + "sigs.k8s.io/controller-runtime/pkg/reconcile" + + seiv1alpha1 "github.com/sei-protocol/sei-k8s-controller/api/v1alpha1" + "github.com/sei-protocol/sei-k8s-controller/internal/planner" +) + +func withUpdateCondition(node *seiv1alpha1.SeiNode, status metav1.ConditionStatus, reason string) *seiv1alpha1.SeiNode { + meta.SetStatusCondition(&node.Status.Conditions, metav1.Condition{ + Type: seiv1alpha1.ConditionNodeUpdateInProgress, Status: status, Reason: reason, + }) + return node +} + +// Spec 012 Req 1.5: only an event that can free or add a slot wakes the +// waiting nodes. +func TestSlotReleased(t *testing.T) { + running := func(status metav1.ConditionStatus, reason string) func() *seiv1alpha1.SeiNode { + return func() *seiv1alpha1.SeiNode { + n := withUpdateCondition(resizeNode("n", false, ""), status, reason) + n.Status.Phase = seiv1alpha1.PhaseRunning + return n + } + } + updating := running(metav1.ConditionTrue, "UpdateStarted") + done := running(metav1.ConditionFalse, "UpdateComplete") + failed := running(metav1.ConditionFalse, "UpdateFailed") + deferred := running(metav1.ConditionFalse, planner.ReasonUpdateDeferred) + cases := []struct { + name string + old, new *seiv1alpha1.SeiNode + want bool + }{ + {"update completes", updating(), done(), true}, + {"update fails", updating(), failed(), true}, + {"update starts", deferred(), updating(), false}, + {"still updating", updating(), updating(), false}, + {"still waiting", deferred(), deferred(), false}, + {"no condition", resizeNode("n", false, ""), resizeNode("n", false, ""), false}, + {"holder paused", updating(), func() *seiv1alpha1.SeiNode { n := updating(); n.Spec.Paused = true; return n }(), true}, + {"node leaves Running", updating(), func() *seiv1alpha1.SeiNode { n := updating(); n.Status.Phase = seiv1alpha1.PhaseFailed; return n }(), true}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + g := NewWithT(t) + g.Expect(slotReleased.Update(event.UpdateEvent{ObjectOld: tc.old, ObjectNew: tc.new})).To(Equal(tc.want)) + }) + } + + g := NewWithT(t) + g.Expect(slotReleased.Delete(event.DeleteEvent{Object: updating()})).To(BeTrue()) + g.Expect(slotReleased.Create(event.CreateEvent{Object: deferred()})).To(BeFalse()) +} + +// Spec 012 Req 1.5: a freed slot wakes every waiting node in the namespace and +// no other node. +func TestDeferredPeers(t *testing.T) { + g := NewWithT(t) + + waitA := withUpdateCondition(resizeNode("wait-a", false, ""), metav1.ConditionFalse, planner.ReasonUpdateDeferred) + waitB := withUpdateCondition(resizeNode("wait-b", true, ""), metav1.ConditionFalse, planner.ReasonUpdateDeferred) + rolling := withUpdateCondition(resizeNode("rolling", false, ""), metav1.ConditionTrue, "UpdateStarted") + settled := withUpdateCondition(resizeNode("settled", false, ""), metav1.ConditionFalse, "UpdateComplete") + plain := resizeNode("plain", false, "") + elsewhere := withUpdateCondition(resizeNode("elsewhere", false, ""), metav1.ConditionFalse, planner.ReasonUpdateDeferred) + elsewhere.Namespace = "other" + + r, _ := newNodeReconciler(t, waitA, waitB, rolling, settled, plain, elsewhere) + + got := r.deferredPeers(context.Background(), settled) + g.Expect(got).To(ConsistOf( + reconcile.Request{NamespacedName: types.NamespacedName{Namespace: testNamespace, Name: "wait-a"}}, + reconcile.Request{NamespacedName: types.NamespacedName{Namespace: testNamespace, Name: "wait-b"}}, + )) +} diff --git a/internal/planner/drift_budget.go b/internal/planner/drift_budget.go index 9b42e65c..c1bc2a58 100644 --- a/internal/planner/drift_budget.go +++ b/internal/planner/drift_budget.go @@ -26,9 +26,12 @@ import ( // for one follow in name order. The list is an uncached read and the node // controller reconciles one node at a time, so each decision sees every earlier // node's persisted status and the slots never overfill. A node without a slot -// keeps its current pod and reports UpdateDeferred. +// keeps its current pod and reports UpdateDeferred. When a slot frees, the +// node controller wakes the waiting nodes at once (deferredPeers). -const reasonUpdateDeferred = "UpdateDeferred" +// ReasonUpdateDeferred is the NodeUpdateInProgress reason of a drifted node +// that waits for a roll slot. +const ReasonUpdateDeferred = "UpdateDeferred" // DriftSlot reports whether node may start a pod-template drift update now. // When it may not, the message names the slot holders. @@ -104,7 +107,7 @@ func (p *NodeResolver) DriftRenderNode(ctx context.Context, node *seiv1alpha1.Se if free { return node, nil } - setNodeUpdateCondition(node, metav1.ConditionFalse, reasonUpdateDeferred, msg) + setNodeUpdateCondition(node, metav1.ConditionFalse, ReasonUpdateDeferred, msg) pinned, err := p.pinnedToRunning(ctx, node) if err != nil || pinned == nil { // With no running image to pin to, the template cannot hold the pod; diff --git a/internal/planner/drift_budget_test.go b/internal/planner/drift_budget_test.go index 1a7ac40f..feeae038 100644 --- a/internal/planner/drift_budget_test.go +++ b/internal/planner/drift_budget_test.go @@ -116,7 +116,7 @@ func TestResolvePlan_DriftWaitsForASlot(t *testing.T) { cond := meta.FindStatusCondition(waiting.Status.Conditions, seiv1alpha1.ConditionNodeUpdateInProgress) g.Expect(cond).NotTo(BeNil()) g.Expect(cond.Status).To(Equal(metav1.ConditionFalse)) - g.Expect(cond.Reason).To(Equal(reasonUpdateDeferred)) + g.Expect(cond.Reason).To(Equal(ReasonUpdateDeferred)) g.Expect(cond.Message).To(ContainSubstring("held by node-a, node-b")) first := nodes[0].DeepCopy() // node-a @@ -214,7 +214,7 @@ func TestDriftRenderNode(t *testing.T) { g.Expect(n.Spec.Image).To(Equal(testImageV2), "the node itself keeps its spec") cond := meta.FindStatusCondition(n.Status.Conditions, seiv1alpha1.ConditionNodeUpdateInProgress) g.Expect(cond).NotTo(BeNil()) - g.Expect(cond.Reason).To(Equal(reasonUpdateDeferred)) + g.Expect(cond.Reason).To(Equal(ReasonUpdateDeferred)) }) } } diff --git a/internal/planner/planner.go b/internal/planner/planner.go index 70535287..c9a50edd 100644 --- a/internal/planner/planner.go +++ b/internal/planner/planner.go @@ -215,13 +215,13 @@ func (p *NodeResolver) ResolvePlan(ctx context.Context, node *seiv1alpha1.SeiNod if plan != nil && p.isDriftUpdatePlan(node, plan) { free, msg, err := p.DriftSlot(ctx, node) if err != nil { - setNodeUpdateCondition(node, metav1.ConditionFalse, reasonUpdateDeferred, err.Error()) + setNodeUpdateCondition(node, metav1.ConditionFalse, ReasonUpdateDeferred, err.Error()) return err } if !free { // BuildPlan stamped UpdateStarted; the plan is not persisted, so // the condition says why the node waits instead. - setNodeUpdateCondition(node, metav1.ConditionFalse, reasonUpdateDeferred, msg) + setNodeUpdateCondition(node, metav1.ConditionFalse, ReasonUpdateDeferred, msg) return nil } }