Skip to content
57 changes: 57 additions & 0 deletions internal/controller/helper.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,11 @@ package controller
import (
"encoding/json"
"maps"
"sort"

corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/types"

readinessv1alpha1 "sigs.k8s.io/node-readiness-controller/api/v1alpha1"
)

Expand Down Expand Up @@ -133,3 +135,58 @@ func filterStatusForExistingNodes(
func labelsEqual(a, b map[string]string) bool {
return maps.Equal(a, b)
}

// nodeStatusDelta captures the per-node NodeEvaluation/NodeFailure changes produced by a single
// processAllNodesForRule sweep, keyed by node name. It excludes AppliedNodes, ObservedGeneration,
// and DryRunResults, which have a single writer and are safe to overwrite directly.
//
// A nil value in failures clears any failure recorded for that node. evaluations only holds
// entries for nodes freshly (re-)evaluated this sweep.
type nodeStatusDelta struct {
evaluations map[string]readinessv1alpha1.NodeEvaluation
failures map[string]*readinessv1alpha1.NodeFailure
}

// sortStatusByNodeName sorts rule.Status.NodeEvaluations and rule.Status.FailedNodes by NodeName.
func sortStatusByNodeName(rule *readinessv1alpha1.NodeReadinessRule) {
sort.Slice(rule.Status.NodeEvaluations, func(i, j int) bool {
return rule.Status.NodeEvaluations[i].NodeName < rule.Status.NodeEvaluations[j].NodeName
})
sort.Slice(rule.Status.FailedNodes, func(i, j int) bool {
return rule.Status.FailedNodes[i].NodeName < rule.Status.FailedNodes[j].NodeName
})
}

// applyNodeStatusDelta merges delta into rule's NodeEvaluations/FailedNodes, replacing only the
// entries for nodes present in delta and leaving every other node's entry untouched.
func applyNodeStatusDelta(rule *readinessv1alpha1.NodeReadinessRule, delta nodeStatusDelta) {
Comment thread
ajaysundark marked this conversation as resolved.
if len(delta.evaluations) > 0 {
merged := make([]readinessv1alpha1.NodeEvaluation, 0, len(rule.Status.NodeEvaluations)+len(delta.evaluations))
for _, eval := range rule.Status.NodeEvaluations {
if _, changed := delta.evaluations[eval.NodeName]; !changed {
merged = append(merged, eval)
}
}
for _, eval := range delta.evaluations {
merged = append(merged, eval)
}
rule.Status.NodeEvaluations = merged
}

if len(delta.failures) > 0 {
merged := make([]readinessv1alpha1.NodeFailure, 0, len(rule.Status.FailedNodes)+len(delta.failures))
for _, failure := range rule.Status.FailedNodes {
if _, changed := delta.failures[failure.NodeName]; !changed {
merged = append(merged, failure)
}
}
for _, failure := range delta.failures {
if failure != nil {
merged = append(merged, *failure)
}
}
rule.Status.FailedNodes = merged
}

sortStatusByNodeName(rule)
}
135 changes: 135 additions & 0 deletions internal/controller/helper_unit_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -168,3 +168,138 @@ func TestGetApplicableRulesForNode_DeepCopy(t *testing.T) {

g.Expect(cachedRule.Status.AppliedNodes).To(Equal([]string{"node-1"}))
}

func TestApplyNodeStatusDelta(t *testing.T) {
g := NewWithT(t)

t.Run("empty delta leaves status untouched", func(t *testing.T) {
rule := &readinessv1alpha1.NodeReadinessRule{
Status: readinessv1alpha1.NodeReadinessRuleStatus{
NodeEvaluations: []readinessv1alpha1.NodeEvaluation{
{NodeName: "node-1"},
},
FailedNodes: []readinessv1alpha1.NodeFailure{
{NodeName: "node-1", Reason: "Err"},
},
},
}

delta := nodeStatusDelta{
evaluations: nil,
failures: nil,
}

applyNodeStatusDelta(rule, delta)
g.Expect(rule.Status.NodeEvaluations).To(HaveLen(1))
g.Expect(rule.Status.NodeEvaluations[0].NodeName).To(Equal("node-1"))
g.Expect(rule.Status.FailedNodes).To(HaveLen(1))
g.Expect(rule.Status.FailedNodes[0].NodeName).To(Equal("node-1"))
})

t.Run("merges evaluation updates and new evaluations in sorted order", func(t *testing.T) {
rule := &readinessv1alpha1.NodeReadinessRule{
Status: readinessv1alpha1.NodeReadinessRuleStatus{
NodeEvaluations: []readinessv1alpha1.NodeEvaluation{
{NodeName: "node-1", TaintStatus: readinessv1alpha1.TaintStatusAbsent},
{NodeName: "node-3", TaintStatus: readinessv1alpha1.TaintStatusAbsent},
},
},
}

delta := nodeStatusDelta{
evaluations: map[string]readinessv1alpha1.NodeEvaluation{
"node-1": {NodeName: "node-1", TaintStatus: readinessv1alpha1.TaintStatusPresent}, // update existing
"node-2": {NodeName: "node-2", TaintStatus: readinessv1alpha1.TaintStatusPresent}, // add new
},
}

applyNodeStatusDelta(rule, delta)

g.Expect(rule.Status.NodeEvaluations).To(HaveLen(3))
g.Expect(rule.Status.NodeEvaluations[0].NodeName).To(Equal("node-1"))
g.Expect(rule.Status.NodeEvaluations[0].TaintStatus).To(Equal(readinessv1alpha1.TaintStatusPresent))
g.Expect(rule.Status.NodeEvaluations[1].NodeName).To(Equal("node-2"))
g.Expect(rule.Status.NodeEvaluations[2].NodeName).To(Equal("node-3"))
g.Expect(rule.Status.NodeEvaluations[2].TaintStatus).To(Equal(readinessv1alpha1.TaintStatusAbsent))
})

t.Run("merges failures and clears failure when nil in delta", func(t *testing.T) {
rule := &readinessv1alpha1.NodeReadinessRule{
Status: readinessv1alpha1.NodeReadinessRuleStatus{
FailedNodes: []readinessv1alpha1.NodeFailure{
{NodeName: "node-1", Reason: "OldError"},
{NodeName: "node-3", Reason: "PersistentError"},
},
},
}

delta := nodeStatusDelta{
failures: map[string]*readinessv1alpha1.NodeFailure{
"node-1": nil, // clear failure for node-1
"node-2": {NodeName: "node-2", Reason: "NewError"}, // add failure for node-2
},
}

applyNodeStatusDelta(rule, delta)

g.Expect(rule.Status.FailedNodes).To(HaveLen(2))
g.Expect(rule.Status.FailedNodes[0].NodeName).To(Equal("node-2"))
g.Expect(rule.Status.FailedNodes[0].Reason).To(Equal("NewError"))
g.Expect(rule.Status.FailedNodes[1].NodeName).To(Equal("node-3"))
g.Expect(rule.Status.FailedNodes[1].Reason).To(Equal("PersistentError"))
})

t.Run("does not add zero-value evaluation when evaluations map has no entry for node", func(t *testing.T) {
rule := &readinessv1alpha1.NodeReadinessRule{
Status: readinessv1alpha1.NodeReadinessRuleStatus{
NodeEvaluations: []readinessv1alpha1.NodeEvaluation{
{NodeName: "other-node", TaintStatus: readinessv1alpha1.TaintStatusPresent},
},
},
}

delta := nodeStatusDelta{
evaluations: make(map[string]readinessv1alpha1.NodeEvaluation),
failures: map[string]*readinessv1alpha1.NodeFailure{
"probe-node": {NodeName: "probe-node", Reason: "EvaluationError"},
},
}

applyNodeStatusDelta(rule, delta)

g.Expect(rule.Status.NodeEvaluations).To(HaveLen(1))
g.Expect(rule.Status.NodeEvaluations[0].NodeName).To(Equal("other-node"))
g.Expect(rule.Status.FailedNodes).To(HaveLen(1))
g.Expect(rule.Status.FailedNodes[0].NodeName).To(Equal("probe-node"))
})

t.Run("clears stale failure when evaluation succeeds and produces new evaluation", func(t *testing.T) {
rule := &readinessv1alpha1.NodeReadinessRule{
Status: readinessv1alpha1.NodeReadinessRuleStatus{
NodeEvaluations: []readinessv1alpha1.NodeEvaluation{
{NodeName: "other-node", TaintStatus: readinessv1alpha1.TaintStatusPresent},
},
FailedNodes: []readinessv1alpha1.NodeFailure{
{NodeName: "node-1", Reason: "EvaluationError", Message: "old failure"},
},
},
}

delta := nodeStatusDelta{
evaluations: map[string]readinessv1alpha1.NodeEvaluation{
"node-1": {NodeName: "node-1", TaintStatus: readinessv1alpha1.TaintStatusAbsent},
},
failures: map[string]*readinessv1alpha1.NodeFailure{
"node-1": nil, // clear failure for node-1 on success
},
}

applyNodeStatusDelta(rule, delta)

g.Expect(rule.Status.NodeEvaluations).To(HaveLen(2))
g.Expect(rule.Status.NodeEvaluations[0].NodeName).To(Equal("node-1"))
g.Expect(rule.Status.NodeEvaluations[0].TaintStatus).To(Equal(readinessv1alpha1.TaintStatusAbsent))
g.Expect(rule.Status.NodeEvaluations[1].NodeName).To(Equal("other-node"))
g.Expect(rule.Status.FailedNodes).To(BeEmpty())
})
}
78 changes: 33 additions & 45 deletions internal/controller/node_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -152,12 +152,13 @@ func (r *RuleReadinessController) processNodeAgainstAllRules(ctx context.Context
"rule", rule.Name,
"ruleResourceVersion", rule.ResourceVersion)

if err := r.evaluateRuleForNode(ctx, rule, node); err != nil {
log.Error(err, "Failed to evaluate rule for node",
evalErr := r.evaluateRuleForNode(ctx, rule, node)
if evalErr != nil {
log.Error(evalErr, "Failed to evaluate rule for node",
"node", node.Name, "rule", rule.Name)
// Continue with other rules even if one fails
r.recordNodeFailure(rule, node.Name, "EvaluationError", err.Error())
errs = append(errs, err)
r.recordNodeFailure(rule, node.Name, "EvaluationError", evalErr.Error())
errs = append(errs, evalErr)
metrics.Failures.WithLabelValues(rule.Name, string(metrics.FailureReasonEvaluationError)).Inc()
}

Expand All @@ -169,58 +170,44 @@ func (r *RuleReadinessController) processNodeAgainstAllRules(ctx context.Context

var successfullyPatchedRule *readinessv1alpha1.NodeReadinessRule

err := retry.RetryOnConflict(retry.DefaultRetry, func() error {
latestRule := &readinessv1alpha1.NodeReadinessRule{}
if err := r.Get(ctx, client.ObjectKey{Name: rule.Name}, latestRule); err != nil {
return err
}

patch := client.MergeFrom(latestRule.DeepCopy())

// update only this specific node evaluation status
currEval := readinessv1alpha1.NodeEvaluation{}
for _, eval := range rule.Status.NodeEvaluations {
if eval.NodeName == node.Name {
currEval = eval
break
}
}
delta := nodeStatusDelta{
evaluations: make(map[string]readinessv1alpha1.NodeEvaluation),
failures: make(map[string]*readinessv1alpha1.NodeFailure),
}

found := false
for i := range latestRule.Status.NodeEvaluations {
if latestRule.Status.NodeEvaluations[i].NodeName == node.Name {
latestRule.Status.NodeEvaluations[i] = currEval
found = true
if evalErr != nil {
// On evaluation failure, record the failure in delta and do not populate delta.evaluations
// so any stale evaluation isn't persisted.
for _, failure := range rule.Status.FailedNodes {
if failure.NodeName == node.Name {
failureCopy := failure
delta.failures[node.Name] = &failureCopy
break
}
}
if !found {
latestRule.Status.NodeEvaluations = append(
latestRule.Status.NodeEvaluations,
currEval,
)
}
} else {
// On evaluation success, clear any previously-recorded failure for this node and record the fresh evaluation.
delta.failures[node.Name] = nil

// handle status.FailedNodes for this node
var updatedFailedNodes []readinessv1alpha1.NodeFailure
for _, failure := range latestRule.Status.FailedNodes {
if failure.NodeName != node.Name {
updatedFailedNodes = append(updatedFailedNodes, failure)
}
}
for _, failure := range rule.Status.FailedNodes {
if failure.NodeName == node.Name {
if failure.NodeName != node.Name {
updatedFailedNodes = append(updatedFailedNodes, failure)
}
}
latestRule.Status.FailedNodes = updatedFailedNodes
rule.Status.FailedNodes = updatedFailedNodes

if err := r.Status().Patch(ctx, latestRule, patch); err != nil {
return err
for _, eval := range rule.Status.NodeEvaluations {
if eval.NodeName == node.Name {
delta.evaluations[node.Name] = eval
break
}
}
}

Comment thread
ajaysundark marked this conversation as resolved.
err := r.patchRuleStatusWithOptimisticLock(ctx, rule.Name, func(latestRule *readinessv1alpha1.NodeReadinessRule) {
applyNodeStatusDelta(latestRule, delta)
successfullyPatchedRule = latestRule
return nil
})

if err != nil {
Expand Down Expand Up @@ -446,7 +433,8 @@ func (r *RuleReadinessController) markBootstrapCompleted(ctx context.Context, no
deferred := false
annotationKey := bootstrapAnnotationKey(rule.GetUID())

// retry to handle conflict with concurrent node updates
// The optimistic lock here protects from a race-condition adding a taint between hasTaintBySpec
// check and mark completed annotation patch from concurrent reconciliations.
err := retry.RetryOnConflict(retry.DefaultRetry, func() error {
node := &corev1.Node{}
if err := r.Get(ctx, client.ObjectKey{Name: nodeName}, node); err != nil {
Expand All @@ -464,15 +452,15 @@ func (r *RuleReadinessController) markBootstrapCompleted(ctx context.Context, no
}
deferred = false

patch := client.MergeFromWithOptions(node.DeepCopy(), client.MergeFromWithOptimisticLock{})
stored := node.DeepCopy()

// Initialize annotations map if nil.
if node.Annotations == nil {
node.Annotations = make(map[string]string)
}

node.Annotations[annotationKey] = bootstrapAnnotationValue(rule.Name)
if err := r.Patch(ctx, node, patch); err != nil {
if err := r.Patch(ctx, node, client.MergeFromWithOptions(stored, client.MergeFromWithOptimisticLock{})); err != nil {
return err
}

Expand Down
Loading