diff --git a/.github/workflows/latest.yaml b/.github/workflows/latest.yaml index 8c5b5d8a..524d579f 100644 --- a/.github/workflows/latest.yaml +++ b/.github/workflows/latest.yaml @@ -25,9 +25,10 @@ jobs: go-version-file: 'go.mod' - name: Lint - uses: golangci/golangci-lint-action@v4 + uses: golangci/golangci-lint-action@v9 with: - args: -p bugs -p unused --timeout=5m + only-new-issues: true + args: --timeout=5m - name: Build and push Docker image run: | diff --git a/.github/workflows/pull_request.yaml b/.github/workflows/pull_request.yaml index 0c91e911..3671ed63 100644 --- a/.github/workflows/pull_request.yaml +++ b/.github/workflows/pull_request.yaml @@ -20,9 +20,10 @@ jobs: go-version-file: 'go.mod' - name: Lint - uses: golangci/golangci-lint-action@v6 + uses: golangci/golangci-lint-action@v9 with: - args: -p bugs -p unused --timeout=5m + only-new-issues: true + args: --timeout=5m - name: Get Kubebuilder Version id: get-kubebuilder-version diff --git a/.github/workflows/release.yaml b/.github/workflows/release.yaml index 45cec546..3110ea38 100644 --- a/.github/workflows/release.yaml +++ b/.github/workflows/release.yaml @@ -26,9 +26,10 @@ jobs: go-version-file: 'go.mod' - name: Lint - uses: golangci/golangci-lint-action@v4 + uses: golangci/golangci-lint-action@v9 with: - args: -p bugs -p unused --timeout=5m + args: --timeout=5m + only-new-issues: true - name: Build and push Docker image run: | diff --git a/.golangci.yaml b/.golangci.yaml new file mode 100644 index 00000000..0b0fdf36 --- /dev/null +++ b/.golangci.yaml @@ -0,0 +1,98 @@ +version: "2" +linters: + enable: + - asasalint + - asciicheck + - bidichk + - bodyclose + - canonicalheader + - containedctx + - contextcheck + - copyloopvar + - decorder + - dogsled + - dupl + - durationcheck + - err113 + - errchkjson + - errname + - errorlint + - exhaustive + - exhaustruct + - exptostd + - forbidigo + - forcetypeassert + - ginkgolinter + - gocheckcompilerdirectives + - gochecknoglobals + - gochecknoinits + - gochecksumtype + - goconst + - gocritic + - godox + - goheader + - gomoddirectives + - gomodguard_v2 + - goprintffuncname + - gosec + - gosmopolitan + - grouper + - iface + - importas + - inamedparam + - interfacebloat + - intrange + - ireturn + - loggercheck + - makezero + - mirror + - misspell + - musttag + - nakedret + - nilerr + - nilnesserr + - nilnil + - nlreturn + - noctx + - nolintlint + - nonamedreturns + - nosprintfhostport + - paralleltest + - predeclared + - promlinter + - protogetter + - reassign + - recvcheck + - rowserrcheck + - sloglint + - spancheck + - sqlclosecheck + - staticcheck + - tagalign + - testifylint + - testpackage + - tparallel + - unconvert + - unparam + - usestdlibvars + - wastedassign + - wrapcheck + - zerologlint + exclusions: + generated: lax + presets: + - comments + - common-false-positives + - legacy + - std-error-handling + paths: + - third_party$ + - builtin$ + - examples$ +formatters: + exclusions: + generated: lax + paths: + - third_party$ + - builtin$ + - examples$ diff --git a/Makefile b/Makefile index a24caab7..8b43f025 100644 --- a/Makefile +++ b/Makefile @@ -251,4 +251,4 @@ localkube-reinstall-postgreslet: localkube-load-image helm upgrade --install postgreslet metal-stack-30/postgreslet --namespace postgreslet-system --values svc-cluster-values.yaml --set-file controlplaneKubeconfig=kubeconfig-ctrl --kubeconfig ./kubeconfig-svc lint: - golangci-lint run -p bugs -p unused --timeout=5m + golangci-lint run --timeout=5m --fix --new-from-merge-base=main diff --git a/api/v1/postgres_types.go b/api/v1/postgres_types.go index 612f213d..8d3cd7d5 100644 --- a/api/v1/postgres_types.go +++ b/api/v1/postgres_types.go @@ -91,7 +91,7 @@ const ( defaultPostgresParamValueWalKeepSize = "1GB" defaultPostgresParamValuePGStatStatementsMax = "500" defaultSelectorDisableValue = "selector-disabled" - defaultPostgresParamValuePasswordEncryption = "scram-sha-256" // nolint + defaultPostgresParamValuePasswordEncryption = "scram-sha-256" //nolint defaultPostgresParamValueLogMinErrorStatement = "WARNING" defaultPostgresParamValueLogErrorVerbosity = "VERBOSE" defaultPostgresParamValueLogLinePrefix = "%m [%p]: [%l-1] db=%d,user=%u,app=%a,client=%h " @@ -319,7 +319,7 @@ func (p *Postgres) HasSourceRanges() bool { // IsBeingDeleted returns true if the deletion-timestamp is set func (p *Postgres) IsBeingDeleted() bool { - return !p.ObjectMeta.DeletionTimestamp.IsZero() + return !p.DeletionTimestamp.IsZero() } // ToCWNP returns CRD ClusterwideNetworkPolicy derived from CRD Postgres @@ -664,6 +664,7 @@ func (p *Postgres) ToDNSName(tlsSubDomain string) string { if len(name) > maxLen { name = name[:maxLen] } + return name + "." + tlsSubDomain } @@ -699,43 +700,43 @@ func (p *Postgres) ToUnstructuredZalandoPostgresql(z *zalando.Postgresql, c *cor z.Spec.DockerImage = image } z.Spec.NumberOfInstances = p.Spec.NumberOfInstances - z.Spec.PostgresqlParam.PgVersion = p.Spec.Version + z.Spec.PgVersion = p.Spec.Version // initialize the parameters - z.Spec.PostgresqlParam.Parameters = map[string]string{} + z.Spec.Parameters = map[string]string{} // enable default audit logs (if not configured otherwise) if p.Spec.AuditLogs == nil || *p.Spec.AuditLogs { - enableAuditLogs(z.Spec.PostgresqlParam.Parameters) + enableAuditLogs(z.Spec.Parameters) } // set some default postgres parameters - setDefaultPostgresParams(z.Spec.PostgresqlParam.Parameters, p.Spec.Version) + setDefaultPostgresParams(z.Spec.Parameters, p.Spec.Version) // now set the given generic parameters (and potentially allow overwriting of default postgres params or audit log params) - setPostgresParams(z.Spec.PostgresqlParam.Parameters, p.Spec.PostgresParams, pgParamBlockList) + setPostgresParams(z.Spec.Parameters, p.Spec.PostgresParams, pgParamBlockList) // finally, overwrite the (special to us) shared buffer parameter - setSharedBufferSize(z.Spec.PostgresqlParam.Parameters, p.Spec.Size.SharedBuffer) + setSharedBufferSize(z.Spec.Parameters, p.Spec.Size.SharedBuffer) z.Spec.Resources = &zalando.Resources{} cpuReq, err := p.calculateCPURequests(p.Spec.Size.CPU, cpuRequestsPercentage) if err != nil { return nil, fmt.Errorf("failed to convert to unstructured zalando postgresql: %w", err) } - z.Spec.Resources.ResourceRequests.CPU = ptr.To(cpuReq) - z.Spec.Resources.ResourceRequests.Memory = ptr.To(p.Spec.Size.Memory) - z.Spec.Resources.ResourceLimits.CPU = ptr.To(p.Spec.Size.CPU) - z.Spec.Resources.ResourceLimits.Memory = ptr.To(p.Spec.Size.Memory) + z.Spec.ResourceRequests.CPU = ptr.To(cpuReq) + z.Spec.ResourceRequests.Memory = ptr.To(p.Spec.Size.Memory) + z.Spec.ResourceLimits.CPU = ptr.To(p.Spec.Size.CPU) + z.Spec.ResourceLimits.Memory = ptr.To(p.Spec.Size.Memory) z.Spec.TeamID = p.generateTeamID() - z.Spec.Volume.Size = p.Spec.Size.StorageSize + z.Spec.Size = p.Spec.Size.StorageSize if p.Spec.StorageClass != nil { - z.Spec.Volume.StorageClass = *p.Spec.StorageClass + z.Spec.StorageClass = *p.Spec.StorageClass } else { - z.Spec.Volume.StorageClass = sc + z.Spec.StorageClass = sc } - z.Spec.Patroni.TTL = patroniTTL - z.Spec.Patroni.LoopWait = patroniLoopWait - z.Spec.Patroni.RetryTimeout = patroniRetryTimeout - z.Spec.Patroni.SynchronousMode = true - z.Spec.Patroni.SynchronousModeStrict = false + z.Spec.TTL = patroniTTL + z.Spec.LoopWait = patroniLoopWait + z.Spec.RetryTimeout = patroniRetryTimeout + z.Spec.SynchronousMode = true + z.Spec.SynchronousModeStrict = false // required with image ermajn/postgres-operator:v1.6.0-20-g1cc71663-dirty // see https://github.com/fi-ts/postgreslet/issues/293 @@ -866,15 +867,15 @@ func (p *Postgres) ToZalandoPostgresqlMatchingLabels() client.MatchingLabels { } func (p *Postgres) HasFinalizer(finalizerName string) bool { - return containsElem(p.ObjectMeta.Finalizers, finalizerName) + return containsElem(p.Finalizers, finalizerName) } func (p *Postgres) AddFinalizer(finalizerName string) { - p.ObjectMeta.Finalizers = append(p.ObjectMeta.Finalizers, finalizerName) + p.Finalizers = append(p.Finalizers, finalizerName) } func (p *Postgres) RemoveFinalizer(finalizerName string) { - p.ObjectMeta.Finalizers = removeElem(p.ObjectMeta.Finalizers, finalizerName) + p.Finalizers = removeElem(p.Finalizers, finalizerName) } func containsElem(ss []string, s string) bool { @@ -883,6 +884,7 @@ func containsElem(ss []string, s string) bool { return true } } + return false } @@ -893,6 +895,7 @@ func removeElem(ss []string, s string) (out []string) { } out = append(out, elem) } + return } @@ -941,6 +944,7 @@ func (p *Postgres) buildSidecars(c *corev1.ConfigMap) []zalando.Sidecar { for j := range sidecars[i].Env { if sidecars[i].Env[j].ValueFrom != nil && sidecars[i].Env[j].ValueFrom.SecretKeyRef != nil { sidecars[i].Env[j].ValueFrom.SecretKeyRef.Name = PostgresConfigMonitoringUsername + "." + p.ToPeripheralResourceName() + ".credentials" + break } } @@ -969,6 +973,7 @@ func (p *Postgres) IsReplicationPrimaryOrStandalone() bool { // nothing is configured, or we are the leader. nothing to do. return true } + return false } @@ -977,6 +982,7 @@ func (p *Postgres) IsReplicationTarget() bool { // sth is configured and we are not the leader return true } + return false } @@ -1125,13 +1131,13 @@ func (p *Postgres) calculateCPURequests(c string, percentage int) (string, error // calculate the percentage value := int64((milliValue / int64(100)) * int64(percentage)) - //return the calculated cpu request, making sure it is not higher than the given input value + // return the calculated cpu request, making sure it is not higher than the given input value return resource.NewMilliQuantity(min(value, milliValue), resource.BinarySI).String(), nil } // sanitize a string so it can be used as a label value. if the string is valid, it // will be returned unmodified. otherwise all illegal character will be replace with "_" -// where multiple "_" will be shrinked to a single one. the string must also start and end +// where multiple "_" will be shrunk to a single one. the string must also start and end // with a alphanumeric character. last but not least, a label value must not be longer than 63 // characters. func sanitizeLabelValue(v string) string { diff --git a/api/v1/postgres_types_test.go b/api/v1/postgres_types_test.go index a7874fd1..005e1c79 100644 --- a/api/v1/postgres_types_test.go +++ b/api/v1/postgres_types_test.go @@ -72,7 +72,6 @@ func Test_setSharedBufferSize(t *testing.T) { }, } for _, tt := range tests { - tt := tt // pin! t.Run(tt.name, func(t *testing.T) { parameters := map[string]string{} @@ -131,9 +130,8 @@ func TestPostgres_generateTeamID(t *testing.T) { }, } for _, tt := range tests { - tt := tt // pin! t.Run(tt.name, func(t *testing.T) { - var dnsRegExp *regexp.Regexp = regexp.MustCompile("^[a-z]([-a-z0-9]*[a-z0-9])?$") + dnsRegExp := regexp.MustCompile("^[a-z]([-a-z0-9]*[a-z0-9])?$") p := &Postgres{ Spec: PostgresSpec{ ProjectID: tt.projectID, @@ -194,9 +192,8 @@ func TestPostgres_ToPeripheralResourceName(t *testing.T) { }, } for _, tt := range tests { - tt := tt // pin! t.Run(tt.name, func(t *testing.T) { - var dnsRegExp *regexp.Regexp = regexp.MustCompile("^[a-z]([-a-z0-9]*[a-z0-9])?$") + dnsRegExp := regexp.MustCompile("^[a-z]([-a-z0-9]*[a-z0-9])?$") p := &Postgres{ ObjectMeta: v1.ObjectMeta{ Name: tt.postgresName, @@ -210,7 +207,7 @@ func TestPostgres_ToPeripheralResourceName(t *testing.T) { if !dnsRegExp.MatchString(got) { t.Errorf("Postgres.ToPeripheralResourceName() got %v, not a valid DNS name", got) } - //This resource name will be used as part of the name of other resources, hence we need to limit it's length, + // This resource name will be used as part of the name of other resources, hence we need to limit it's length, // e.g. "postgres.bce25ade7552494c-33d21de46d284ea6bec0.credentials" maxLen := 37 if len(got) > maxLen { @@ -378,7 +375,6 @@ func TestPostgresRestoreTimestamp_ToUnstructuredZalandoPostgresql(t *testing.T) }, } for _, tt := range tests { - tt := tt // pin! t.Run(tt.name, func(t *testing.T) { p := &Postgres{ Spec: tt.spec, @@ -445,7 +441,6 @@ func Test_calculateCPURequests(t *testing.T) { }, } for _, tt := range tests { - tt := tt // pin! t.Run(tt.name, func(t *testing.T) { p := &Postgres{ Spec: PostgresSpec{ diff --git a/controllers/postgres_controller.go b/controllers/postgres_controller.go index 230c9e44..9b526061 100644 --- a/controllers/postgres_controller.go +++ b/controllers/postgres_controller.go @@ -133,17 +133,19 @@ type PatroniConfig struct { // +kubebuilder:rbac:groups=acid.zalan.do,resources=postgresqls,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups=acid.zalan.do,resources=postgresqls/status,verbs=get;list;watch func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { - log := r.Log.WithValues("pgID", req.NamespacedName.Name) + log := r.Log.WithValues("pgID", req.Name) instance := &pg.Postgres{} if err := r.CtrlClient.Get(ctx, req.NamespacedName, instance); err != nil { if apierrors.IsNotFound(err) { // the instance was updated, but does not exist anymore -> do nothing, it was probably deleted log.Info("postgres already deleted") + return ctrl.Result{}, nil } r.recorder.Eventf(instance, "Warning", "Error", "failed to get resource: %v", err) + return ctrl.Result{}, err } log.V(debugLogLevel).Info("postgres fetched", "postgres", instance) @@ -152,6 +154,7 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c if !r.isManagedByUs(instance) { log.V(debugLogLevel).Info("object should be managed by another postgreslet, ignored.") + return ctrl.Result{}, nil } @@ -162,6 +165,7 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c instance.Status.Description = "Terminating" if err := r.CtrlClient.Status().Update(ctx, instance); err != nil { log.Error(err, "failed to update owner object") + return ctrl.Result{}, err } log.Info("instance being deleted") @@ -171,17 +175,20 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c if err := r.deleteCWNP(log, ctx, instance); client.IgnoreNotFound(err) != nil { // todo: remove ignorenotfound r.recorder.Event(instance, "Warning", "Error", "failed to delete ClusterwideNetworkPolicy") + return ctrl.Result{}, err } log.V(debugLogLevel).Info("corresponding CRD ClusterwideNetworkPolicy deleted") - if err := r.LBManager.DeleteSharedSvcLB(ctx, instance); err != nil { + if err := r.DeleteSharedSvcLB(ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to delete Service with shared ip: %v", err) + return ctrl.Result{}, err } - if err := r.LBManager.DeleteDedicatedSvcLB(ctx, instance); err != nil { + if err := r.DeleteDedicatedSvcLB(ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to delete Service with dedicated ip: %v", err) + return ctrl.Result{}, err } log.V(debugLogLevel).Info("corresponding Service(s) of type LoadBalancer deleted") @@ -193,6 +200,7 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c if err := r.deleteZPostgresqlByLabels(log, ctx, matchingLabels, namespace); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to delete Zalando resource: %v", err) + return ctrl.Result{}, err } log.V(debugLogLevel).Info("owned zalando postgresql deleted") @@ -211,18 +219,21 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c log.V(debugLogLevel).Info("finalizer from storage encryption secret removed") } - deletable, err := r.OperatorManager.IsOperatorDeletable(ctx, namespace, instance.ToPeripheralResourceName()) + deletable, err := r.IsOperatorDeletable(ctx, namespace, instance.ToPeripheralResourceName()) if err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to check if the operator is idle: %v", err) + return ctrl.Result{}, fmt.Errorf("error while checking if the operator is idle: %w", err) } if !deletable { r.recorder.Event(instance, "Warning", "Self-Reconciliation", "operator not yet deletable, requeuing") log.Info("operator not yet deletable, requeuing") + return ctrl.Result{Requeue: true}, nil } - if err := r.OperatorManager.UninstallOperator(ctx, namespace); err != nil { + if err := r.UninstallOperator(ctx, namespace); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to uninstall operator: %v", err) + return ctrl.Result{}, fmt.Errorf("error while uninstalling operator: %w", err) } log.V(debugLogLevel).Info("corresponding operator deleted") @@ -235,10 +246,12 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c instance.RemoveFinalizer(pg.PostgresFinalizerName) if err := r.CtrlClient.Update(ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Self-Reconciliation", "failed to remove finalizer: %v", err) + return ctrl.Result{}, fmt.Errorf("failed to update finalizers: %w", err) } log.V(debugLogLevel).Info("finalizers removed") log.Info("postgres deletion reconciled") + return ctrl.Result{}, nil } @@ -247,6 +260,7 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c instance.AddFinalizer(pg.PostgresFinalizerName) if err := r.CtrlClient.Update(ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Self-Reconciliation", "failed to add finalizer: %v", err) + return ctrl.Result{}, fmt.Errorf("error while adding finalizer: %w", err) } log.V(debugLogLevel).Info("finalizer added") @@ -255,12 +269,14 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c backupConfig, err := r.getBackupConfig(ctx, instance.Namespace, instance.Spec.BackupSecretRef) if err != nil { r.recorder.Eventf(instance, "Warning", "Self-Reconciliation", "failed to fetch backupConfig: %v", err) + return ctrl.Result{}, fmt.Errorf("failed to fetch backupConfig: %w", err) } // Check if zalando dependencies are installed. If not, install them. if err := r.ensureZalandoDependencies(log, ctx, instance, backupConfig); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to install operator: %v", err) + return ctrl.Result{}, fmt.Errorf("error while ensuring Zalando dependencies: %w", err) } @@ -268,11 +284,13 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c if r.EnableNetPol { if err := r.createOrUpdateNetPol(ctx, instance, r.EtcdHost); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to create netpol: %v", err) + return ctrl.Result{}, fmt.Errorf("error while creating netpol: %w", err) } } else { if err := r.deleteNetPol(ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to delete netpol: %v", err) + return ctrl.Result{}, fmt.Errorf("error while deleting netpol: %w", err) } } @@ -280,12 +298,14 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c // Request certificate, if necessary if err := r.createOrUpdateCertificate(log, ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to create certificate request: %v", err) + return ctrl.Result{}, fmt.Errorf("error while creating certificate request: %w", err) } // Make sure the postgres secrets exist, if necessary if err := r.ensurePostgresSecrets(log, ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to create postgres secrets: %v", err) + return ctrl.Result{}, fmt.Errorf("error while creating postgres secrets: %w", err) } @@ -295,6 +315,7 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c // create standby egress rule first, so the standby can actually connect to the primary if err := r.createOrUpdateEgressCWNP(ctx, instance); err != nil { r.recorder.Event(instance, "Warning", "Error", "failed to create or update egress ClusterwideNetworkPolicy") + return ctrl.Result{}, fmt.Errorf("unable to create or update egress ClusterwideNetworkPolicy: %w", err) } @@ -323,47 +344,56 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c // Add pod monitor if err := r.createOrUpdatePatroniPodMonitor(ctx, namespace, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to create podmonitor: %v", err) + return ctrl.Result{}, fmt.Errorf("error while creating podmonitor %v: %w", namespace, err) } // Make sure the storage secret exist, if necessary if err := r.ensureStorageEncryptionSecret(log, ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to create storage secret: %v", err) + return ctrl.Result{}, fmt.Errorf("error while creating storage secret: %w", err) } if err := r.createOrUpdateZalandoPostgresql(ctx, instance, log, globalSidecarsCM, r.PatroniTTL, r.PatroniLoopWait, r.PatroniRetryTimeout); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to create Zalando resource: %v", err) + return ctrl.Result{}, fmt.Errorf("failed to create or update zalando postgresql: %w", err) } if err := r.ensureInitDBJob(log, ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to create initDB job resource: %v", err) + return ctrl.Result{}, fmt.Errorf("failed to create or update initdb job: %w", err) } - if err := r.LBManager.ReconcileSvcLBs(ctx, instance); err != nil { + if err := r.ReconcileSvcLBs(ctx, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to create Service: %v", err) + return ctrl.Result{}, err } if r.EnableWalGExporter { if err := r.createOrUpdateWalGExporterDeployment(log, ctx, namespace, instance, backupConfig); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to deploy wal-g-exporter: %v", err) + return ctrl.Result{}, fmt.Errorf("error while deploying wal-g-exporter %v: %w", namespace, err) } if err := r.createOrUpdateWalGExporterPodMonitor(log, ctx, namespace, instance); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to deploy wal-g-exporter podMonitor: %v", err) + return ctrl.Result{}, fmt.Errorf("error while deploying wal-g-exporter podMonitor %v: %w", namespace, err) } } else { // remove wal-g-exporter when disabled if err := r.deleteWalGExporterDeployment(ctx, namespace); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to delete wal-g-exporter: %v", err) + return ctrl.Result{}, fmt.Errorf("error while deleting wal-g-exporter: %w", err) } if err := r.deleteWalGExporterPodMonitor(ctx, namespace); err != nil { r.recorder.Eventf(instance, "Warning", "Error", "failed to delete wal-g-exporter podMonitor: %v", err) + return ctrl.Result{}, fmt.Errorf("error while deleting wal-g-exporter podMonitor: %w", err) } } @@ -373,12 +403,14 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c if port == 0 { r.recorder.Event(instance, "Warning", "Self-Reconciliation", "socket port not ready") log.Info("socket port not ready, requeuing") + return requeue, nil } // Update status will be handled by the StatusReconciler, based on the Zalando Status if err := r.createOrUpdateIngressCWNP(log, ctx, instance, int(port)); err != nil { r.recorder.Event(instance, "Warning", "Error", "failed to create or update ingress ClusterwideNetworkPolicy") + return ctrl.Result{}, fmt.Errorf("unable to create or update ingress ClusterwideNetworkPolicy: %w", err) } @@ -386,23 +418,27 @@ func (r *PostgresReconciler) Reconcile(ctx context.Context, req ctrl.Request) (c // we try again in the next loop, hoping things will settle if patroniConfigChangeErr != nil { log.Info("Requeuing after getting/setting patroni replication config failed") + return ctrl.Result{Requeue: true, RequeueAfter: 10 * time.Second}, patroniConfigChangeErr } // if the config isn't in the expected state yet (we only add values to an existing config, we do not perform the actual switch), we simply requeue. // on the next reconciliation loop, postgres-operator should have caught up and the config should hopefully be correct already so we can continue with adding our values. if requeueAfterReconcile { log.Info("Requeuing after patroni replication hasn't returned the expected state (yet)") + return ctrl.Result{Requeue: true, RequeueAfter: r.ReplicationChangeRequeueDuration}, nil } log.Info("postgres reconciled") r.recorder.Event(instance, "Normal", "Reconciled", "postgres up to date") + return ctrl.Result{}, nil } // SetupWithManager informs mgr when this reconciler should be called. func (r *PostgresReconciler) SetupWithManager(mgr ctrl.Manager) error { r.recorder = mgr.GetEventRecorderFor("PostgresController") + return ctrl.NewControllerManagedBy(mgr). For(&pg.Postgres{}). WithEventFilter(predicate.GenerationChangedPredicate{}). @@ -479,6 +515,7 @@ func (r *PostgresReconciler) deleteUserPasswordsSecret(ctx context.Context, inst if err := r.CtrlClient.Delete(ctx, secret); client.IgnoreNotFound(err) != nil { msgWithFormat := "failed to delete user passwords secret: %w" r.recorder.Eventf(instance, "Warning", "Error", msgWithFormat, err) + return fmt.Errorf(msgWithFormat, err) } @@ -488,13 +525,13 @@ func (r *PostgresReconciler) deleteUserPasswordsSecret(ctx context.Context, inst // ensureZalandoDependencies makes sure Zalando resources are installed in the service-cluster. func (r *PostgresReconciler) ensureZalandoDependencies(log logr.Logger, ctx context.Context, p *pg.Postgres, b *pg.BackupConfig) error { namespace := p.ToPeripheralResourceNamespace() - isInstalled, err := r.OperatorManager.IsOperatorInstalled(ctx, namespace) + isInstalled, err := r.IsOperatorInstalled(ctx, namespace) if err != nil { return fmt.Errorf("error while querying if zalando dependencies are installed: %w", err) } if !isInstalled { - if err := r.OperatorManager.InstallOrUpdateOperator(ctx, namespace); err != nil { + if err := r.InstallOrUpdateOperator(ctx, namespace); err != nil { return fmt.Errorf("error while installing zalando dependencies: %w", err) } } @@ -513,6 +550,7 @@ func (r *PostgresReconciler) ensureZalandoDependencies(log logr.Logger, ctx cont func (r *PostgresReconciler) updatePodEnvironmentConfigMap(log logr.Logger, ctx context.Context, p *pg.Postgres, b *pg.BackupConfig) error { if b == nil { log.Info("No backupConfig found, skipping configuration of postgres backup") + return nil } @@ -589,10 +627,10 @@ func (r *PostgresReconciler) updatePodEnvironmentConfigMap(log logr.Logger, ctx if err := r.SvcClient.Get(ctx, ns, cm); err != nil { // when updating from v0.7.0 straight to v0.10.0, we neither have that ConfigMap (as we use a Secret in version // v0.7.0) nor do we create it (the new labels aren't there yet, so the selector does not match and - // operatormanager.OperatorManager.UpdateAllManagedOperators does not call InstallOrUpdateOperator) + // operatormanager.UpdateAllManagedOperators does not call InstallOrUpdateOperator) // we previously aborted here (before the postgresql resource was updated with the new labels), meaning we would // simply restart the loop without solving the problem. - if cm, err = r.OperatorManager.CreatePodEnvironmentConfigMap(ctx, ns.Namespace); err != nil { + if cm, err = r.CreatePodEnvironmentConfigMap(ctx, ns.Namespace); err != nil { return fmt.Errorf("error while creating the missing Pod Environment ConfigMap %v: %w", ns.Namespace, err) } log.Info("missing Pod Environment ConfigMap created!") @@ -608,6 +646,7 @@ func (r *PostgresReconciler) updatePodEnvironmentConfigMap(log logr.Logger, ctx func (r *PostgresReconciler) updatePodEnvironmentSecret(log logr.Logger, ctx context.Context, p *pg.Postgres) error { if p.Spec.BackupSecretRef == "" { log.Info("No configured backupSecretRef found, skipping configuration of postgres backup") + return nil } @@ -664,7 +703,7 @@ func (r *PostgresReconciler) updatePodEnvironmentSecret(log logr.Logger, ctx con var s *corev1.Secret ns := p.ToPeripheralResourceNamespace() - if s, err = r.OperatorManager.CreateOrGetPodEnvironmentSecret(ctx, ns); err != nil { + if s, err = r.CreateOrGetPodEnvironmentSecret(ctx, ns); err != nil { return fmt.Errorf("error while accessing the pod environment secret %v: %w", ns, err) } @@ -693,22 +732,26 @@ func (r *PostgresReconciler) getStandbyEnvs(ctx context.Context, p *pg.Postgres) } r.recorder.Eventf(primary, "Warning", "Error", "failed to get referenced primary postgres: %v", err) + return standbyEnvs } if primary.Spec.BackupSecretRef == "" { r.recorder.Eventf(primary, "Warning", "Error", "No backupSecretRef for primary postgres found, skipping configuration of wal_e bootstrapping") + return standbyEnvs } primaryBackupConfig, err := r.getBackupConfig(ctx, primary.Namespace, primary.Spec.BackupSecretRef) if err != nil { r.recorder.Eventf(primary, "Warning", "Error", "failed to get referenced primary backup config, skipping configuration of wal_e bootstrapping: %v", err) + return standbyEnvs } primaryS3url, err := url.Parse(primaryBackupConfig.S3Endpoint) if err != nil { r.recorder.Eventf(primary, "Warning", "Error", "error while parsing the s3 endpoint url in the backup secret: %w", err) + return standbyEnvs } @@ -793,6 +836,7 @@ func (r *PostgresReconciler) createOrUpdateIngressCWNP(log logr.Logger, ctx cont key := &firewall.ClusterwideNetworkPolicy{ObjectMeta: metav1.ObjectMeta{Name: policy.Name, Namespace: policy.Namespace}} if _, err := controllerutil.CreateOrUpdate(ctx, r.SvcClient, key, func() error { key.Spec.Ingress = policy.Spec.Ingress + return nil }); err != nil { return fmt.Errorf("unable to deploy CRD ClusterwideNetworkPolicy: %w", err) @@ -814,6 +858,7 @@ func (r *PostgresReconciler) createOrUpdateIngressCWNP(log logr.Logger, ctx cont key2 := &firewall.ClusterwideNetworkPolicy{ObjectMeta: metav1.ObjectMeta{Name: standbyIngressCWNP.Name, Namespace: standbyIngressCWNP.Namespace}} if _, err := controllerutil.CreateOrUpdate(ctx, r.SvcClient, key2, func() error { key2.Spec.Ingress = standbyIngressCWNP.Spec.Ingress + return nil }); err != nil { return fmt.Errorf("unable to deploy standby ingress ClusterwideNetworkPolicy: %w", err) @@ -843,6 +888,7 @@ func (r *PostgresReconciler) createOrUpdateEgressCWNP(ctx context.Context, in *p key3 := &firewall.ClusterwideNetworkPolicy{ObjectMeta: metav1.ObjectMeta{Name: standbyEgressCWNP.Name, Namespace: standbyEgressCWNP.Namespace}} if _, err := controllerutil.CreateOrUpdate(ctx, r.SvcClient, key3, func() error { key3.Spec.Egress = standbyEgressCWNP.Spec.Egress + return nil }); err != nil { return fmt.Errorf("unable to deploy standby egress ClusterwideNetworkPolicy: %w", err) @@ -958,6 +1004,7 @@ func (r *PostgresReconciler) ensureStandbySecrets(log logr.Logger, ctx context.C if err == nil { log.V(debugLogLevel).Info("local monitoring secret found, no action needed") + return nil } else if !apierrors.IsNotFound(err) { // we got an error other than not found, so we cannot continue! @@ -967,9 +1014,10 @@ func (r *PostgresReconciler) ensureStandbySecrets(log logr.Logger, ctx context.C log.Info("not all expected local secrets found, continuing to create them") remoteSecretNamespacedName := types.NamespacedName{ - Namespace: instance.ObjectMeta.Namespace, + Namespace: instance.Namespace, Name: instance.Spec.PostgresConnection.ConnectionSecretName, } + return r.copySecrets(log, ctx, remoteSecretNamespacedName, instance, false) } @@ -994,6 +1042,7 @@ func (r *PostgresReconciler) ensureCloneSecrets(log logr.Logger, ctx context.Con if err == nil { log.V(debugLogLevel).Info("local postgres secret found, no action needed") + return nil } @@ -1006,9 +1055,10 @@ func (r *PostgresReconciler) ensureCloneSecrets(log logr.Logger, ctx context.Con remoteSecretName := strings.Replace(instance.ToUserPasswordsSecretName(), instance.Name, instance.Spec.PostgresRestore.SourcePostgresID, 1) // TODO this is hacky-wacky... remoteSecretNamespacedName := types.NamespacedName{ - Namespace: instance.ObjectMeta.Namespace, + Namespace: instance.Namespace, Name: remoteSecretName, } + return r.copySecrets(log, ctx, remoteSecretNamespacedName, instance, true) } @@ -1047,8 +1097,10 @@ func (r *PostgresReconciler) copySecrets(log logr.Logger, ctx context.Context, s if err := r.SvcClient.Create(ctx, postgresSecret); err != nil { if apierrors.IsAlreadyExists(err) { log.Info("local postgres secret already exists, skipping", "name", currentSecretName) + continue } + return fmt.Errorf("error while creating local secrets in service cluster: %w", err) } } @@ -1069,6 +1121,7 @@ func (r *PostgresReconciler) checkAndUpdatePatroniReplicationConfig(log logr.Log leaderPods, err := r.findLeaderPods(log, ctx, instance) if err != nil { log.V(debugLogLevel).Info("could not query pods, requeuing") + return requeueAfterReconcile, err } @@ -1082,6 +1135,7 @@ func (r *PostgresReconciler) checkAndUpdatePatroniReplicationConfig(log logr.Log // If there is no connected postgres, we still need to possibly clean up a former synchronous primary if instance.Spec.PostgresConnection == nil { log.V(debugLogLevel).Info("single instance, updating with empty config and requeing") + return allDone, r.httpPatchPatroni(log, ctx, instance, leaderIP, nil) } @@ -1089,16 +1143,19 @@ func (r *PostgresReconciler) checkAndUpdatePatroniReplicationConfig(log logr.Log resp, err = r.httpGetPatroniConfig(log, ctx, leaderIP) if err != nil { log.V(debugLogLevel).Info("could not query patroni, requeuing") + return requeueAfterReconcile, err } if resp == nil { log.V(debugLogLevel).Info("got nil response from patroni, requeuing") + return requeueAfterReconcile, nil } if instance.IsReplicationPrimaryOrStandalone() { if resp.StandbyCluster != nil { log.V(debugLogLevel).Info("standby_cluster mismatch, requeing", "response", resp) + return requeueAfterReconcile, nil } if instance.Spec.PostgresConnection.SynchronousReplication { @@ -1119,14 +1176,17 @@ func (r *PostgresReconciler) checkAndUpdatePatroniReplicationConfig(log logr.Log // compare the actual value with the expected value if synchronousStandbyApplicationName == nil { log.V(debugLogLevel).Info("could not fetch synchronous_nodes_additional, disabling sync replication and requeing", "response", resp) + return requeueAfterReconcile, r.httpPatchPatroni(log, ctx, instance, leaderIP, nil) } else if resp.SynchronousNodesAdditional == nil || *resp.SynchronousNodesAdditional != *synchronousStandbyApplicationName { log.V(debugLogLevel).Info("synchronous_nodes_additional mismatch, updating and requeing", "response", resp) + return requeueAfterReconcile, r.httpPatchPatroni(log, ctx, instance, leaderIP, synchronousStandbyApplicationName) } } else { if resp.SynchronousNodesAdditional != nil { log.V(debugLogLevel).Info("synchronous_nodes_additional mismatch, updating and requeing", "response", resp) + return requeueAfterReconcile, r.httpPatchPatroni(log, ctx, instance, leaderIP, nil) } } @@ -1134,31 +1194,38 @@ func (r *PostgresReconciler) checkAndUpdatePatroniReplicationConfig(log logr.Log } else { if resp.StandbyCluster == nil { log.V(debugLogLevel).Info("standby_cluster mismatch, requeing", "response", resp) + return requeueAfterReconcile, nil } if resp.StandbyCluster.CreateReplicaMethods == nil { log.V(debugLogLevel).Info("create_replica_methods mismatch, updating and requeing", "response", resp) + return requeueAfterReconcile, r.httpPatchPatroni(log, ctx, instance, leaderIP, nil) } if resp.StandbyCluster.Host != instance.Spec.PostgresConnection.ConnectionIP { log.V(debugLogLevel).Info("host mismatch, updating and requeing", "updating", resp) + return requeueAfterReconcile, r.httpPatchPatroni(log, ctx, instance, leaderIP, nil) } if resp.StandbyCluster.Port != int(instance.Spec.PostgresConnection.ConnectionPort) { log.V(debugLogLevel).Info("port mismatch, updating and requeing", "updating", resp) + return requeueAfterReconcile, r.httpPatchPatroni(log, ctx, instance, leaderIP, nil) } if resp.StandbyCluster.ApplicationName != instance.ToPeripheralResourceName() { log.V(debugLogLevel).Info("application_name mismatch, updating and requeing", "response", resp) + return requeueAfterReconcile, r.httpPatchPatroni(log, ctx, instance, leaderIP, nil) } if resp.SynchronousNodesAdditional != nil { log.V(debugLogLevel).Info("synchronous_nodes_additional mismatch, updating and requeing", "response", resp) + return requeueAfterReconcile, r.httpPatchPatroni(log, ctx, instance, leaderIP, nil) } } log.V(debugLogLevel).Info("replication config from Patroni API up to date") + return allDone, nil } @@ -1167,6 +1234,7 @@ func (r *PostgresReconciler) findLeaderPods(log logr.Logger, ctx context.Context roleReq, err := labels.NewRequirement(pg.SpiloRoleLabelName, selection.In, []string{pg.SpiloRoleLabelValueMaster, pg.SpiloRoleLabelValueStandbyLeader}) if err != nil { log.V(debugLogLevel).Info("could not create requirements for label selector to query pods, requeuing") + return leaderPods, err } leaderSelector := labels.NewSelector().Add(*roleReq) @@ -1174,6 +1242,7 @@ func (r *PostgresReconciler) findLeaderPods(log logr.Logger, ctx context.Context client.InNamespace(instance.ToPeripheralResourceNamespace()), client.MatchingLabelsSelector{Selector: leaderSelector}, } + return leaderPods, r.SvcClient.List(ctx, leaderPods, opts...) } @@ -1185,11 +1254,13 @@ func (r *PostgresReconciler) updatePatroniReplicationConfigOnAllPods(log logr.Lo } if err := r.SvcClient.List(ctx, pods, opts...); err != nil { log.V(debugLogLevel).Info("could not query pods, requeuing") + return err } if len(pods.Items) == 0 { log.V(debugLogLevel).Info("no spilo pods found at all, requeuing") + return errors.New("no spilo pods found at all") } else if len(pods.Items) < int(instance.Spec.NumberOfInstances) { log.V(debugLogLevel).Info("unexpected number of pods (might be ok if it is still creating)") @@ -1198,7 +1269,6 @@ func (r *PostgresReconciler) updatePatroniReplicationConfigOnAllPods(log logr.Lo // iterate all spilo pods var lastErr error for _, pod := range pods.Items { - pod := pod // pin! podIP := pod.Status.PodIP if err := r.httpPatchPatroni(log, ctx, instance, podIP, nil); err != nil { lastErr = err @@ -1207,9 +1277,11 @@ func (r *PostgresReconciler) updatePatroniReplicationConfigOnAllPods(log logr.Lo } if lastErr != nil { log.V(debugLogLevel).Info("updating patroni config failed, got one or more errors") + return lastErr } log.V(debugLogLevel).Info("updating patroni config succeeded") + return nil } @@ -1266,6 +1338,7 @@ func (r *PostgresReconciler) httpPatchPatroni(log logr.Logger, ctx context.Conte jsonReq, err := json.Marshal(request) if err != nil { log.V(debugLogLevel).Info("could not create config") + return err } @@ -1275,6 +1348,7 @@ func (r *PostgresReconciler) httpPatchPatroni(log logr.Logger, ctx context.Conte req, err := http.NewRequestWithContext(ctx, http.MethodPatch, url, bytes.NewBuffer(jsonReq)) if err != nil { log.Error(err, "could not create PATCH request") + return err } req.Header.Set("Content-Type", "application/json") @@ -1282,6 +1356,7 @@ func (r *PostgresReconciler) httpPatchPatroni(log logr.Logger, ctx context.Conte resp, err := httpClient.Do(req) if err != nil { log.Error(err, "could not perform PATCH request") + return err } defer resp.Body.Close() @@ -1289,6 +1364,7 @@ func (r *PostgresReconciler) httpPatchPatroni(log logr.Logger, ctx context.Conte if resp.StatusCode/100 != 2 { err = fmt.Errorf("received unexpected return code %d", resp.StatusCode) log.Error(err, "could not perform PATCH request") + return err } @@ -1318,9 +1394,10 @@ func (r *PostgresReconciler) httpGetPatroniConfig(log logr.Logger, ctx context.C httpClient := &http.Client{} url := "http://" + podIP + ":" + podPort + "/" + path - req, err := http.NewRequestWithContext(ctx, "GET", url, nil) + req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) if err != nil { log.Error(err, "could not create GET request") + return nil, err } req.Header.Set("Content-Type", "application/json") @@ -1328,6 +1405,7 @@ func (r *PostgresReconciler) httpGetPatroniConfig(log logr.Logger, ctx context.C resp, err := httpClient.Do(req) if err != nil { log.Error(err, "could not perform GET request") + return nil, err } @@ -1336,12 +1414,14 @@ func (r *PostgresReconciler) httpGetPatroniConfig(log logr.Logger, ctx context.C body, err := io.ReadAll(resp.Body) if err != nil { log.Info("could not read body") + return nil, err } var jsonResp PatroniConfig err = json.Unmarshal(body, &jsonResp) if err != nil { log.V(debugLogLevel).Info("could not parse config response") + return nil, err } @@ -1370,6 +1450,7 @@ func (r *PostgresReconciler) getBackupConfig(ctx context.Context, ns, name strin if err != nil { return nil, fmt.Errorf("unable to unmarshal backupconfig:%w", err) } + return &backupConfig, nil } @@ -1502,6 +1583,7 @@ func (r *PostgresReconciler) createOrUpdateNetPol(ctx context.Context, instance np := &networkingv1.NetworkPolicy{ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: namespace}} if _, err := controllerutil.CreateOrUpdate(ctx, r.SvcClient, np, func() error { np.Spec = spec + return nil }); err != nil { return fmt.Errorf("unable to deploy NetworkPolicy: %w", err) @@ -1545,7 +1627,7 @@ func (r *PostgresReconciler) createOrUpdateExporterSidecarServices(log logr.Logg pes.Spec.Ports = []corev1.ServicePort{ { Name: postgresExporterServicePortName, - Port: int32(pesPort), //nolint + Port: int32(pesPort), Protocol: corev1.ProtocolTCP, TargetPort: intstr.FromInt(int(pesTargetPort)), }, @@ -1565,11 +1647,12 @@ func (r *PostgresReconciler) createOrUpdateExporterSidecarServices(log logr.Logg if err := r.SvcClient.Get(ctx, ns, old); err == nil { // service exists, overwriting it (but using the same clusterip) pes.Spec.ClusterIP = old.Spec.ClusterIP - pes.ObjectMeta.ResourceVersion = old.GetObjectMeta().GetResourceVersion() + pes.ResourceVersion = old.GetObjectMeta().GetResourceVersion() if err := r.SvcClient.Update(ctx, pes); err != nil { return fmt.Errorf("error while updating the postgres-exporter service: %w", err) } log.V(debugLogLevel).Info("postgres-exporter service updated") + return nil } // todo: handle errors other than `NotFound` @@ -1592,8 +1675,10 @@ func (r *PostgresReconciler) deleteNetPol(ctx context.Context, instance *pg.Post if apierrors.IsNotFound(err) { return nil } + return fmt.Errorf("unable to delete NetworkPolicy %v: %w", netpol.Name, err) } + return nil } @@ -1648,11 +1733,12 @@ func (r *PostgresReconciler) createOrUpdateExporterSidecarServiceMonitor(log log old := &coreosv1.ServiceMonitor{} if err := r.SvcClient.Get(ctx, ns, old); err == nil { // Copy the resource version - pesm.ObjectMeta.ResourceVersion = old.ObjectMeta.ResourceVersion + pesm.ResourceVersion = old.ResourceVersion if err := r.SvcClient.Update(ctx, pesm); err != nil { return fmt.Errorf("error while updating the postgres-exporter servicemonitor: %w", err) } log.V(debugLogLevel).Info("postgres-exporter servicemonitor updated") + return nil } // todo: handle errors other than `NotFound` @@ -1713,11 +1799,12 @@ func (r *PostgresReconciler) createOrUpdatePatroniPodMonitor(ctx context.Context old := &coreosv1.PodMonitor{} if err := r.SvcClient.Get(ctx, ns, old); err == nil { // Copy the resource version - pm.ObjectMeta.ResourceVersion = old.ObjectMeta.ResourceVersion + pm.ResourceVersion = old.ResourceVersion if err := r.SvcClient.Update(ctx, pm); err != nil { return fmt.Errorf("error while updating the podmonitor: %w", err) } log.Info("pod monitor updated") + return nil } @@ -1767,6 +1854,7 @@ func (r *PostgresReconciler) ensureStorageEncryptionSecret(log logr.Logger, ctx if !r.EnableRandomStorageEncryptionSecret { log.V(debugLogLevel).Info("storage secret disabled, no action needed") + return nil } @@ -1778,6 +1866,7 @@ func (r *PostgresReconciler) ensureStorageEncryptionSecret(log logr.Logger, ctx err := r.SvcClient.Get(ctx, types.NamespacedName{Namespace: ns, Name: n}, s) if err == nil { log.V(debugLogLevel).Info("storage secret found, no action needed") + return nil } @@ -1819,7 +1908,7 @@ func (r *PostgresReconciler) ensureStorageEncryptionSecret(log logr.Logger, ctx func (r *PostgresReconciler) generateRandomString() (string, error) { const chars string = "!\"#$%&'()*+,-./0123456789:;<=>?@ABCDEFGHIJKLMNOPQRSTUVWXYZ[\\]^_`abcdefghijklmnopqrstuvwxyz{|}~" - var size *big.Int = big.NewInt(int64(len(chars))) + size := big.NewInt(int64(len(chars))) b := make([]byte, 64) for i := range b { x, err := rand.Int(rand.Reader, size) @@ -1828,6 +1917,7 @@ func (r *PostgresReconciler) generateRandomString() (string, error) { } b[i] = chars[x.Int64()] } + return string(b), nil } @@ -1842,6 +1932,7 @@ func (r *PostgresReconciler) removeStorageEncryptionSecretFinalizer(log logr.Log if err != nil { if apierrors.IsNotFound(err) { log.V(debugLogLevel).Info("storage secret not found, nothing to do", "name", n) + return nil } // this would be blocking if we couldn't remove the finalizer, so we should keep trying! @@ -1849,12 +1940,13 @@ func (r *PostgresReconciler) removeStorageEncryptionSecretFinalizer(log logr.Log } // Remove finalizer - s.ObjectMeta.Finalizers = removeElem(s.ObjectMeta.Finalizers, storageEncryptionKeyFinalizerName) + s.Finalizers = removeElem(s.Finalizers, storageEncryptionKeyFinalizerName) if err := r.SvcClient.Update(ctx, s); err != nil { return fmt.Errorf("error while removing finalizer from storage secret in service cluster: %w", err) } log.V(debugLogLevel).Info("finalizer removed from storage secret", "name", n) + return nil } @@ -1865,6 +1957,7 @@ func removeElem(ss []string, s string) (out []string) { } out = append(out, elem) } + return } @@ -1877,6 +1970,7 @@ func (r *PostgresReconciler) ensureInitDBJob(log logr.Logger, ctx context.Contex if err := r.SvcClient.Get(ctx, ns, cm); err == nil { // configmap already exists, nothing to do here log.V(debugLogLevel).Info("initdb ConfigMap already exists") + return nil } @@ -1912,6 +2006,7 @@ func (r *PostgresReconciler) ensureInitDBJob(log logr.Logger, ctx context.Contex if instance.IsReplicationTarget() || instance.Spec.PostgresRestore != nil { log.V(debugLogLevel).Info("initdb job not required") + return nil } @@ -1921,6 +2016,7 @@ func (r *PostgresReconciler) ensureInitDBJob(log logr.Logger, ctx context.Contex if err := r.SvcClient.Get(ctx, ns, j); err == nil { // job already exists, nothing to do here log.V(debugLogLevel).Info("initdb Job already exists") + return nil // TODO return or update? } @@ -2013,6 +2109,7 @@ func (r *PostgresReconciler) ensureInitDBJob(log logr.Logger, ctx context.Contex func (r *PostgresReconciler) createOrUpdateCertificate(log logr.Logger, ctx context.Context, instance *pg.Postgres) error { if r.TLSClusterIssuer == "" { log.V(debugLogLevel).Info("certificate skipped") + return nil } @@ -2032,12 +2129,14 @@ func (r *PostgresReconciler) createOrUpdateCertificate(log logr.Logger, ctx cont Name: r.TLSClusterIssuer, }, } + return nil }); err != nil { return fmt.Errorf("unable to create or update certificate: %w", err) } log.V(debugLogLevel).Info("certificate created or updated") + return nil } @@ -2045,6 +2144,7 @@ func (r *PostgresReconciler) createOrUpdateCertificate(log logr.Logger, ctx cont func (r *PostgresReconciler) createOrUpdateWalGExporterDeployment(log logr.Logger, ctx context.Context, namespace string, instance *pg.Postgres, b *pg.BackupConfig) error { if b == nil { log.Info("No backupConfig found, skipping configuration of wa-l-exporter") + return nil } @@ -2196,11 +2296,12 @@ func (r *PostgresReconciler) createOrUpdateWalGExporterDeployment(log logr.Logge old := &appsv1.Deployment{} if err := r.SvcClient.Get(ctx, ns, old); err == nil { // Copy the resource version - deploy.ObjectMeta.ResourceVersion = old.ObjectMeta.ResourceVersion + deploy.ResourceVersion = old.ResourceVersion if err := r.SvcClient.Update(ctx, deploy); err != nil { return fmt.Errorf("error while updating the wal-g-exporter deployment: %w", err) } log.Info("wal-g-exporter deployment updated") + return nil } @@ -2222,8 +2323,10 @@ func (r *PostgresReconciler) deleteWalGExporterDeployment(ctx context.Context, n if apierrors.IsNotFound(err) { return nil } + return fmt.Errorf("error while deleting the wal-g-exporter deployment: %w", err) } + return nil } @@ -2269,11 +2372,12 @@ func (r *PostgresReconciler) createOrUpdateWalGExporterPodMonitor(log logr.Logge old := &coreosv1.PodMonitor{} if err := r.SvcClient.Get(ctx, ns, old); err == nil { // podMonitor exists, overwriting it - s.ObjectMeta.ResourceVersion = old.GetObjectMeta().GetResourceVersion() + s.ResourceVersion = old.GetObjectMeta().GetResourceVersion() if err := r.SvcClient.Update(ctx, s); err != nil { return fmt.Errorf("error while updating the wal-g-exporter podMonitor: %w", err) } log.V(debugLogLevel).Info("wal-g-exporter podMonitor updated") + return nil } // todo: handle errors other than `NotFound` @@ -2296,7 +2400,9 @@ func (r *PostgresReconciler) deleteWalGExporterPodMonitor(ctx context.Context, n if apierrors.IsNotFound(err) { return nil } + return fmt.Errorf("error while deleting the wal-g-exporter podMonitor: %w", err) } + return nil } diff --git a/controllers/postgres_controller_test.go b/controllers/postgres_controller_test.go index 00a77b48..bf7b0d0c 100644 --- a/controllers/postgres_controller_test.go +++ b/controllers/postgres_controller_test.go @@ -32,6 +32,7 @@ var _ = Describe("postgres controller", func() { if len(instance.Finalizers) == 0 { return false } + return instance.Finalizers[0] == pg.PostgresFinalizerName }, timeout, interval).Should(BeTrue()) }) diff --git a/controllers/status_controller.go b/controllers/status_controller.go index ffb61309..def74b1b 100644 --- a/controllers/status_controller.go +++ b/controllers/status_controller.go @@ -42,7 +42,7 @@ type StatusReconciler struct { // +kubebuilder:rbac:groups=acid.zalan.do,resources=postgresqls,verbs=get;list;watch;create;update;patch;delete // +kubebuilder:rbac:groups=acid.zalan.do,resources=postgresqls/status,verbs=get;update;patch func (r *StatusReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { - log := r.Log.WithValues("ns", req.NamespacedName.Namespace) + log := r.Log.WithValues("ns", req.Namespace) log.V(debugLogLevel).Info("fetching postgresql") instance := &zalando.Postgresql{} @@ -51,6 +51,7 @@ func (r *StatusReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctr return ctrl.Result{}, err } log.Info("status changed to Deleted") + return ctrl.Result{}, nil } @@ -82,13 +83,15 @@ func (r *StatusReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctr // update the status of the remote object owner.Status.Description = instance.Status.PostgresClusterStatus // update the reference to the zalando instance in the remote object - owner.Status.ChildName = instance.ObjectMeta.Name + owner.Status.ChildName = instance.Name log.V(debugLogLevel).Info("Updating owner", "owner", owner.UID) if err := r.CtrlClient.Status().Update(ctx, owner); err != nil { log.Error(err, "failed to update owner object") + return err } + return nil }) @@ -184,6 +187,7 @@ func (r *StatusReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctr if len(secrets.Items) == 0 { log.Info("no local secrets found yet, requeuing", "status", owner.Status) + return ctrl.Result{Requeue: true, RequeueAfter: 2 * time.Second}, nil } @@ -230,6 +234,7 @@ func (r *StatusReconciler) createOrUpdateSecret(ctx context.Context, in *pg.Post // todo: update to CreateOrPatch() result, err := controllerutil.CreateOrUpdate(ctx, r.CtrlClient, fetched, func() error { fetched.Data = secret.Data + return nil }) if err != nil { @@ -243,9 +248,10 @@ func (r *StatusReconciler) createOrUpdateSecret(ctx context.Context, in *pg.Post // Extract the UID of the owner object by reading the value of a certain label func deriveOwnerName(instance *zalando.Postgresql) (string, error) { - value, ok := instance.ObjectMeta.Labels[pg.NameLabelName] + value, ok := instance.Labels[pg.NameLabelName] if !ok { return "", fmt.Errorf("could not derive owner reference") } + return value, nil } diff --git a/main.go b/main.go index b42e7b46..8a93f717 100644 --- a/main.go +++ b/main.go @@ -344,10 +344,7 @@ func main() { enableSuperUserForDBO = viper.GetBool(enableSuperUserForDBOFlg) tlsClusterIssuer = viper.GetString(tlsClusterIssuerFlg) - enableCustomTLSCert := false - if tlsClusterIssuer != "" { - enableCustomTLSCert = true - } + enableCustomTLSCert := tlsClusterIssuer != "" tlsSubDomain = viper.GetString(tlsSubDomainFlg) viper.SetDefault(enablePatroniFailsafeModeFlg, true) @@ -480,7 +477,7 @@ func main() { os.Exit(1) } - var etcdMgrOpts etcdmanager.Options = etcdmanager.Options{ + etcdMgrOpts := etcdmanager.Options{ EtcdImage: etcdImage, EtcdBackupSidecarImage: etcdBackupSidecarImage, SecretKeyRefName: etcdBackupSecretName, @@ -505,7 +502,7 @@ func main() { } } - var opMgrOpts operatormanager.Options = operatormanager.Options{ + opMgrOpts := operatormanager.Options{ PspName: pspName, OperatorImage: operatorImage, DockerImage: postgresImage, @@ -528,7 +525,7 @@ func main() { os.Exit(1) } - var lbMgrOpts lbmanager.Options = lbmanager.Options{ + lbMgrOpts := lbmanager.Options{ LBIP: lbIP, PortRangeStart: portRangeStart, PortRangeSize: portRangeSize, diff --git a/pkg/etcdmanager/etcdmanager.go b/pkg/etcdmanager/etcdmanager.go index 505fb482..766d5c02 100644 --- a/pkg/etcdmanager/etcdmanager.go +++ b/pkg/etcdmanager/etcdmanager.go @@ -82,6 +82,7 @@ func New(confRest *rest.Config, fileName string, scheme *runtime.Scheme, log log } log.Info("new `EtcdManager` created") + return &EtcdManager{ metadataAccessor: meta.NewAccessor(), client: client, @@ -121,6 +122,7 @@ func (m *EtcdManager) InstallOrUpdateEtcd() error { } m.log.Info("etcd installed") + return nil } @@ -170,19 +172,18 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje case *corev1.ServiceAccount: m.log.Info("handling ServiceAccount") - v.ObjectMeta.Name = saName + v.Name = saName // Use the updated name to get the resource - key.Name = v.ObjectMeta.Name + key.Name = v.Name err = m.client.Get(ctx, key, &corev1.ServiceAccount{}) case *rbacv1.Role: m.log.Info("handling Role") - v.ObjectMeta.Name = roleName + v.Name = roleName m.log.Info("Updating psp") for i := range v.Rules { - i := i if !slices.Contains(v.Rules[i].APIGroups, "extensions") { continue @@ -198,12 +199,12 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje } // Use the updated name to get the resource - key.Name = v.ObjectMeta.Name + key.Name = v.Name err = m.client.Get(ctx, key, &rbacv1.Role{}) case *rbacv1.RoleBinding: m.log.Info("handling RoleBinding") - v.ObjectMeta.Name = rbName + v.Name = rbName m.log.Info("Updating roleRef") v.RoleRef.Name = roleName @@ -217,12 +218,12 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje } // Use the updated name to get the resource - key.Name = v.ObjectMeta.Name + key.Name = v.Name err = m.client.Get(ctx, key, &rbacv1.RoleBinding{}) case *corev1.ConfigMap: m.log.Info("handling ConfigMap") - v.ObjectMeta.Name = cmName + v.Name = cmName var configYaml strings.Builder configYaml.WriteString("db: etcd\n") @@ -236,16 +237,15 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje v.Data["config.yaml"] = configYaml.String() // Use the updated name to get the resource - key.Name = v.ObjectMeta.Name + key.Name = v.Name err = m.client.Get(ctx, key, &corev1.ConfigMap{}) case *appsv1.StatefulSet: m.log.Info("handling StatefulSet") - v.ObjectMeta.Name = stsName + v.Name = stsName m.log.Info("Updating containers") for i := range v.Spec.Template.Spec.Containers { - i := i // Patch EtcdImage if m.options.EtcdImage != "" { @@ -256,8 +256,7 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje m.log.Info("Updating envs") // Patch Env for j, env := range v.Spec.Template.Spec.Containers[i].Env { - j := j - env := env + switch env.Name { case "BACKUP_RESTORE_SIDECAR_S3_BUCKET_NAME": if env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil { @@ -293,7 +292,6 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje if m.options.EtcdBackupSidecarImage != "" { m.log.Info("Updating initContainers") for i := range v.Spec.Template.Spec.InitContainers { - i := i m.log.Info("Updating etcd backup sidecar image") v.Spec.Template.Spec.InitContainers[i].Image = m.options.EtcdBackupSidecarImage @@ -302,7 +300,6 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje m.log.Info("Updating configMap volume") for i := range v.Spec.Template.Spec.Volumes { - i := i if v.Spec.Template.Spec.Volumes[i].Name != "backup-restore-sidecar-config" { continue @@ -312,10 +309,10 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje m.log.Info("Updating labels") // Add partition ID label - v.Spec.Template.ObjectMeta.Labels[pg.PartitionIDLabelName] = m.options.PartitionID - v.Spec.Template.ObjectMeta.Labels[pg.ManagedByLabelName] = m.options.PostgresletFullname - v.Spec.Template.ObjectMeta.Labels[etcdComponentLabelName] = etcdComponentLabelValue - v.Spec.Template.ObjectMeta.Labels[pg.NameLabelName] = stsName + v.Spec.Template.Labels[pg.PartitionIDLabelName] = m.options.PartitionID + v.Spec.Template.Labels[pg.ManagedByLabelName] = m.options.PostgresletFullname + v.Spec.Template.Labels[etcdComponentLabelName] = etcdComponentLabelValue + v.Spec.Template.Labels[pg.NameLabelName] = stsName m.log.Info("Updating selector") // spec.selector.matchLabels @@ -334,29 +331,29 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje got := appsv1.StatefulSet{} // Use the updated name to get the resource - key.Name = v.ObjectMeta.Name + key.Name = v.Name err = m.client.Get(ctx, key, &got) if err == nil { // Copy the ResourceVersion m.log.Info("Copying existing resource version") - v.ObjectMeta.ResourceVersion = got.ObjectMeta.ResourceVersion + v.ResourceVersion = got.ResourceVersion } case *corev1.Service: m.log.Info("handling Service") - switch v.ObjectMeta.Name { + switch v.Name { case "backup-restore-sidecar-svc": - v.ObjectMeta.Name = svcSidecarName + v.Name = svcSidecarName case "etcd-psql-headless": - v.ObjectMeta.Name = svcHeadlessName + v.Name = svcHeadlessName case "etcd-psql": - v.ObjectMeta.Name = svcName + v.Name = svcName default: return fmt.Errorf("unknown service name: %v", v.ObjectMeta.Name) } m.log.Info("Updating labels") - v.ObjectMeta.Labels[pg.NameLabelName] = v.ObjectMeta.Name + v.Labels[pg.NameLabelName] = v.Name m.log.Info("Updating selector") v.Spec.Selector[pg.PartitionIDLabelName] = m.options.PartitionID @@ -364,11 +361,11 @@ func (m *EtcdManager) createNewClientObject(ctx context.Context, obj client.Obje v.Spec.Selector[pg.NameLabelName] = stsName got := corev1.Service{} - key.Name = v.ObjectMeta.Name + key.Name = v.Name err = m.client.Get(ctx, key, &got) if err == nil { // Copy the ResourceVersion - v.ObjectMeta.ResourceVersion = got.ObjectMeta.ResourceVersion + v.ResourceVersion = got.ResourceVersion // Copy the ClusterIP v.Spec.ClusterIP = got.Spec.ClusterIP } @@ -425,6 +422,7 @@ func (m *EtcdManager) toObjectKey(obj runtime.Object, namespace string) (client. if err != nil { return client.ObjectKey{}, fmt.Errorf("error while extracting the name of the k8s resource: %w", err) } + return client.ObjectKey{ Namespace: namespace, Name: name, @@ -475,11 +473,12 @@ func (m *EtcdManager) createOrUpdateServiceMonitor(ctx context.Context, targetNa old := &coreosv1.ServiceMonitor{} if err := m.client.Get(ctx, nsn, old); err == nil { // Copy the resource version - sm.ObjectMeta.ResourceVersion = old.ObjectMeta.ResourceVersion + sm.ResourceVersion = old.ResourceVersion if err := m.client.Update(ctx, sm); err != nil { return fmt.Errorf("error while updating the servicemonitor: %w", err) } m.log.Info("servicemonitor updated") + return nil } diff --git a/pkg/lbmanager/lbmanager.go b/pkg/lbmanager/lbmanager.go index 32fac43b..6cc805b2 100644 --- a/pkg/lbmanager/lbmanager.go +++ b/pkg/lbmanager/lbmanager.go @@ -54,6 +54,7 @@ func (m *LBManager) ReconcileSvcLBs(ctx context.Context, in *api.Postgres) error if len(errs) > 0 { return errors.Join(errs...) } + return nil } @@ -65,6 +66,7 @@ func (m *LBManager) CreateOrUpdateSharedSvcLB(ctx context.Context, in *api.Postg if err != nil { m.log.Info("could not delete dedicated loadbalancer", "ns", in.Namespace, "pgID", in.Name) } + return nil } @@ -101,6 +103,7 @@ func (m *LBManager) CreateOrUpdateSharedSvcLB(ctx context.Context, in *api.Postg if err := m.client.Create(ctx, svc); err != nil { return fmt.Errorf("failed to create Service of type LoadBalancer: %w", err) } + return nil } @@ -116,7 +119,7 @@ func (m *LBManager) CreateOrUpdateSharedSvcLB(ctx context.Context, in *api.Postg svc.Spec.LoadBalancerSourceRanges = []string{} } // also update the annotations for our custom tls certs - svc.ObjectMeta.Annotations = updated.ObjectMeta.Annotations + svc.Annotations = updated.Annotations if err := m.client.Update(ctx, svc); err != nil { return fmt.Errorf("failed to update Service of type LoadBalancer (shared): %w", err) @@ -133,6 +136,7 @@ func (m *LBManager) CreateOrUpdateDedicatedSvcLB(ctx context.Context, in *api.Po if err != nil { m.log.Info("could not delete dedicated loadbalancer", "ns", in.Namespace, "pgID", in.Name) } + return nil } @@ -140,7 +144,7 @@ func (m *LBManager) CreateOrUpdateDedicatedSvcLB(ctx context.Context, in *api.Po if in.Spec.DedicatedLoadBalancerPort != nil && *in.Spec.DedicatedLoadBalancerPort != 0 { nextFreePort = *in.Spec.DedicatedLoadBalancerPort } - var lbIPToUse string = *in.Spec.DedicatedLoadBalancerIP + lbIPToUse := *in.Spec.DedicatedLoadBalancerIP sharedSvcLbAlsoEnabled := in.EnableSharedSVCLB(m.options.EnableForceSharedIP) @@ -162,6 +166,7 @@ func (m *LBManager) CreateOrUpdateDedicatedSvcLB(ctx context.Context, in *api.Po if err := m.client.Create(ctx, new); err != nil { return fmt.Errorf("failed to create Service of type LoadBalancer: %w", err) } + return nil } @@ -169,7 +174,7 @@ func (m *LBManager) CreateOrUpdateDedicatedSvcLB(ctx context.Context, in *api.Po existing.Spec = new.Spec // also update the annotations for our custom tls certs - existing.ObjectMeta.Annotations = new.ObjectMeta.Annotations + existing.Annotations = new.Annotations if err := m.client.Update(ctx, existing); err != nil { return fmt.Errorf("failed to update Service of type LoadBalancer (dedicated): %w", err) @@ -186,6 +191,7 @@ func (m *LBManager) DeleteSharedSvcLB(ctx context.Context, in *api.Postgres) err if err := m.client.Delete(ctx, lb); client.IgnoreNotFound(err) != nil { return err } + return nil } @@ -197,6 +203,7 @@ func (m *LBManager) DeleteDedicatedSvcLB(ctx context.Context, in *api.Postgres) if err := m.client.Delete(ctx, lb); client.IgnoreNotFound(err) != nil { return err } + return nil } @@ -254,5 +261,6 @@ func containsElem(s []int32, v int32) bool { return true } } + return false } diff --git a/pkg/lbmanager/lbmanager_test.go b/pkg/lbmanager/lbmanager_test.go index f3103a70..7c30f514 100644 --- a/pkg/lbmanager/lbmanager_test.go +++ b/pkg/lbmanager/lbmanager_test.go @@ -89,11 +89,12 @@ func TestLBManager_nextFreePort(t *testing.T) { } for _, tt := range tests { - tt := tt + t.Run(tt.name, func(t *testing.T) { _, portGot, err := tt.lbMgr.nextFreeSocket(context.Background()) if (err != nil) != tt.wantErr { t.Errorf("LBManager.nextFreePort() error = %v, wantErr %v", err, tt.wantErr) + return } if portGot != tt.portWant { @@ -116,6 +117,7 @@ func svcListWithPorts(ports ...int32) *corev1.ServiceList { for _, port := range ports { svcList.Items = append(svcList.Items, *svcWithPort(port)) } + return svcList } @@ -128,5 +130,6 @@ func svcWithPort(port int32) *corev1.Service { Port: port, }, } + return &svc } diff --git a/pkg/operatormanager/operatormanager.go b/pkg/operatormanager/operatormanager.go index 0feeda13..87ae0868 100644 --- a/pkg/operatormanager/operatormanager.go +++ b/pkg/operatormanager/operatormanager.go @@ -110,6 +110,7 @@ func New(confRest *rest.Config, fileName string, scheme *runtime.Scheme, log log } log.Info("new `OperatorManager` created") + return &OperatorManager{ metadataAccessor: meta.NewAccessor(), client: client, @@ -161,6 +162,7 @@ func (m *OperatorManager) InstallOrUpdateOperator(ctx context.Context, namespace } m.log.Info("operator installed", "ns", namespace) + return nil } @@ -174,6 +176,7 @@ func (m *OperatorManager) IsOperatorDeletable(ctx context.Context, namespace str } if len(setList.Items) != 0 { log.Info("statefulset still running") + return false, nil } @@ -185,10 +188,12 @@ func (m *OperatorManager) IsOperatorDeletable(ctx context.Context, namespace str err := m.client.Get(ctx, ns, &corev1.Service{}) if !errors.IsNotFound(err) { log.Info("service still running") + return false, nil } log.Info("operator deletable") + return true, nil } @@ -206,6 +211,7 @@ func (m *OperatorManager) IsOperatorInstalled(ctx context.Context, namespace str return false, nil } m.log.Info("operator is installed", "ns", namespace) + return true, nil } @@ -254,6 +260,7 @@ func (m *OperatorManager) UninstallOperator(ctx context.Context, namespace strin if errors.IsNotFound(err) { return nil } + return fmt.Errorf("error while deleting %v: %w", v, err) } } @@ -370,7 +377,7 @@ func (m *OperatorManager) createNewClientObject(ctx context.Context, obj client. err = m.client.Get(ctx, key, &got) if err == nil { // Copy the ResourceVersion - v.ObjectMeta.ResourceVersion = got.ObjectMeta.ResourceVersion + v.ResourceVersion = got.ResourceVersion // Copy the ClusterIP v.Spec.ClusterIP = got.Spec.ClusterIP } @@ -505,14 +512,14 @@ func (m *OperatorManager) createOrUpdateNamespace(ctx context.Context, namespace // Create the namespace. nsObj := &corev1.Namespace{} nsObj.Name = namespace - nsObj.ObjectMeta.Labels = labels + nsObj.Labels = labels if err := m.client.Create(ctx, nsObj); err != nil { return fmt.Errorf("error while creating namespace %v: %w", namespace, err) } log.Info("namespace created") } else { // update namespace - ns.ObjectMeta.Labels = labels + ns.Labels = labels if err := m.client.Update(ctx, &ns); err != nil { return fmt.Errorf("error while updating namespace: %w", err) } @@ -535,6 +542,7 @@ func (m *OperatorManager) CreatePodEnvironmentConfigMap(ctx context.Context, nam // configmap already exists, nothing to do here // we will update the configmap with the correct S3 config in the postgres controller log.Info("Pod Environment ConfigMap already exists") + return cm, nil } @@ -566,6 +574,7 @@ func (m *OperatorManager) CreateOrGetPodEnvironmentSecret(ctx context.Context, n // secret already exists, nothing to do here // we will update the secret with the correct S3 config in the postgres controller log.Info("Pod Environment Secret already exists") + return s, nil } @@ -594,6 +603,7 @@ func (m *OperatorManager) createOrUpdateSidecarsConfig(ctx context.Context, name if err := m.client.Get(ctx, cns, globalSidecarsCM); err != nil { // configmap with configuration does not exists, nothing we can do here... m.log.Error(err, "could not fetch global config for sidecars", "ns", namespace) + return err } @@ -643,6 +653,7 @@ func (m *OperatorManager) createOrUpdateSidecarsConfigMap(ctx context.Context, n return fmt.Errorf("error while updating the new Sidecars ConfigMap: %w", err) } log.Info("Sidecars ConfigMap updated") + return nil } // todo: handle errors other than `NotFound` @@ -668,6 +679,7 @@ func (m *OperatorManager) deletePodEnvironmentConfigMap(ctx context.Context, nam return fmt.Errorf("error while deleting the Pod Environment ConfigMap: %w", err) } m.log.Info("Pod Environment ConfigMap deleted", "ns", namespace) + return nil } @@ -682,6 +694,7 @@ func (m *OperatorManager) toObjectKey(obj runtime.Object, namespace string) (cli if err != nil { return client.ObjectKey{}, fmt.Errorf("error while extracting the name of the k8s resource: %w", err) } + return client.ObjectKey{ Namespace: namespace, Name: name, @@ -711,5 +724,6 @@ func (m *OperatorManager) UpdateAllManagedOperators(ctx context.Context) error { } m.log.Info("Done updating postgres operators in managed namespaces") + return nil } diff --git a/pkg/webhooks/spiloPodMutator.go b/pkg/webhooks/spiloPodMutator.go index 16071c23..7ef2729e 100644 --- a/pkg/webhooks/spiloPodMutator.go +++ b/pkg/webhooks/spiloPodMutator.go @@ -37,6 +37,7 @@ func (a *SpiloPodMutator) Handle(ctx context.Context, req admission.Request) adm err := a.Decoder.Decode(req, pod) if err != nil { log.Error(err, "failed to decode request") + return admission.Errored(http.StatusBadRequest, err) } @@ -62,13 +63,13 @@ func (a *SpiloPodMutator) Handle(ctx context.Context, req admission.Request) adm LabelSelector: &metav1.LabelSelector{ MatchLabels: map[string]string{ "application": "spilo", - "cluster-name": pod.ObjectMeta.Labels["cluster-name"], - pg.NameLabelName: pod.ObjectMeta.Labels[pg.NameLabelName], - pg.PartitionIDLabelName: pod.ObjectMeta.Labels[pg.PartitionIDLabelName], - pg.ProjectIDLabelName: pod.ObjectMeta.Labels[pg.ProjectIDLabelName], - pg.TenantLabelName: pod.ObjectMeta.Labels[pg.TenantLabelName], - pg.UIDLabelName: pod.ObjectMeta.Labels[pg.UIDLabelName], - "team": pod.ObjectMeta.Labels["team"], + "cluster-name": pod.Labels["cluster-name"], + pg.NameLabelName: pod.Labels[pg.NameLabelName], + pg.PartitionIDLabelName: pod.Labels[pg.PartitionIDLabelName], + pg.ProjectIDLabelName: pod.Labels[pg.ProjectIDLabelName], + pg.TenantLabelName: pod.Labels[pg.TenantLabelName], + pg.UIDLabelName: pod.Labels[pg.UIDLabelName], + "team": pod.Labels["team"], }, }, TopologyKey: a.PodTopologySpreadConstraintTopologyKey, @@ -88,9 +89,11 @@ func (a *SpiloPodMutator) Handle(ctx context.Context, req admission.Request) adm marshaledPod, err := json.Marshal(pod) if err != nil { log.Error(err, "failed to marshal response") + return admission.Errored(http.StatusInternalServerError, err) } log.V(1).Info("done") + return admission.PatchResponseFromRaw(req.Object.Raw, marshaledPod) }