From 63119819abd70b09433b10e1526607ba0b29ddb1 Mon Sep 17 00:00:00 2001 From: yerkennz Date: Fri, 21 Aug 2026 15:49:24 +0500 Subject: [PATCH 1/7] feat: [CPCAP-9364] automatic pvc extension --- operator/pkg/helper/resource_management.go | 230 +++++++++++++++ operator/pkg/patroni/patroni.go | 27 ++ operator/pkg/reconciler/backup_daemon.go | 88 +++++- operator/pkg/reconciler/patroni.go | 326 +++++++++++++++++++-- 4 files changed, 650 insertions(+), 21 deletions(-) diff --git a/operator/pkg/helper/resource_management.go b/operator/pkg/helper/resource_management.go index b335bac0..ede246ac 100644 --- a/operator/pkg/helper/resource_management.go +++ b/operator/pkg/helper/resource_management.go @@ -38,6 +38,7 @@ import ( "k8s.io/apimachinery/pkg/api/equality" "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" + "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" "k8s.io/apimachinery/pkg/util/wait" @@ -536,6 +537,98 @@ func (rm *ResourceManager) CreatePvcIfNotExists(pvc *corev1.PersistentVolumeClai return nil } +func (rm *ResourceManager) CreateOrResizePvc(pvc *corev1.PersistentVolumeClaim) (bool, error) { + foundPvc := &corev1.PersistentVolumeClaim{} + + err := rm.kubeClient.Get( + context.TODO(), + types.NamespacedName{ + Name: pvc.Name, + Namespace: pvc.Namespace, + }, + foundPvc, + ) + + if errors.IsNotFound(err) { + logger.Info(fmt.Sprintf("Creating %s PVC", pvc.Name)) + + if err := rm.kubeClient.Create(context.TODO(), pvc); err != nil { + return false, err + } + + return false, nil + } + + if err != nil { + return false, err + } + + currentSize := foundPvc.Spec.Resources.Requests[corev1.ResourceStorage] + desiredSize := pvc.Spec.Resources.Requests[corev1.ResourceStorage] + + resizeInProgress := false + changed := false + + switch desiredSize.Cmp(currentSize) { + case 1: + logger.Info(fmt.Sprintf("Expanding PVC %s from %s to %s", pvc.Name, currentSize.String(), desiredSize.String())) + + foundPvc.Spec.Resources.Requests[corev1.ResourceStorage] = desiredSize + + resizeInProgress = true + changed = true + + case 0: + capacity := foundPvc.Status.Capacity[corev1.ResourceStorage] + + if !capacity.IsZero() && capacity.Cmp(desiredSize) < 0 { + logger.Info(fmt.Sprintf("PVC %s resize is still in progress: requested=%s capacity=%s", pvc.Name, desiredSize.String(), capacity.String())) + + resizeInProgress = true + } + + case -1: + return false, fmt.Errorf( + "PVC %s shrinking from %s to %s is not supported", + pvc.Name, + currentSize.String(), + desiredSize.String(), + ) + } + + // Preserve existing PVC behavior. + if len(foundPvc.OwnerReferences) > 0 { + foundPvc.OwnerReferences = nil + changed = true + } + + // Apply desired annotations only when they differ. + if pvc.Annotations != nil { + if foundPvc.Annotations == nil { + foundPvc.Annotations = make(map[string]string) + } + + for key, value := range pvc.Annotations { + if foundPvc.Annotations[key] != value { + foundPvc.Annotations[key] = value + changed = true + } + } + } + + // Nothing in PVC spec/metadata changed. + // Do not perform unnecessary Kubernetes Update(). + if !changed { + return resizeInProgress, nil + } + + if err := rm.kubeClient.Update(context.TODO(), foundPvc); err != nil { + return false, err + } + + return resizeInProgress, nil +} + func (rm *ResourceManager) CreateSecretIfNotExists(secret *corev1.Secret) error { foundSecret := &corev1.Secret{} err := rm.kubeClient.Get(context.TODO(), types.NamespacedName{ @@ -1075,3 +1168,140 @@ func (rm *ResourceManager) commonLabels(name string) map[string]string { return labels } + +func (rm *ResourceManager) WaitForPvcResizeState(pvcName string, namespace string, desiredSize resource.Quantity) (bool, error) { + restartRequired := false + + err := wait.PollUntilContextTimeout(context.Background(), time.Second, 2*time.Minute, true, + func(ctx context.Context) (bool, error) { + currentPvc := &corev1.PersistentVolumeClaim{} + + if err := rm.kubeClient.Get( + ctx, + types.NamespacedName{ + Name: pvcName, + Namespace: namespace, + }, + currentPvc, + ); err != nil { + return false, err + } + + capacity := currentPvc.Status.Capacity[corev1.ResourceStorage] + + if capacity.Cmp(desiredSize) >= 0 { + return true, nil + } + + for _, condition := range currentPvc.Status.Conditions { + if condition.Type == corev1.PersistentVolumeClaimFileSystemResizePending && + condition.Status == corev1.ConditionTrue { + + restartRequired = true + return true, nil + } + } + + return false, nil + }, + ) + + return restartRequired, err +} + +func (rm *ResourceManager) WaitForPodDeletion(podName string, timeout time.Duration) error { + return wait.PollUntilContextTimeout(context.Background(), time.Second, timeout, true, + func(ctx context.Context) (bool, error) { + pod := &corev1.Pod{} + + err := rm.kubeClient.Get( + ctx, + types.NamespacedName{ + Name: podName, + Namespace: util.GetNameSpace(), + }, + pod, + ) + + if errors.IsNotFound(err) { + return true, nil + } + + if err != nil { + return false, err + } + + return false, nil + }, + ) +} + +func (rm *ResourceManager) WaitForPvcCapacity(pvcName string, namespace string, desiredSize resource.Quantity, timeout time.Duration) (bool, error) { + resized := false + + err := wait.PollUntilContextTimeout(context.Background(), time.Second, timeout, true, + func(ctx context.Context) (bool, error) { + pvc := &corev1.PersistentVolumeClaim{} + + if err := rm.kubeClient.Get( + ctx, + types.NamespacedName{ + Name: pvcName, + Namespace: namespace, + }, + pvc, + ); err != nil { + logger.Warn( + fmt.Sprintf("Failed to get PVC %s while waiting for resize, retrying", pvcName), + zap.Error(err), + ) + + return false, nil + } + + capacity := pvc.Status.Capacity[corev1.ResourceStorage] + + if capacity.Cmp(desiredSize) >= 0 { + logger.Info(fmt.Sprintf("PVC %s resize completed: capacity=%s", pvcName, capacity.String())) + + resized = true + return true, nil + } + + logger.Info(fmt.Sprintf("Waiting for PVC %s capacity: current=%s desired=%s", pvcName, capacity.String(), desiredSize.String())) + + return false, nil + }, + ) + + if err == context.DeadlineExceeded { + return false, nil + } + + if err != nil { + return false, err + } + + return resized, nil +} + +func (rm *ResourceManager) ScaleStatefulSet(name string, replicas int32) error { + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + sts := &appsv1.StatefulSet{} + + if err := rm.kubeClient.Get( + context.TODO(), + types.NamespacedName{ + Name: name, + Namespace: namespace, + }, + sts, + ); err != nil { + return err + } + + sts.Spec.Replicas = &replicas + + return rm.kubeClient.Update(context.TODO(), sts) + }) +} diff --git a/operator/pkg/patroni/patroni.go b/operator/pkg/patroni/patroni.go index 82747f57..a8ffb41f 100644 --- a/operator/pkg/patroni/patroni.go +++ b/operator/pkg/patroni/patroni.go @@ -599,3 +599,30 @@ func GenerateLDAPConfig(cr *patroniv1.PatroniCore) []string { ), } } + +func Switchover(patroniURL, leader, candidate string) error { + body := map[string]string{"leader": leader, "candidate": candidate} + + data, err := json.Marshal(body) + if err != nil { + return err + } + + resp, err := http.Post(patroniURL+"switchover", "application/json", bytes.NewBuffer(data)) + if err != nil { + return err + } + defer resp.Body.Close() + + if resp.StatusCode < http.StatusOK || + resp.StatusCode >= http.StatusMultipleChoices { + + responseBody, _ := io.ReadAll(resp.Body) + + return fmt.Errorf("Patroni switchover failed: %s, response: %s", resp.Status, string(responseBody)) + } + + logger.Info(fmt.Sprintf("Patroni switchover from %s to %s requested successfully", leader, candidate)) + + return nil +} diff --git a/operator/pkg/reconciler/backup_daemon.go b/operator/pkg/reconciler/backup_daemon.go index 8e5ac705..04eb0663 100644 --- a/operator/pkg/reconciler/backup_daemon.go +++ b/operator/pkg/reconciler/backup_daemon.go @@ -18,6 +18,7 @@ import ( "fmt" "strconv" "strings" + "time" qubershipv1 "github.com/Netcracker/pgskipper-operator/api/apps/v1" commonv1 "github.com/Netcracker/pgskipper-operator/api/common/v1" @@ -61,12 +62,40 @@ func NewBackupDaemonReconciler(cr *qubershipv1.PatroniServices, helper *helper.H func (r *BackupDaemonReconciler) Reconcile() error { cr := r.cr bdSpec := cr.Spec.BackupDaemon + var backupPvc *corev1.PersistentVolumeClaim + backupPvcRestartRequired := false + if bdSpec.Storage.Type != "ephemeral" && bdSpec.Storage.Type != "s3" { - backupPvc := storage.NewPvc("postgres-backup-pvc", &bdSpec.Storage, 1) - if err := r.helper.CreatePvcIfNotExists(backupPvc); err != nil { - logger.Error(fmt.Sprintf("Cannot create pvc %s", backupPvc.Name), zap.Error(err)) + backupPvc = storage.NewPvc( + "postgres-backup-pvc", + &bdSpec.Storage, + 1, + ) + + resizeInProgress, err := r.helper.CreateOrResizePvc(backupPvc) + if err != nil { + logger.Error( + fmt.Sprintf("Cannot create or resize pvc %s", backupPvc.Name), + zap.Error(err), + ) return err } + + if resizeInProgress { + logger.Info(fmt.Sprintf( + "Waiting for PVC %s resize state", + backupPvc.Name, + )) + + backupPvcRestartRequired, err = r.helper.WaitForPvcResizeState( + backupPvc.Name, + backupPvc.Namespace, + backupPvc.Spec.Resources.Requests[corev1.ResourceStorage], + ) + if err != nil { + return err + } + } } if bdSpec.ExternalPv != nil { logger.Info("External Pv for PostgreSQL Backup Daemon is not empty, start configuration") @@ -256,10 +285,63 @@ func (r *BackupDaemonReconciler) Reconcile() error { backupDaemonDeployment.Spec.Template.Spec.Containers[0].Env = append(backupDaemonDeployment.Spec.Template.Spec.Containers[0].Env, envValue...) } + if backupPvcRestartRequired { + logger.Info(fmt.Sprintf("Restarting Backup Daemon deployment %s to complete PVC resize", backupDaemonDeployment.Name)) + + backupPods, err := r.helper.GetNamespacePodListBySelectors( + map[string]string{"app": "postgres-backup-daemon"}, + ) + if err != nil { + return err + } + + if err := r.helper.DeleteDeployment(backupDaemonDeployment); err != nil { + return err + } + + if err := r.helper.WaitTillDeploymentDeleted(backupDaemonDeployment); err != nil { + return err + } + + for _, pod := range backupPods.Items { + if err := r.helper.WaitForPodDeletion(pod.Name, 2*time.Minute); err != nil { + return err + } + } + + time.Sleep(10 * time.Second) + } + if err := r.helper.CreateOrUpdateDeploymentForce(backupDaemonDeployment, true); err != nil { logger.Error(fmt.Sprintf("Cannot create or update deployment %s", backupDaemonDeployment.Name), zap.Error(err)) return err } + if backupPvcRestartRequired { + desiredSize := backupPvc.Spec.Resources.Requests[corev1.ResourceStorage] + + resized, err := r.helper.WaitForPvcCapacity( + backupPvc.Name, + backupPvc.Namespace, + desiredSize, + 2*time.Minute, + ) + if err != nil { + return err + } + + if !resized { + return fmt.Errorf( + "Backup Daemon PVC %s filesystem resize is still pending", + backupPvc.Name, + ) + } + + logger.Info(fmt.Sprintf( + "Backup Daemon PVC %s successfully resized to %s", + backupPvc.Name, + desiredSize.String(), + )) + } if err := util.WaitForBackupDaemon(); err != nil { logger.Error("Failed to wait for backup daemon, exiting", zap.Error(err)) return err diff --git a/operator/pkg/reconciler/patroni.go b/operator/pkg/reconciler/patroni.go index 4cbe2a0e..52f41a88 100644 --- a/operator/pkg/reconciler/patroni.go +++ b/operator/pkg/reconciler/patroni.go @@ -41,6 +41,7 @@ import ( "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/util/wait" ) var wg sync.WaitGroup @@ -156,7 +157,9 @@ func (r *PatroniReconciler) Reconcile() error { } } - + if err := r.processPgBackRestPvc(cr); err != nil { + return err + } // find possible deployments by pods // try to get master pod masterPod, err = r.helper.GetPodsByLabel(r.cluster.PatroniMasterSelectors) @@ -493,26 +496,17 @@ func (r *PatroniReconciler) processPatroniStatefulset(cr *v1.PatroniCore, deploy } patroniSpec := cr.Spec.Patroni - pvc := storage.NewPvc(fmt.Sprintf("%s-data-%v", opUtil.GetPatroniClusterName(cr.Spec.Patroni.ClusterName), deploymentIdx), patroniSpec.Storage, deploymentIdx) - if err := r.helper.CreatePvcIfNotExists(pvc); err != nil { - logger.Error(fmt.Sprintf("Cannot create pvc %s", pvc.Name), zap.Error(err)) - return err + + patroniPvcs := []*corev1.PersistentVolumeClaim{ + storage.NewPvc(fmt.Sprintf("%s-data-%v", opUtil.GetPatroniClusterName(cr.Spec.Patroni.ClusterName), deploymentIdx), patroniSpec.Storage, deploymentIdx), } + if patroniSpec.PgWalStorage != nil { - pvc := storage.NewPvc(fmt.Sprintf("%s-wals-data-%v", opUtil.GetPatroniClusterName(cr.Spec.Patroni.ClusterName), deploymentIdx), patroniSpec.PgWalStorage, deploymentIdx) - if err := r.helper.CreatePvcIfNotExists(pvc); err != nil { - logger.Error(fmt.Sprintf("Cannot create pvc %s", pvc.Name), zap.Error(err)) - return err - } + patroniPvcs = append(patroniPvcs, storage.NewPvc(fmt.Sprintf("%s-wals-data-%v", opUtil.GetPatroniClusterName(cr.Spec.Patroni.ClusterName), deploymentIdx), patroniSpec.PgWalStorage, deploymentIdx)) } - if cr.Spec.PgBackRest != nil && strings.ToLower(cr.Spec.PgBackRest.RepoType) == "rwx" { - pgBackrestStorage := cr.Spec.PgBackRest.Rwx - pgBackrestStorage.AccessModes = []string{"ReadWriteMany"} - pvc = storage.NewPvc("pgbackrest-backups", pgBackrestStorage, 1) - if err := r.helper.CreatePvcIfNotExists(pvc); err != nil { - logger.Error(fmt.Sprintf("Cannot create pvc %s", pvc.Name), zap.Error(err)) - return err - } + + if err := r.processPatroniPvcResize(patroniPvcs, deploymentIdx); err != nil { + return err } // check deployments @@ -543,6 +537,302 @@ func (r *PatroniReconciler) processPatroniStatefulset(cr *v1.PatroniCore, deploy return nil } +func (r *PatroniReconciler) processPgBackRestPvc(cr *v1.PatroniCore) error { + if cr.Spec.PgBackRest == nil || strings.ToLower(cr.Spec.PgBackRest.RepoType) != "rwx" { + return nil + } + + pgBackrestStorage := cr.Spec.PgBackRest.Rwx + pgBackrestStorage.AccessModes = []string{"ReadWriteMany"} + + pvc := storage.NewPvc("pgbackrest-backups", pgBackrestStorage, 1) + + resizeInProgress, err := r.helper.CreateOrResizePvc(pvc) + if err != nil { + logger.Error(fmt.Sprintf("Cannot create or resize pvc %s", pvc.Name), zap.Error(err)) + return err + } + + if !resizeInProgress { + return nil + } + + logger.Info(fmt.Sprintf("Waiting for PVC %s resize state", pvc.Name)) + + restartRequired, err := r.helper.WaitForPvcResizeState( + pvc.Name, + pvc.Namespace, + pvc.Spec.Resources.Requests[corev1.ResourceStorage], + ) + if err != nil { + return err + } + + if !restartRequired { + return nil + } + + desiredSize := pvc.Spec.Resources.Requests[corev1.ResourceStorage] + + replicaPods, err := r.helper.GetPodsByLabel(r.cluster.PatroniReplicasSelector) + if err != nil && !errors.IsNotFound(err) { + return err + } + + var podName string + + if len(replicaPods.Items) > 0 { + podName = replicaPods.Items[0].Name + logger.Info(fmt.Sprintf("Restarting Patroni replica %s to complete pgBackRest PVC resize", podName)) + } else { + masterPods, err := r.helper.GetPodsByLabel(r.cluster.PatroniMasterSelectors) + if err != nil { + return err + } + + if len(masterPods.Items) != 1 { + return fmt.Errorf("expected exactly one Patroni master, found %d", len(masterPods.Items)) + } + + podName = masterPods.Items[0].Name + + logger.Warn(fmt.Sprintf("No Patroni replica available, restarting master %s to complete pgBackRest PVC resize", podName)) + } + + statefulSetName := strings.TrimSuffix(podName, "-0") + + const maxResizeAttempts = 3 + + for attempt := 1; attempt <= maxResizeAttempts; attempt++ { + logger.Info(fmt.Sprintf("Restarting StatefulSet %s to complete pgBackRest PVC resize, attempt %d/%d", statefulSetName, attempt, maxResizeAttempts)) + + if err := r.helper.ScaleStatefulSet(statefulSetName, 0); err != nil { + return err + } + + if err := r.helper.WaitForPodDeletion( + podName, + 2*time.Minute, + ); err != nil { + return err + } + + time.Sleep(10 * time.Second) + + if err := r.helper.ScaleStatefulSet(statefulSetName, 1); err != nil { + return err + } + + resized, err := r.helper.WaitForPvcCapacity( + pvc.Name, + pvc.Namespace, + desiredSize, + 30*time.Second, + ) + if err != nil { + return err + } + + if resized { + logger.Info(fmt.Sprintf("pgBackRest PVC %s successfully resized to %s", pvc.Name, desiredSize.String())) + return nil + } + + if attempt == maxResizeAttempts { + return fmt.Errorf("pgBackRest PVC %s filesystem resize is still pending after %d restart attempts", pvc.Name, maxResizeAttempts) + } + + logger.Warn(fmt.Sprintf("pgBackRest PVC %s resize is still pending, retrying StatefulSet restart", pvc.Name)) + } + + return nil +} + +func (r *PatroniReconciler) processPatroniPvcResize(pvcs []*corev1.PersistentVolumeClaim, deploymentIdx int) error { + var resizingPvcs []*corev1.PersistentVolumeClaim + restartRequired := false + + for _, pvc := range pvcs { + resizeInProgress, err := r.helper.CreateOrResizePvc(pvc) + if err != nil { + return err + } + + if resizeInProgress { + resizingPvcs = append(resizingPvcs, pvc) + } + } + + for _, pvc := range resizingPvcs { + logger.Info(fmt.Sprintf("Waiting for PVC %s resize state", pvc.Name)) + + pvcRestartRequired, err := r.helper.WaitForPvcResizeState( + pvc.Name, + pvc.Namespace, + pvc.Spec.Resources.Requests[corev1.ResourceStorage], + ) + if err != nil { + return err + } + + if pvcRestartRequired { + restartRequired = true + } + } + + if !restartRequired { + return nil + } + + statefulSetName := fmt.Sprintf("pg-%s-node%d", r.cluster.ClusterName, deploymentIdx) + + podName := fmt.Sprintf("%s-0", statefulSetName) + + if err := r.switchoverIfMaster(podName); err != nil { + return err + } + + const maxResizeAttempts = 3 + + for attempt := 1; attempt <= maxResizeAttempts; attempt++ { + logger.Info(fmt.Sprintf("Restarting StatefulSet %s to complete PVC resize, attempt %d/%d", statefulSetName, attempt, maxResizeAttempts)) + + if err := r.helper.ScaleStatefulSet( + statefulSetName, + 0, + ); err != nil { + return err + } + + if err := r.helper.WaitForPodDeletion( + podName, + 2*time.Minute, + ); err != nil { + return err + } + + time.Sleep(10 * time.Second) + + if err := r.helper.ScaleStatefulSet( + statefulSetName, + 1, + ); err != nil { + return err + } + + allResized := true + + for _, pvc := range resizingPvcs { + desiredSize := pvc.Spec.Resources.Requests[corev1.ResourceStorage] + + resized, err := r.helper.WaitForPvcCapacity( + pvc.Name, + pvc.Namespace, + desiredSize, + 30*time.Second, + ) + if err != nil { + return err + } + + if !resized { + allResized = false + continue + } + + logger.Info(fmt.Sprintf( + "PVC %s successfully resized to %s", + pvc.Name, + desiredSize.String(), + )) + } + + if allResized { + return nil + } + + if attempt == maxResizeAttempts { + return fmt.Errorf("Patroni PVC filesystem resize is still pending after %d restart attempts for StatefulSet %s", maxResizeAttempts, statefulSetName) + } + + logger.Warn(fmt.Sprintf("Patroni PVC resize is still pending for StatefulSet %s, retrying restart", statefulSetName)) + } + + return nil +} + +func (r *PatroniReconciler) switchoverIfMaster(podName string) error { + masterPods, err := r.helper.GetPodsByLabel(r.cluster.PatroniMasterSelectors) + if err != nil { + return err + } + + if len(masterPods.Items) != 1 { + return fmt.Errorf("expected exactly one Patroni master, found %d", len(masterPods.Items)) + } + + currentMaster := masterPods.Items[0].Name + + // Target node is already replica. + // Nothing to do. + if currentMaster != podName { + return nil + } + + logger.Info(fmt.Sprintf("Patroni node %s is master, switchover is required before restart", podName)) + + replicaPods, err := r.helper.GetPodsByLabel( + r.cluster.PatroniReplicasSelector, + ) + if err != nil { + return err + } + + if len(replicaPods.Items) == 0 { + return fmt.Errorf("cannot restart Patroni master %s: no replica available for switchover", podName) + } + + candidate := replicaPods.Items[0].Name + + logger.Info(fmt.Sprintf("Switching Patroni master from %s to %s", currentMaster, candidate)) + + if err := patroni.Switchover( + r.cluster.PatroniUrl, + currentMaster, + candidate, + ); err != nil { + return err + } + + // Wait specifically until our selected replica becomes master. + if err := wait.PollUntilContextTimeout( + context.Background(), + time.Second, + 2*time.Minute, + true, + func(ctx context.Context) (bool, error) { + masters, err := r.helper.GetPodsByLabel( + r.cluster.PatroniMasterSelectors, + ) + if err != nil { + return false, nil + } + + if len(masters.Items) != 1 { + return false, nil + } + + return masters.Items[0].Name == candidate, nil + }, + ); err != nil { + return fmt.Errorf("timeout waiting for %s to become Patroni master: %w", candidate, err) + } + + logger.Info(fmt.Sprintf("Patroni switchover completed, new master is %s", candidate)) + + return nil +} + func (r *PatroniReconciler) createEndpointsForEtcdAsDcs() error { pgEndpointSlice := reconcileEndpointSlice(r.cluster.PostgresServiceName, r.cluster.PatroniLabels) if err := r.helper.CreateEndpointSliceIfNotExists(pgEndpointSlice); err != nil { From 953e6d453e944517a15fff1da2c2582b9fca9570 Mon Sep 17 00:00:00 2001 From: yerkennz Date: Wed, 26 Aug 2026 15:22:09 +0500 Subject: [PATCH 2/7] fix: remove pvc extension for pgbackrest and refactor code --- operator/pkg/helper/resource_management.go | 48 +++---- operator/pkg/patroni/patroni.go | 2 +- operator/pkg/reconciler/backup_daemon.go | 45 ++----- operator/pkg/reconciler/patroni.go | 143 ++------------------- 4 files changed, 40 insertions(+), 198 deletions(-) diff --git a/operator/pkg/helper/resource_management.go b/operator/pkg/helper/resource_management.go index ede246ac..b4d74cdb 100644 --- a/operator/pkg/helper/resource_management.go +++ b/operator/pkg/helper/resource_management.go @@ -1209,33 +1209,6 @@ func (rm *ResourceManager) WaitForPvcResizeState(pvcName string, namespace strin return restartRequired, err } -func (rm *ResourceManager) WaitForPodDeletion(podName string, timeout time.Duration) error { - return wait.PollUntilContextTimeout(context.Background(), time.Second, timeout, true, - func(ctx context.Context) (bool, error) { - pod := &corev1.Pod{} - - err := rm.kubeClient.Get( - ctx, - types.NamespacedName{ - Name: podName, - Namespace: util.GetNameSpace(), - }, - pod, - ) - - if errors.IsNotFound(err) { - return true, nil - } - - if err != nil { - return false, err - } - - return false, nil - }, - ) -} - func (rm *ResourceManager) WaitForPvcCapacity(pvcName string, namespace string, desiredSize resource.Quantity, timeout time.Duration) (bool, error) { resized := false @@ -1305,3 +1278,24 @@ func (rm *ResourceManager) ScaleStatefulSet(name string, replicas int32) error { return rm.kubeClient.Update(context.TODO(), sts) }) } + +func (rm *ResourceManager) ScaleDeployment(name string, replicas int32) error { + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + deployment := &appsv1.Deployment{} + + if err := rm.kubeClient.Get( + context.TODO(), + types.NamespacedName{ + Name: name, + Namespace: namespace, + }, + deployment, + ); err != nil { + return err + } + + deployment.Spec.Replicas = &replicas + + return rm.kubeClient.Update(context.TODO(), deployment) + }) +} diff --git a/operator/pkg/patroni/patroni.go b/operator/pkg/patroni/patroni.go index a8ffb41f..d9779433 100644 --- a/operator/pkg/patroni/patroni.go +++ b/operator/pkg/patroni/patroni.go @@ -608,7 +608,7 @@ func Switchover(patroniURL, leader, candidate string) error { return err } - resp, err := http.Post(patroniURL+"switchover", "application/json", bytes.NewBuffer(data)) + resp, err := patroniPost(http.DefaultClient, patroniURL+"switchover", bytes.NewBuffer(data)) if err != nil { return err } diff --git a/operator/pkg/reconciler/backup_daemon.go b/operator/pkg/reconciler/backup_daemon.go index 04eb0663..b3db3d07 100644 --- a/operator/pkg/reconciler/backup_daemon.go +++ b/operator/pkg/reconciler/backup_daemon.go @@ -29,6 +29,7 @@ import ( "github.com/Netcracker/pgskipper-operator/pkg/patroni" "github.com/Netcracker/pgskipper-operator/pkg/storage" "github.com/Netcracker/pgskipper-operator/pkg/util" + opUtil "github.com/Netcracker/pgskipper-operator/pkg/util" "github.com/Netcracker/pgskipper-operator/pkg/util/constants" "github.com/Netcracker/qubership-credential-manager/pkg/manager" "go.uber.org/zap" @@ -66,26 +67,15 @@ func (r *BackupDaemonReconciler) Reconcile() error { backupPvcRestartRequired := false if bdSpec.Storage.Type != "ephemeral" && bdSpec.Storage.Type != "s3" { - backupPvc = storage.NewPvc( - "postgres-backup-pvc", - &bdSpec.Storage, - 1, - ) - + backupPvc = storage.NewPvc("postgres-backup-pvc", &bdSpec.Storage, 1) resizeInProgress, err := r.helper.CreateOrResizePvc(backupPvc) if err != nil { - logger.Error( - fmt.Sprintf("Cannot create or resize pvc %s", backupPvc.Name), - zap.Error(err), - ) + logger.Error(fmt.Sprintf("Cannot create or resize pvc %s", backupPvc.Name), zap.Error(err)) return err } if resizeInProgress { - logger.Info(fmt.Sprintf( - "Waiting for PVC %s resize state", - backupPvc.Name, - )) + logger.Info(fmt.Sprintf("Waiting for PVC %s resize state", backupPvc.Name)) backupPvcRestartRequired, err = r.helper.WaitForPvcResizeState( backupPvc.Name, @@ -294,17 +284,12 @@ func (r *BackupDaemonReconciler) Reconcile() error { if err != nil { return err } - - if err := r.helper.DeleteDeployment(backupDaemonDeployment); err != nil { - return err - } - - if err := r.helper.WaitTillDeploymentDeleted(backupDaemonDeployment); err != nil { + if err := r.helper.ScaleDeployment(backupDaemonDeployment.Name, 0); err != nil { return err } for _, pod := range backupPods.Items { - if err := r.helper.WaitForPodDeletion(pod.Name, 2*time.Minute); err != nil { + if err := opUtil.WaitDeletePod(&pod); err != nil { return err } } @@ -319,28 +304,16 @@ func (r *BackupDaemonReconciler) Reconcile() error { if backupPvcRestartRequired { desiredSize := backupPvc.Spec.Resources.Requests[corev1.ResourceStorage] - resized, err := r.helper.WaitForPvcCapacity( - backupPvc.Name, - backupPvc.Namespace, - desiredSize, - 2*time.Minute, - ) + resized, err := r.helper.WaitForPvcCapacity(backupPvc.Name, backupPvc.Namespace, desiredSize, 2*time.Minute) if err != nil { return err } if !resized { - return fmt.Errorf( - "Backup Daemon PVC %s filesystem resize is still pending", - backupPvc.Name, - ) + return fmt.Errorf("Backup Daemon PVC %s filesystem resize is still pending", backupPvc.Name) } - logger.Info(fmt.Sprintf( - "Backup Daemon PVC %s successfully resized to %s", - backupPvc.Name, - desiredSize.String(), - )) + logger.Info(fmt.Sprintf("Backup Daemon PVC %s successfully resized to %s", backupPvc.Name, desiredSize.String())) } if err := util.WaitForBackupDaemon(); err != nil { logger.Error("Failed to wait for backup daemon, exiting", zap.Error(err)) diff --git a/operator/pkg/reconciler/patroni.go b/operator/pkg/reconciler/patroni.go index 52f41a88..64b58562 100644 --- a/operator/pkg/reconciler/patroni.go +++ b/operator/pkg/reconciler/patroni.go @@ -157,9 +157,6 @@ func (r *PatroniReconciler) Reconcile() error { } } - if err := r.processPgBackRestPvc(cr); err != nil { - return err - } // find possible deployments by pods // try to get master pod masterPod, err = r.helper.GetPodsByLabel(r.cluster.PatroniMasterSelectors) @@ -537,117 +534,6 @@ func (r *PatroniReconciler) processPatroniStatefulset(cr *v1.PatroniCore, deploy return nil } -func (r *PatroniReconciler) processPgBackRestPvc(cr *v1.PatroniCore) error { - if cr.Spec.PgBackRest == nil || strings.ToLower(cr.Spec.PgBackRest.RepoType) != "rwx" { - return nil - } - - pgBackrestStorage := cr.Spec.PgBackRest.Rwx - pgBackrestStorage.AccessModes = []string{"ReadWriteMany"} - - pvc := storage.NewPvc("pgbackrest-backups", pgBackrestStorage, 1) - - resizeInProgress, err := r.helper.CreateOrResizePvc(pvc) - if err != nil { - logger.Error(fmt.Sprintf("Cannot create or resize pvc %s", pvc.Name), zap.Error(err)) - return err - } - - if !resizeInProgress { - return nil - } - - logger.Info(fmt.Sprintf("Waiting for PVC %s resize state", pvc.Name)) - - restartRequired, err := r.helper.WaitForPvcResizeState( - pvc.Name, - pvc.Namespace, - pvc.Spec.Resources.Requests[corev1.ResourceStorage], - ) - if err != nil { - return err - } - - if !restartRequired { - return nil - } - - desiredSize := pvc.Spec.Resources.Requests[corev1.ResourceStorage] - - replicaPods, err := r.helper.GetPodsByLabel(r.cluster.PatroniReplicasSelector) - if err != nil && !errors.IsNotFound(err) { - return err - } - - var podName string - - if len(replicaPods.Items) > 0 { - podName = replicaPods.Items[0].Name - logger.Info(fmt.Sprintf("Restarting Patroni replica %s to complete pgBackRest PVC resize", podName)) - } else { - masterPods, err := r.helper.GetPodsByLabel(r.cluster.PatroniMasterSelectors) - if err != nil { - return err - } - - if len(masterPods.Items) != 1 { - return fmt.Errorf("expected exactly one Patroni master, found %d", len(masterPods.Items)) - } - - podName = masterPods.Items[0].Name - - logger.Warn(fmt.Sprintf("No Patroni replica available, restarting master %s to complete pgBackRest PVC resize", podName)) - } - - statefulSetName := strings.TrimSuffix(podName, "-0") - - const maxResizeAttempts = 3 - - for attempt := 1; attempt <= maxResizeAttempts; attempt++ { - logger.Info(fmt.Sprintf("Restarting StatefulSet %s to complete pgBackRest PVC resize, attempt %d/%d", statefulSetName, attempt, maxResizeAttempts)) - - if err := r.helper.ScaleStatefulSet(statefulSetName, 0); err != nil { - return err - } - - if err := r.helper.WaitForPodDeletion( - podName, - 2*time.Minute, - ); err != nil { - return err - } - - time.Sleep(10 * time.Second) - - if err := r.helper.ScaleStatefulSet(statefulSetName, 1); err != nil { - return err - } - - resized, err := r.helper.WaitForPvcCapacity( - pvc.Name, - pvc.Namespace, - desiredSize, - 30*time.Second, - ) - if err != nil { - return err - } - - if resized { - logger.Info(fmt.Sprintf("pgBackRest PVC %s successfully resized to %s", pvc.Name, desiredSize.String())) - return nil - } - - if attempt == maxResizeAttempts { - return fmt.Errorf("pgBackRest PVC %s filesystem resize is still pending after %d restart attempts", pvc.Name, maxResizeAttempts) - } - - logger.Warn(fmt.Sprintf("pgBackRest PVC %s resize is still pending, retrying StatefulSet restart", pvc.Name)) - } - - return nil -} - func (r *PatroniReconciler) processPatroniPvcResize(pvcs []*corev1.PersistentVolumeClaim, deploymentIdx int) error { var resizingPvcs []*corev1.PersistentVolumeClaim restartRequired := false @@ -704,10 +590,12 @@ func (r *PatroniReconciler) processPatroniPvcResize(pvcs []*corev1.PersistentVol return err } - if err := r.helper.WaitForPodDeletion( - podName, - 2*time.Minute, - ); err != nil { + if err := opUtil.WaitDeletePod(&corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: podName, + Namespace: r.cr.Namespace, + }, + }); err != nil { return err } @@ -725,12 +613,7 @@ func (r *PatroniReconciler) processPatroniPvcResize(pvcs []*corev1.PersistentVol for _, pvc := range resizingPvcs { desiredSize := pvc.Spec.Resources.Requests[corev1.ResourceStorage] - resized, err := r.helper.WaitForPvcCapacity( - pvc.Name, - pvc.Namespace, - desiredSize, - 30*time.Second, - ) + resized, err := r.helper.WaitForPvcCapacity(pvc.Name, pvc.Namespace, desiredSize, 30*time.Second) if err != nil { return err } @@ -740,11 +623,7 @@ func (r *PatroniReconciler) processPatroniPvcResize(pvcs []*corev1.PersistentVol continue } - logger.Info(fmt.Sprintf( - "PVC %s successfully resized to %s", - pvc.Name, - desiredSize.String(), - )) + logger.Info(fmt.Sprintf("PVC %s successfully resized to %s", pvc.Name, desiredSize.String())) } if allResized { @@ -805,11 +684,7 @@ func (r *PatroniReconciler) switchoverIfMaster(podName string) error { } // Wait specifically until our selected replica becomes master. - if err := wait.PollUntilContextTimeout( - context.Background(), - time.Second, - 2*time.Minute, - true, + if err := wait.PollUntilContextTimeout(context.Background(), time.Second, 2*time.Minute, true, func(ctx context.Context) (bool, error) { masters, err := r.helper.GetPodsByLabel( r.cluster.PatroniMasterSelectors, From 6b0c5ee62e8bc0d60d9569c62c753191d11f3d3e Mon Sep 17 00:00:00 2001 From: yerkennz Date: Thu, 27 Aug 2026 13:56:08 +0500 Subject: [PATCH 3/7] fix: wait for detaching in exponent time --- operator/pkg/reconciler/backup_daemon.go | 87 ++++++++++++++++-------- operator/pkg/reconciler/patroni.go | 40 +++++------ 2 files changed, 76 insertions(+), 51 deletions(-) diff --git a/operator/pkg/reconciler/backup_daemon.go b/operator/pkg/reconciler/backup_daemon.go index b3db3d07..64d161a8 100644 --- a/operator/pkg/reconciler/backup_daemon.go +++ b/operator/pkg/reconciler/backup_daemon.go @@ -15,6 +15,7 @@ package reconciler import ( + "context" "fmt" "strconv" "strings" @@ -276,44 +277,72 @@ func (r *BackupDaemonReconciler) Reconcile() error { } if backupPvcRestartRequired { - logger.Info(fmt.Sprintf("Restarting Backup Daemon deployment %s to complete PVC resize", backupDaemonDeployment.Name)) + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute) + defer cancel() - backupPods, err := r.helper.GetNamespacePodListBySelectors( - map[string]string{"app": "postgres-backup-daemon"}, - ) - if err != nil { - return err - } - if err := r.helper.ScaleDeployment(backupDaemonDeployment.Name, 0); err != nil { - return err - } + attempt := 1 + + for { + select { + case <-ctx.Done(): + return fmt.Errorf("timeout waiting for Backup Daemon PVC resize") + default: + } + + delay := 10 * time.Second * time.Duration(1<<(attempt-1)) + if delay > time.Minute { + delay = time.Minute + } + + logger.Info(fmt.Sprintf("Restarting Backup Daemon deployment %s to complete PVC resize, attempt %d, waiting %s", backupDaemonDeployment.Name, attempt, delay)) - for _, pod := range backupPods.Items { - if err := opUtil.WaitDeletePod(&pod); err != nil { + backupPods, err := r.helper.GetNamespacePodListBySelectors( + map[string]string{"app": "postgres-backup-daemon"}, + ) + if err != nil { return err } - } - time.Sleep(10 * time.Second) - } + if err := r.helper.ScaleDeployment(backupDaemonDeployment.Name, 0); err != nil { + return err + } - if err := r.helper.CreateOrUpdateDeploymentForce(backupDaemonDeployment, true); err != nil { - logger.Error(fmt.Sprintf("Cannot create or update deployment %s", backupDaemonDeployment.Name), zap.Error(err)) - return err - } - if backupPvcRestartRequired { - desiredSize := backupPvc.Spec.Resources.Requests[corev1.ResourceStorage] + for _, pod := range backupPods.Items { + if err := opUtil.WaitDeletePod(&pod); err != nil { + return err + } + } - resized, err := r.helper.WaitForPvcCapacity(backupPvc.Name, backupPvc.Namespace, desiredSize, 2*time.Minute) - if err != nil { - return err - } + select { + case <-time.After(delay): + case <-ctx.Done(): + return fmt.Errorf("timeout waiting for Backup Daemon PVC resize") + } - if !resized { - return fmt.Errorf("Backup Daemon PVC %s filesystem resize is still pending", backupPvc.Name) - } + if err := r.helper.CreateOrUpdateDeploymentForce(backupDaemonDeployment, true); err != nil { + logger.Error(fmt.Sprintf("Cannot create or update deployment %s", backupDaemonDeployment.Name), zap.Error(err)) + return err + } - logger.Info(fmt.Sprintf("Backup Daemon PVC %s successfully resized to %s", backupPvc.Name, desiredSize.String())) + desiredSize := backupPvc.Spec.Resources.Requests[corev1.ResourceStorage] + + resized, err := r.helper.WaitForPvcCapacity(backupPvc.Name, backupPvc.Namespace, desiredSize, 30*time.Second) + if err != nil { + return err + } + + if resized { + logger.Info(fmt.Sprintf("Backup Daemon PVC %s successfully resized to %s", backupPvc.Name, desiredSize.String())) + break + } + + attempt++ + } + } else { + if err := r.helper.CreateOrUpdateDeploymentForce(backupDaemonDeployment, true); err != nil { + logger.Error(fmt.Sprintf("Cannot create or update deployment %s", backupDaemonDeployment.Name), zap.Error(err)) + return err + } } if err := util.WaitForBackupDaemon(); err != nil { logger.Error("Failed to wait for backup daemon, exiting", zap.Error(err)) diff --git a/operator/pkg/reconciler/patroni.go b/operator/pkg/reconciler/patroni.go index 2c1e104f..91f722b4 100644 --- a/operator/pkg/reconciler/patroni.go +++ b/operator/pkg/reconciler/patroni.go @@ -583,15 +583,23 @@ func (r *PatroniReconciler) processPatroniPvcResize(pvcs []*corev1.PersistentVol return err } - const maxResizeAttempts = 3 + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute) + defer cancel() - for attempt := 1; attempt <= maxResizeAttempts; attempt++ { - logger.Info(fmt.Sprintf("Restarting StatefulSet %s to complete PVC resize, attempt %d/%d", statefulSetName, attempt, maxResizeAttempts)) + attempt := 1 - if err := r.helper.ScaleStatefulSet( - statefulSetName, - 0, - ); err != nil { + for { + select { + case <-ctx.Done(): + return fmt.Errorf("timeout waiting for Patroni PVC resize") + default: + } + + delay := 10 * time.Second * time.Duration(1<<(attempt-1)) + + logger.Info(fmt.Sprintf("Restarting StatefulSet %s to complete PVC resize, attempt %d, waiting %s", statefulSetName, attempt, delay)) + + if err := r.helper.ScaleStatefulSet(statefulSetName, 0); err != nil { return err } @@ -604,12 +612,9 @@ func (r *PatroniReconciler) processPatroniPvcResize(pvcs []*corev1.PersistentVol return err } - time.Sleep(10 * time.Second) + time.Sleep(delay) - if err := r.helper.ScaleStatefulSet( - statefulSetName, - 1, - ); err != nil { + if err := r.helper.ScaleStatefulSet(statefulSetName, 1); err != nil { return err } @@ -625,24 +630,15 @@ func (r *PatroniReconciler) processPatroniPvcResize(pvcs []*corev1.PersistentVol if !resized { allResized = false - continue } - - logger.Info(fmt.Sprintf("PVC %s successfully resized to %s", pvc.Name, desiredSize.String())) } if allResized { return nil } - if attempt == maxResizeAttempts { - return fmt.Errorf("Patroni PVC filesystem resize is still pending after %d restart attempts for StatefulSet %s", maxResizeAttempts, statefulSetName) - } - - logger.Warn(fmt.Sprintf("Patroni PVC resize is still pending for StatefulSet %s, retrying restart", statefulSetName)) + attempt++ } - - return nil } func (r *PatroniReconciler) switchoverIfMaster(podName string) error { From 5b6fd3ffac59b576b064c14ddb19e87aa22554f9 Mon Sep 17 00:00:00 2001 From: yerkennz Date: Thu, 27 Aug 2026 18:58:54 +0500 Subject: [PATCH 4/7] fix: revert pgbackrest pvc creation && use const backupDeamon label --- operator/pkg/reconciler/backup_daemon.go | 2 +- operator/pkg/reconciler/patroni.go | 11 +++++++++++ 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/operator/pkg/reconciler/backup_daemon.go b/operator/pkg/reconciler/backup_daemon.go index 64d161a8..cab6589c 100644 --- a/operator/pkg/reconciler/backup_daemon.go +++ b/operator/pkg/reconciler/backup_daemon.go @@ -297,7 +297,7 @@ func (r *BackupDaemonReconciler) Reconcile() error { logger.Info(fmt.Sprintf("Restarting Backup Daemon deployment %s to complete PVC resize, attempt %d, waiting %s", backupDaemonDeployment.Name, attempt, delay)) backupPods, err := r.helper.GetNamespacePodListBySelectors( - map[string]string{"app": "postgres-backup-daemon"}, + deployment.BackupDaemonLabels, ) if err != nil { return err diff --git a/operator/pkg/reconciler/patroni.go b/operator/pkg/reconciler/patroni.go index 91f722b4..62330995 100644 --- a/operator/pkg/reconciler/patroni.go +++ b/operator/pkg/reconciler/patroni.go @@ -507,6 +507,17 @@ func (r *PatroniReconciler) processPatroniStatefulset(cr *v1.PatroniCore, deploy patroniPvcs = append(patroniPvcs, storage.NewPvc(fmt.Sprintf("%s-wals-data-%v", opUtil.GetPatroniClusterName(cr.Spec.Patroni.ClusterName), deploymentIdx), patroniSpec.PgWalStorage, deploymentIdx)) } + if cr.Spec.PgBackRest != nil && strings.ToLower(cr.Spec.PgBackRest.RepoType) == "rwx" { + pgBackrestStorage := cr.Spec.PgBackRest.Rwx + pgBackrestStorage.AccessModes = []string{"ReadWriteMany"} + + pgBackrestPvc := storage.NewPvc("pgbackrest-backups", pgBackrestStorage, 1) + if err := r.helper.CreatePvcIfNotExists(pgBackrestPvc); err != nil { + logger.Error(fmt.Sprintf("Cannot create pvc %s", pgBackrestPvc.Name), zap.Error(err)) + return err + } + } + if err := r.processPatroniPvcResize(patroniPvcs, deploymentIdx); err != nil { return err } From 8627f53622bd145a378deb60775d34a2c6599247 Mon Sep 17 00:00:00 2001 From: yerkennz Date: Fri, 28 Aug 2026 14:44:40 +0500 Subject: [PATCH 5/7] fix: CreateOrUpdateDeploymentForce and prevent waiting for deployment stability --- operator/pkg/helper/resource_management.go | 4 ++-- operator/pkg/reconciler/backup_daemon.go | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/operator/pkg/helper/resource_management.go b/operator/pkg/helper/resource_management.go index b4d74cdb..1aae89e5 100644 --- a/operator/pkg/helper/resource_management.go +++ b/operator/pkg/helper/resource_management.go @@ -345,7 +345,7 @@ func (rm *ResourceManager) CreateOrUpdateService(service *corev1.Service) error // This method performs delete and re-create deployment in case update was failed func (rm *ResourceManager) CreateOrUpdateDeploymentForce(deployment *appsv1.Deployment, waitStability bool) error { - if err := rm.CreateOrUpdateDeployment(deployment, true); err != nil { + if err := rm.CreateOrUpdateDeployment(deployment, waitStability); err != nil { logger.Error(fmt.Sprintf("Cannot create deployment %s", deployment.Name), zap.Error(err)) if err = rm.DeleteDeployment(deployment.Name); err != nil { @@ -354,7 +354,7 @@ func (rm *ResourceManager) CreateOrUpdateDeploymentForce(deployment *appsv1.Depl if err = rm.WaitTillDeploymentDeleted(deployment); err != nil { logger.Error(fmt.Sprintf("Deployment: %s was not deleted in time", deployment.Name), zap.Error(err)) } - if err = rm.CreateOrUpdateDeployment(deployment, true); err != nil { + if err = rm.CreateOrUpdateDeployment(deployment, waitStability); err != nil { logger.Error(fmt.Sprintf("Cannot create deployment after delete %s", deployment.Name), zap.Error(err)) } } diff --git a/operator/pkg/reconciler/backup_daemon.go b/operator/pkg/reconciler/backup_daemon.go index cab6589c..22ce2ed3 100644 --- a/operator/pkg/reconciler/backup_daemon.go +++ b/operator/pkg/reconciler/backup_daemon.go @@ -319,7 +319,7 @@ func (r *BackupDaemonReconciler) Reconcile() error { return fmt.Errorf("timeout waiting for Backup Daemon PVC resize") } - if err := r.helper.CreateOrUpdateDeploymentForce(backupDaemonDeployment, true); err != nil { + if err := r.helper.CreateOrUpdateDeploymentForce(backupDaemonDeployment, false); err != nil { logger.Error(fmt.Sprintf("Cannot create or update deployment %s", backupDaemonDeployment.Name), zap.Error(err)) return err } From 0b481316a97d0933b0459ad33c5724a5aac1151c Mon Sep 17 00:00:00 2001 From: yerkennz Date: Fri, 28 Aug 2026 16:03:13 +0500 Subject: [PATCH 6/7] fix: rm condition --- operator/pkg/reconciler/backup_daemon.go | 3 --- 1 file changed, 3 deletions(-) diff --git a/operator/pkg/reconciler/backup_daemon.go b/operator/pkg/reconciler/backup_daemon.go index 22ce2ed3..e10af318 100644 --- a/operator/pkg/reconciler/backup_daemon.go +++ b/operator/pkg/reconciler/backup_daemon.go @@ -290,9 +290,6 @@ func (r *BackupDaemonReconciler) Reconcile() error { } delay := 10 * time.Second * time.Duration(1<<(attempt-1)) - if delay > time.Minute { - delay = time.Minute - } logger.Info(fmt.Sprintf("Restarting Backup Daemon deployment %s to complete PVC resize, attempt %d, waiting %s", backupDaemonDeployment.Name, attempt, delay)) From 38cdccc26b25d55ceeedf0e415a9d2fb2fbd0c54 Mon Sep 17 00:00:00 2001 From: yerkennz Date: Fri, 28 Aug 2026 17:22:23 +0500 Subject: [PATCH 7/7] fix: switch to available replica --- operator/pkg/patroni/patroni.go | 10 +++++----- operator/pkg/reconciler/patroni.go | 24 ++++++++++++------------ 2 files changed, 17 insertions(+), 17 deletions(-) diff --git a/operator/pkg/patroni/patroni.go b/operator/pkg/patroni/patroni.go index 46c21118..67ad13c4 100644 --- a/operator/pkg/patroni/patroni.go +++ b/operator/pkg/patroni/patroni.go @@ -623,8 +623,8 @@ func GenerateLDAPConfig(cr *patroniv1.PatroniCore) []string { } } -func Switchover(patroniURL, leader, candidate string) error { - body := map[string]string{"leader": leader, "candidate": candidate} +func Switchover(patroniURL, leader string) error { + body := map[string]string{"leader": leader} data, err := json.Marshal(body) if err != nil { @@ -645,9 +645,9 @@ func Switchover(patroniURL, leader, candidate string) error { return fmt.Errorf("Patroni switchover failed: %s, response: %s", resp.Status, string(responseBody)) } - logger.Info(fmt.Sprintf("Patroni switchover from %s to %s requested successfully", leader, candidate)) - - return nil + logger.Info(fmt.Sprintf("Patroni switchover from %s requested successfully", leader)) + + return nil } func ReloadPatroniConfig(patroniUrl string) error { diff --git a/operator/pkg/reconciler/patroni.go b/operator/pkg/reconciler/patroni.go index 62330995..04a4ffd9 100644 --- a/operator/pkg/reconciler/patroni.go +++ b/operator/pkg/reconciler/patroni.go @@ -683,19 +683,14 @@ func (r *PatroniReconciler) switchoverIfMaster(podName string) error { return fmt.Errorf("cannot restart Patroni master %s: no replica available for switchover", podName) } - candidate := replicaPods.Items[0].Name + logger.Info(fmt.Sprintf("Switching Patroni master from %s", currentMaster)) - logger.Info(fmt.Sprintf("Switching Patroni master from %s to %s", currentMaster, candidate)) - - if err := patroni.Switchover( - r.cluster.PatroniUrl, - currentMaster, - candidate, - ); err != nil { + if err := patroni.Switchover(r.cluster.PatroniUrl, currentMaster); err != nil { return err } - // Wait specifically until our selected replica becomes master. + var newMaster string + if err := wait.PollUntilContextTimeout(context.Background(), time.Second, 2*time.Minute, true, func(ctx context.Context) (bool, error) { masters, err := r.helper.GetPodsByLabel( @@ -709,13 +704,18 @@ func (r *PatroniReconciler) switchoverIfMaster(podName string) error { return false, nil } - return masters.Items[0].Name == candidate, nil + if masters.Items[0].Name == currentMaster { + return false, nil + } + + newMaster = masters.Items[0].Name + return true, nil }, ); err != nil { - return fmt.Errorf("timeout waiting for %s to become Patroni master: %w", candidate, err) + return fmt.Errorf("timeout waiting for new Patroni master after switchover from %s: %w", currentMaster, err) } - logger.Info(fmt.Sprintf("Patroni switchover completed, new master is %s", candidate)) + logger.Info(fmt.Sprintf("Patroni switchover completed, new master is %s", newMaster)) return nil }