diff --git a/pkg/health/health.go b/pkg/health/health.go index a5579aea0..c615deea9 100644 --- a/pkg/health/health.go +++ b/pkg/health/health.go @@ -4,6 +4,7 @@ import ( "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime/schema" + "github.com/argoproj/gitops-engine/pkg/sync/hook" "github.com/argoproj/gitops-engine/pkg/utils/kube" ) @@ -65,7 +66,7 @@ func IsWorse(current, new HealthStatusCode) bool { // GetResourceHealth returns the health of a k8s resource func GetResourceHealth(obj *unstructured.Unstructured, healthOverride HealthOverride) (health *HealthStatus, err error) { - if obj.GetDeletionTimestamp() != nil { + if obj.GetDeletionTimestamp() != nil && !hook.HasHookFinalizer(obj) { return &HealthStatus{ Status: HealthStatusProgressing, Message: "Pending deletion", diff --git a/pkg/sync/hook/hook.go b/pkg/sync/hook/hook.go index 7e0c38752..66dfc26e5 100644 --- a/pkg/sync/hook/hook.go +++ b/pkg/sync/hook/hook.go @@ -8,6 +8,21 @@ import ( resourceutil "github.com/argoproj/gitops-engine/pkg/sync/resource" ) +const ( + // HookFinalizer is the finalizer added to hooks to ensure they are deleted only after the sync phase is completed. + HookFinalizer = "argocd.argoproj.io/hook-finalizer" +) + +func HasHookFinalizer(obj *unstructured.Unstructured) bool { + finalizers := obj.GetFinalizers() + for _, finalizer := range finalizers { + if finalizer == HookFinalizer { + return true + } + } + return false +} + func IsHook(obj *unstructured.Unstructured) bool { _, ok := obj.GetAnnotations()[common.AnnotationKeyHook] if ok { diff --git a/pkg/sync/sync_context.go b/pkg/sync/sync_context.go index 9d024ed12..7d43899c9 100644 --- a/pkg/sync/sync_context.go +++ b/pkg/sync/sync_context.go @@ -295,12 +295,12 @@ const ( crdReadinessTimeout = time.Duration(3) * time.Second ) -// getOperationPhase returns a hook status from an _live_ unstructured object -func (sc *syncContext) getOperationPhase(hook *unstructured.Unstructured) (common.OperationPhase, string, error) { +// getOperationPhase returns a health status from a _live_ unstructured object +func (sc *syncContext) getOperationPhase(obj *unstructured.Unstructured) (common.OperationPhase, string, error) { phase := common.OperationSucceeded - message := hook.GetName() + " created" + message := obj.GetName() + " created" - resHealth, err := health.GetResourceHealth(hook, sc.healthOverride) + resHealth, err := health.GetResourceHealth(obj, sc.healthOverride) if err != nil { return "", "", err } @@ -475,6 +475,15 @@ func (sc *syncContext) Sync() { return } + hooksCompleted := tasks.Filter(func(task *syncTask) bool { + return task.isHook() && task.completed() + }) + for _, task := range hooksCompleted { + if err := sc.removeHookFinalizer(task); err != nil { + sc.setResourceResult(task, task.syncStatus, common.OperationError, fmt.Sprintf("Failed to remove hook finalizer: %v", err)) + } + } + // collect all completed hooks which have appropriate delete policy hooksPendingDeletionSuccessful := tasks.Filter(func(task *syncTask) bool { return task.isHook() && task.liveObj != nil && !task.running() && task.deleteOnPhaseSuccessful() @@ -576,6 +585,64 @@ func (sc *syncContext) filterOutOfSyncTasks(tasks syncTasks) syncTasks { }) } +func (sc *syncContext) removeHookFinalizer(task *syncTask) error { + if task.liveObj == nil { + return nil + } + removeFinalizerMutation := func(obj *unstructured.Unstructured) bool { + finalizers := obj.GetFinalizers() + for i, finalizer := range finalizers { + if finalizer == hook.HookFinalizer { + obj.SetFinalizers(append(finalizers[:i], finalizers[i+1:]...)) + return true + } + } + return false + } + + // The cached live object may be stale in the controller cache, and the actual object may have been updated in the meantime, + // and Kubernetes API will return a conflict error on the Update call. + // In that case, we need to get the latest version of the object and retry the update. + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + mutated := removeFinalizerMutation(task.liveObj) + if !mutated { + return nil + } + + updateErr := sc.updateResource(task) + if apierrors.IsConflict(updateErr) { + sc.log.WithValues("task", task).V(1).Info("Retrying hook finalizer removal due to conflict on update") + resIf, err := sc.getResourceIf(task, "get") + if err != nil { + return err + } + liveObj, err := resIf.Get(context.TODO(), task.liveObj.GetName(), metav1.GetOptions{}) + if apierrors.IsNotFound(err) { + sc.log.WithValues("task", task).V(1).Info("Resource is already deleted") + return nil + } else if err != nil { + return err + } + task.liveObj = liveObj + } else if apierrors.IsNotFound(updateErr) { + // If the resource is already deleted, it is a no-op + sc.log.WithValues("task", task).V(1).Info("Resource is already deleted") + return nil + } + return updateErr + }) +} + +func (sc *syncContext) updateResource(task *syncTask) error { + sc.log.WithValues("task", task).V(1).Info("Updating resource") + resIf, err := sc.getResourceIf(task, "update") + if err != nil { + return err + } + _, err = resIf.Update(context.TODO(), task.liveObj, metav1.UpdateOptions{}) + return err +} + func (sc *syncContext) deleteHooks(hooksPendingDeletion syncTasks) { for _, task := range hooksPendingDeletion { err := sc.deleteResource(task) @@ -597,8 +664,8 @@ func (sc *syncContext) GetState() (common.OperationPhase, string, []common.Resou } func (sc *syncContext) setOperationFailed(syncFailTasks, syncFailedTasks syncTasks, message string) { - errorMessageFactory := func(_ []*syncTask, message string) string { - messages := syncFailedTasks.Map(func(task *syncTask) string { + errorMessageFactory := func(tasks syncTasks, message string) string { + messages := tasks.Map(func(task *syncTask) string { return task.message }) if len(messages) > 0 { @@ -619,7 +686,9 @@ func (sc *syncContext) setOperationFailed(syncFailTasks, syncFailedTasks syncTas // the phase, so we make sure we have at least one more sync sc.log.WithValues("syncFailTasks", syncFailTasks).V(1).Info("Running sync fail tasks") if sc.runTasks(syncFailTasks, false) == failed { - sc.setOperationPhase(common.OperationFailed, errorMessage) + failedSyncFailTasks := syncFailTasks.Filter(func(t *syncTask) bool { return t.syncStatus == common.ResultCodeSyncFailed }) + syncFailTasksMessage := errorMessageFactory(failedSyncFailTasks, "one or more SyncFail hooks failed") + sc.setOperationPhase(common.OperationFailed, fmt.Sprintf("%s\n%s", errorMessage, syncFailTasksMessage)) } } else { sc.setOperationPhase(common.OperationFailed, errorMessage) @@ -679,7 +748,9 @@ func (sc *syncContext) getSyncTasks() (_ syncTasks, successful bool) { generateName := obj.GetGenerateName() targetObj.SetName(fmt.Sprintf("%s%s", generateName, postfix)) } - + if !hook.HasHookFinalizer(targetObj) { + targetObj.SetFinalizers(append(targetObj.GetFinalizers(), hook.HookFinalizer)) + } hookTasks = append(hookTasks, &syncTask{phase: phase, targetObj: targetObj}) } } @@ -1085,6 +1156,11 @@ func (sc *syncContext) Terminate() { if !task.isHook() || task.liveObj == nil { continue } + if err := sc.removeHookFinalizer(task); err != nil { + sc.setResourceResult(task, task.syncStatus, common.OperationError, fmt.Sprintf("Failed to remove hook finalizer: %v", err)) + terminateSuccessful = false + continue + } phase, msg, err := sc.getOperationPhase(task.liveObj) if err != nil { sc.setOperationPhase(common.OperationError, fmt.Sprintf("Failed to get hook health: %v", err)) @@ -1092,7 +1168,7 @@ func (sc *syncContext) Terminate() { } if phase == common.OperationRunning { err := sc.deleteResource(task) - if err != nil { + if err != nil && !apierrors.IsNotFound(err) { sc.setResourceResult(task, "", common.OperationFailed, fmt.Sprintf("Failed to delete: %v", err)) terminateSuccessful = false } else { diff --git a/pkg/sync/sync_context_test.go b/pkg/sync/sync_context_test.go index 1ec904ad8..21f9730bb 100644 --- a/pkg/sync/sync_context_test.go +++ b/pkg/sync/sync_context_test.go @@ -29,6 +29,7 @@ import ( "github.com/argoproj/gitops-engine/pkg/diff" "github.com/argoproj/gitops-engine/pkg/health" synccommon "github.com/argoproj/gitops-engine/pkg/sync/common" + "github.com/argoproj/gitops-engine/pkg/sync/hook" "github.com/argoproj/gitops-engine/pkg/utils/kube" "github.com/argoproj/gitops-engine/pkg/utils/kube/kubetest" testingutils "github.com/argoproj/gitops-engine/pkg/utils/testing" @@ -1285,15 +1286,19 @@ func (r resourceNameHealthOverride) GetResourceHealth(obj *unstructured.Unstruct } func TestRunSync_HooksNotDeletedIfPhaseNotCompleted(t *testing.T) { - completedHook := newHook(synccommon.HookTypePreSync) - completedHook.SetName("completed-hook") - completedHook.SetNamespace(testingutils.FakeArgoCDNamespace) - _ = testingutils.Annotate(completedHook, synccommon.AnnotationKeyHookDeletePolicy, "HookSucceeded") - - inProgressHook := newHook(synccommon.HookTypePreSync) - inProgressHook.SetNamespace(testingutils.FakeArgoCDNamespace) - inProgressHook.SetName("in-progress-hook") - _ = testingutils.Annotate(inProgressHook, synccommon.AnnotationKeyHookDeletePolicy, "HookSucceeded") + hook1 := newHook(synccommon.HookTypePreSync) + hook1.SetName("completed-hook") + hook1.SetNamespace(testingutils.FakeArgoCDNamespace) + _ = testingutils.Annotate(hook1, synccommon.AnnotationKeyHookDeletePolicy, string(synccommon.HookDeletePolicyHookSucceeded)) + completedHook := hook1.DeepCopy() + completedHook.SetFinalizers(append(completedHook.GetFinalizers(), hook.HookFinalizer)) + + hook2 := newHook(synccommon.HookTypePreSync) + hook2.SetNamespace(testingutils.FakeArgoCDNamespace) + hook2.SetName("in-progress-hook") + _ = testingutils.Annotate(hook2, synccommon.AnnotationKeyHookDeletePolicy, string(synccommon.HookDeletePolicyHookSucceeded)) + inProgressHook := hook2.DeepCopy() + inProgressHook.SetFinalizers(append(inProgressHook.GetFinalizers(), hook.HookFinalizer)) syncCtx := newTestSyncCtx(nil, WithHealthOverride(resourceNameHealthOverride(map[string]health.HealthStatusCode{ @@ -1312,6 +1317,12 @@ func TestRunSync_HooksNotDeletedIfPhaseNotCompleted(t *testing.T) { )) fakeDynamicClient := fake.NewSimpleDynamicClient(runtime.NewScheme()) syncCtx.dynamicIf = fakeDynamicClient + updatedCount := 0 + fakeDynamicClient.PrependReactor("update", "*", func(_ testcore.Action) (handled bool, ret runtime.Object, err error) { + // Removing the finalizers + updatedCount++ + return true, nil, nil + }) deletedCount := 0 fakeDynamicClient.PrependReactor("delete", "*", func(_ testcore.Action) (handled bool, ret runtime.Object, err error) { deletedCount++ @@ -1321,7 +1332,7 @@ func TestRunSync_HooksNotDeletedIfPhaseNotCompleted(t *testing.T) { Live: []*unstructured.Unstructured{completedHook, inProgressHook}, Target: []*unstructured.Unstructured{nil, nil}, }) - syncCtx.hooks = []*unstructured.Unstructured{completedHook, inProgressHook} + syncCtx.hooks = []*unstructured.Unstructured{hook1, hook2} syncCtx.kubectl = &kubetest.MockKubectlCmd{ Commands: map[string]kubetest.KubectlOutput{}, @@ -1330,19 +1341,24 @@ func TestRunSync_HooksNotDeletedIfPhaseNotCompleted(t *testing.T) { syncCtx.Sync() assert.Equal(t, synccommon.OperationRunning, syncCtx.phase) + assert.Equal(t, 0, updatedCount) assert.Equal(t, 0, deletedCount) } func TestRunSync_HooksDeletedAfterPhaseCompleted(t *testing.T) { - completedHook1 := newHook(synccommon.HookTypePreSync) - completedHook1.SetName("completed-hook1") - completedHook1.SetNamespace(testingutils.FakeArgoCDNamespace) - _ = testingutils.Annotate(completedHook1, synccommon.AnnotationKeyHookDeletePolicy, "HookSucceeded") - - completedHook2 := newHook(synccommon.HookTypePreSync) - completedHook2.SetNamespace(testingutils.FakeArgoCDNamespace) - completedHook2.SetName("completed-hook2") - _ = testingutils.Annotate(completedHook2, synccommon.AnnotationKeyHookDeletePolicy, "HookSucceeded") + hook1 := newHook(synccommon.HookTypePreSync) + hook1.SetName("completed-hook1") + hook1.SetNamespace(testingutils.FakeArgoCDNamespace) + _ = testingutils.Annotate(hook1, synccommon.AnnotationKeyHookDeletePolicy, string(synccommon.HookDeletePolicyHookSucceeded)) + completedHook1 := hook1.DeepCopy() + completedHook1.SetFinalizers(append(completedHook1.GetFinalizers(), hook.HookFinalizer)) + + hook2 := newHook(synccommon.HookTypePreSync) + hook2.SetNamespace(testingutils.FakeArgoCDNamespace) + hook2.SetName("completed-hook2") + _ = testingutils.Annotate(hook2, synccommon.AnnotationKeyHookDeletePolicy, string(synccommon.HookDeletePolicyHookSucceeded)) + completedHook2 := hook2.DeepCopy() + completedHook2.SetFinalizers(append(completedHook1.GetFinalizers(), hook.HookFinalizer)) syncCtx := newTestSyncCtx(nil, WithInitialState(synccommon.OperationRunning, "", []synccommon.ResourceSyncResult{{ @@ -1358,6 +1374,12 @@ func TestRunSync_HooksDeletedAfterPhaseCompleted(t *testing.T) { )) fakeDynamicClient := fake.NewSimpleDynamicClient(runtime.NewScheme()) syncCtx.dynamicIf = fakeDynamicClient + updatedCount := 0 + fakeDynamicClient.PrependReactor("update", "*", func(_ testcore.Action) (handled bool, ret runtime.Object, err error) { + // Removing the finalizers + updatedCount++ + return true, nil, nil + }) deletedCount := 0 fakeDynamicClient.PrependReactor("delete", "*", func(_ testcore.Action) (handled bool, ret runtime.Object, err error) { deletedCount++ @@ -1367,7 +1389,7 @@ func TestRunSync_HooksDeletedAfterPhaseCompleted(t *testing.T) { Live: []*unstructured.Unstructured{completedHook1, completedHook2}, Target: []*unstructured.Unstructured{nil, nil}, }) - syncCtx.hooks = []*unstructured.Unstructured{completedHook1, completedHook2} + syncCtx.hooks = []*unstructured.Unstructured{hook1, hook2} syncCtx.kubectl = &kubetest.MockKubectlCmd{ Commands: map[string]kubetest.KubectlOutput{}, @@ -1376,19 +1398,24 @@ func TestRunSync_HooksDeletedAfterPhaseCompleted(t *testing.T) { syncCtx.Sync() assert.Equal(t, synccommon.OperationSucceeded, syncCtx.phase) + assert.Equal(t, 2, updatedCount) assert.Equal(t, 2, deletedCount) } func TestRunSync_HooksDeletedAfterPhaseCompletedFailed(t *testing.T) { - completedHook1 := newHook(synccommon.HookTypeSync) - completedHook1.SetName("completed-hook1") - completedHook1.SetNamespace(testingutils.FakeArgoCDNamespace) - _ = testingutils.Annotate(completedHook1, synccommon.AnnotationKeyHookDeletePolicy, "HookFailed") - - completedHook2 := newHook(synccommon.HookTypeSync) - completedHook2.SetNamespace(testingutils.FakeArgoCDNamespace) - completedHook2.SetName("completed-hook2") - _ = testingutils.Annotate(completedHook2, synccommon.AnnotationKeyHookDeletePolicy, "HookFailed") + hook1 := newHook(synccommon.HookTypeSync) + hook1.SetName("completed-hook1") + hook1.SetNamespace(testingutils.FakeArgoCDNamespace) + _ = testingutils.Annotate(hook1, synccommon.AnnotationKeyHookDeletePolicy, string(synccommon.HookDeletePolicyHookFailed)) + completedHook1 := hook1.DeepCopy() + completedHook1.SetFinalizers(append(completedHook1.GetFinalizers(), hook.HookFinalizer)) + + hook2 := newHook(synccommon.HookTypeSync) + hook2.SetNamespace(testingutils.FakeArgoCDNamespace) + hook2.SetName("completed-hook2") + _ = testingutils.Annotate(hook2, synccommon.AnnotationKeyHookDeletePolicy, string(synccommon.HookDeletePolicyHookFailed)) + completedHook2 := hook2.DeepCopy() + completedHook2.SetFinalizers(append(completedHook1.GetFinalizers(), hook.HookFinalizer)) syncCtx := newTestSyncCtx(nil, WithInitialState(synccommon.OperationRunning, "", []synccommon.ResourceSyncResult{{ @@ -1404,6 +1431,12 @@ func TestRunSync_HooksDeletedAfterPhaseCompletedFailed(t *testing.T) { )) fakeDynamicClient := fake.NewSimpleDynamicClient(runtime.NewScheme()) syncCtx.dynamicIf = fakeDynamicClient + updatedCount := 0 + fakeDynamicClient.PrependReactor("update", "*", func(_ testcore.Action) (handled bool, ret runtime.Object, err error) { + // Removing the finalizers + updatedCount++ + return true, nil, nil + }) deletedCount := 0 fakeDynamicClient.PrependReactor("delete", "*", func(_ testcore.Action) (handled bool, ret runtime.Object, err error) { deletedCount++ @@ -1413,7 +1446,7 @@ func TestRunSync_HooksDeletedAfterPhaseCompletedFailed(t *testing.T) { Live: []*unstructured.Unstructured{completedHook1, completedHook2}, Target: []*unstructured.Unstructured{nil, nil}, }) - syncCtx.hooks = []*unstructured.Unstructured{completedHook1, completedHook2} + syncCtx.hooks = []*unstructured.Unstructured{hook1, hook2} syncCtx.kubectl = &kubetest.MockKubectlCmd{ Commands: map[string]kubetest.KubectlOutput{}, @@ -1422,6 +1455,7 @@ func TestRunSync_HooksDeletedAfterPhaseCompletedFailed(t *testing.T) { syncCtx.Sync() assert.Equal(t, synccommon.OperationFailed, syncCtx.phase) + assert.Equal(t, 2, updatedCount) assert.Equal(t, 2, deletedCount) }