diff --git a/internal/controller/helper.go b/internal/controller/helper.go index 959741a3..4799012a 100644 --- a/internal/controller/helper.go +++ b/internal/controller/helper.go @@ -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" ) @@ -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) { + 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) +} diff --git a/internal/controller/helper_unit_test.go b/internal/controller/helper_unit_test.go index 50d7215d..1d6123d8 100644 --- a/internal/controller/helper_unit_test.go +++ b/internal/controller/helper_unit_test.go @@ -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()) + }) +} diff --git a/internal/controller/node_controller.go b/internal/controller/node_controller.go index 2d1b1c12..b4209e82 100644 --- a/internal/controller/node_controller.go +++ b/internal/controller/node_controller.go @@ -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() } @@ -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 + } } + } + err := r.patchRuleStatusWithOptimisticLock(ctx, rule.Name, func(latestRule *readinessv1alpha1.NodeReadinessRule) { + applyNodeStatusDelta(latestRule, delta) successfullyPatchedRule = latestRule - return nil }) if err != nil { @@ -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 { @@ -464,7 +452,7 @@ 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 { @@ -472,7 +460,7 @@ func (r *RuleReadinessController) markBootstrapCompleted(ctx context.Context, no } 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 } diff --git a/internal/controller/nodereadinessrule_controller.go b/internal/controller/nodereadinessrule_controller.go index 3b79df67..f67c9d0d 100644 --- a/internal/controller/nodereadinessrule_controller.go +++ b/internal/controller/nodereadinessrule_controller.go @@ -25,6 +25,7 @@ import ( "github.com/prometheus/client_golang/prometheus" corev1 "k8s.io/api/core/v1" + apiequality "k8s.io/apimachinery/pkg/api/equality" apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" @@ -139,6 +140,7 @@ func (r *RuleReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl. r.Controller.updateRuleCache(ctx, rule) // Handle dry run + var delta nodeStatusDelta if rule.Spec.DryRun { if err := r.Controller.processDryRun(ctx, rule, nodeList); err != nil { log.Error(err, "Failed to process dry run", "rule", rule.Name) @@ -149,14 +151,16 @@ func (r *RuleReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl. rule.Status.DryRunResults = readinessv1alpha1.DryRunResults{} // Process all applicable nodes for this rule - if err := r.Controller.processAllNodesForRule(ctx, rule, nodeList); err != nil { + var err error + delta, err = r.Controller.processAllNodesForRule(ctx, rule, nodeList) + if err != nil { log.Error(err, "Failed to process nodes for rule", "rule", rule.Name) return ctrl.Result{RequeueAfter: time.Minute}, err } } // Update rule status - if err := r.Controller.updateRuleStatus(ctx, rule); err != nil { + if err := r.Controller.updateRuleStatus(ctx, rule, delta); err != nil { log.Error(err, "Failed to update rule status", "rule", rule.Name) return ctrl.Result{RequeueAfter: time.Minute}, err } @@ -198,9 +202,19 @@ func (r *RuleReconciler) reconcileDelete(ctx context.Context, rule *readinessv1a r.Controller.removeRuleFromCache(ctx, rule.Name) log.V(3).Info("Removing the finalizer from the rule") - patch := client.MergeFrom(rule.DeepCopy()) - controllerutil.RemoveFinalizer(rule, finalizerName) - err := r.Patch(ctx, rule, patch) + err := retry.RetryOnConflict(retry.DefaultRetry, func() error { + latest := &readinessv1alpha1.NodeReadinessRule{} + if err := r.Get(ctx, client.ObjectKey{Name: rule.Name}, latest); err != nil { + return client.IgnoreNotFound(err) + } + if !controllerutil.ContainsFinalizer(latest, finalizerName) { + return nil + } + + stored := latest.DeepCopy() + controllerutil.RemoveFinalizer(latest, finalizerName) + return r.Patch(ctx, latest, client.MergeFromWithOptions(stored, client.MergeFromWithOptimisticLock{})) + }) if err != nil { return ctrl.Result{}, err } @@ -251,56 +265,70 @@ func (r *RuleReadinessController) cleanupDeletedNodes(ctx context.Context, rule "before", len(rule.Status.NodeEvaluations), "after", len(newNodeEvaluations)) - // Use retry on conflict to update status to avoid race conditions from node updates - return retry.RetryOnConflict(retry.DefaultRetry, func() error { - fresh := &readinessv1alpha1.NodeReadinessRule{} - if err := r.Get(ctx, client.ObjectKey{Name: rule.Name}, fresh); err != nil { - return err - } - + // Use an optimistic-locked patch to avoid race conditions from concurrent node updates. + return r.patchRuleStatusWithOptimisticLock(ctx, rule.Name, func(fresh *readinessv1alpha1.NodeReadinessRule) { freshNodeEvaluations, freshFailedNodes := filterStatusForExistingNodes( existingNodes, fresh.Status.NodeEvaluations, fresh.Status.FailedNodes, ) - if len(freshNodeEvaluations) == len(fresh.Status.NodeEvaluations) && - len(freshFailedNodes) == len(fresh.Status.FailedNodes) { - return nil - } - - patch := client.MergeFrom(fresh.DeepCopy()) fresh.Status.NodeEvaluations = freshNodeEvaluations fresh.Status.FailedNodes = freshFailedNodes - return r.Status().Patch(ctx, fresh, patch) }) } -// processAllNodesForRule processes all nodes when a rule changes. +// processAllNodesForRule processes all nodes when a rule changes. It mutates rule.Status in place +// and additionally returns a nodeStatusDelta describing exactly which nodes' status are changed. +// so updateRuleStatus can merge those changes into the latest stored status +// instead of replacing NodeEvaluations/FailedNodes wholesale. // //nolint:unparam // Keep error return for future extensibility and API stability. -func (r *RuleReadinessController) processAllNodesForRule(ctx context.Context, rule *readinessv1alpha1.NodeReadinessRule, nodeList *corev1.NodeList) error { +func (r *RuleReadinessController) processAllNodesForRule(ctx context.Context, rule *readinessv1alpha1.NodeReadinessRule, nodeList *corev1.NodeList) (nodeStatusDelta, error) { log := ctrl.LoggerFrom(ctx) log.Info("Processing all nodes for rule", "rule", rule.Name, "totalNodes", len(nodeList.Items)) + delta := nodeStatusDelta{ + evaluations: make(map[string]readinessv1alpha1.NodeEvaluation), + failures: make(map[string]*readinessv1alpha1.NodeFailure), + } + var appliedNodes []string for _, node := range nodeList.Items { - if r.ruleAppliesTo(ctx, rule, &node) { - log.Info("Processing node for rule", "rule", rule.Name, "node", node.Name) - if err := r.evaluateRuleForNode(ctx, rule, &node); err != nil { - log.Error(err, "Failed to evaluate node for rule", "rule", rule.Name, "node", node.Name) - r.recordNodeFailure(rule, node.Name, "EvaluationError", err.Error()) - metrics.Failures.WithLabelValues(rule.Name, string(metrics.FailureReasonEvaluationError)).Inc() - } else { - appliedNodes = append(appliedNodes, node.Name) - var updatedFailedNodes []readinessv1alpha1.NodeFailure - for _, f := range rule.Status.FailedNodes { - if f.NodeName != node.Name { - updatedFailedNodes = append(updatedFailedNodes, f) - } + if !r.ruleAppliesTo(ctx, rule, &node) { + continue + } + + log.Info("Processing node for rule", "rule", rule.Name, "node", node.Name) + if err := r.evaluateRuleForNode(ctx, rule, &node); err != nil { + log.Error(err, "Failed to evaluate node for rule", "rule", rule.Name, "node", node.Name) + r.recordNodeFailure(rule, node.Name, "EvaluationError", err.Error()) + metrics.Failures.WithLabelValues(rule.Name, string(metrics.FailureReasonEvaluationError)).Inc() + + for _, f := range rule.Status.FailedNodes { + if f.NodeName == node.Name { + failure := f + delta.failures[node.Name] = &failure + break + } + } + } else { + appliedNodes = append(appliedNodes, node.Name) + var updatedFailedNodes []readinessv1alpha1.NodeFailure + for _, f := range rule.Status.FailedNodes { + if f.NodeName != node.Name { + updatedFailedNodes = append(updatedFailedNodes, f) + } + } + rule.Status.FailedNodes = updatedFailedNodes + delta.failures[node.Name] = nil // clear any previously-recorded failure + + for _, eval := range rule.Status.NodeEvaluations { + if eval.NodeName == node.Name { + delta.evaluations[node.Name] = eval + break } - rule.Status.FailedNodes = updatedFailedNodes } } } @@ -314,7 +342,7 @@ func (r *RuleReadinessController) processAllNodesForRule(ctx context.Context, ru } log.Info("Completed processing nodes for rule", "rule", rule.Name, "processedCount", len(appliedNodes)) - return nil + return delta, nil } // evaluateRuleForNode evaluates a single rule against a single node. @@ -570,8 +598,43 @@ func (r *RuleReadinessController) removeRuleFromCache(ctx context.Context, ruleN log.Info("Removed rule from cache", "rule", ruleName, "totalRules", len(r.ruleCache)) } -// updateRuleStatus updates the status of a NodeReadinessRule. -func (r *RuleReadinessController) updateRuleStatus(ctx context.Context, rule *readinessv1alpha1.NodeReadinessRule) error { +// patchRuleStatusWithOptimisticLock fetches the latest NodeReadinessRule, and apply mutate status +// changes to it. It then patches the result to API with an optimistic-locked JSON merge patch. mutate +// should return false if it made no changes, to skip an unnecessary Patch call. +// +// We use client.MergeFromWithOptimisticLock here for a JSON merge patch replaces slice fields +// (NodeEvaluations, AppliedNodes, FailedNodes) wholesale rather than merging them, so without a +// resourceVersion precondition retry.RetryOnConflict can never observe a genuine conflict and a +// concurrent status write from the other reconciler (RuleReconciler and NodeReconciler both patch +// NodeReadinessRule.Status independently) can be silently overwritten. +func (r *RuleReadinessController) patchRuleStatusWithOptimisticLock( + ctx context.Context, + ruleName string, + mutate func(latest *readinessv1alpha1.NodeReadinessRule), +) error { + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + latestRule := &readinessv1alpha1.NodeReadinessRule{} + if err := r.Get(ctx, client.ObjectKey{Name: ruleName}, latestRule); err != nil { + return err + } + + stored := latestRule.DeepCopy() + mutate(latestRule) + + if apiequality.Semantic.DeepEqual(stored.Status, latestRule.Status) { + return nil + } + + return r.Status().Patch(ctx, latestRule, client.MergeFromWithOptions(stored, client.MergeFromWithOptimisticLock{})) + }) +} + +// updateRuleStatus updates the status of a NodeReadinessRule. delta carries the per-node +// NodeEvaluations/FailedNodes changes processAllNodesForRule actually produced this reconcile; +// it is merged into the latest stored status by node name (see applyNodeStatusDelta) rather than +// replacing those fields wholesale, so a concurrent per-node update from NodeReconciler +// (processNodeAgainstAllRules) for a node outside this sweep isn't silently discarded. +func (r *RuleReadinessController) updateRuleStatus(ctx context.Context, rule *readinessv1alpha1.NodeReadinessRule, delta nodeStatusDelta) error { log := ctrl.LoggerFrom(ctx) log.V(1).Info("Updating rule status", @@ -579,30 +642,19 @@ func (r *RuleReadinessController) updateRuleStatus(ctx context.Context, rule *re "nodeEvaluations", len(rule.Status.NodeEvaluations), "appliedNodes", len(rule.Status.AppliedNodes)) - return 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()) - - latestRule.Status.NodeEvaluations = rule.Status.NodeEvaluations + err := r.patchRuleStatusWithOptimisticLock(ctx, rule.Name, func(latestRule *readinessv1alpha1.NodeReadinessRule) { + applyNodeStatusDelta(latestRule, delta) latestRule.Status.AppliedNodes = rule.Status.AppliedNodes - latestRule.Status.FailedNodes = rule.Status.FailedNodes latestRule.Status.ObservedGeneration = rule.Status.ObservedGeneration latestRule.Status.DryRunResults = rule.Status.DryRunResults - - if err := r.Status().Patch(ctx, latestRule, patch); err != nil { - log.V(1).Info("Status patch conflict, will retry", - "rule", rule.Name, - "error", err.Error()) - return err - } - - log.V(1).Info("Successfully patched rule status", "rule", rule.Name) - return nil }) + if err != nil { + log.V(1).Info("Failed to patch rule status", "rule", rule.Name, "error", err.Error()) + return err + } + + log.V(1).Info("Successfully patched rule status", "rule", rule.Name) + return nil } // processDryRun processes dry run for a rule. @@ -725,13 +777,30 @@ func (r *RuleReconciler) ensureFinalizer(ctx context.Context, rule *readinessv1a return false, nil } - patch := client.MergeFrom(rule.DeepCopy()) - controllerutil.AddFinalizer(rule, finalizer) - err = r.Patch(ctx, rule, patch) + added := false + err = retry.RetryOnConflict(retry.DefaultRetry, func() error { + latest := &readinessv1alpha1.NodeReadinessRule{} + if err := r.Get(ctx, client.ObjectKey{Name: rule.Name}, latest); err != nil { + return err + } + if controllerutil.ContainsFinalizer(latest, finalizer) { + return nil + } + + stored := latest.DeepCopy() + controllerutil.AddFinalizer(latest, finalizer) + if err := r.Patch(ctx, latest, client.MergeFromWithOptions(stored, client.MergeFromWithOptimisticLock{})); err != nil { + return err + } + + *rule = *latest + added = true + return nil + }) if err != nil { return false, err } - return true, nil + return added, nil } // getPreviousNodeEvaluation retrieves the previous evaluation result for a specific node from the rule status. diff --git a/internal/controller/nodereadinessrule_controller_test.go b/internal/controller/nodereadinessrule_controller_test.go index 07fb4b30..000889b5 100644 --- a/internal/controller/nodereadinessrule_controller_test.go +++ b/internal/controller/nodereadinessrule_controller_test.go @@ -19,6 +19,7 @@ package controller import ( "context" "fmt" + "sync/atomic" "time" . "github.com/onsi/ginkgo/v2" @@ -33,6 +34,8 @@ import ( "k8s.io/client-go/kubernetes/fake" "k8s.io/client-go/tools/events" "sigs.k8s.io/controller-runtime/pkg/client" + fakeclient "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" "sigs.k8s.io/controller-runtime/pkg/reconcile" nodereadinessiov1alpha1 "sigs.k8s.io/node-readiness-controller/api/v1alpha1" @@ -68,6 +71,15 @@ func (c *errorInjectingClient) Patch(ctx context.Context, obj client.Object, pat return c.Client.Patch(ctx, obj, patch, opts...) } +func findNodeEvaluation(evals []nodereadinessiov1alpha1.NodeEvaluation, nodeName string) *nodereadinessiov1alpha1.NodeEvaluation { + for i := range evals { + if evals[i].NodeName == nodeName { + return &evals[i] + } + } + return nil +} + var _ = Describe("NodeReadinessRule Controller", func() { var ( ctx context.Context @@ -2361,7 +2373,8 @@ var _ = Describe("NodeReadinessRule Controller", func() { } nodeList := &corev1.NodeList{Items: []corev1.Node{*failNode}} - Expect(failController.processAllNodesForRule(ctx, rule, nodeList)).To(Succeed()) + _, err := failController.processAllNodesForRule(ctx, rule, nodeList) + Expect(err).NotTo(HaveOccurred()) Expect(rule.Status.AppliedNodes).NotTo(ContainElement("fail-path-node")) @@ -2418,7 +2431,8 @@ var _ = Describe("NodeReadinessRule Controller", func() { } nodeList := &corev1.NodeList{Items: []corev1.Node{*successNode}} - Expect(successController.processAllNodesForRule(ctx, rule, nodeList)).To(Succeed()) + _, err := successController.processAllNodesForRule(ctx, rule, nodeList) + Expect(err).NotTo(HaveOccurred()) Expect(rule.Status.AppliedNodes).To(ContainElement("stale-recovery-node")) @@ -2535,4 +2549,143 @@ var _ = Describe("NodeReadinessRule Controller", func() { Expect(readinessController.hasTaintBySpec(anyOfNode, rule.Spec.Taint)).To(BeFalse()) }) }) + + Context("concurrent status writes between RuleReconciler and NodeReconciler", func() { + var testScheme *runtime.Scheme + + BeforeEach(func() { + testScheme = runtime.NewScheme() + Expect(corev1.AddToScheme(testScheme)).To(Succeed()) + Expect(nodereadinessiov1alpha1.AddToScheme(testScheme)).To(Succeed()) + }) + + It("should not discard a NodeReconciler-written evaluation for a node outside the RuleReconciler sweep", func() { + nodeA := &corev1.Node{ + ObjectMeta: metav1.ObjectMeta{Name: "concurrent-node-a", Labels: map[string]string{"role": "worker"}}, + Status: corev1.NodeStatus{Conditions: []corev1.NodeCondition{{Type: "Ready", Status: corev1.ConditionTrue}}}, + } + nodeB := &corev1.Node{ + ObjectMeta: metav1.ObjectMeta{Name: "concurrent-node-b", Labels: map[string]string{"role": "worker"}}, + Status: corev1.NodeStatus{Conditions: []corev1.NodeCondition{{Type: "Ready", Status: corev1.ConditionTrue}}}, + } + + rule := &nodereadinessiov1alpha1.NodeReadinessRule{ + ObjectMeta: metav1.ObjectMeta{Name: "concurrent-status-rule"}, + Spec: nodereadinessiov1alpha1.NodeReadinessRuleSpec{ + Conditions: []nodereadinessiov1alpha1.ConditionRequirement{ + {Type: "Ready", RequiredStatus: corev1.ConditionTrue}, + }, + Taint: corev1.Taint{Key: "readiness.k8s.io/concurrent-test", Effect: corev1.TaintEffectNoSchedule}, + NodeSelector: metav1.LabelSelector{MatchLabels: map[string]string{"role": "worker"}}, + EnforcementMode: nodereadinessiov1alpha1.EnforcementModeContinuous, + }, + } + + fc := fakeclient.NewClientBuilder(). + WithScheme(testScheme). + WithObjects(nodeA, nodeB, rule). + WithStatusSubresource(rule). + Build() + + controller := &RuleReadinessController{ + Client: fc, + Scheme: testScheme, + clientset: fake.NewSimpleClientset(), + ruleCache: make(map[string]*nodereadinessiov1alpha1.NodeReadinessRule), + EventRecorder: events.NewFakeRecorder(10), + } + controller.updateRuleCache(ctx, rule) + + // Simulate RuleReconciler starting a sweep from a nodeList snapshot that + // only knows about node-a (e.g. node-b joined after the snapshot was taken). + staleNodeList := &corev1.NodeList{Items: []corev1.Node{*nodeA}} + delta, err := controller.processAllNodesForRule(ctx, rule, staleNodeList) + Expect(err).NotTo(HaveOccurred()) + + // While that sweep is "in flight", NodeReconciler independently reconciles + // node-b and persists its own evaluation into the same rule's status. + Expect(controller.processNodeAgainstAllRules(ctx, nodeB)).To(Succeed()) + + afterConcurrentWrite := &nodereadinessiov1alpha1.NodeReadinessRule{} + Expect(fc.Get(ctx, types.NamespacedName{Name: rule.Name}, afterConcurrentWrite)).To(Succeed()) + Expect(findNodeEvaluation(afterConcurrentWrite.Status.NodeEvaluations, "concurrent-node-b")). + NotTo(BeNil(), "NodeReconciler should have persisted an evaluation for node-b") + + // RuleReconciler now finishes its (stale) sweep and patches the rule status. + Expect(controller.updateRuleStatus(ctx, rule, delta)).To(Succeed()) + + final := &nodereadinessiov1alpha1.NodeReadinessRule{} + Expect(fc.Get(ctx, types.NamespacedName{Name: rule.Name}, final)).To(Succeed()) + + Expect(findNodeEvaluation(final.Status.NodeEvaluations, "concurrent-node-a")).NotTo(BeNil(), + "node-a's freshly computed evaluation should be present") + Expect(findNodeEvaluation(final.Status.NodeEvaluations, "concurrent-node-b")).NotTo(BeNil(), + "node-b's evaluation, written concurrently by NodeReconciler outside this sweep, must not be discarded") + }) + + It("should retry updateRuleStatus on a genuine conflict instead of silently clobbering it", func() { + node := &corev1.Node{ + ObjectMeta: metav1.ObjectMeta{Name: "conflict-node", Labels: map[string]string{"role": "worker"}}, + Status: corev1.NodeStatus{Conditions: []corev1.NodeCondition{{Type: "Ready", Status: corev1.ConditionTrue}}}, + } + + rule := &nodereadinessiov1alpha1.NodeReadinessRule{ + ObjectMeta: metav1.ObjectMeta{Name: "conflict-rule"}, + Spec: nodereadinessiov1alpha1.NodeReadinessRuleSpec{ + Conditions: []nodereadinessiov1alpha1.ConditionRequirement{ + {Type: "Ready", RequiredStatus: corev1.ConditionTrue}, + }, + Taint: corev1.Taint{Key: "readiness.k8s.io/conflict-test", Effect: corev1.TaintEffectNoSchedule}, + NodeSelector: metav1.LabelSelector{MatchLabels: map[string]string{"role": "worker"}}, + EnforcementMode: nodereadinessiov1alpha1.EnforcementModeContinuous, + }, + } + + var patchCount atomic.Int32 + + fc := fakeclient.NewClientBuilder(). + WithScheme(testScheme). + WithObjects(node, rule). + WithStatusSubresource(rule). + WithInterceptorFuncs(interceptor.Funcs{ + SubResourcePatch: func(ctx context.Context, c client.Client, subResourceName string, obj client.Object, patch client.Patch, opts ...client.SubResourcePatchOption) error { + if r, ok := obj.(*nodereadinessiov1alpha1.NodeReadinessRule); ok && r.Name == "conflict-rule" && patchCount.Add(1) == 1 { + // Simulate a concurrent NodeReconciler write for a different + // node landing in between our Get and our Patch. + current := &nodereadinessiov1alpha1.NodeReadinessRule{} + Expect(c.Get(ctx, types.NamespacedName{Name: r.Name}, current)).To(Succeed()) + current.Status.NodeEvaluations = append(current.Status.NodeEvaluations, nodereadinessiov1alpha1.NodeEvaluation{ + NodeName: "conflict-other-node", + TaintStatus: nodereadinessiov1alpha1.TaintStatusAbsent, + }) + Expect(c.Status().Update(ctx, current)).To(Succeed()) + } + return c.SubResource(subResourceName).Patch(ctx, obj, patch, opts...) + }, + }). + Build() + + controller := &RuleReadinessController{ + Client: fc, + Scheme: testScheme, + clientset: fake.NewSimpleClientset(), + ruleCache: make(map[string]*nodereadinessiov1alpha1.NodeReadinessRule), + EventRecorder: events.NewFakeRecorder(10), + } + + nodeList := &corev1.NodeList{Items: []corev1.Node{*node}} + delta, err := controller.processAllNodesForRule(ctx, rule, nodeList) + Expect(err).NotTo(HaveOccurred()) + + Expect(controller.updateRuleStatus(ctx, rule, delta)).To(Succeed()) + Expect(patchCount.Load()).To(BeNumerically(">=", 2), + "the first patch attempt should hit a conflict and force a retry") + + final := &nodereadinessiov1alpha1.NodeReadinessRule{} + Expect(fc.Get(ctx, types.NamespacedName{Name: rule.Name}, final)).To(Succeed()) + Expect(findNodeEvaluation(final.Status.NodeEvaluations, "conflict-node")).NotTo(BeNil()) + Expect(findNodeEvaluation(final.Status.NodeEvaluations, "conflict-other-node")).NotTo(BeNil(), + "the concurrently-written evaluation must survive the retried patch") + }) + }) })