diff --git a/operator/pkg/helper/resource_management.go b/operator/pkg/helper/resource_management.go index b335bac0..b4d74cdb 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,134 @@ 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) 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) + }) +} + +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 94975444..46c21118 100644 --- a/operator/pkg/patroni/patroni.go +++ b/operator/pkg/patroni/patroni.go @@ -623,6 +623,33 @@ 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 := patroniPost(http.DefaultClient, patroniURL+"switchover", 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 +} + func ReloadPatroniConfig(patroniUrl string) error { hosts, err := getPatroniHosts(patroniUrl) if err != nil { diff --git a/operator/pkg/reconciler/backup_daemon.go b/operator/pkg/reconciler/backup_daemon.go index 8e5ac705..cab6589c 100644 --- a/operator/pkg/reconciler/backup_daemon.go +++ b/operator/pkg/reconciler/backup_daemon.go @@ -15,9 +15,11 @@ package reconciler import ( + "context" "fmt" "strconv" "strings" + "time" qubershipv1 "github.com/Netcracker/pgskipper-operator/api/apps/v1" commonv1 "github.com/Netcracker/pgskipper-operator/api/common/v1" @@ -28,6 +30,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" @@ -61,12 +64,29 @@ 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,9 +276,73 @@ func (r *BackupDaemonReconciler) Reconcile() error { backupDaemonDeployment.Spec.Template.Spec.Containers[0].Env = append(backupDaemonDeployment.Spec.Template.Spec.Containers[0].Env, envValue...) } - 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 { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute) + defer cancel() + + 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)) + + backupPods, err := r.helper.GetNamespacePodListBySelectors( + deployment.BackupDaemonLabels, + ) + if err != nil { + return err + } + + if err := r.helper.ScaleDeployment(backupDaemonDeployment.Name, 0); err != nil { + return err + } + + for _, pod := range backupPods.Items { + if err := opUtil.WaitDeletePod(&pod); err != nil { + return err + } + } + + select { + case <-time.After(delay): + case <-ctx.Done(): + return fmt.Errorf("timeout waiting for Backup Daemon PVC resize") + } + + 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 + } + + 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 3962cf1b..62330995 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 @@ -161,7 +162,6 @@ func (r *PatroniReconciler) Reconcile() error { } } - // find possible deployments by pods // try to get master pod masterPod, err = r.helper.GetPodsByLabel(r.cluster.PatroniMasterSelectors) @@ -498,28 +498,30 @@ 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)) + + 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 + } + // check deployments patroniSfs := deployment.NewPatroniStatefulset(cr, deploymentIdx, r.cluster.ClusterName, r.cluster.PatroniTemplate, r.cluster.PostgreSQLUserConf, r.cluster.PatroniLabels) @@ -548,6 +550,176 @@ func (r *PatroniReconciler) processPatroniStatefulset(cr *v1.PatroniCore, deploy 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 + } + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute) + defer cancel() + + attempt := 1 + + 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 + } + + if err := opUtil.WaitDeletePod(&corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: podName, + Namespace: r.cr.Namespace, + }, + }); err != nil { + return err + } + + time.Sleep(delay) + + 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 + } + } + + if allResized { + return nil + } + + attempt++ + } +} + +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 {