Skip to content

Commit 80e3bf3

Browse files
authored
fix: address an issue in AppliedWork object processing when the agent leaves then re-joins (kubefleet-dev#421)
1 parent 5e7013c commit 80e3bf3

2 files changed

Lines changed: 244 additions & 27 deletions

File tree

pkg/controllers/workapplier/controller.go

Lines changed: 37 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -551,23 +551,40 @@ func (r *Reconciler) forgetWorkAndRemoveFinalizer(ctx context.Context, work *fle
551551
// ensureAppliedWork makes sure that an associated appliedWork and a finalizer on the work resource exists on the cluster.
552552
func (r *Reconciler) ensureAppliedWork(ctx context.Context, work *fleetv1beta1.Work) (*fleetv1beta1.AppliedWork, error) {
553553
workRef := klog.KObj(work)
554-
appliedWork := &fleetv1beta1.AppliedWork{}
555-
hasFinalizer := false
556-
if controllerutil.ContainsFinalizer(work, fleetv1beta1.WorkFinalizer) {
557-
hasFinalizer = true
558-
err := r.spokeClient.Get(ctx, types.NamespacedName{Name: work.Name}, appliedWork)
559-
switch {
560-
case apierrors.IsNotFound(err):
561-
klog.ErrorS(err, "AppliedWork finalizer resource does not exist even with the finalizer, it will be recreated", "appliedWork", workRef.Name)
562-
case err != nil:
563-
klog.ErrorS(err, "Failed to retrieve the appliedWork ", "appliedWork", workRef.Name)
564-
return nil, controller.NewAPIServerError(true, err)
565-
default:
566-
return appliedWork, nil
554+
555+
// Add a finalizer to the Work object.
556+
if !controllerutil.ContainsFinalizer(work, fleetv1beta1.WorkFinalizer) {
557+
work.Finalizers = append(work.Finalizers, fleetv1beta1.WorkFinalizer)
558+
559+
if err := r.hubClient.Update(ctx, work, &client.UpdateOptions{}); err != nil {
560+
klog.ErrorS(err, "Failed to add the cleanup finalizer to the work", "work", workRef)
561+
return nil, controller.NewAPIServerError(false, err)
567562
}
563+
klog.V(2).InfoS("Added the cleanup finalizer to the Work object", "work", workRef)
564+
}
565+
566+
// Check if an AppliedWork object already exists for the Work object.
567+
//
568+
// Since we only create an AppliedWork object after adding the finalizer to the Work object,
569+
// usually it is safe for us to assume that if the finalizer is absent, the AppliedWork object should
570+
// not exist. This is not the case with the work applier though, as the controller features a
571+
// Leave method that will strip all Work objects off their finalizers, which is called when the
572+
// member cluster leaves the fleet. If the member cluster chooses to re-join the fleet, the controller
573+
// will see a Work object with no finalizer but with an AppliedWork object. Because of this, here we always
574+
// check for the existence of the AppliedWork object, with or without the finalizer.
575+
appliedWork := &fleetv1beta1.AppliedWork{}
576+
err := r.spokeClient.Get(ctx, types.NamespacedName{Name: work.Name}, appliedWork)
577+
switch {
578+
case err == nil:
579+
// The AppliedWork already exists; no further action is needed.
580+
klog.V(2).InfoS("Found an AppliedWork for the Work object", "work", workRef, "appliedWork", klog.KObj(appliedWork))
581+
return appliedWork, nil
582+
case !apierrors.IsNotFound(err):
583+
klog.ErrorS(err, "Failed to retrieve the appliedWork object", "appliedWork", workRef.Name)
584+
return nil, controller.NewAPIServerError(true, err)
568585
}
569586

570-
// we create the appliedWork before setting the finalizer, so it should always exist unless it's deleted behind our back
587+
// The AppliedWork object does not exist; create one.
571588
appliedWork = &fleetv1beta1.AppliedWork{
572589
ObjectMeta: metav1.ObjectMeta{
573590
Name: work.Name,
@@ -577,20 +594,14 @@ func (r *Reconciler) ensureAppliedWork(ctx context.Context, work *fleetv1beta1.W
577594
WorkNamespace: work.Namespace,
578595
},
579596
}
580-
if err := r.spokeClient.Create(ctx, appliedWork); err != nil && !apierrors.IsAlreadyExists(err) {
581-
klog.ErrorS(err, "AppliedWork create failed", "appliedWork", workRef.Name)
597+
if err := r.spokeClient.Create(ctx, appliedWork); err != nil {
598+
// Note: the controller must retry on AppliedWork AlreadyExists errors; otherwise the
599+
// controller will run the reconciliation loop with an AppliedWork that has no UID,
600+
// which might lead to takeover failures in later steps.
601+
klog.ErrorS(err, "Failed to create an AppliedWork object for the Work object", "appliedWork", klog.KObj(appliedWork), "work", workRef)
582602
return nil, controller.NewAPIServerError(false, err)
583603
}
584-
if !hasFinalizer {
585-
klog.InfoS("Add the finalizer to the work", "work", workRef)
586-
work.Finalizers = append(work.Finalizers, fleetv1beta1.WorkFinalizer)
587-
588-
if err := r.hubClient.Update(ctx, work, &client.UpdateOptions{}); err != nil {
589-
klog.ErrorS(err, "Failed to add the finalizer to the work", "work", workRef)
590-
return nil, controller.NewAPIServerError(false, err)
591-
}
592-
}
593-
klog.InfoS("Recreated the appliedWork resource", "appliedWork", workRef.Name)
604+
klog.V(2).InfoS("Created an AppliedWork for the Work object", "work", workRef, "appliedWork", klog.KObj(appliedWork))
594605
return appliedWork, nil
595606
}
596607

pkg/controllers/workapplier/controller_test.go

Lines changed: 207 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ limitations under the License.
1717
package workapplier
1818

1919
import (
20+
"context"
2021
"fmt"
2122
"log"
2223
"os"
@@ -32,8 +33,10 @@ import (
3233
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
3334
"k8s.io/apimachinery/pkg/runtime"
3435
"k8s.io/apimachinery/pkg/runtime/schema"
36+
"k8s.io/apimachinery/pkg/types"
3537
"k8s.io/client-go/kubernetes/scheme"
3638
"k8s.io/utils/ptr"
39+
"sigs.k8s.io/controller-runtime/pkg/client/fake"
3740

3841
fleetv1beta1 "github.com/kubefleet-dev/kubefleet/apis/placement/v1beta1"
3942
)
@@ -208,7 +211,8 @@ var (
208211
)
209212

210213
var (
211-
ignoreFieldTypeMetaInNamespace = cmpopts.IgnoreFields(corev1.Namespace{}, "TypeMeta")
214+
ignoreFieldTypeMetaInNamespace = cmpopts.IgnoreFields(corev1.Namespace{}, "TypeMeta")
215+
ignoreFieldObjectMetaresourceVersion = cmpopts.IgnoreFields(metav1.ObjectMeta{}, "ResourceVersion")
212216

213217
lessFuncAppliedResourceMeta = func(i, j fleetv1beta1.AppliedResourceMeta) bool {
214218
iStr := fmt.Sprintf("%s/%s/%s/%s/%s", i.Group, i.Version, i.Kind, i.Namespace, i.Name)
@@ -262,6 +266,14 @@ func TestMain(m *testing.M) {
262266
os.Exit(m.Run())
263267
}
264268

269+
func fakeClientScheme(t *testing.T) *runtime.Scheme {
270+
scheme := runtime.NewScheme()
271+
if err := fleetv1beta1.AddToScheme(scheme); err != nil {
272+
t.Fatalf("Failed to add placement v1beta1 scheme: %v", err)
273+
}
274+
return scheme
275+
}
276+
265277
func initializeVariables() {
266278
var err error
267279

@@ -332,3 +344,197 @@ func TestPrepareManifestProcessingBundles(t *testing.T) {
332344
t.Errorf("prepareManifestProcessingBundles() mismatches (-got +want):\n%s", diff)
333345
}
334346
}
347+
348+
// TestEnsureAppliedWork tests the ensureAppliedWork method.
349+
func TestEnsureAppliedWork(t *testing.T) {
350+
ctx := context.Background()
351+
352+
fakeUID := types.UID("foo")
353+
testCases := []struct {
354+
name string
355+
work *fleetv1beta1.Work
356+
appliedWork *fleetv1beta1.AppliedWork
357+
wantWork *fleetv1beta1.Work
358+
wantAppliedWork *fleetv1beta1.AppliedWork
359+
}{
360+
{
361+
name: "with work cleanup finalizer present, but no corresponding AppliedWork exists",
362+
work: &fleetv1beta1.Work{
363+
ObjectMeta: metav1.ObjectMeta{
364+
Name: workName,
365+
Namespace: memberReservedNSName1,
366+
Finalizers: []string{
367+
fleetv1beta1.WorkFinalizer,
368+
},
369+
},
370+
},
371+
wantWork: &fleetv1beta1.Work{
372+
ObjectMeta: metav1.ObjectMeta{
373+
Name: workName,
374+
Namespace: memberReservedNSName1,
375+
Finalizers: []string{
376+
fleetv1beta1.WorkFinalizer,
377+
},
378+
},
379+
},
380+
wantAppliedWork: &fleetv1beta1.AppliedWork{
381+
ObjectMeta: metav1.ObjectMeta{
382+
Name: workName,
383+
},
384+
Spec: fleetv1beta1.AppliedWorkSpec{
385+
WorkName: workName,
386+
WorkNamespace: memberReservedNSName1,
387+
},
388+
},
389+
},
390+
{
391+
name: "with work cleanup finalizer present, and corresponding AppliedWork exists",
392+
work: &fleetv1beta1.Work{
393+
ObjectMeta: metav1.ObjectMeta{
394+
Name: workName,
395+
Namespace: memberReservedNSName1,
396+
Finalizers: []string{
397+
fleetv1beta1.WorkFinalizer,
398+
},
399+
},
400+
},
401+
appliedWork: &fleetv1beta1.AppliedWork{
402+
ObjectMeta: metav1.ObjectMeta{
403+
Name: workName,
404+
// Add the UID field to track if the method returns the existing object.
405+
UID: fakeUID,
406+
},
407+
Spec: fleetv1beta1.AppliedWorkSpec{
408+
WorkName: workName,
409+
WorkNamespace: memberReservedNSName1,
410+
},
411+
},
412+
wantWork: &fleetv1beta1.Work{
413+
ObjectMeta: metav1.ObjectMeta{
414+
Name: workName,
415+
Namespace: memberReservedNSName1,
416+
Finalizers: []string{
417+
fleetv1beta1.WorkFinalizer,
418+
},
419+
},
420+
},
421+
wantAppliedWork: &fleetv1beta1.AppliedWork{
422+
ObjectMeta: metav1.ObjectMeta{
423+
Name: workName,
424+
UID: fakeUID,
425+
},
426+
Spec: fleetv1beta1.AppliedWorkSpec{
427+
WorkName: workName,
428+
WorkNamespace: memberReservedNSName1,
429+
},
430+
},
431+
},
432+
{
433+
name: "without work cleanup finalizer, but corresponding AppliedWork exists",
434+
work: &fleetv1beta1.Work{
435+
ObjectMeta: metav1.ObjectMeta{
436+
Name: workName,
437+
Namespace: memberReservedNSName1,
438+
},
439+
},
440+
appliedWork: &fleetv1beta1.AppliedWork{
441+
ObjectMeta: metav1.ObjectMeta{
442+
Name: workName,
443+
// Add the UID field to track if the method returns the existing object.
444+
UID: fakeUID,
445+
},
446+
Spec: fleetv1beta1.AppliedWorkSpec{
447+
WorkName: workName,
448+
WorkNamespace: memberReservedNSName1,
449+
},
450+
},
451+
wantWork: &fleetv1beta1.Work{
452+
ObjectMeta: metav1.ObjectMeta{
453+
Name: workName,
454+
Namespace: memberReservedNSName1,
455+
Finalizers: []string{
456+
fleetv1beta1.WorkFinalizer,
457+
},
458+
},
459+
},
460+
wantAppliedWork: &fleetv1beta1.AppliedWork{
461+
ObjectMeta: metav1.ObjectMeta{
462+
Name: workName,
463+
UID: fakeUID,
464+
},
465+
Spec: fleetv1beta1.AppliedWorkSpec{
466+
WorkName: workName,
467+
WorkNamespace: memberReservedNSName1,
468+
},
469+
},
470+
},
471+
{
472+
name: "without work cleanup finalizer, and no corresponding AppliedWork exists",
473+
work: &fleetv1beta1.Work{
474+
ObjectMeta: metav1.ObjectMeta{
475+
Name: workName,
476+
Namespace: memberReservedNSName1,
477+
},
478+
},
479+
wantWork: &fleetv1beta1.Work{
480+
ObjectMeta: metav1.ObjectMeta{
481+
Name: workName,
482+
Namespace: memberReservedNSName1,
483+
Finalizers: []string{
484+
fleetv1beta1.WorkFinalizer,
485+
},
486+
},
487+
},
488+
wantAppliedWork: &fleetv1beta1.AppliedWork{
489+
ObjectMeta: metav1.ObjectMeta{
490+
Name: workName,
491+
},
492+
Spec: fleetv1beta1.AppliedWorkSpec{
493+
WorkName: workName,
494+
WorkNamespace: memberReservedNSName1,
495+
},
496+
},
497+
},
498+
}
499+
500+
for _, tc := range testCases {
501+
t.Run(tc.name, func(t *testing.T) {
502+
hubClientScheme := fakeClientScheme(t)
503+
fakeHubClient := fake.NewClientBuilder().
504+
WithScheme(hubClientScheme).
505+
WithObjects(tc.work).
506+
Build()
507+
508+
memberClientScheme := fakeClientScheme(t)
509+
fakeMemberClientBuilder := fake.NewClientBuilder().WithScheme(memberClientScheme)
510+
if tc.appliedWork != nil {
511+
fakeMemberClientBuilder = fakeMemberClientBuilder.WithObjects(tc.appliedWork)
512+
}
513+
fakeMemberClient := fakeMemberClientBuilder.Build()
514+
515+
r := &Reconciler{
516+
hubClient: fakeHubClient,
517+
spokeClient: fakeMemberClient,
518+
}
519+
520+
gotAppliedWork, err := r.ensureAppliedWork(ctx, tc.work)
521+
if err != nil {
522+
t.Fatalf("ensureAppliedWork() = %v, want no error", err)
523+
}
524+
525+
// Verify the Work object.
526+
gotWork := &fleetv1beta1.Work{}
527+
if err := fakeHubClient.Get(ctx, types.NamespacedName{Name: tc.work.Name, Namespace: tc.work.Namespace}, gotWork); err != nil {
528+
t.Fatalf("failed to get Work object from fake hub client: %v", err)
529+
}
530+
if diff := cmp.Diff(gotWork, tc.wantWork, ignoreFieldObjectMetaresourceVersion); diff != "" {
531+
t.Errorf("Work objects diff (-got +want):\n%s", diff)
532+
}
533+
534+
// Verify the AppliedWork object.
535+
if diff := cmp.Diff(gotAppliedWork, tc.wantAppliedWork, ignoreFieldObjectMetaresourceVersion); diff != "" {
536+
t.Errorf("AppliedWork objects diff (-got +want):\n%s", diff)
537+
}
538+
})
539+
}
540+
}

0 commit comments

Comments
 (0)