Skip to content

Commit cb75223

Browse files
fix(controller): use optimistic locking on rule status patches
RuleReconciler and NodeReconciler both patch NodeReadinessRule.Status concurrently, but every status/finalizer patch used a plain client.MergeFrom with no resourceVersion precondition, wrapped in retry.RetryOnConflict. Since a JSON merge patch never carries that precondition unless MergeFromWithOptimisticLock is used, the API server never returns a conflict and the retry wrapper never actually retries -- the same bug fixed earlier for node taint patches, just left open on the rule-status side. Worse, updateRuleStatus replaced NodeEvaluations/FailedNodes wholesale from a snapshot computed at the start of a RuleReconciler sweep, so it could silently discard a concurrent NodeReconciler per-node update for a node outside that sweep. Fix this by having processAllNodesForRule return a delta of exactly the per-node changes it made, and merging that delta by node name instead of overwriting the whole slice. Also add the missing optimistic lock to ensureFinalizer, the finalizer removal in reconcileDelete, cleanupDeletedNodes, and markBootstrapCompleted's node annotation patch, matching the pattern already used by addTaintBySpec/removeTaintBySpec.
1 parent 55dbc33 commit cb75223

4 files changed

Lines changed: 347 additions & 75 deletions

File tree

‎internal/controller/helper.go‎

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,12 @@ package controller
1919
import (
2020
"encoding/json"
2121
"maps"
22+
"sort"
2223

2324
corev1 "k8s.io/api/core/v1"
2425
"k8s.io/apimachinery/pkg/types"
26+
27+
readinessv1alpha1 "sigs.k8s.io/node-readiness-controller/api/v1alpha1"
2528
)
2629

