Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion docs/specs/012-drift-roll-budget/spec.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand All @@ -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

Expand All @@ -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
Expand Down
1 change: 1 addition & 0 deletions internal/controller/node/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down
66 changes: 66 additions & 0 deletions internal/controller/node/driftwake.go
Original file line number Diff line number Diff line change
@@ -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
}
85 changes: 85 additions & 0 deletions internal/controller/node/driftwake_test.go
Original file line number Diff line number Diff line change
@@ -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"}},
))
}
9 changes: 6 additions & 3 deletions internal/planner/drift_budget.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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;
Expand Down
4 changes: 2 additions & 2 deletions internal/planner/drift_budget_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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))
})
}
}
Expand Down
4 changes: 2 additions & 2 deletions internal/planner/planner.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
}
Expand Down
Loading