diff --git a/internal/controller/idempotency_test.go b/internal/controller/idempotency_test.go new file mode 100644 index 0000000..4d4ca8d --- /dev/null +++ b/internal/controller/idempotency_test.go @@ -0,0 +1,279 @@ +package controller + +import ( + "context" + "sync/atomic" + "testing" + + barmancloudv1 "github.com/cloudnative-pg/plugin-barman-cloud/api/v1" + appsv1 "k8s.io/api/apps/v1" + batchv1 "k8s.io/api/batch/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/utils/ptr" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" + + supabasev1alpha1 "github.com/GuionAI/cloudnative-supabase/api/v1alpha1" + cnpgresources "github.com/GuionAI/cloudnative-supabase/internal/resources/cnpg" + deploymentresources "github.com/GuionAI/cloudnative-supabase/internal/resources/deployments" + cnpgv1 "github.com/cloudnative-pg/cloudnative-pg/api/v1" +) + +type updateCountingClient struct { + client.Client + updates atomic.Int32 + statusUpdates atomic.Int32 +} + +func (c *updateCountingClient) Update(ctx context.Context, object client.Object, options ...client.UpdateOption) error { + c.updates.Add(1) + return c.Client.Update(ctx, object, options...) +} + +func (c *updateCountingClient) Status() client.SubResourceWriter { + return &statusUpdateCounter{SubResourceWriter: c.Client.Status(), updates: &c.statusUpdates} +} + +type statusUpdateCounter struct { + client.SubResourceWriter + updates *atomic.Int32 +} + +func (w *statusUpdateCounter) Update(ctx context.Context, object client.Object, options ...client.SubResourceUpdateOption) error { + w.updates.Add(1) + return w.SubResourceWriter.Update(ctx, object, options...) +} + +func TestOwnedResourceHelpersSkipNoOpUpdates(t *testing.T) { + t.Parallel() + + scheme := newIdempotencyTestScheme(t) + project := &supabasev1alpha1.SupabaseProject{ + TypeMeta: metav1.TypeMeta{APIVersion: supabasev1alpha1.GroupVersion.String(), Kind: "SupabaseProject"}, + ObjectMeta: metav1.ObjectMeta{Name: "app", Namespace: "default", UID: "project-uid"}, + } + configMap := &corev1.ConfigMap{ObjectMeta: metav1.ObjectMeta{Name: "config", Namespace: "default"}, Data: map[string]string{"key": "value"}} + deployment := &appsv1.Deployment{ObjectMeta: metav1.ObjectMeta{Name: "deployment", Namespace: "default"}, Spec: appsv1.DeploymentSpec{Replicas: ptr.To[int32](1)}} + service := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "service", Namespace: "default"}, Spec: corev1.ServiceSpec{Ports: []corev1.ServicePort{{Name: "http", Port: 80}}}} + cronJob := &batchv1.CronJob{ObjectMeta: metav1.ObjectMeta{Name: "cron", Namespace: "default"}, Spec: batchv1.CronJobSpec{Schedule: "0 3 * * *"}} + for _, object := range []client.Object{configMap, deployment, service, cronJob} { + if err := setTestControllerReference(project, object, scheme); err != nil { + t.Fatal(err) + } + } + objectStore := &barmancloudv1.ObjectStore{ObjectMeta: metav1.ObjectMeta{Name: "store", Namespace: "default"}, Spec: barmancloudv1.ObjectStoreSpec{RetentionPolicy: "30d"}} + scheduledBackup := &cnpgv1.ScheduledBackup{ObjectMeta: metav1.ObjectMeta{Name: "backup", Namespace: "default"}, Spec: cnpgv1.ScheduledBackupSpec{Schedule: "0 0 2 * * *"}} + + objects := []client.Object{ + configMap.DeepCopy(), deployment.DeepCopy(), service.DeepCopy(), cronJob.DeepCopy(), objectStore.DeepCopy(), scheduledBackup.DeepCopy(), + } + countingClient := &updateCountingClient{Client: fake.NewClientBuilder().WithScheme(scheme).WithObjects(objects...).Build()} + reconciler := &SupabaseProjectReconciler{Client: countingClient, Scheme: scheme} + + ctx := context.Background() + for _, call := range []func() error{ + func() error { return reconciler.createOrUpdateConfigMap(ctx, project, configMap.DeepCopy()) }, + func() error { return reconciler.createOrUpdateDeployment(ctx, project, deployment.DeepCopy()) }, + func() error { return reconciler.createOrUpdateService(ctx, project, service.DeepCopy()) }, + func() error { return reconciler.createOrUpdateCronJob(ctx, project, cronJob.DeepCopy()) }, + func() error { return reconciler.createOrUpdateObjectStore(ctx, objectStore.DeepCopy()) }, + func() error { return reconciler.createOrUpdateScheduledBackup(ctx, scheduledBackup.DeepCopy()) }, + } { + if err := call(); err != nil { + t.Fatal(err) + } + } + if got := countingClient.updates.Load(); got != 0 { + t.Fatalf("no-op owned resource updates = %d, want 0", got) + } + + configMap.Data["key"] = "changed" + deployment.Spec.Replicas = ptr.To[int32](2) + service.Spec.Ports[0].Port = 8080 + cronJob.Spec.Schedule = "0 4 * * *" + objectStore.Spec.RetentionPolicy = "60d" + scheduledBackup.Spec.Schedule = "0 0 4 * * *" + for _, call := range []func() error{ + func() error { return reconciler.createOrUpdateConfigMap(ctx, project, configMap.DeepCopy()) }, + func() error { return reconciler.createOrUpdateDeployment(ctx, project, deployment.DeepCopy()) }, + func() error { return reconciler.createOrUpdateService(ctx, project, service.DeepCopy()) }, + func() error { return reconciler.createOrUpdateCronJob(ctx, project, cronJob.DeepCopy()) }, + func() error { return reconciler.createOrUpdateObjectStore(ctx, objectStore.DeepCopy()) }, + func() error { return reconciler.createOrUpdateScheduledBackup(ctx, scheduledBackup.DeepCopy()) }, + } { + if err := call(); err != nil { + t.Fatal(err) + } + } + if got := countingClient.updates.Load(); got != 6 { + t.Fatalf("drift updates = %d, want 6", got) + } +} + +func TestOwnedResourceHelpersRepairClearedFields(t *testing.T) { + t.Parallel() + + scheme := newIdempotencyTestScheme(t) + project := &supabasev1alpha1.SupabaseProject{ + TypeMeta: metav1.TypeMeta{APIVersion: supabasev1alpha1.GroupVersion.String(), Kind: "SupabaseProject"}, + ObjectMeta: metav1.ObjectMeta{Name: "app", Namespace: "default", UID: "project-uid"}, + Spec: supabasev1alpha1.SupabaseProjectSpec{Auth: supabasev1alpha1.AuthSpec{ + EmailHook: &supabasev1alpha1.EmailHookSpec{Enabled: true, URI: "https://hook.example.com"}, + }}, + } + secretNames := supabasev1alpha1.SecretNamesStatus{ + JWT: "app-jwt", AuthAdmin: "app-auth-admin-password", + } + existingDeployment := deploymentresources.BuildAuthDeployment(project, &secretNames) + project.Spec.Auth.EmailHook.Enabled = false + deployment := deploymentresources.BuildAuthDeployment(project, &secretNames) + service := &corev1.Service{ObjectMeta: metav1.ObjectMeta{Name: "service", Namespace: "default"}} + cronJob := &batchv1.CronJob{ObjectMeta: metav1.ObjectMeta{Name: "cron", Namespace: "default"}} + for _, object := range []client.Object{deployment, service, cronJob} { + if err := setTestControllerReference(project, object, scheme); err != nil { + t.Fatal(err) + } + } + objectStore := &barmancloudv1.ObjectStore{ObjectMeta: metav1.ObjectMeta{Name: "store", Namespace: "default"}} + scheduledBackup := &cnpgv1.ScheduledBackup{ObjectMeta: metav1.ObjectMeta{Name: "backup", Namespace: "default"}} + + existingService := service.DeepCopy() + existingService.Spec.Ports = []corev1.ServicePort{{Name: "stale", Port: 80}} + existingCronJob := cronJob.DeepCopy() + existingCronJob.Spec.Schedule = "0 3 * * *" + existingObjectStore := objectStore.DeepCopy() + existingObjectStore.Spec.RetentionPolicy = "30d" + existingObjectStore.Spec.InstanceSidecarConfiguration.LogLevel = "info" + existingScheduledBackup := scheduledBackup.DeepCopy() + existingScheduledBackup.Spec.Schedule = "0 0 2 * * *" + + countingClient := &updateCountingClient{Client: fake.NewClientBuilder().WithScheme(scheme).WithObjects( + existingDeployment, existingService, existingCronJob, existingObjectStore, existingScheduledBackup, + ).Build()} + reconciler := &SupabaseProjectReconciler{Client: countingClient, Scheme: scheme} + ctx := context.Background() + + for _, call := range []func() error{ + func() error { return reconciler.createOrUpdateDeployment(ctx, project, deployment.DeepCopy()) }, + func() error { return reconciler.createOrUpdateService(ctx, project, service.DeepCopy()) }, + func() error { return reconciler.createOrUpdateCronJob(ctx, project, cronJob.DeepCopy()) }, + func() error { return reconciler.createOrUpdateObjectStore(ctx, objectStore.DeepCopy()) }, + func() error { return reconciler.createOrUpdateScheduledBackup(ctx, scheduledBackup.DeepCopy()) }, + } { + if err := call(); err != nil { + t.Fatal(err) + } + } + if got := countingClient.updates.Load(); got != 5 { + t.Fatalf("cleared-field updates = %d, want 5", got) + } + + updatedDeployment := &appsv1.Deployment{} + if err := countingClient.Get(ctx, client.ObjectKeyFromObject(deployment), updatedDeployment); err != nil { + t.Fatal(err) + } + for _, env := range updatedDeployment.Spec.Template.Spec.Containers[0].Env { + if env.Name == "GOTRUE_HOOK_SEND_EMAIL_ENABLED" || env.Name == "GOTRUE_HOOK_SEND_EMAIL_URI" { + t.Fatalf("disabled email hook env was not removed: %#v", env) + } + } + updatedService := &corev1.Service{} + if err := countingClient.Get(ctx, client.ObjectKeyFromObject(service), updatedService); err != nil { + t.Fatal(err) + } + if got := updatedService.Spec.Ports; len(got) != 0 { + t.Fatalf("service ports were not cleared: %#v", got) + } + updatedObjectStore := &barmancloudv1.ObjectStore{} + if err := countingClient.Get(ctx, client.ObjectKeyFromObject(objectStore), updatedObjectStore); err != nil { + t.Fatal(err) + } + if updatedObjectStore.Spec.RetentionPolicy != "" { + t.Fatalf("retention policy = %q, want empty", updatedObjectStore.Spec.RetentionPolicy) + } + if updatedObjectStore.Spec.InstanceSidecarConfiguration.LogLevel != "info" { + t.Fatalf("API-managed sidecar settings were overwritten: %#v", updatedObjectStore.Spec.InstanceSidecarConfiguration) + } +} + +func TestReconcileSecretsSkipsUnchangedStatusUpdate(t *testing.T) { + t.Parallel() + + scheme := newIdempotencyTestScheme(t) + project := &supabasev1alpha1.SupabaseProject{ + ObjectMeta: metav1.ObjectMeta{Name: "app", Namespace: "default", UID: "project-uid"}, + Status: supabasev1alpha1.SupabaseProjectStatus{ + Phase: supabasev1alpha1.PhaseRunning, + }, + } + objects := make([]client.Object, 0, 5) + objects = append(objects, project) + for _, name := range []string{"app-jwt", "app-supabase-admin-password", "app-authenticator-password", "app-auth-admin-password"} { + objects = append(objects, &corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "default"}}) + } + countingClient := &updateCountingClient{ + Client: fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(project).WithObjects(objects...).Build(), + } + reconciler := &SupabaseProjectReconciler{Client: countingClient, Scheme: scheme} + + ctx := context.Background() + if err := reconciler.reconcileSecrets(ctx, project); err != nil { + t.Fatal(err) + } + countingClient.statusUpdates.Store(0) + if err := reconciler.reconcileSecrets(ctx, project); err != nil { + t.Fatal(err) + } + if got := countingClient.statusUpdates.Load(); got != 0 { + t.Fatalf("unchanged status updates = %d, want 0", got) + } +} + +func TestReconcilePowerSyncPublicationRemovesStaleParameters(t *testing.T) { + t.Parallel() + + scheme := newIdempotencyTestScheme(t) + project := &supabasev1alpha1.SupabaseProject{ + TypeMeta: metav1.TypeMeta{APIVersion: supabasev1alpha1.GroupVersion.String(), Kind: "SupabaseProject"}, + ObjectMeta: metav1.ObjectMeta{Name: "app", Namespace: "default", UID: "project-uid"}, + } + existing := cnpgresources.BuildPowerSyncPublication(project) + existing.Spec.Parameters = map[string]string{"publish": "insert"} + existing.Status.Applied = ptr.To(true) + if err := setTestControllerReference(project, existing, scheme); err != nil { + t.Fatal(err) + } + countingClient := &updateCountingClient{Client: fake.NewClientBuilder().WithScheme(scheme).WithObjects(existing).Build()} + reconciler := &SupabaseProjectReconciler{Client: countingClient, Scheme: scheme} + ctx := context.Background() + + if _, err := reconciler.reconcilePowerSyncPublication(ctx, project); err != nil { + t.Fatal(err) + } + if got := countingClient.updates.Load(); got != 1 { + t.Fatalf("publication updates = %d, want 1", got) + } + updated := &cnpgv1.Publication{} + if err := countingClient.Get(ctx, client.ObjectKeyFromObject(existing), updated); err != nil { + t.Fatal(err) + } + if len(updated.Spec.Parameters) != 0 { + t.Fatalf("publication parameters were not cleared: %#v", updated.Spec.Parameters) + } +} + +func newIdempotencyTestScheme(t *testing.T) *runtime.Scheme { + t.Helper() + scheme := newPowerSyncTestScheme(t) + if err := barmancloudv1.AddToScheme(scheme); err != nil { + t.Fatal(err) + } + return scheme +} + +func setTestControllerReference(project *supabasev1alpha1.SupabaseProject, object client.Object, scheme *runtime.Scheme) error { + return controllerutil.SetControllerReference(project, object, scheme) +} diff --git a/internal/controller/parent_watch_test.go b/internal/controller/parent_watch_test.go new file mode 100644 index 0000000..095ba42 --- /dev/null +++ b/internal/controller/parent_watch_test.go @@ -0,0 +1,139 @@ +package controller + +import ( + "context" + "sync/atomic" + "time" + + cnpgv1 "github.com/cloudnative-pg/cloudnative-pg/api/v1" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server" + + supabasev1alpha1 "github.com/GuionAI/cloudnative-supabase/api/v1alpha1" + deploymentresources "github.com/GuionAI/cloudnative-supabase/internal/resources/deployments" +) + +type projectGetCountingClient struct { + client.Client + projectGets atomic.Int32 +} + +func (c *projectGetCountingClient) Get(ctx context.Context, key client.ObjectKey, object client.Object, options ...client.GetOption) error { + if _, ok := object.(*supabasev1alpha1.SupabaseProject); ok { + c.projectGets.Add(1) + } + return c.Client.Get(ctx, key, object, options...) +} + +var _ = Describe("SupabaseProject parent watch", func() { + It("reconciles removed deployment fields once despite API defaults", func() { + project := &supabasev1alpha1.SupabaseProject{ + TypeMeta: metav1.TypeMeta{APIVersion: supabasev1alpha1.GroupVersion.String(), Kind: "SupabaseProject"}, + ObjectMeta: metav1.ObjectMeta{Name: "defaults", Namespace: "default", UID: "defaults-project-uid"}, + Spec: supabasev1alpha1.SupabaseProjectSpec{Auth: supabasev1alpha1.AuthSpec{ + EmailHook: &supabasev1alpha1.EmailHookSpec{Enabled: true, URI: "https://hook.example.com"}, + }}, + } + secretNames := &supabasev1alpha1.SecretNamesStatus{ + JWT: "defaults-jwt", AuthAdmin: "defaults-auth-admin", + } + desired := deploymentresources.BuildAuthDeployment(project, secretNames) + Expect(setTestControllerReference(project, desired, k8sClient.Scheme())).To(Succeed()) + Expect(k8sClient.Create(ctx, desired)).To(Succeed()) + DeferCleanup(func() { Expect(client.IgnoreNotFound(k8sClient.Delete(ctx, desired))).To(Succeed()) }) + + freshDesired := deploymentresources.BuildAuthDeployment(project, secretNames) + countingClient := &updateCountingClient{Client: k8sClient} + reconciler := &SupabaseProjectReconciler{Client: countingClient, Scheme: k8sClient.Scheme()} + Expect(reconciler.createOrUpdateDeployment(ctx, project, freshDesired)).To(Succeed()) + Expect(countingClient.updates.Load()).To(BeZero()) + + project.Spec.Auth.EmailHook.Enabled = false + disabled := deploymentresources.BuildAuthDeployment(project, secretNames) + Expect(reconciler.createOrUpdateDeployment(ctx, project, disabled)).To(Succeed()) + Expect(countingClient.updates.Load()).To(Equal(int32(1))) + Expect(reconciler.createOrUpdateDeployment(ctx, project, disabled)).To(Succeed()) + Expect(countingClient.updates.Load()).To(Equal(int32(1))) + }) + + It("does not update a CronJob only because the API server defaulted it", func() { + project := &supabasev1alpha1.SupabaseProject{ + TypeMeta: metav1.TypeMeta{APIVersion: supabasev1alpha1.GroupVersion.String(), Kind: "SupabaseProject"}, + ObjectMeta: metav1.ObjectMeta{Name: "cron-defaults", Namespace: "default", UID: "cron-defaults-project-uid"}, + Spec: supabasev1alpha1.SupabaseProjectSpec{Powersync: &supabasev1alpha1.PowersyncSpec{ + Compact: supabasev1alpha1.PowersyncCompactSpec{Enabled: true}, + }}, + } + secretNames := &supabasev1alpha1.SecretNamesStatus{ + JWT: "cron-defaults-jwt", PowersyncStoragePassword: "cron-defaults-storage", + PowersyncReplicationPassword: "cron-defaults-replication", + } + desired := deploymentresources.BuildPowersyncCompactCronJob(project, secretNames) + Expect(setTestControllerReference(project, desired, k8sClient.Scheme())).To(Succeed()) + Expect(k8sClient.Create(ctx, desired)).To(Succeed()) + DeferCleanup(func() { Expect(client.IgnoreNotFound(k8sClient.Delete(ctx, desired))).To(Succeed()) }) + + freshDesired := deploymentresources.BuildPowersyncCompactCronJob(project, secretNames) + countingClient := &updateCountingClient{Client: k8sClient} + reconciler := &SupabaseProjectReconciler{Client: countingClient, Scheme: k8sClient.Scheme()} + Expect(reconciler.createOrUpdateCronJob(ctx, project, freshDesired)).To(Succeed()) + Expect(countingClient.updates.Load()).To(BeZero()) + }) + + It("ignores status-only updates but reconciles spec changes", func() { + manager, err := ctrl.NewManager(cfg, ctrl.Options{ + Scheme: k8sClient.Scheme(), + Metrics: metricsserver.Options{BindAddress: "0"}, + HealthProbeBindAddress: "0", + }) + Expect(err).NotTo(HaveOccurred()) + + countingClient := &projectGetCountingClient{ + Client: fake.NewClientBuilder().WithScheme(k8sClient.Scheme()).Build(), + } + reconciler := &SupabaseProjectReconciler{Client: countingClient, Scheme: k8sClient.Scheme()} + Expect(reconciler.SetupWithManager(manager)).To(Succeed()) + + managerCtx, cancel := context.WithCancel(ctx) + defer cancel() + go func() { + defer GinkgoRecover() + Expect(manager.Start(managerCtx)).To(Succeed()) + }() + Expect(manager.GetCache().WaitForCacheSync(managerCtx)).To(BeTrue()) + + project := &supabasev1alpha1.SupabaseProject{ + ObjectMeta: metav1.ObjectMeta{Name: "parent-watch", Namespace: "default"}, + Spec: supabasev1alpha1.SupabaseProjectSpec{ + Database: supabasev1alpha1.DatabaseSpec{ + Instances: 1, + Storage: cnpgv1.StorageConfiguration{Size: "1Gi"}, + }, + Auth: supabasev1alpha1.AuthSpec{ + SiteURL: "https://app.example.com", + ExternalURL: "https://auth.example.com", + }, + }, + } + Expect(k8sClient.Create(ctx, project)).To(Succeed()) + DeferCleanup(func() { Expect(client.IgnoreNotFound(k8sClient.Delete(ctx, project))).To(Succeed()) }) + + Eventually(countingClient.projectGets.Load, 10*time.Second).Should(BeNumerically(">=", 1)) + baseline := countingClient.projectGets.Load() + + Expect(k8sClient.Get(ctx, client.ObjectKeyFromObject(project), project)).To(Succeed()) + project.Status.Phase = supabasev1alpha1.PhaseRunning + Expect(k8sClient.Status().Update(ctx, project)).To(Succeed()) + Consistently(countingClient.projectGets.Load, time.Second).Should(Equal(baseline)) + + Expect(k8sClient.Get(ctx, client.ObjectKeyFromObject(project), project)).To(Succeed()) + project.Spec.Auth.SiteURL = "https://new.example.com" + Expect(k8sClient.Update(ctx, project)).To(Succeed()) + Eventually(countingClient.projectGets.Load, 10*time.Second).Should(BeNumerically(">", baseline)) + }) +}) diff --git a/internal/controller/powersync_lifecycle_test.go b/internal/controller/powersync_lifecycle_test.go index a869d2b..dcab464 100644 --- a/internal/controller/powersync_lifecycle_test.go +++ b/internal/controller/powersync_lifecycle_test.go @@ -88,6 +88,36 @@ func TestPowerSyncStatusNeedsCleanup(t *testing.T) { } } +func TestReconcileSecretsDoesNotDowngradeRunningPhase(t *testing.T) { + t.Parallel() + + scheme := newPowerSyncTestScheme(t) + project := &supabasev1alpha1.SupabaseProject{ + ObjectMeta: metav1.ObjectMeta{Name: "app", Namespace: "default", UID: "project-uid"}, + Status: supabasev1alpha1.SupabaseProjectStatus{Phase: supabasev1alpha1.PhaseRunning}, + } + objects := []client.Object{project} + for _, name := range []string{ + "app-jwt", + "app-supabase-admin-password", + "app-authenticator-password", + "app-auth-admin-password", + } { + objects = append(objects, &corev1.Secret{ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "default"}}) + } + reconciler := &SupabaseProjectReconciler{ + Client: fake.NewClientBuilder().WithScheme(scheme).WithStatusSubresource(project).WithObjects(objects...).Build(), + Scheme: scheme, + } + + if err := reconciler.reconcileSecrets(context.Background(), project); err != nil { + t.Fatal(err) + } + if project.Status.Phase != supabasev1alpha1.PhaseRunning { + t.Fatalf("phase = %q, want %q", project.Status.Phase, supabasev1alpha1.PhaseRunning) + } +} + func TestCleanupPowerSyncCompactDeletesOwnedCronJob(t *testing.T) { t.Parallel() diff --git a/internal/controller/suite_test.go b/internal/controller/suite_test.go index cb8637c..f6d41a1 100644 --- a/internal/controller/suite_test.go +++ b/internal/controller/suite_test.go @@ -19,10 +19,13 @@ package controller import ( "context" "os" + "os/exec" "path/filepath" + "strings" "testing" cnpgv1 "github.com/cloudnative-pg/cloudnative-pg/api/v1" + barmancloudv1 "github.com/cloudnative-pg/plugin-barman-cloud/api/v1" . "github.com/onsi/ginkgo/v2" . "github.com/onsi/gomega" @@ -64,12 +67,20 @@ var _ = BeforeSuite(func() { Expect(err).NotTo(HaveOccurred()) err = cnpgv1.AddToScheme(scheme.Scheme) Expect(err).NotTo(HaveOccurred()) + err = barmancloudv1.AddToScheme(scheme.Scheme) + Expect(err).NotTo(HaveOccurred()) // +kubebuilder:scaffold:scheme By("bootstrapping test environment") + cnpgModuleDir := goModuleDir("github.com/cloudnative-pg/cloudnative-pg") + barmanModuleDir := goModuleDir("github.com/cloudnative-pg/plugin-barman-cloud") testEnv = &envtest.Environment{ - CRDDirectoryPaths: []string{filepath.Join("..", "..", "config", "crd", "bases")}, + CRDDirectoryPaths: []string{ + filepath.Join("..", "..", "config", "crd", "bases"), + filepath.Join(cnpgModuleDir, "config", "crd", "bases"), + filepath.Join(barmanModuleDir, "config", "crd", "bases"), + }, ErrorIfCRDPathMissing: true, } @@ -88,6 +99,12 @@ var _ = BeforeSuite(func() { Expect(k8sClient).NotTo(BeNil()) }) +func goModuleDir(module string) string { + output, err := exec.Command("go", "list", "-m", "-f={{.Dir}}", module).Output() + Expect(err).NotTo(HaveOccurred()) + return strings.TrimSpace(string(output)) +} + var _ = AfterSuite(func() { By("tearing down the test environment") cancel() diff --git a/internal/controller/supabaseproject_controller.go b/internal/controller/supabaseproject_controller.go index 78ca3d7..41df8cb 100644 --- a/internal/controller/supabaseproject_controller.go +++ b/internal/controller/supabaseproject_controller.go @@ -34,11 +34,14 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" + "k8s.io/utils/ptr" ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/builder" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller/controllerutil" "sigs.k8s.io/controller-runtime/pkg/handler" logf "sigs.k8s.io/controller-runtime/pkg/log" + "sigs.k8s.io/controller-runtime/pkg/predicate" "sigs.k8s.io/controller-runtime/pkg/reconcile" "sigs.k8s.io/yaml" @@ -184,7 +187,7 @@ func (r *SupabaseProjectReconciler) Reconcile(ctx context.Context, req ctrl.Requ project.Status.ObservedGeneration = project.Generation r.setCondition(project, supabasev1alpha1.ConditionTypeReady, metav1.ConditionTrue, "AllComponentsReady", "All Supabase components are running") - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return ctrl.Result{}, err } @@ -198,7 +201,7 @@ func (r *SupabaseProjectReconciler) reconcilePowerSyncRulesValidation(ctx contex return rules, nil } r.setCondition(project, supabasev1alpha1.ConditionTypePowersyncReady, metav1.ConditionFalse, "InvalidSyncRules", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return nil, statusErr } return nil, err @@ -244,14 +247,14 @@ func (r *SupabaseProjectReconciler) waitForPowerSyncManagedRoles(ctx context.Con ready, err := powerSyncManagedRolesReady(cluster) if err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeCDCReady, metav1.ConditionFalse, "ManagedRolesFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err } if !ready { r.setCondition(project, supabasev1alpha1.ConditionTypeCDCReady, metav1.ConditionFalse, "ManagedRolesPending", "Waiting for CloudNativePG to reconcile PowerSync roles") - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return ctrl.Result{}, err } return ctrl.Result{RequeueAfter: RequeueDelay}, nil @@ -279,7 +282,9 @@ func (r *SupabaseProjectReconciler) reconcileSecrets(ctx context.Context, projec log := logf.FromContext(ctx) log.Info("Reconciling secrets") - project.Status.Phase = supabasev1alpha1.PhaseProvisioning + if project.Status.Phase == supabasev1alpha1.PhasePending { + project.Status.Phase = supabasev1alpha1.PhaseProvisioning + } // Check if user-specified secrets mode is enabled if project.Spec.Secrets != nil && !project.Spec.Secrets.AutoGenerate { @@ -303,7 +308,7 @@ func (r *SupabaseProjectReconciler) reconcileUserSpecifiedSecrets(ctx context.Co if apierrors.IsNotFound(err) { r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionFalse, "SecretNotFound", fmt.Sprintf("JWT secret %s not found", secretNames.JWT)) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return fmt.Errorf("JWT secret %s not found", secretNames.JWT) @@ -312,7 +317,7 @@ func (r *SupabaseProjectReconciler) reconcileUserSpecifiedSecrets(ctx context.Co } if err := secrets.ValidateJWTSecret(jwtSecret); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionFalse, "InvalidSecret", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -335,7 +340,7 @@ func (r *SupabaseProjectReconciler) reconcileUserSpecifiedSecrets(ctx context.Co if apierrors.IsNotFound(err) { r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionFalse, "SecretNotFound", fmt.Sprintf("Secret %s for role %s not found", secretName, roleName)) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return fmt.Errorf("secret %s for role %s not found", secretName, roleName) @@ -344,7 +349,7 @@ func (r *SupabaseProjectReconciler) reconcileUserSpecifiedSecrets(ctx context.Co } if err := secrets.ValidateRoleSecret(secret, secretName); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionFalse, "InvalidSecret", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -354,7 +359,7 @@ func (r *SupabaseProjectReconciler) reconcileUserSpecifiedSecrets(ctx context.Co // All secrets validated successfully project.Status.SecretNames = secretNames r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionTrue, "SecretsValidated", "All user-specified secrets are valid") - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return err } @@ -410,7 +415,7 @@ func (r *SupabaseProjectReconciler) reconcileAutoGeneratedSecrets(ctx context.Co project.Status.SecretNames = secretNames r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionTrue, "SecretsExist", "All secrets exist") - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return err } return nil @@ -424,7 +429,7 @@ func (r *SupabaseProjectReconciler) reconcileAutoGeneratedSecrets(ctx context.Co generatedSecrets, secretNames, err := secrets.GenerateSecrets(project) if err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionFalse, "GenerationFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -434,7 +439,7 @@ func (r *SupabaseProjectReconciler) reconcileAutoGeneratedSecrets(ctx context.Co for _, secret := range generatedSecrets { if err := r.createOrUpdateSecret(ctx, project, secret); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionFalse, "CreateFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -452,7 +457,7 @@ func (r *SupabaseProjectReconciler) reconcileAutoGeneratedSecrets(ctx context.Co project.Status.SecretNames = secretNames r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionTrue, "SecretsCreated", "All secrets have been created") - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return err } @@ -488,7 +493,7 @@ func (r *SupabaseProjectReconciler) reconcilePowersyncSecrets(ctx context.Contex psSecrets, err := secrets.GeneratePowersyncSecrets(project) if err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionFalse, "PowersyncSecretsFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -497,7 +502,7 @@ func (r *SupabaseProjectReconciler) reconcilePowersyncSecrets(ctx context.Contex for _, secret := range psSecrets { if err := r.createOrUpdateSecret(ctx, project, secret); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeSecretsReady, metav1.ConditionFalse, "CreateFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -538,14 +543,14 @@ func (r *SupabaseProjectReconciler) reconcileCNPGCluster(ctx context.Context, pr log.Info("Creating CNPG Cluster", "name", cluster.Name) if err := r.Create(ctx, cluster); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeDatabaseReady, metav1.ConditionFalse, "CreateFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err } r.setCondition(project, supabasev1alpha1.ConditionTypeDatabaseReady, metav1.ConditionFalse, "Creating", "CNPG Cluster is being created") - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return ctrl.Result{}, err } @@ -595,7 +600,7 @@ func (r *SupabaseProjectReconciler) waitForDatabase(ctx context.Context, project r.setCondition(project, supabasev1alpha1.ConditionTypeDatabaseReady, metav1.ConditionFalse, "WaitingForInstances", fmt.Sprintf("Waiting for database instances: %d/%d ready", cluster.Status.ReadyInstances, cluster.Spec.Instances)) - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return ctrl.Result{}, err } @@ -616,7 +621,7 @@ func (r *SupabaseProjectReconciler) waitForDatabase(ctx context.Context, project } r.setCondition(project, supabasev1alpha1.ConditionTypeDatabaseReady, metav1.ConditionTrue, "Ready", "Database is ready") - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return ctrl.Result{}, err } @@ -668,7 +673,7 @@ func (r *SupabaseProjectReconciler) reconcileServiceComponent(ctx context.Contex log.V(1).Info(fmt.Sprintf("Built %s deployment", config.name), config.logFields...) if err := r.createOrUpdateDeployment(ctx, project, deployment); err != nil { r.setCondition(project, config.conditionType, metav1.ConditionFalse, "DeploymentFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -678,7 +683,7 @@ func (r *SupabaseProjectReconciler) reconcileServiceComponent(ctx context.Contex service := config.buildService() if err := r.createOrUpdateService(ctx, project, service); err != nil { r.setCondition(project, config.conditionType, metav1.ConditionFalse, "ServiceFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -764,7 +769,7 @@ func (r *SupabaseProjectReconciler) reconcileKong(ctx context.Context, project * log.V(1).Info("Built Kong ConfigMap", "name", kongConfig.Name) if err := r.createOrUpdateConfigMap(ctx, project, kongConfig); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeKongReady, metav1.ConditionFalse, "ConfigMapFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -776,7 +781,7 @@ func (r *SupabaseProjectReconciler) reconcileKong(ctx context.Context, project * "image", fmt.Sprintf("%s:%s", "kong", project.Spec.Kong.ImageTag)) if err := r.createOrUpdateDeployment(ctx, project, deployment); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeKongReady, metav1.ConditionFalse, "DeploymentFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -786,7 +791,7 @@ func (r *SupabaseProjectReconciler) reconcileKong(ctx context.Context, project * service := services.BuildKongService(project) if err := r.createOrUpdateService(ctx, project, service); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeKongReady, metav1.ConditionFalse, "ServiceFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -839,6 +844,9 @@ func (r *SupabaseProjectReconciler) createOrUpdateConfigMap(ctx context.Context, } // Update existing + if apiequality.Semantic.DeepEqual(existing.Data, configMap.Data) { + return nil + } existing.Data = configMap.Data return r.Update(ctx, existing) } @@ -866,13 +874,104 @@ func (r *SupabaseProjectReconciler) createOrUpdateDeployment(ctx context.Context return err } - // Update existing - only update operator-owned fields + desired := deployment.DeepCopy() + actual := existing.DeepCopy() + normalizePodTemplateDefaults(&desired.Spec.Template) + normalizePodTemplateDefaults(&actual.Spec.Template) + if apiequality.Semantic.DeepEqual(actual.Spec.Replicas, desired.Spec.Replicas) && + apiequality.Semantic.DeepEqual(actual.Spec.Template, desired.Spec.Template) { + return nil + } + + // Update existing - only update operator-owned fields. log.V(1).Info("Updating deployment", "name", deployment.Name) - existing.Spec.Replicas = deployment.Spec.Replicas - existing.Spec.Template = deployment.Spec.Template + existing.Spec.Replicas = desired.Spec.Replicas + existing.Spec.Template = desired.Spec.Template return r.Update(ctx, existing) } +func normalizePodTemplateDefaults(template *corev1.PodTemplateSpec) { + spec := &template.Spec + if spec.RestartPolicy == "" { + spec.RestartPolicy = corev1.RestartPolicyAlways + } + if spec.TerminationGracePeriodSeconds == nil { + spec.TerminationGracePeriodSeconds = ptr.To[int64](30) + } + if spec.DNSPolicy == "" { + spec.DNSPolicy = corev1.DNSClusterFirst + } + if spec.SecurityContext == nil { + spec.SecurityContext = &corev1.PodSecurityContext{} + } + if spec.SchedulerName == "" { + spec.SchedulerName = corev1.DefaultSchedulerName + } + for i := range spec.Volumes { + normalizeVolumeDefaults(&spec.Volumes[i]) + } + for i := range spec.InitContainers { + normalizeContainerDefaults(&spec.InitContainers[i]) + } + for i := range spec.Containers { + normalizeContainerDefaults(&spec.Containers[i]) + } +} + +func normalizeVolumeDefaults(volume *corev1.Volume) { + const defaultMode int32 = 0o644 + if volume.ConfigMap != nil && volume.ConfigMap.DefaultMode == nil { + volume.ConfigMap.DefaultMode = ptr.To(defaultMode) + } + if volume.Secret != nil && volume.Secret.DefaultMode == nil { + volume.Secret.DefaultMode = ptr.To(defaultMode) + } + if volume.DownwardAPI != nil && volume.DownwardAPI.DefaultMode == nil { + volume.DownwardAPI.DefaultMode = ptr.To(defaultMode) + } + if volume.Projected != nil && volume.Projected.DefaultMode == nil { + volume.Projected.DefaultMode = ptr.To(defaultMode) + } +} + +func normalizeContainerDefaults(container *corev1.Container) { + if container.TerminationMessagePath == "" { + container.TerminationMessagePath = corev1.TerminationMessagePathDefault + } + if container.TerminationMessagePolicy == "" { + container.TerminationMessagePolicy = corev1.TerminationMessageReadFile + } + for i := range container.Ports { + if container.Ports[i].Protocol == "" { + container.Ports[i].Protocol = corev1.ProtocolTCP + } + } + normalizeProbeDefaults(container.LivenessProbe) + normalizeProbeDefaults(container.ReadinessProbe) + normalizeProbeDefaults(container.StartupProbe) +} + +func normalizeProbeDefaults(probe *corev1.Probe) { + if probe == nil { + return + } + if probe.TimeoutSeconds == 0 { + probe.TimeoutSeconds = 1 + } + if probe.PeriodSeconds == 0 { + probe.PeriodSeconds = 10 + } + if probe.SuccessThreshold == 0 { + probe.SuccessThreshold = 1 + } + if probe.FailureThreshold == 0 { + probe.FailureThreshold = 3 + } + if probe.HTTPGet != nil && probe.HTTPGet.Scheme == "" { + probe.HTTPGet.Scheme = corev1.URISchemeHTTP + } +} + func (r *SupabaseProjectReconciler) createOrUpdateService(ctx context.Context, project *supabasev1alpha1.SupabaseProject, service *corev1.Service) error { if service == nil { logf.FromContext(ctx).V(1).Info("Skipping nil service - service may be disabled") @@ -895,6 +994,10 @@ func (r *SupabaseProjectReconciler) createOrUpdateService(ctx context.Context, p } // Update existing (preserve ClusterIP) + if apiequality.Semantic.DeepEqual(existing.Spec.Ports, service.Spec.Ports) && + apiequality.Semantic.DeepEqual(existing.Spec.Selector, service.Spec.Selector) { + return nil + } existing.Spec.Ports = service.Spec.Ports existing.Spec.Selector = service.Spec.Selector return r.Update(ctx, existing) @@ -948,7 +1051,7 @@ func (r *SupabaseProjectReconciler) reconcileBackup(ctx context.Context, project // Success r.setCondition(project, supabasev1alpha1.ConditionTypeBackupReady, metav1.ConditionTrue, "BackupConfigured", "Backup infrastructure is ready") - return r.Status().Update(ctx, project) + return r.updateProjectStatus(ctx, project) } // cleanupBackupResources removes ObjectStore and ScheduledBackup when backup is disabled @@ -988,7 +1091,7 @@ func (r *SupabaseProjectReconciler) cleanupBackupResources(ctx context.Context, // Remove stale condition after successful cleanup meta.RemoveStatusCondition(&project.Status.Conditions, supabasev1alpha1.ConditionTypeBackupReady) - return r.Status().Update(ctx, project) + return r.updateProjectStatus(ctx, project) } // reconcileRecovery handles recovery ObjectStore creation @@ -1021,7 +1124,7 @@ func (r *SupabaseProjectReconciler) reconcileRecovery(ctx context.Context, proje // Success r.setCondition(project, supabasev1alpha1.ConditionTypeRecoveryReady, metav1.ConditionTrue, "RecoveryConfigured", "Recovery infrastructure is ready") - return r.Status().Update(ctx, project) + return r.updateProjectStatus(ctx, project) } // cleanupRecoveryResources removes recovery ObjectStore when recovery is disabled @@ -1048,13 +1151,13 @@ func (r *SupabaseProjectReconciler) cleanupRecoveryResources(ctx context.Context // Remove stale condition after successful cleanup meta.RemoveStatusCondition(&project.Status.Conditions, supabasev1alpha1.ConditionTypeRecoveryReady) - return r.Status().Update(ctx, project) + return r.updateProjectStatus(ctx, project) } // failRecovery sets RecoveryReady=False condition and updates status func (r *SupabaseProjectReconciler) failRecovery(ctx context.Context, project *supabasev1alpha1.SupabaseProject, reason string, err error) error { r.setCondition(project, supabasev1alpha1.ConditionTypeRecoveryReady, metav1.ConditionFalse, reason, err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -1084,7 +1187,7 @@ func (r *SupabaseProjectReconciler) validateS3Secret(ctx context.Context, namesp // failBackup sets BackupReady=False condition and updates status func (r *SupabaseProjectReconciler) failBackup(ctx context.Context, project *supabasev1alpha1.SupabaseProject, reason string, err error) error { r.setCondition(project, supabasev1alpha1.ConditionTypeBackupReady, metav1.ConditionFalse, reason, err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return statusErr } return err @@ -1107,8 +1210,15 @@ func (r *SupabaseProjectReconciler) createOrUpdateObjectStore(ctx context.Contex return fmt.Errorf("failed to get ObjectStore: %w", err) } - // Update existing ObjectStore - existing.Spec = objectStore.Spec + if apiequality.Semantic.DeepEqual(existing.Spec.Configuration, objectStore.Spec.Configuration) && + existing.Spec.RetentionPolicy == objectStore.Spec.RetentionPolicy { + return nil + } + + // Update only fields owned by this controller. The ObjectStore CRD defaults + // instanceSidecarConfiguration, so preserve that API-managed field. + existing.Spec.Configuration = objectStore.Spec.Configuration + existing.Spec.RetentionPolicy = objectStore.Spec.RetentionPolicy if err := r.Update(ctx, existing); err != nil { return fmt.Errorf("failed to update ObjectStore: %w", err) } @@ -1132,7 +1242,11 @@ func (r *SupabaseProjectReconciler) createOrUpdateScheduledBackup(ctx context.Co return fmt.Errorf("failed to get ScheduledBackup: %w", err) } - // Update existing ScheduledBackup + if apiequality.Semantic.DeepEqual(existing.Spec, scheduledBackup.Spec) { + return nil + } + + // Update existing ScheduledBackup. existing.Spec = scheduledBackup.Spec if err := r.Update(ctx, existing); err != nil { return fmt.Errorf("failed to update ScheduledBackup: %w", err) @@ -1153,7 +1267,7 @@ func (r *SupabaseProjectReconciler) reconcileCDCPermissions(ctx context.Context, configMap := jobs.BuildCDCMigrationsConfigMap(project) if err := r.createOrUpdateConfigMap(ctx, project, configMap); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeCDCReady, metav1.ConditionFalse, "ConfigMapFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err @@ -1167,7 +1281,7 @@ func (r *SupabaseProjectReconciler) reconcileCDCPermissions(ctx context.Context, completed, err := r.createOrCheckJob(ctx, project, job, scriptHash) if err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypeCDCReady, metav1.ConditionFalse, "JobFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err @@ -1175,7 +1289,7 @@ func (r *SupabaseProjectReconciler) reconcileCDCPermissions(ctx context.Context, if !completed { r.setCondition(project, supabasev1alpha1.ConditionTypeCDCReady, metav1.ConditionFalse, "JobRunning", "CDC permissions Job is still running") - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{RequeueAfter: RequeueDelay}, nil @@ -1208,21 +1322,33 @@ func (r *SupabaseProjectReconciler) reconcilePowerSyncPublication(ctx context.Co return ctrl.Result{}, err } r.setCondition(project, supabasev1alpha1.ConditionTypeCDCReady, metav1.ConditionFalse, "PublicationPending", "Waiting for the PowerSync publication") - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return ctrl.Result{}, err } return ctrl.Result{RequeueAfter: RequeueDelay}, nil } - before := existing.DeepCopy() - existing.Spec = desired.Spec - existing.Labels = desired.Labels + needsUpdate := !apiequality.Semantic.DeepEqual(desired.Spec, existing.Spec) + if needsUpdate { + existing.Spec = desired.Spec + } + if existing.Labels == nil && len(desired.Labels) > 0 { + existing.Labels = make(map[string]string, len(desired.Labels)) + } + for key, value := range desired.Labels { + if existing.Labels[key] != value { + existing.Labels[key] = value + needsUpdate = true + } + } + ownerReferences := append([]metav1.OwnerReference(nil), existing.OwnerReferences...) if err := controllerutil.SetControllerReference(project, existing, r.Scheme); err != nil { return ctrl.Result{}, err } - if !apiequality.Semantic.DeepEqual(before.Spec, existing.Spec) || - !apiequality.Semantic.DeepEqual(before.Labels, existing.Labels) || - !apiequality.Semantic.DeepEqual(before.OwnerReferences, existing.OwnerReferences) { + if !apiequality.Semantic.DeepEqual(ownerReferences, existing.OwnerReferences) { + needsUpdate = true + } + if needsUpdate { if err := r.Update(ctx, existing); err != nil { return ctrl.Result{}, fmt.Errorf("updating PowerSync publication: %w", err) } @@ -1235,7 +1361,7 @@ func (r *SupabaseProjectReconciler) reconcilePowerSyncPublication(ctx context.Co message = existing.Status.Message } r.setCondition(project, supabasev1alpha1.ConditionTypeCDCReady, metav1.ConditionFalse, "PublicationPending", message) - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return ctrl.Result{}, err } return ctrl.Result{RequeueAfter: RequeueDelay}, nil @@ -1260,7 +1386,7 @@ func (r *SupabaseProjectReconciler) reconcilePowersync(ctx context.Context, proj psConfig := configmaps.BuildPowersyncConfigMap(project) if err := r.createOrUpdateConfigMap(ctx, project, psConfig); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypePowersyncReady, metav1.ConditionFalse, "ConfigMapFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err @@ -1271,7 +1397,7 @@ func (r *SupabaseProjectReconciler) reconcilePowersync(ctx context.Context, proj if syncRulesConfigMap != nil { if err := r.createOrUpdateConfigMap(ctx, project, syncRulesConfigMap); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypePowersyncReady, metav1.ConditionFalse, "SyncRulesConfigMapFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err @@ -1283,7 +1409,7 @@ func (r *SupabaseProjectReconciler) reconcilePowersync(ctx context.Context, proj applyPowerSyncConfigHash(apiDeployment, psConfig.Data["config.json"], syncRulesContent) if err := r.createOrUpdateDeployment(ctx, project, apiDeployment); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypePowersyncReady, metav1.ConditionFalse, "APIDeploymentFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err @@ -1293,7 +1419,7 @@ func (r *SupabaseProjectReconciler) reconcilePowersync(ctx context.Context, proj apiService := services.BuildPowersyncAPIService(project) if err := r.createOrUpdateService(ctx, project, apiService); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypePowersyncReady, metav1.ConditionFalse, "APIServiceFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err @@ -1304,7 +1430,7 @@ func (r *SupabaseProjectReconciler) reconcilePowersync(ctx context.Context, proj applyPowerSyncConfigHash(replDeployment, psConfig.Data["config.json"], syncRulesContent) if err := r.createOrUpdateDeployment(ctx, project, replDeployment); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypePowersyncReady, metav1.ConditionFalse, "ReplicationDeploymentFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err @@ -1315,7 +1441,7 @@ func (r *SupabaseProjectReconciler) reconcilePowersync(ctx context.Context, proj if compactCronJob != nil { if err := r.createOrUpdateCronJob(ctx, project, compactCronJob); err != nil { r.setCondition(project, supabasev1alpha1.ConditionTypePowersyncReady, metav1.ConditionFalse, "CronJobFailed", err.Error()) - if statusErr := r.Status().Update(ctx, project); statusErr != nil { + if statusErr := r.updateProjectStatus(ctx, project); statusErr != nil { return ctrl.Result{}, statusErr } return ctrl.Result{}, err @@ -1336,7 +1462,7 @@ func (r *SupabaseProjectReconciler) reconcilePowersync(ctx context.Context, proj project.Status.Services.PowersyncReplication = supabasev1alpha1.ServiceStatus{Ready: replicationReady, AvailableReplicas: replicationAvailable} if !apiReady || !replicationReady { r.setCondition(project, supabasev1alpha1.ConditionTypePowersyncReady, metav1.ConditionFalse, "DeploymentsPending", "Waiting for PowerSync deployments to become ready") - if err := r.Status().Update(ctx, project); err != nil { + if err := r.updateProjectStatus(ctx, project); err != nil { return ctrl.Result{}, err } return ctrl.Result{RequeueAfter: RequeueDelay}, nil @@ -1388,7 +1514,7 @@ func (r *SupabaseProjectReconciler) cleanupPowerSync(ctx context.Context, projec project.Status.Services.PowersyncReplication = supabasev1alpha1.ServiceStatus{} meta.RemoveStatusCondition(&project.Status.Conditions, supabasev1alpha1.ConditionTypeCDCReady) meta.RemoveStatusCondition(&project.Status.Conditions, supabasev1alpha1.ConditionTypePowersyncReady) - return r.Status().Update(ctx, project) + return r.updateProjectStatus(ctx, project) } func powerSyncStatusNeedsCleanup(project *supabasev1alpha1.SupabaseProject) bool { @@ -1441,6 +1567,17 @@ func powersyncDeploymentIsReady(deployment *appsv1.Deployment) bool { deployment.Status.UnavailableReplicas == 0 } +func (r *SupabaseProjectReconciler) updateProjectStatus(ctx context.Context, project *supabasev1alpha1.SupabaseProject) error { + current := &supabasev1alpha1.SupabaseProject{} + if err := r.Get(ctx, client.ObjectKeyFromObject(project), current); err != nil { + return err + } + if apiequality.Semantic.DeepEqual(current.Status, project.Status) { + return nil + } + return r.Status().Update(ctx, project) +} + // createOrUpdateCronJob creates or updates a CronJob resource func (r *SupabaseProjectReconciler) createOrUpdateCronJob(ctx context.Context, project *supabasev1alpha1.SupabaseProject, cronJob *batchv1.CronJob) error { log := logf.FromContext(ctx) @@ -1459,11 +1596,26 @@ func (r *SupabaseProjectReconciler) createOrUpdateCronJob(ctx context.Context, p return err } - // Update existing - existing.Spec = cronJob.Spec + desired := cronJob.DeepCopy() + actual := existing.DeepCopy() + normalizeCronJobDefaults(&desired.Spec) + normalizeCronJobDefaults(&actual.Spec) + if apiequality.Semantic.DeepEqual(actual.Spec, desired.Spec) { + return nil + } + + // Update existing. + existing.Spec = desired.Spec return r.Update(ctx, existing) } +func normalizeCronJobDefaults(spec *batchv1.CronJobSpec) { + if spec.Suspend == nil { + spec.Suspend = ptr.To(false) + } + normalizePodTemplateDefaults(&spec.JobTemplate.Spec.Template) +} + const cdcScriptHashAnnotation = "supabase.guion.dev/cdc-script-hash" // createOrCheckJob creates a Job if it doesn't exist, or checks status of an existing Job. @@ -1562,7 +1714,7 @@ func (r *SupabaseProjectReconciler) mapPowerSyncConfigMapToProjects(ctx context. // SetupWithManager sets up the controller with the Manager. func (r *SupabaseProjectReconciler) SetupWithManager(mgr ctrl.Manager) error { return ctrl.NewControllerManagedBy(mgr). - For(&supabasev1alpha1.SupabaseProject{}). + For(&supabasev1alpha1.SupabaseProject{}, builder.WithPredicates(predicate.GenerationChangedPredicate{})). Owns(&corev1.Secret{}). Owns(&corev1.ConfigMap{}). Watches(&corev1.ConfigMap{}, handler.EnqueueRequestsFromMapFunc(r.mapPowerSyncConfigMapToProjects)).