diff --git a/cmd/postgres-operator/main.go b/cmd/postgres-operator/main.go index 8ec39c398e..6cae8a70e8 100644 --- a/cmd/postgres-operator/main.go +++ b/cmd/postgres-operator/main.go @@ -155,6 +155,7 @@ func addControllersToManager(ctx context.Context, mgr manager.Manager) error { r := &postgrescluster.Reconciler{ Client: mgr.GetClient(), + APIReader: mgr.GetAPIReader(), // K8SPG-992 Scheme: mgr.GetScheme(), Owner: postgrescluster.ControllerName, Recorder: mgr.GetEventRecorderFor(postgrescluster.ControllerName), diff --git a/internal/controller/postgrescluster/controller.go b/internal/controller/postgrescluster/controller.go index a816d08b34..3a56378162 100644 --- a/internal/controller/postgrescluster/controller.go +++ b/internal/controller/postgrescluster/controller.go @@ -70,7 +70,9 @@ const ( // Reconciler holds resources for the PostgresCluster reconciler type Reconciler struct { - Client client.Client + Client client.Client + // K8SPG-992: APIReader reads directly from the API server, bypassing the cache + APIReader client.Reader Scheme *k8sruntime.Scheme DiscoveryClient *discovery.DiscoveryClient IsOpenShift bool @@ -87,6 +89,13 @@ type Reconciler struct { newPGBouncerAdmin func(opts pgbruntime.AdminClientOptions) (pgbruntime.AdminClient, error) } +func (r *Reconciler) apiReader() client.Reader { + if r.APIReader != nil { + return r.APIReader + } + return r.Client +} + // +kubebuilder:rbac:groups="",resources="events",verbs={create,patch} // +kubebuilder:rbac:groups="upstream.pgv2.percona.com",resources="postgresclusters",verbs={get,list,watch} // +kubebuilder:rbac:groups="upstream.pgv2.percona.com",resources="postgresclusters/status",verbs={patch} diff --git a/internal/controller/postgrescluster/pgbackrest.go b/internal/controller/postgrescluster/pgbackrest.go index d94d76b6fd..ab3b8a97b7 100644 --- a/internal/controller/postgrescluster/pgbackrest.go +++ b/internal/controller/postgrescluster/pgbackrest.go @@ -2568,6 +2568,29 @@ func (r *Reconciler) reconcileDedicatedRepoHost(ctx context.Context, return repoHost, nil } +// K8SPG-992: manualBackupJobNeeded returns true if a backup job should be created +// for the given PerconaPGBackup. +func manualBackupJobNeeded(ctx context.Context, cl client.Reader, cluster *v1beta1.PostgresCluster, backupName string) (bool, error) { + pgBackup := new(v2.PerconaPGBackup) + err := cl.Get(ctx, client.ObjectKey{Namespace: cluster.Namespace, Name: backupName}, pgBackup) + if apierrors.IsNotFound(err) { + return false, nil + } + if err != nil { + return false, errors.WithStack(err) + } + + if pgBackup.Spec.PGCluster != cluster.Name { + return false, nil + } + + if !pgBackup.DeletionTimestamp.IsZero() || pgBackup.Status.State.IsTerminal() { + return false, nil + } + + return pgBackup.Status.JobName == "", nil +} + // +kubebuilder:rbac:groups="batch",resources="jobs",verbs={create,patch,delete} // reconcileManualBackup is responsible for reconciling pgBackRest backups that are initiated @@ -2697,6 +2720,17 @@ func (r *Reconciler) reconcileManualBackup(ctx context.Context, return nil } + // K8SPG-992: we need this check to eliminate possibility of creating a second backup Job for the same PerconaPGBackup. + if currentBackupJob == nil { + needsJob, err := manualBackupJobNeeded(ctx, r.apiReader(), postgresCluster, manualAnnotation) + if err != nil { + return err + } + if !needsJob { + return nil + } + } + // determine if the dedicated repository host is ready using the repo host ready // condition, and return if not repoCondition := meta.FindStatusCondition(postgresCluster.Status.Conditions, ConditionRepoHostReady) diff --git a/internal/controller/postgrescluster/pgbackrest_test.go b/internal/controller/postgrescluster/pgbackrest_test.go index 0daf32be78..f6c5e2436f 100644 --- a/internal/controller/postgrescluster/pgbackrest_test.go +++ b/internal/controller/postgrescluster/pgbackrest_test.go @@ -31,6 +31,7 @@ import ( "k8s.io/apimachinery/pkg/util/rand" "k8s.io/apimachinery/pkg/util/wait" "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" "sigs.k8s.io/controller-runtime/pkg/manager" "sigs.k8s.io/controller-runtime/pkg/reconcile" @@ -1134,6 +1135,7 @@ func TestReconcileManualBackup(t *testing.T) { Tracer: otel.Tracer(ControllerName), Owner: ControllerName, } + r.APIReader = mgr.GetAPIReader() }) t.Cleanup(func() { teardownManager(cancel, t) }) @@ -1144,9 +1146,11 @@ func TestReconcileManualBackup(t *testing.T) { fakeJob := func(clusterName, repoName string) *batchv1.Job { return &batchv1.Job{ ObjectMeta: metav1.ObjectMeta{ - Name: "manual-backup-" + rand.String(4), - Namespace: ns.GetName(), - Annotations: map[string]string{naming.PGBackRestBackup: defaultBackupId}, + Name: "manual-backup-" + rand.String(4), + Namespace: ns.GetName(), + Annotations: map[string]string{ + naming.PGBackRestBackup: clusterName + "-" + defaultBackupId, + }, Labels: naming.PGBackRestBackupJobLabels(clusterName, repoName, naming.BackupManual), }, @@ -1186,6 +1190,8 @@ func TestReconcileManualBackup(t *testing.T) { status *v1beta1.PostgresClusterStatus // the ID used to populate the "backup" annotation for the test (can be empty) backupId string + // K8SPG-992: whether to skip creating the PerconaPGBackup that owns the backup ID + noPGBackup bool // the manual backup field to define in the postgrescluster spec for the test manual *v1beta1.PGBackRestManualBackup // whether or not the test should expect a Job to be reconciled @@ -1392,6 +1398,23 @@ func TestReconcileManualBackup(t *testing.T) { expectCurrentJobDeletion: false, expectReconcile: true, verifyStaleJobConflict: true, + }, { + testDesc: "no PerconaPGBackup for backup id should not reconcile", + createCurrentJob: false, + noPGBackup: true, + clusterConditions: map[string]metav1.ConditionStatus{ + ConditionRepoHostReady: metav1.ConditionTrue, + ConditionReplicaCreate: metav1.ConditionTrue, + }, + status: &v1beta1.PostgresClusterStatus{ + PGBackRest: &v1beta1.PGBackRestStatus{ + Repos: []v1beta1.RepoStatus{{Name: "repo1", StanzaCreated: true}}, + }, + }, + backupId: backupId, + manual: &v1beta1.PGBackRestManualBackup{RepoName: "repo1"}, + expectCurrentJobDeletion: false, + expectReconcile: false, }, { testDesc: "reconcile new job when in-progress job exists for another id", createCurrentJob: true, @@ -1460,12 +1483,33 @@ func TestReconcileManualBackup(t *testing.T) { ctx := context.Background() + if !tc.noPGBackup { + for _, id := range []string{defaultBackupId, backupId} { + assert.NilError(t, tClient.Create(ctx, &v2.PerconaPGBackup{ + ObjectMeta: metav1.ObjectMeta{ + Name: clusterName + "-" + id, + Namespace: ns.GetName(), + }, + Spec: v2.PerconaPGBackupSpec{ + PGCluster: clusterName, + RepoName: new("repo1"), + }, + })) + } + } + postgresCluster := fakePostgresCluster(clusterName, ns.GetName(), "", dedicated) postgresCluster.Spec.Backups.PGBackRest.Manual = tc.manual - postgresCluster.Annotations = map[string]string{naming.PGBackRestBackup: tc.backupId} + postgresCluster.Annotations = map[string]string{} + if tc.backupId != "" { + postgresCluster.Annotations[naming.PGBackRestBackup] = clusterName + "-" + tc.backupId + } assert.NilError(t, tClient.Create(ctx, postgresCluster)) - postgresCluster.Status = *tc.status + postgresCluster.Status = *tc.status.DeepCopy() + if mb := postgresCluster.Status.PGBackRest.ManualBackup; mb != nil { + mb.ID = clusterName + "-" + mb.ID + } for condition, status := range tc.clusterConditions { meta.SetStatusCondition(&postgresCluster.Status.Conditions, metav1.Condition{ Type: condition, Reason: "testing", Status: status, @@ -4793,3 +4837,86 @@ func TestPgBackRestCACert(t *testing.T) { assert.ErrorContains(t, err, "did not return a CA certificate") }) } + +// K8SPG-992 +func TestManualBackupJobNeeded(t *testing.T) { + backup := func(f func(*v2.PerconaPGBackup)) *v2.PerconaPGBackup { + b := &v2.PerconaPGBackup{ + ObjectMeta: metav1.ObjectMeta{Name: "backup1", Namespace: "postgres-operator"}, + Spec: v2.PerconaPGBackupSpec{ + PGCluster: "hippo", + RepoName: new("repo1"), + }, + } + if f != nil { + f(b) + } + return b + } + + for _, tt := range []struct { + name string + pgBackup *v2.PerconaPGBackup + expected bool + }{{ + name: "no PerconaPGBackup", + expected: false, + }, { + name: "new backup without a Job", + pgBackup: backup(nil), + expected: true, + }, { + name: "starting backup without a Job", + pgBackup: backup(func(b *v2.PerconaPGBackup) { + b.Status.State = v2.BackupStarting + }), + expected: true, + }, { + name: "backup that already had a Job", + pgBackup: backup(func(b *v2.PerconaPGBackup) { + b.Status.State = v2.BackupRunning + b.Status.JobName = "hippo-backup-abcd" + }), + expected: false, + }, { + name: "backup for another cluster", + pgBackup: backup(func(b *v2.PerconaPGBackup) { + b.Spec.PGCluster = "elephant" + }), + expected: false, + }, { + name: "backup being deleted", + pgBackup: backup(func(b *v2.PerconaPGBackup) { + b.DeletionTimestamp = new(metav1.Now()) + b.Finalizers = []string{"internal.percona.com/delete-backup"} + }), + expected: false, + }, { + name: "failed backup", + pgBackup: backup(func(b *v2.PerconaPGBackup) { + b.Status.State = v2.BackupFailed + }), + expected: false, + }, { + name: "succeeded backup", + pgBackup: backup(func(b *v2.PerconaPGBackup) { + b.Status.State = v2.BackupSucceeded + }), + expected: false, + }} { + t.Run(tt.name, func(t *testing.T) { + builder := fake.NewClientBuilder().WithScheme(runtime.Scheme) + if tt.pgBackup != nil { + builder = builder.WithObjects(tt.pgBackup) + } + + cluster := &v1beta1.PostgresCluster{ + ObjectMeta: metav1.ObjectMeta{Name: "hippo", Namespace: "postgres-operator"}, + } + + needed, err := manualBackupJobNeeded(t.Context(), builder.Build(), cluster, "backup1") + assert.NilError(t, err) + assert.Equal(t, needed, tt.expected) + }) + } +} diff --git a/internal/controller/postgrescluster/suite_test.go b/internal/controller/postgrescluster/suite_test.go index 1e0d828ee1..dc368d82f4 100644 --- a/internal/controller/postgrescluster/suite_test.go +++ b/internal/controller/postgrescluster/suite_test.go @@ -14,18 +14,25 @@ import ( . "github.com/onsi/gomega" "k8s.io/apimachinery/pkg/util/version" "k8s.io/client-go/discovery" - - // Google Kubernetes Engine / Google Cloud Platform authentication provider _ "k8s.io/client-go/plugin/pkg/client/auth/gcp" "k8s.io/client-go/rest" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/log" "sigs.k8s.io/controller-runtime/pkg/manager" + "github.com/percona/percona-postgresql-operator/v2/internal/controller/runtime" "github.com/percona/percona-postgresql-operator/v2/internal/logging" "github.com/percona/percona-postgresql-operator/v2/internal/testing/require" + v2 "github.com/percona/percona-postgresql-operator/v2/pkg/apis/pgv2.percona.com/v2" ) +// K8SPG-992 +func init() { + if err := v2.AddToScheme(runtime.Scheme); err != nil { + panic(err) + } +} + var suite struct { Client client.Client Config *rest.Config @@ -63,5 +70,4 @@ var _ = BeforeSuite(func() { }) var _ = AfterSuite(func() { - }) diff --git a/percona/controller/pgbackup/controller.go b/percona/controller/pgbackup/controller.go index c409dc5008..a0e09451ee 100644 --- a/percona/controller/pgbackup/controller.go +++ b/percona/controller/pgbackup/controller.go @@ -343,16 +343,8 @@ func (r *PGBackupReconciler) Reconcile(ctx context.Context, request reconcile.Re case v2.BackupSucceeded: job, err := findBackupJob(ctx, r.Client, pgBackup) if err == nil && controllerutil.ContainsFinalizer(job, pNaming.FinalizerKeepJob) { - if err := retry.RetryOnConflict(retry.DefaultBackoff, func() error { - j := new(batchv1.Job) - if err := r.Client.Get(ctx, client.ObjectKeyFromObject(job), j); err != nil { - return errors.Wrap(err, "get job") - } - controllerutil.RemoveFinalizer(j, pNaming.FinalizerKeepJob) - - return r.Client.Update(ctx, j) - }); err != nil { - return reconcile.Result{}, errors.Wrap(err, "update PGBackup status") + if err := controller.RemoveKeepJobFinalizer(ctx, r.Client, job); err != nil { + return reconcile.Result{}, errors.Wrap(err, "remove keep-job finalizer") } } @@ -409,16 +401,51 @@ func ensureFinalizers(ctx context.Context, cl client.Client, pgBackup *v2.Percon func deleteBackupFinalizer(c client.Client, pg *v2.PerconaPGCluster) func(ctx context.Context, pgBackup *v2.PerconaPGBackup) error { return func(ctx context.Context, pgBackup *v2.PerconaPGBackup) error { if pg == nil { - return nil - } + // If the cluster was previously deleted, we cannot call finishBackup to remove + // annotations from the PGCluster. + // In that case, we need to + // - delete the job so that it stops trying to create pods + // - remove the keep-job finalizer + job, err := findBackupJob(ctx, c, pgBackup) + if errors.Is(err, ErrBackupJobNotFound) || k8serrors.IsNotFound(err) { + return nil + } + if err != nil { + return errors.Wrap(err, "find backup job") + } - job := new(batchv1.Job) - err := c.Get(ctx, types.NamespacedName{Name: pgBackup.Status.JobName, Namespace: pgBackup.Namespace}, job) - if client.IgnoreNotFound(err) != nil { - return errors.Wrap(err, "get backup job") + if job.DeletionTimestamp.IsZero() { + // If the backup is deleted too early, the job doesn't have an owner reference + // pointing to the backup, so we should delete the job manually + if err := c.Delete(ctx, job, client.PropagationPolicy(metav1.DeletePropagationForeground)); client.IgnoreNotFound(err) != nil { + return errors.Wrap(err, "delete backup job") + } + } + + if err := controller.RemoveKeepJobFinalizer(ctx, c, job); err != nil { + return errors.Wrap(err, "remove keep-job finalizer") + } + + return controller.ErrFinalizerPending } - if k8serrors.IsNotFound(err) { + + job, err := findBackupJob(ctx, c, pgBackup) + if errors.Is(err, ErrBackupJobNotFound) || k8serrors.IsNotFound(err) { job = nil + } else if err != nil { + return errors.Wrap(err, "find backup job") + } + + if job != nil && pgBackup.Status.JobName == "" { + // If the backup is delete too early `pgBackup.Status.JobName` will be empty. + // finishBackup will remove cluster annotations. Without them findBackupJob won't find the job. + // We should set `pgBackup.Status.JobName` in that case + if err := pgBackup.UpdateStatus(ctx, c, func(bcp *v2.PerconaPGBackup) { + bcp.Status.JobName = job.Name + }); err != nil { + return errors.Wrap(err, "update backup job name") + } + return controller.ErrFinalizerPending } rr, err := finishBackup(ctx, c, pgBackup, job) @@ -428,6 +455,16 @@ func deleteBackupFinalizer(c client.Client, pg *v2.PerconaPGCluster) func(ctx co if rr != nil && rr.RequeueAfter != 0 { return controller.ErrFinalizerPending } + if !pgBackup.DeletionTimestamp.IsZero() && job != nil { + // If the backup is deleted too early, the job doesn't have an owner reference + // pointing to the backup, so we should delete the job manually + if err := c.Delete(ctx, job, client.PropagationPolicy(metav1.DeletePropagationForeground)); client.IgnoreNotFound(err) != nil { + return errors.Wrap(err, "delete backup job") + } + if err := controller.RemoveKeepJobFinalizer(ctx, c, job); err != nil { + return errors.Wrap(err, "remove keep-job finalizer") + } + } return nil } } diff --git a/percona/controller/pgbackup/controller_test.go b/percona/controller/pgbackup/controller_test.go index 7af6c08386..136644feae 100644 --- a/percona/controller/pgbackup/controller_test.go +++ b/percona/controller/pgbackup/controller_test.go @@ -7,6 +7,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + batchv1 "k8s.io/api/batch/v1" coordinationv1 "k8s.io/api/coordination/v1" corev1 "k8s.io/api/core/v1" k8serrors "k8s.io/apimachinery/pkg/api/errors" @@ -31,30 +32,30 @@ func TestReconcileFailsBackupWhenClusterIsUnavailable(t *testing.T) { ctx := feature.NewContext(t.Context(), gate) tests := []struct { - name string - deleteCluster bool - snapshot bool - expectedError string + name string + deleteBackupsFinalizer bool + snapshot bool + expectedError string }{ { - name: "cluster is deleted", - deleteCluster: true, + name: "cluster without delete-backups finalizer is deleted", expectedError: "PerconaPGCluster test-cluster is not found", }, { - name: "cluster is terminating", - expectedError: "PerconaPGCluster test-cluster is being deleted", + name: "cluster with delete-backups finalizer is terminating", + deleteBackupsFinalizer: true, + expectedError: "PerconaPGCluster test-cluster is being deleted", }, { - name: "cluster is deleted during snapshot", - deleteCluster: true, + name: "cluster without delete-backups finalizer is deleted during snapshot", snapshot: true, expectedError: "PerconaPGCluster test-cluster is not found", }, { - name: "cluster is terminating during snapshot", - snapshot: true, - expectedError: "PerconaPGCluster test-cluster is being deleted", + name: "cluster with delete-backups finalizer is terminating during snapshot", + deleteBackupsFinalizer: true, + snapshot: true, + expectedError: "PerconaPGCluster test-cluster is being deleted", }, } @@ -62,7 +63,9 @@ func TestReconcileFailsBackupWhenClusterIsUnavailable(t *testing.T) { t.Run(tt.name, func(t *testing.T) { cluster, err := readDefaultCR("test-cluster", "test-namespace") require.NoError(t, err) - cluster.Finalizers = []string{pNaming.FinalizerDeleteBackups} + if tt.deleteBackupsFinalizer { + cluster.Finalizers = []string{pNaming.FinalizerDeleteBackups} + } if tt.snapshot { cluster.Spec.Backups.VolumeSnapshots = &v2.VolumeSnapshots{ Mode: v2.VolumeSnapshotModeOffline, @@ -80,7 +83,10 @@ func TestReconcileFailsBackupWhenClusterIsUnavailable(t *testing.T) { PGCluster: cluster.Name, RepoName: new("repo1"), }, - Status: v2.PerconaPGBackupStatus{State: v2.BackupRunning}, + Status: v2.PerconaPGBackupStatus{ + State: v2.BackupRunning, + JobName: "test-backup-job", + }, } if tt.snapshot { backup.Spec.Method = new(v2.BackupMethodVolumeSnapshot) @@ -93,6 +99,13 @@ func TestReconcileFailsBackupWhenClusterIsUnavailable(t *testing.T) { } objects := []client.Object{backup} + if !tt.snapshot { + objects = append(objects, &batchv1.Job{ObjectMeta: metav1.ObjectMeta{ + Name: backup.Status.JobName, + Namespace: backup.Namespace, + Finalizers: []string{pNaming.FinalizerKeepJob}, + }}) + } if tt.snapshot { holder := backupLeaseHolder(backup) objects = append(objects, &coordinationv1.Lease{ @@ -108,11 +121,6 @@ func TestReconcileFailsBackupWhenClusterIsUnavailable(t *testing.T) { cl, err := buildFakeClient(ctx, cluster, objects...) require.NoError(t, err) require.NoError(t, cl.Delete(ctx, cluster)) - if tt.deleteCluster { - require.NoError(t, cl.Get(ctx, client.ObjectKeyFromObject(cluster), cluster)) - cluster.Finalizers = nil - require.NoError(t, cl.Update(ctx, cluster)) - } r := &PGBackupReconciler{Client: cl} _, err = r.Reconcile(ctx, reconcile.Request{NamespacedName: client.ObjectKeyFromObject(backup)}) @@ -133,12 +141,177 @@ func TestReconcileFailsBackupWhenClusterIsUnavailable(t *testing.T) { }, new(coordinationv1.Lease)) assert.True(t, k8serrors.IsNotFound(err)) } else { + assert.Contains(t, updated.Finalizers, pNaming.FinalizerDeleteBackup) + + err = cl.Get(ctx, client.ObjectKey{ + Name: backup.Status.JobName, Namespace: backup.Namespace, + }, new(batchv1.Job)) + assert.True(t, k8serrors.IsNotFound(err)) + + _, err = r.Reconcile(ctx, reconcile.Request{NamespacedName: client.ObjectKeyFromObject(backup)}) + require.NoError(t, err) + require.NoError(t, cl.Get(ctx, client.ObjectKeyFromObject(backup), updated)) assert.NotContains(t, updated.Finalizers, pNaming.FinalizerDeleteBackup) } }) } } +func TestReconcileDeletesStartingBackupJobWhenClusterIsDeleted(t *testing.T) { + ctx := t.Context() + tests := []struct { + name string + deleteBackupsFinalizer bool + expectedError string + }{ + { + name: "cluster without delete-backups finalizer is deleted", + expectedError: "PerconaPGCluster test-cluster is not found", + }, + { + name: "cluster with delete-backups finalizer is terminating", + deleteBackupsFinalizer: true, + expectedError: "PerconaPGCluster test-cluster is being deleted", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + cluster, err := readDefaultCR("test-cluster", "test-namespace") + require.NoError(t, err) + if tt.deleteBackupsFinalizer { + cluster.Finalizers = []string{pNaming.FinalizerDeleteBackups} + } + + backup := &v2.PerconaPGBackup{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-backup", + Namespace: cluster.Namespace, + Finalizers: []string{pNaming.FinalizerDeleteBackup}, + }, + Spec: v2.PerconaPGBackupSpec{ + PGCluster: cluster.Name, + RepoName: new("repo1"), + }, + Status: v2.PerconaPGBackupStatus{State: v2.BackupStarting}, + } + job := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{ + Name: "test-cluster-backup-xff6", + Namespace: backup.Namespace, + Labels: naming.PGBackRestBackupJobLabels( + cluster.Name, *backup.Spec.RepoName, naming.BackupManual, + ), + Annotations: map[string]string{ + naming.PGBackRestBackup: backup.Name, + }, + Finalizers: []string{pNaming.FinalizerKeepJob}, + }} + + cl, err := buildFakeClient(ctx, cluster, backup, job) + require.NoError(t, err) + require.NoError(t, cl.Delete(ctx, cluster)) + + r := &PGBackupReconciler{Client: cl} + request := reconcile.Request{NamespacedName: client.ObjectKeyFromObject(backup)} + _, err = r.Reconcile(ctx, request) + require.NoError(t, err) + + updated := new(v2.PerconaPGBackup) + require.NoError(t, cl.Get(ctx, client.ObjectKeyFromObject(backup), updated)) + assert.Equal(t, v2.BackupFailed, updated.Status.State) + assert.Equal(t, tt.expectedError, updated.Status.Error) + assert.Empty(t, updated.Status.JobName) + assert.Contains(t, updated.Finalizers, pNaming.FinalizerDeleteBackup) + + err = cl.Get(ctx, client.ObjectKeyFromObject(job), new(batchv1.Job)) + assert.True(t, k8serrors.IsNotFound(err)) + + _, err = r.Reconcile(ctx, request) + require.NoError(t, err) + require.NoError(t, cl.Get(ctx, client.ObjectKeyFromObject(backup), updated)) + assert.NotContains(t, updated.Finalizers, pNaming.FinalizerDeleteBackup) + }) + } +} + +func TestReconcileDeletesJobWhenStartingBackupIsDeleted(t *testing.T) { + ctx := t.Context() + cluster, err := readDefaultCR("test-cluster", "test-namespace") + require.NoError(t, err) + + backup := &v2.PerconaPGBackup{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-backup", + Namespace: cluster.Namespace, + Finalizers: []string{pNaming.FinalizerDeleteBackup}, + }, + Spec: v2.PerconaPGBackupSpec{ + PGCluster: cluster.Name, + RepoName: new("repo1"), + }, + Status: v2.PerconaPGBackupStatus{State: v2.BackupStarting}, + } + cluster.Annotations[naming.PGBackRestBackup] = backup.Name + cluster.Annotations[pNaming.AnnotationBackupInProgress] = backup.Name + + job := &batchv1.Job{ObjectMeta: metav1.ObjectMeta{ + Name: "test-cluster-backup-xff6", + Namespace: backup.Namespace, + Labels: naming.PGBackRestBackupJobLabels( + cluster.Name, *backup.Spec.RepoName, naming.BackupManual, + ), + Annotations: map[string]string{ + naming.PGBackRestBackup: backup.Name, + }, + Finalizers: []string{pNaming.FinalizerKeepJob}, + }} + + cl, err := buildFakeClient(ctx, cluster, backup, job) + require.NoError(t, err) + require.NoError(t, cl.Delete(ctx, backup)) + + r := &PGBackupReconciler{Client: cl} + request := reconcile.Request{NamespacedName: client.ObjectKeyFromObject(backup)} + _, err = r.Reconcile(ctx, request) + require.NoError(t, err) + + updated := new(v2.PerconaPGBackup) + require.NoError(t, cl.Get(ctx, client.ObjectKeyFromObject(backup), updated)) + assert.False(t, updated.DeletionTimestamp.IsZero()) + assert.Contains(t, updated.Finalizers, pNaming.FinalizerDeleteBackup) + assert.Equal(t, job.Name, updated.Status.JobName) + + currentJob := new(batchv1.Job) + require.NoError(t, cl.Get(ctx, client.ObjectKeyFromObject(job), currentJob)) + assert.True(t, currentJob.DeletionTimestamp.IsZero()) + assert.Contains(t, currentJob.Finalizers, pNaming.FinalizerKeepJob) + + // The next reconciliation waits for Crunchy to observe the annotation cleanup, + // so the Job is left untouched. + _, err = r.Reconcile(ctx, request) + require.NoError(t, err) + require.NoError(t, cl.Get(ctx, client.ObjectKeyFromObject(job), currentJob)) + assert.True(t, currentJob.DeletionTimestamp.IsZero()) + assert.Contains(t, currentJob.Finalizers, pNaming.FinalizerKeepJob) + + currentCluster := new(v2.PerconaPGCluster) + require.NoError(t, cl.Get(ctx, client.ObjectKeyFromObject(cluster), currentCluster)) + delete(currentCluster.Annotations, naming.PGBackRestBackup) + delete(currentCluster.Annotations, pNaming.AnnotationBackupInProgress) + require.NoError(t, cl.Update(ctx, currentCluster)) + + crunchyCluster := new(v1beta1.PostgresCluster) + require.NoError(t, cl.Get(ctx, client.ObjectKeyFromObject(cluster), crunchyCluster)) + delete(crunchyCluster.Annotations, pNaming.ToCrunchyAnnotation(naming.PGBackRestBackup)) + delete(crunchyCluster.Annotations, pNaming.ToCrunchyAnnotation(pNaming.AnnotationBackupInProgress)) + require.NoError(t, cl.Update(ctx, crunchyCluster)) + + _, err = r.Reconcile(ctx, request) + require.NoError(t, err) + err = cl.Get(ctx, client.ObjectKeyFromObject(job), new(batchv1.Job)) + assert.True(t, k8serrors.IsNotFound(err)) +} + func TestReconcileNotUpdatingOldBackup(t *testing.T) { ctx := t.Context() cluster, err := readDefaultCR("test-cluster", "test-namespace") diff --git a/percona/controller/pgcluster/backup.go b/percona/controller/pgcluster/backup.go index bb7cda10a2..1e05ed2a2d 100644 --- a/percona/controller/pgcluster/backup.go +++ b/percona/controller/pgcluster/backup.go @@ -136,6 +136,30 @@ func (r *PGClusterReconciler) cleanupOutdatedBackups(ctx context.Context, cr *v2 return nil } +// removeStaleBackupAnnotation removes the manual backup annotation when the PerconaPGBackup is gone. +// Otherwise crunchy would start a backup job that no PerconaPGBackup owns. +func (r *PGClusterReconciler) removeStaleBackupAnnotation(ctx context.Context, cr *v2.PerconaPGCluster) error { + backupName := cr.Annotations[naming.PGBackRestBackup] + if backupName == "" { + return nil + } + + pgBackup := new(v2.PerconaPGBackup) + if err := r.Client.Get(ctx, types.NamespacedName{Name: backupName, Namespace: cr.Namespace}, pgBackup); err != nil { + if k8serrors.IsNotFound(err) { + delete(cr.Annotations, naming.PGBackRestBackup) + return nil + } + return errors.Wrapf(err, "get PerconaPGBackup %s", backupName) + } + + if pgBackup.Spec.PGCluster != cr.Name { + delete(cr.Annotations, naming.PGBackRestBackup) + } + + return nil +} + func (r *PGClusterReconciler) reconcileBackupJobs(ctx context.Context, cr *v2.PerconaPGCluster) error { for _, repo := range cr.Spec.Backups.PGBackRest.Repos { backupJobs, err := listBackupJobs(ctx, r.Client, cr, repo.Name) diff --git a/percona/controller/pgcluster/backup_test.go b/percona/controller/pgcluster/backup_test.go index 35ba0a7558..45e24d885a 100644 --- a/percona/controller/pgcluster/backup_test.go +++ b/percona/controller/pgcluster/backup_test.go @@ -350,3 +350,67 @@ func createFakePodsForStatefulsets(ctx context.Context, cl client.Client, stsLis } return nil } + +func TestRemoveStaleBackupAnnotation(t *testing.T) { + const crName = "some-cluster" + const ns = crName + const backupName = "some-backup" + + newBackupForCluster := func(cluster string, deleting bool) *v2.PerconaPGBackup { + pb := &v2.PerconaPGBackup{ + ObjectMeta: metav1.ObjectMeta{Name: backupName, Namespace: ns}, + Spec: v2.PerconaPGBackupSpec{ + PGCluster: cluster, + RepoName: new("repo1"), + }, + } + if deleting { + pb.Finalizers = []string{pNaming.FinalizerDeleteBackup} + pb.DeletionTimestamp = new(metav1.Now()) + } + return pb + } + newBackup := func(deleting bool) *v2.PerconaPGBackup { + return newBackupForCluster(crName, deleting) + } + + tests := []struct { + name string + backup *v2.PerconaPGBackup + expectKept bool + }{ + // The annotation must not reach crunchy cluster once the backup is gone, + // otherwise it starts a backup Job that no PerconaPGBackup owns. + {name: "backup is gone", expectKept: false}, + {name: "backup exists", backup: newBackup(false), expectKept: true}, + // The delete-backup finalizer still reads the annotation while it cleans up. + {name: "backup is terminating", backup: newBackup(true), expectKept: true}, + // A backup with the same name in the same namespace can belong to another cluster. + {name: "backup belongs to another cluster", backup: newBackupForCluster("other-cluster", false), expectKept: false}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ctx := t.Context() + cr, err := readDefaultCR(crName, ns) + require.NoError(t, err) + cr.Annotations[naming.PGBackRestBackup] = backupName + + objs := []client.Object{} + if tt.backup != nil { + objs = append(objs, tt.backup) + } + cl, err := buildFakeClient(ctx, cr, objs...) + require.NoError(t, err) + + r := &PGClusterReconciler{Client: cl} + require.NoError(t, r.removeStaleBackupAnnotation(ctx, cr)) + + if tt.expectKept { + assert.Equal(t, backupName, cr.Annotations[naming.PGBackRestBackup]) + } else { + assert.NotContains(t, cr.Annotations, naming.PGBackRestBackup) + } + }) + } +} diff --git a/percona/controller/pgcluster/controller.go b/percona/controller/pgcluster/controller.go index e272cc5569..619148fa6f 100644 --- a/percona/controller/pgcluster/controller.go +++ b/percona/controller/pgcluster/controller.go @@ -437,6 +437,10 @@ func (r *PGClusterReconciler) Reconcile(ctx context.Context, request reconcile.R return reconcile.Result{}, errors.Wrap(err, "reconcile replication main site annotation") } + if err := r.removeStaleBackupAnnotation(ctx, cr); err != nil { + return reconcile.Result{}, errors.Wrap(err, "remove stale backup annotation") + } + if cr.Spec.Pause != nil && *cr.Spec.Pause { backupRunning, err := isBackupRunning(ctx, r.Client, cr) if err != nil { diff --git a/percona/controller/pgcluster/testutils_test.go b/percona/controller/pgcluster/testutils_test.go index 0ae826d020..e4eed392af 100644 --- a/percona/controller/pgcluster/testutils_test.go +++ b/percona/controller/pgcluster/testutils_test.go @@ -65,6 +65,7 @@ func reconciler(cr *v2.PerconaPGCluster) *PGClusterReconciler { func crunchyReconciler() *postgrescluster.Reconciler { return &postgrescluster.Reconciler{ Client: k8sClient, + APIReader: k8sClient, // K8SPG-992 Owner: postgrescluster.ControllerName, Recorder: new(record.FakeRecorder), Tracer: otel.Tracer("test"), diff --git a/percona/controller/utils.go b/percona/controller/utils.go index 796fa15f08..d355db10ce 100644 --- a/percona/controller/utils.go +++ b/percona/controller/utils.go @@ -8,6 +8,7 @@ import ( batchv1 "k8s.io/api/batch/v1" corev1 "k8s.io/api/core/v1" k8serrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/client-go/util/retry" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" @@ -15,6 +16,7 @@ import ( "github.com/percona/percona-postgresql-operator/v2/internal/logging" "github.com/percona/percona-postgresql-operator/v2/internal/naming" + pNaming "github.com/percona/percona-postgresql-operator/v2/percona/naming" ) // jobCompleted returns "true" if the Job provided completed successfully. Otherwise it returns @@ -29,6 +31,21 @@ func JobCompleted(job *batchv1.Job) bool { return false } +func RemoveKeepJobFinalizer(ctx context.Context, cl client.Client, job *batchv1.Job) error { + return retry.RetryOnConflict(retry.DefaultBackoff, func() error { + j := new(batchv1.Job) + if err := cl.Get(ctx, client.ObjectKeyFromObject(job), j); err != nil { + return client.IgnoreNotFound(err) + } + + if !controllerutil.RemoveFinalizer(j, pNaming.FinalizerKeepJob) { + return nil + } + + return cl.Update(ctx, j) + }) +} + // jobFailed returns "true" if the Job provided has failed. Otherwise it returns "false". func JobFailed(job *batchv1.Job) bool { conditions := job.Status.Conditions