2730
//nolint:godot
@@ -109,3 +112,56 @@ func taintsEqual(a, b []corev1.Taint) bool {
109112
func labelsEqual(a, b map[string]string) bool {
110113
return maps.Equal(a, b)
111114
}
115+
116+
// nodeStatusDelta captures the per-node NodeEvaluation/NodeFailure changes a single
117+
// processAllNodesForRule sweep actually produced, keyed by node name. It intentionally does not
118+
// carry AppliedNodes/ObservedGeneration/DryRunResults: those fields only ever have one writer
119+
// (RuleReconciler), so a plain overwrite of them is safe.
120+
//
121+
// A nil value in failures means "clear any failure recorded for this node" (the node evaluated
122+
// successfully). evaluations only ever holds entries for nodes that were freshly (re-)evaluated
123+
// this sweep; a failed evaluation leaves the node's prior NodeEvaluation untouched.
124+
type nodeStatusDelta struct {
125+
evaluations map[string]readinessv1alpha1.NodeEvaluation
126+
failures map[string]*readinessv1alpha1.NodeFailure
127+
}
128+
129+
// applyNodeStatusDelta merges delta into rule's NodeEvaluations/FailedNodes, replacing only the
130+
// entries for nodes present in delta and leaving every other node's entry untouched.
131+
//
132+
// This is the crux of fixing the lost-update bug described in #341: a naive full-slice
133+
// replacement of NodeEvaluations/FailedNodes (computed from a nodeList snapshot taken at the
134+
// start of a RuleReconciler sweep) would silently discard any per-node status update written
135+
// concurrently by NodeReconciler for a node this particular sweep didn't touch. Merging by node
136+
// name instead means each writer only ever overwrites the entries it just recomputed.
137+
func applyNodeStatusDelta(rule *readinessv1alpha1.NodeReadinessRule, delta nodeStatusDelta) {
138+
if len(delta.evaluations) > 0 {
139+
merged := make([]readinessv1alpha1.NodeEvaluation, 0, len(rule.Status.NodeEvaluations)+len(delta.evaluations))
140+
for _, eval := range rule.Status.NodeEvaluations {
141+
if _, changed := delta.evaluations[eval.NodeName]; !changed {
142+
merged = append(merged, eval)
143+
}
144+
}
145+
for _, eval := range delta.evaluations {
146+
merged = append(merged, eval)
147+
}
148+
sort.Slice(merged, func(i, j int) bool { return merged[i].NodeName < merged[j].NodeName })
149+
rule.Status.NodeEvaluations = merged
150+
}
151+
152+
if len(delta.failures) > 0 {
153+
merged := make([]readinessv1alpha1.NodeFailure, 0, len(rule.Status.FailedNodes)+len(delta.failures))
154+
for _, failure := range rule.Status.FailedNodes {
155+
if _, changed := delta.failures[failure.NodeName]; !changed {
156+
merged = append(merged, failure)
157+
}
158+
}
159+
for _, failure := range delta.failures {
160+
if failure != nil {
161+
merged = append(merged, *failure)
162+
}
163+
}
164+
sort.Slice(merged, func(i, j int) bool { return merged[i].NodeName < merged[j].NodeName })
165+
rule.Status.FailedNodes = merged
166+
}
167+
}

‎internal/controller/node_controller.go‎

Lines changed: 4 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -169,14 +169,7 @@ func (r *RuleReadinessController) processNodeAgainstAllRules(ctx context.Context
169169

170170
var successfullyPatchedRule *readinessv1alpha1.NodeReadinessRule
171171

172-
err := retry.RetryOnConflict(retry.DefaultRetry, func() error {
173-
latestRule := &readinessv1alpha1.NodeReadinessRule{}
174-
if err := r.Get(ctx, client.ObjectKey{Name: rule.Name}, latestRule); err != nil {
175-
return err
176-
}
177-
178-
patch := client.MergeFrom(latestRule.DeepCopy())
179-
172+
err := r.patchRuleStatusWithOptimisticLock(ctx, rule.Name, func(latestRule *readinessv1alpha1.NodeReadinessRule) bool {
180173
// update only this specific node evaluation status
181174
currEval := readinessv1alpha1.NodeEvaluation{}
182175
for _, eval := range rule.Status.NodeEvaluations {
@@ -215,12 +208,8 @@ func (r *RuleReadinessController) processNodeAgainstAllRules(ctx context.Context
215208
}
216209
latestRule.Status.FailedNodes = updatedFailedNodes
217210

218-
if err := r.Status().Patch(ctx, latestRule, patch); err != nil {
219-
return err
220-
}
221-
222211
successfullyPatchedRule = latestRule
223-
return nil
212+
return true
224213
})
225214

226215
if err != nil {
@@ -374,15 +363,15 @@ func (r *RuleReadinessController) markBootstrapCompleted(ctx context.Context, no
374363
return nil
375364
}
376365

377-
patch := client.MergeFrom(node.DeepCopy())
366+
stored := node.DeepCopy()
378367

379368
// Initialize annotations map if nil.
380369
if node.Annotations == nil {
381370
node.Annotations = make(map[string]string)
382371
}
383372

384373
node.Annotations[annotationKey] = bootstrapAnnotationValue(ruleName)
385-
if err := r.Patch(ctx, node, patch); err != nil {
374+
if err := r.Patch(ctx, node, client.MergeFromWithOptions(stored, client.MergeFromWithOptimisticLock{})); err != nil {
386375
return err
387376
}
388377

‎internal/controller/nodereadinessrule_controller.go‎

Lines changed: 132 additions & 58 deletions
Original file line numberDiff line numberDiff line change
@@ -139,6 +139,7 @@ func (r *RuleReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.
139139
r.Controller.updateRuleCache(ctx, rule)
140140

141141
// Handle dry run
142+
var delta nodeStatusDelta
142143
if rule.Spec.DryRun {
143144
if err := r.Controller.processDryRun(ctx, rule, nodeList); err != nil {
144145
log.Error(err, "Failed to process dry run", "rule", rule.Name)
@@ -149,14 +150,16 @@ func (r *RuleReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.
149150
rule.Status.DryRunResults = readinessv1alpha1.DryRunResults{}
150151

151152
// Process all applicable nodes for this rule
152-
if err := r.Controller.processAllNodesForRule(ctx, rule, nodeList); err != nil {
153+
var err error
154+
delta, err = r.Controller.processAllNodesForRule(ctx, rule, nodeList)
155+
if err != nil {
153156
log.Error(err, "Failed to process nodes for rule", "rule", rule.Name)
154157
return ctrl.Result{RequeueAfter: time.Minute}, err
155158
}
156159
}
157160

158161
// Update rule status
159-
if err := r.Controller.updateRuleStatus(ctx, rule); err != nil {
162+
if err := r.Controller.updateRuleStatus(ctx, rule, delta); err != nil {
160163
log.Error(err, "Failed to update rule status", "rule", rule.Name)
161164
return ctrl.Result{RequeueAfter: time.Minute}, err
162165
}
@@ -198,9 +201,19 @@ func (r *RuleReconciler) reconcileDelete(ctx context.Context, rule *readinessv1a
198201
r.Controller.removeRuleFromCache(ctx, rule.Name)
199202

200203
log.V(3).Info("Removing the finalizer from the rule")
201-
patch := client.MergeFrom(rule.DeepCopy())
202-
controllerutil.RemoveFinalizer(rule, finalizerName)
203-
err := r.Patch(ctx, rule, patch)
204+
err := retry.RetryOnConflict(retry.DefaultRetry, func() error {
205+
latest := &readinessv1alpha1.NodeReadinessRule{}
206+
if err := r.Get(ctx, client.ObjectKey{Name: rule.Name}, latest); err != nil {
207+
return client.IgnoreNotFound(err)
208+
}
209+
if !controllerutil.ContainsFinalizer(latest, finalizerName) {
210+
return nil
211+
}
212+
213+
stored := latest.DeepCopy()
214+
controllerutil.RemoveFinalizer(latest, finalizerName)
215+
return r.Patch(ctx, latest, client.MergeFromWithOptions(stored, client.MergeFromWithOptimisticLock{}))
216+
})
204217
if err != nil {
205218
return ctrl.Result{}, err
206219
}
@@ -251,13 +264,8 @@ func (r *RuleReadinessController) cleanupDeletedNodes(ctx context.Context, rule
251264
"before", len(rule.Status.NodeEvaluations),
252265
"after", len(newNodeEvaluations))
253266

254-
// Use retry on conflict to update status to avoid race conditions from node updates
255-
return retry.RetryOnConflict(retry.DefaultRetry, func() error {
256-
fresh := &readinessv1alpha1.NodeReadinessRule{}
257-
if err := r.Get(ctx, client.ObjectKey{Name: rule.Name}, fresh); err != nil {
258-
return err
259-
}
260-
267+
// Use an optimistic-locked patch to avoid race conditions from concurrent node updates.
268+
return r.patchRuleStatusWithOptimisticLock(ctx, rule.Name, func(fresh *readinessv1alpha1.NodeReadinessRule) bool {
261269
var freshNodeEvaluations []readinessv1alpha1.NodeEvaluation
262270
for _, evaluation := range fresh.Status.NodeEvaluations {
263271
if existingNodes[evaluation.NodeName] {
@@ -266,40 +274,65 @@ func (r *RuleReadinessController) cleanupDeletedNodes(ctx context.Context, rule
266274
}
267275

268276
if len(freshNodeEvaluations) == len(fresh.Status.NodeEvaluations) {
269-
return nil
277+
return false
270278
}
271279

272-
patch := client.MergeFrom(fresh.DeepCopy())
273280
fresh.Status.NodeEvaluations = freshNodeEvaluations
274-
return r.Status().Patch(ctx, fresh, patch)
281+
return true
275282
})
276283
}
277284

278-
// processAllNodesForRule processes all nodes when a rule changes.
285+
// processAllNodesForRule processes all nodes when a rule changes. It mutates rule.Status in place
286+
// (as before) and additionally returns a nodeStatusDelta describing exactly which nodes' status
287+
// this sweep changed, so updateRuleStatus can merge those changes into the latest stored status
288+
// instead of replacing NodeEvaluations/FailedNodes wholesale.
279289
//
280290
//nolint:unparam // Keep error return for future extensibility and API stability.
281-
func (r *RuleReadinessController) processAllNodesForRule(ctx context.Context, rule *readinessv1alpha1.NodeReadinessRule, nodeList *corev1.NodeList) error {
291+
func (r *RuleReadinessController) processAllNodesForRule(ctx context.Context, rule *readinessv1alpha1.NodeReadinessRule, nodeList *corev1.NodeList) (nodeStatusDelta, error) {
282292
log := ctrl.LoggerFrom(ctx)
283293

284294
log.Info("Processing all nodes for rule", "rule", rule.Name, "totalNodes", len(nodeList.Items))
285295

296+
delta := nodeStatusDelta{
297+
evaluations: make(map[string]readinessv1alpha1.NodeEvaluation),
298+
failures: make(map[string]*readinessv1alpha1.NodeFailure),
299+
}
300+
286301
var appliedNodes []string
287302
for _, node := range nodeList.Items {
288-
if r.ruleAppliesTo(ctx, rule, &node) {
289-
log.Info("Processing node for rule", "rule", rule.Name, "node", node.Name)
290-
if err := r.evaluateRuleForNode(ctx, rule, &node); err != nil {
291-
log.Error(err, "Failed to evaluate node for rule", "rule", rule.Name, "node", node.Name)
292-
r.recordNodeFailure(rule, node.Name, "EvaluationError", err.Error())
293-
metrics.Failures.WithLabelValues(rule.Name, "EvaluationError").Inc()
294-
} else {
295-
appliedNodes = append(appliedNodes, node.Name)
296-
var updatedFailedNodes []readinessv1alpha1.NodeFailure
297-
for _, f := range rule.Status.FailedNodes {
298-
if f.NodeName != node.Name {
299-
updatedFailedNodes = append(updatedFailedNodes, f)
300-
}
303+
if !r.ruleAppliesTo(ctx, rule, &node) {
304+
continue
305+
}
306+
307+
log.Info("Processing node for rule", "rule", rule.Name, "node", node.Name)
308+
if err := r.evaluateRuleForNode(ctx, rule, &node); err != nil {
309+
log.Error(err, "Failed to evaluate node for rule", "rule", rule.Name, "node", node.Name)
310+
r.recordNodeFailure(rule, node.Name, "EvaluationError", err.Error())
311+
metrics.Failures.WithLabelValues(rule.Name, "EvaluationError").Inc()
312+
313+
for _, f := range rule.Status.FailedNodes {
314+
if f.NodeName == node.Name {
315+
failure := f
316+
delta.failures[node.Name] = &failure
317+
break
318+
}
319+
}
320+
} else {
321+
appliedNodes = append(appliedNodes, node.Name)
322+
var updatedFailedNodes []readinessv1alpha1.NodeFailure
323+
for _, f := range rule.Status.FailedNodes {
324+
if f.NodeName != node.Name {
325+
updatedFailedNodes = append(updatedFailedNodes, f)
326+
}
327+
}
328+
rule.Status.FailedNodes = updatedFailedNodes
329+
delta.failures[node.Name] = nil // clear any previously-recorded failure
330+
331+
for _, eval := range rule.Status.NodeEvaluations {
332+
if eval.NodeName == node.Name {
333+
delta.evaluations[node.Name] = eval
334+
break
301335
}
302-
rule.Status.FailedNodes = updatedFailedNodes
303336
}
304337
}
305338
}
@@ -313,7 +346,7 @@ func (r *RuleReadinessController) processAllNodesForRule(ctx context.Context, ru
313346
}
314347

315348
log.Info("Completed processing nodes for rule", "rule", rule.Name, "processedCount", len(appliedNodes))
316-
return nil
349+
return delta, nil
317350
}
318351

319352
// evaluateRuleForNode evaluates a single rule against a single node.
@@ -557,39 +590,63 @@ func (r *RuleReadinessController) removeRuleFromCache(ctx context.Context, ruleN
557590
log.Info("Removed rule from cache", "rule", ruleName, "totalRules", len(r.ruleCache))
558591
}
559592

560-
// updateRuleStatus updates the status of a NodeReadinessRule.
561-
func (r *RuleReadinessController) updateRuleStatus(ctx context.Context, rule *readinessv1alpha1.NodeReadinessRule) error {
593+
// patchRuleStatusWithOptimisticLock fetches the latest NodeReadinessRule, lets mutate apply status
594+
// changes to it, and patches the result back with an optimistic-locked JSON merge patch. mutate
595+
// should return false if it made no changes, to skip an unnecessary Patch call.
596+
//
597+
// We use client.MergeFromWithOptimisticLock here for the same reason addTaintBySpec/
598+
// removeTaintBySpec do (see node_controller.go): a JSON merge patch replaces slice fields
599+
// (NodeEvaluations, AppliedNodes, FailedNodes) wholesale rather than merging them, so without a
600+
// resourceVersion precondition retry.RetryOnConflict can never observe a genuine conflict and a
601+
// concurrent status write from the other reconciler (RuleReconciler and NodeReconciler both patch
602+
// NodeReadinessRule.Status independently) can be silently overwritten.
603+
func (r *RuleReadinessController) patchRuleStatusWithOptimisticLock(
604+
ctx context.Context,
605+
ruleName string,
606+
mutate func(latest *readinessv1alpha1.NodeReadinessRule) (changed bool),
607+
) error {
608+
return retry.RetryOnConflict(retry.DefaultRetry, func() error {
609+
latestRule := &readinessv1alpha1.NodeReadinessRule{}
610+
if err := r.Get(ctx, client.ObjectKey{Name: ruleName}, latestRule); err != nil {
611+
return err
612+
}
613+
614+
stored := latestRule.DeepCopy()
615+
if !mutate(latestRule) {
616+
return nil
617+
}
618+
619+
return r.Status().Patch(ctx, latestRule, client.MergeFromWithOptions(stored, client.MergeFromWithOptimisticLock{}))
620+
})
621+
}
622+
623+
// updateRuleStatus updates the status of a NodeReadinessRule. delta carries the per-node
624+
// NodeEvaluations/FailedNodes changes processAllNodesForRule actually produced this reconcile;
625+
// it is merged into the latest stored status by node name (see applyNodeStatusDelta) rather than
626+
// replacing those fields wholesale, so a concurrent per-node update from NodeReconciler
627+
// (processNodeAgainstAllRules) for a node outside this sweep isn't silently discarded.
628+
func (r *RuleReadinessController) updateRuleStatus(ctx context.Context, rule *readinessv1alpha1.NodeReadinessRule, delta nodeStatusDelta) error {
562629
log := ctrl.LoggerFrom(ctx)
563630

564631
log.V(1).Info("Updating rule status",
565632
"rule", rule.Name,
566633
"nodeEvaluations", len(rule.Status.NodeEvaluations),
567634
"appliedNodes", len(rule.Status.AppliedNodes))
568635

569-
return retry.RetryOnConflict(retry.DefaultRetry, func() error {
570-
latestRule := &readinessv1alpha1.NodeReadinessRule{}
571-
if err := r.Get(ctx, client.ObjectKey{Name: rule.Name}, latestRule); err != nil {
572-
return err
573-
}
574-
575-
patch := client.MergeFrom(latestRule.DeepCopy())
576-
577-
latestRule.Status.NodeEvaluations = rule.Status.NodeEvaluations
636+
err := r.patchRuleStatusWithOptimisticLock(ctx, rule.Name, func(latestRule *readinessv1alpha1.NodeReadinessRule) bool {
637+
applyNodeStatusDelta(latestRule, delta)
578638
latestRule.Status.AppliedNodes = rule.Status.AppliedNodes
579-
latestRule.Status.FailedNodes = rule.Status.FailedNodes
580639
latestRule.Status.ObservedGeneration = rule.Status.ObservedGeneration
581640
latestRule.Status.DryRunResults = rule.Status.DryRunResults
582-
583-
if err := r.Status().Patch(ctx, latestRule, patch); err != nil {
584-
log.V(1).Info("Status patch conflict, will retry",
585-
"rule", rule.Name,
586-
"error", err.Error())
587-
return err
588-
}
589-
590-
log.V(1).Info("Successfully patched rule status", "rule", rule.Name)
591-
return nil
641+
return true
592642
})
643+
if err != nil {
644+
log.V(1).Info("Failed to patch rule status", "rule", rule.Name, "error", err.Error())
645+
return err
646+
}
647+
648+
log.V(1).Info("Successfully patched rule status", "rule", rule.Name)
649+
return nil
593650
}
594651

595652
// processDryRun processes dry run for a rule.
@@ -705,13 +762,30 @@ func (r *RuleReconciler) ensureFinalizer(ctx context.Context, rule *readinessv1a
705762
return false, nil
706763
}
707764

708-
patch := client.MergeFrom(rule.DeepCopy())
709-
controllerutil.AddFinalizer(rule, finalizer)
710-
err = r.Patch(ctx, rule, patch)
765+
added := false
766+
err = retry.RetryOnConflict(retry.DefaultRetry, func() error {
767+
latest := &readinessv1alpha1.NodeReadinessRule{}
768+
if err := r.Get(ctx, client.ObjectKey{Name: rule.Name}, latest); err != nil {
769+
return err
770+
}
771+
if controllerutil.ContainsFinalizer(latest, finalizer) {
772+
return nil
773+
}
774+
775+
stored := latest.DeepCopy()
776+
controllerutil.AddFinalizer(latest, finalizer)
777+
if err := r.Patch(ctx, latest, client.MergeFromWithOptions(stored, client.MergeFromWithOptimisticLock{})); err != nil {
778+
return err
779+
}
780+
781+
*rule = *latest
782+
added = true
783+
return nil
784+
})
711785
if err != nil {
712786
return false, err
713787
}
714-
return true, nil
788+
return added, nil
715789
}
716790

717791
// getPreviousNodeEvaluation retrieves the previous evaluation result for a specific node from the rule status.

0 commit comments

Comments
 (0)