Skip to content
Open
1 change: 1 addition & 0 deletions cmd/postgres-operator/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
11 changes: 10 additions & 1 deletion internal/controller/postgrescluster/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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}
Expand Down
34 changes: 34 additions & 0 deletions internal/controller/postgrescluster/pgbackrest.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down
137 changes: 132 additions & 5 deletions internal/controller/postgrescluster/pgbackrest_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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) })

Expand All @@ -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),
},
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)
})
}
}
12 changes: 9 additions & 3 deletions internal/controller/postgrescluster/suite_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -63,5 +70,4 @@ var _ = BeforeSuite(func() {
})

var _ = AfterSuite(func() {

})
71 changes: 54 additions & 17 deletions percona/controller/pgbackup/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
}

Expand Down Expand Up @@ -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)
Expand All @@ -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
}
}
Expand Down
Loading
Loading