From 86ef2929d8b6b4902186ce3cf9e99143078efdc2 Mon Sep 17 00:00:00 2001 From: Mats Kindahl Date: Wed, 12 Aug 2026 10:35:10 +0200 Subject: [PATCH 1/3] refactor(shard): remove dormant DRAINED stand-in machinery Quarantine remediation (delete pod + wipe data PVC + re-bootstrap from backup) supersedes the older "stand-in replica" model, in which a pod whose topology role was DRAINED was kept alive for investigation while a replacement replica was provisioned at a higher index. GetPoolerStatus only ever emits PRIMARY, REPLICA, or QUARANTINED, so nothing sets the DRAINED role anymore and that machinery is dead code. Remove it: - countDrainedPods (topology-role counter) and every effectiveReplicas = replicas + drainedCount offset that keyed on it. createMissingResources and handleScaleDown retain an effectiveReplicas parameter, but it now reflects only maintenance surges (replicas + maintenanceSurges); the DRAINED term is gone. - the drainedCount offsets that maintenance-surge indexing (maintenance_surge.go) and disruption-budget accounting (disruption.go) had picked up; both now use the user-desired replica count directly. These were no-ops because countDrainedPods was always zero. - syncDrainedLabels and the multigres.com/role (LabelPodRole) label it maintained, plus the DRAINED branch in cleanupDrainedPod that orphaned a pod's PVC unconditionally. - the DRAINED short-circuits in isDrainStale, isPoolHealthy, and the maintenance-surge base-stability check (isPoolHealthy still excludes QUARANTINED pods). No behavior change for live paths: rolling updates, scale-down, maintenance surges, and quarantine remediation are unaffected. Tests that exercised the removed stand-in behavior are deleted or refolded onto the index-vs-replicas semantics that remain. Part of MUL-1009. Signed-off-by: Mats Kindahl Co-Authored-By: Claude Opus 4.8 --- pkg/data-handler/topo/pooler.go | 3 +- .../controller/shard/disruption.go | 2 +- .../controller/shard/drain_helpers.go | 15 +- .../controller/shard/maintenance_surge.go | 5 +- .../controller/shard/reconcile_data_plane.go | 8 - .../controller/shard/reconcile_pool_pods.go | 83 ++-------- .../controller/shard/reconcile_quarantine.go | 6 - .../shard/shard_controller_internal_test.go | 151 +++--------------- .../controller/shard/shard_controller_test.go | 114 +------------ pkg/util/metadata/labels.go | 6 - 10 files changed, 40 insertions(+), 353 deletions(-) diff --git a/pkg/data-handler/topo/pooler.go b/pkg/data-handler/topo/pooler.go index e0b03350..bbea5450 100644 --- a/pkg/data-handler/topo/pooler.go +++ b/pkg/data-handler/topo/pooler.go @@ -57,8 +57,7 @@ func GetPoolerStatus( } // Quarantined poolers are unrecoverable (postgres cannot start). // They are surfaced with a distinct QUARANTINED role so they are - // visible in Shard.Status.PodRoles, but they are not routed and do - // not drive the stand-in-replica path (which keyed on DRAINED). + // visible in Shard.Status.PodRoles, but they are not routed. // The operator replaces them via quarantine remediation (delete pod // + wipe data PVC + re-bootstrap from backup); GetQuarantinedPods // carries the reason for that. diff --git a/pkg/resource-handler/controller/shard/disruption.go b/pkg/resource-handler/controller/shard/disruption.go index caefdd53..6db18622 100644 --- a/pkg/resource-handler/controller/shard/disruption.go +++ b/pkg/resource-handler/controller/shard/disruption.go @@ -144,7 +144,7 @@ func (r *ShardReconciler) selectShardScaleDownPod( for poolName, pool := range shard.Spec.Pools { for _, cell := range pool.Cells { group := groups[string(poolName)+"/"+string(cell)] - replicas := poolReplicas(pool) + countDrainedPods(shard, group) + replicas := poolReplicas(pool) for _, pod := range group { index, ok := resolvePodIndex(pod.Name) if ok && index >= int(replicas) && !isMaintenanceSurge(pod) { diff --git a/pkg/resource-handler/controller/shard/drain_helpers.go b/pkg/resource-handler/controller/shard/drain_helpers.go index e3a49d19..246a6646 100644 --- a/pkg/resource-handler/controller/shard/drain_helpers.go +++ b/pkg/resource-handler/controller/shard/drain_helpers.go @@ -13,8 +13,8 @@ import ( "github.com/multigres/multigres-operator/pkg/util/metadata" ) -// resolvePodRole returns the role (e.g. "PRIMARY", "REPLICA", "DRAINED") for a -// pod by checking shard.Status.PodRoles. It checks both the exact pod name and +// resolvePodRole returns the role (e.g. "PRIMARY", "REPLICA", "QUARANTINED") for +// a pod by checking shard.Status.PodRoles. It checks both the exact pod name and // FQDN prefix (podName.subdomain...) since the data-handler may store either. func resolvePodRole(shard *multigresv1alpha1.Shard, podName string) string { if shard.Status.PodRoles == nil { @@ -31,17 +31,6 @@ func resolvePodRole(shard *multigresv1alpha1.Shard, podName string) string { return "" } -// countDrainedPods returns the number of pods whose topology role is DRAINED. -func countDrainedPods(shard *multigresv1alpha1.Shard, existingPods map[string]*corev1.Pod) int32 { - var count int32 - for _, pod := range existingPods { - if resolvePodRole(shard, pod.Name) == "DRAINED" { - count++ - } - } - return count -} - // clearDrainAnnotations removes all drain annotations from a pod via merge patch, // cancelling a drain that is no longer needed (e.g. scale-down reversed). func clearDrainAnnotations(ctx context.Context, k8sClient client.Client, pod *corev1.Pod) error { diff --git a/pkg/resource-handler/controller/shard/maintenance_surge.go b/pkg/resource-handler/controller/shard/maintenance_surge.go index 614436b1..98602218 100644 --- a/pkg/resource-handler/controller/shard/maintenance_surge.go +++ b/pkg/resource-handler/controller/shard/maintenance_surge.go @@ -99,9 +99,6 @@ func (r *ShardReconciler) reconcileCellMaintenanceSurge( baseUnsettled = true continue } - if resolvePodRole(shard, pod.Name) == "DRAINED" { - continue - } stable := isAvailablePooler(pod) && pod.Annotations[metadata.AnnotationDrainState] == "" if !stable { @@ -220,7 +217,7 @@ func (r *ShardReconciler) createOrAdoptMaintenanceSurge( replicas int32, ) error { logger := log.FromContext(ctx) - index := replicas + countDrainedPods(shard, existingPods) + index := replicas podName := BuildPoolPodName(shard, poolName, cellName, int(index)) pvcName := BuildPoolDataPVCName(shard, poolName, cellName, int(index)) diff --git a/pkg/resource-handler/controller/shard/reconcile_data_plane.go b/pkg/resource-handler/controller/shard/reconcile_data_plane.go index 183d9a05..57ac0301 100644 --- a/pkg/resource-handler/controller/shard/reconcile_data_plane.go +++ b/pkg/resource-handler/controller/shard/reconcile_data_plane.go @@ -586,14 +586,6 @@ func (r *ShardReconciler) isDrainStale( return false } - // A drain on a DRAINED pod comes from external deletion (kubectl delete), - // which also sets DeletionTimestamp (handled above). If we somehow reach - // here with a DRAINED pod in requested state without a DeletionTimestamp, - // the drain should still complete — it should never be cancelled. - if resolvePodRole(shard, pod.Name) == "DRAINED" { - return false - } - poolName := pod.Labels[metadata.LabelMultigresPool] cellName := pod.Labels[metadata.LabelMultigresCell] if poolName == "" || cellName == "" { diff --git a/pkg/resource-handler/controller/shard/reconcile_pool_pods.go b/pkg/resource-handler/controller/shard/reconcile_pool_pods.go index 2f7b32a2..eba431ce 100644 --- a/pkg/resource-handler/controller/shard/reconcile_pool_pods.go +++ b/pkg/resource-handler/controller/shard/reconcile_pool_pods.go @@ -78,11 +78,7 @@ func (r *ShardReconciler) reconcilePoolPods( existingPVCs[pvc.Name] = pvc } - // Phase 0: Sync DRAINED labels and reconcile temporary maintenance capacity. - if err := r.syncDrainedLabels(ctx, shard, existingPods); err != nil { - return err - } - drainedCount := countDrainedPods(shard, existingPods) + // Phase 0: Reconcile temporary maintenance capacity. maintenanceSurges, surgeAction, err := r.reconcileCellMaintenanceSurge( ctx, shard, @@ -101,9 +97,8 @@ func (r *ShardReconciler) reconcilePoolPods( return nil } - // DRAINED pods stay alive for investigation; stand-in replicas compensate. // Active maintenance surges remain desired until the cell has settled. - effectiveReplicas := replicas + drainedCount + maintenanceSurges + effectiveReplicas := replicas + maintenanceSurges // Phase 1: Create missing resources and handle terminal/deleted pods driftedCount, actionTaken, err := r.createMissingResources( @@ -155,7 +150,7 @@ func (r *ShardReconciler) reconcilePoolPods( // createMissingResources creates PVCs and Pods that should exist but don't. // It also handles terminal pods (Failed/Succeeded) and externally-deleted pods. -// effectiveReplicas includes stand-in pods for DRAINED pods (replicas + drainedCount). +// effectiveReplicas includes temporary maintenance-surge pods (replicas + maintenanceSurges). // Returns the number of drifted pods and whether an action was taken this reconcile. func (r *ShardReconciler) createMissingResources( ctx context.Context, @@ -370,10 +365,10 @@ func isPodReady(pod *corev1.Pod) bool { // isPoolHealthy returns true if the pool/cell has at least effectiveReplicas // pods, and all of them — except extras (index >= effectiveReplicas) and -// DRAINED/QUARANTINED pods — are Ready, with none draining or terminating. +// QUARANTINED pods — are Ready, with none draining or terminating. // Extra pods are excluded so an unhealthy extra pod does not block its own -// removal. DRAINED and QUARANTINED pods are excluded because they are -// expected to be unhealthy and should not block scale-down of stand-in pods. +// removal. QUARANTINED pods are excluded because they are expected to be +// unhealthy and are being replaced by quarantine remediation. // // The count check matters because a pod drained all the way to deletion // disappears from existingPods entirely — there is nothing left for the @@ -393,7 +388,7 @@ func isPoolHealthy( if idx, ok := resolvePodIndex(pod.Name); !ok || idx >= int(effectiveReplicas) { continue } - if role := resolvePodRole(shard, pod.Name); role == "DRAINED" || role == "QUARANTINED" { + if resolvePodRole(shard, pod.Name) == "QUARANTINED" { continue } if !pod.DeletionTimestamp.IsZero() { @@ -443,8 +438,7 @@ func (r *ShardReconciler) isShardHealthy( replicas = *pool.ReplicasPerCell } group := podsByPoolCell[string(poolName)+"/"+string(cell)] - effectiveReplicas := replicas + countDrainedPods(shard, group) - if !isPoolHealthy(group, effectiveReplicas, shard) { + if !isPoolHealthy(group, replicas, shard) { return false, nil } } @@ -494,7 +488,7 @@ func (r *ShardReconciler) handleExternalDeletion( // handleScaleDown processes pods that need removal: ready-for-deletion cleanup // and draining extra pods beyond the effective replica count. -// replicas is the user-desired count; effectiveReplicas = replicas + drainedCount. +// replicas is the user-desired count; effectiveReplicas = replicas + maintenanceSurges. // Returns whether an action was taken and whether any drain is in progress. func (r *ShardReconciler) handleScaleDown( ctx context.Context, @@ -936,43 +930,6 @@ func (r *ShardReconciler) selectPodToDrain( return bestPod } -// syncDrainedLabels ensures pods with topology role DRAINED have the -// multigres.com/role=DRAINED label, and pods no longer DRAINED have it removed. -// The label is the durable signal for DRAINED PVC cleanup — PodRoles may be -// cleared by the data-handler during drain before cleanup runs. -func (r *ShardReconciler) syncDrainedLabels( - ctx context.Context, - shard *multigresv1alpha1.Shard, - existingPods map[string]*corev1.Pod, -) error { - for _, pod := range existingPods { - role := resolvePodRole(shard, pod.Name) - currentLabel := pod.Labels[metadata.LabelPodRole] - - if role == "DRAINED" && currentLabel != "DRAINED" { - patch := client.MergeFrom(pod.DeepCopy()) - if pod.Labels == nil { - pod.Labels = make(map[string]string) - } - pod.Labels[metadata.LabelPodRole] = "DRAINED" - if err := r.Patch(ctx, pod, patch); err != nil { - return fmt.Errorf("failed to set DRAINED label on pod %s: %w", pod.Name, err) - } - r.Recorder.Eventf(shard, "Warning", "PodDrained", - "Pod %s detected as DRAINED — provisioning stand-in replica", pod.Name) - } else if role != "DRAINED" && currentLabel == "DRAINED" { - patch := client.MergeFrom(pod.DeepCopy()) - delete(pod.Labels, metadata.LabelPodRole) - if err := r.Patch(ctx, pod, patch); err != nil { - return fmt.Errorf("failed to remove DRAINED label from pod %s: %w", pod.Name, err) - } - r.Recorder.Eventf(shard, "Normal", "PodRecovered", - "Pod %s is no longer DRAINED", pod.Name) - } - } - return nil -} - func (r *ShardReconciler) cleanupDrainedPod( ctx context.Context, shard *multigresv1alpha1.Shard, @@ -983,27 +940,7 @@ func (r *ShardReconciler) cleanupDrainedPod( ) error { logger := log.FromContext(ctx) - // DRAINED pods always get their PVC marked orphan — data is known-bad. - // The multigres-gc CronJob deletes the PVC after the retention window. - // We check the pod label (not PodRoles) because the data-handler clears - // the topology entry during drain before this cleanup point. - if pod.Labels[metadata.LabelPodRole] == "DRAINED" { - if err := r.cleanupPodPVC( - ctx, - shard, - pod, - poolName, - "DRAINED (data known-bad)", - ); err != nil { - return err - } - logger.Info("Drained pod cleanup complete", "pod", pod.Name) - r.Recorder.Eventf(shard, "Normal", "DrainCompleted", - "Completed drain for DRAINED pod %s — PVC cleanup queued", pod.Name) - return nil - } - - // For non-DRAINED pods, respect WhenScaled policy + // Respect the WhenScaled PVC-deletion policy. mergedPolicy := multigresv1alpha1.MergePVCDeletionPolicy( poolSpec.PVCDeletionPolicy, shard.Spec.PVCDeletionPolicy, diff --git a/pkg/resource-handler/controller/shard/reconcile_quarantine.go b/pkg/resource-handler/controller/shard/reconcile_quarantine.go index d14fc6eb..db05103f 100644 --- a/pkg/resource-handler/controller/shard/reconcile_quarantine.go +++ b/pkg/resource-handler/controller/shard/reconcile_quarantine.go @@ -57,12 +57,6 @@ const ( // // Returns true when it took a (destructive) action so the caller requeues and // skips other disruptive work this cycle. -// -// NOTE(review): this supersedes, for quarantined poolers, the older "stand-in -// replica" model (GetPoolerStatus previously mapped quarantined -> DRAINED, -// which provisioned a replacement at a new index and kept the bad pod). That -// DRAINED machinery in reconcile_pool_pods.go is now dormant for the quarantine -// case; a follow-up can remove it if we settle on wipe-in-place. func (r *ShardReconciler) reconcileQuarantineRemediation( ctx context.Context, store topoclient.Store, diff --git a/pkg/resource-handler/controller/shard/shard_controller_internal_test.go b/pkg/resource-handler/controller/shard/shard_controller_internal_test.go index 2dab036e..344f3ce7 100644 --- a/pkg/resource-handler/controller/shard/shard_controller_internal_test.go +++ b/pkg/resource-handler/controller/shard/shard_controller_internal_test.go @@ -883,7 +883,7 @@ func TestUpdateStatus_FieldOwner(t *testing.T) { // TestHandleScaleDown_ConcurrentDrainPrevention verifies that handleScaleDown // respects the inProgress flag: when any pod already has a drain annotation // (DrainStateRequested, DrainStateDraining, or DrainStateAcknowledged), no new -// drains are initiated for either DRAINED replacement or extra-pod scale-down. +// drains are initiated for extra-pod scale-down. func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { scheme := runtime.NewScheme() _ = multigresv1alpha1.AddToScheme(scheme) @@ -933,14 +933,13 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { tests := map[string]struct { replicas int32 pods []*corev1.Pod - podRoles map[string]string actionTaken bool wantAction bool wantInProgress bool wantNoDrains bool wantDrainedPod string }{ - "drain in progress (DrainStateRequested) blocks DRAINED replacement": { + "drain in progress (DrainStateRequested) blocks a new drain": { replicas: 2, pods: []*corev1.Pod{ makePod(podName0, map[string]string{ @@ -948,14 +947,11 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { }), makePod(podName1, nil), }, - podRoles: map[string]string{ - podName1: "DRAINED", - }, wantAction: false, wantInProgress: true, wantNoDrains: true, }, - "drain in progress (DrainStateDraining) blocks DRAINED replacement": { + "drain in progress (DrainStateDraining) blocks a new drain": { replicas: 2, pods: []*corev1.Pod{ makePod(podName0, map[string]string{ @@ -963,14 +959,11 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { }), makePod(podName1, nil), }, - podRoles: map[string]string{ - podName1: "DRAINED", - }, wantAction: false, wantInProgress: true, wantNoDrains: true, }, - "drain in progress (DrainStateAcknowledged) blocks DRAINED replacement": { + "drain in progress (DrainStateAcknowledged) blocks a new drain": { replicas: 2, pods: []*corev1.Pod{ makePod(podName0, map[string]string{ @@ -978,9 +971,6 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { }), makePod(podName1, nil), }, - podRoles: map[string]string{ - podName1: "DRAINED", - }, wantAction: false, wantInProgress: true, wantNoDrains: true, @@ -1009,15 +999,12 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { wantInProgress: true, wantNoDrains: true, }, - "no drain in progress with DRAINED pod does not auto-drain": { + "no drain in progress with matching replicas does not auto-drain": { replicas: 2, pods: []*corev1.Pod{ makePod(podName0, nil), makePod(podName1, nil), }, - podRoles: map[string]string{ - podName1: "DRAINED", - }, wantAction: false, wantNoDrains: true, }, @@ -1031,15 +1018,12 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { wantInProgress: false, wantDrainedPod: podName1, }, - "actionTaken from earlier phase blocks DRAINED replacement": { + "actionTaken from earlier phase blocks a new drain": { replicas: 2, pods: []*corev1.Pod{ makePod(podName0, nil), makePod(podName1, nil), }, - podRoles: map[string]string{ - podName1: "DRAINED", - }, actionTaken: true, wantAction: true, wantInProgress: false, @@ -1056,7 +1040,7 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { wantInProgress: false, wantNoDrains: true, }, - "pod with DeletionTimestamp sets inProgress and blocks DRAINED replacement": { + "pod with DeletionTimestamp sets inProgress and blocks a new drain": { replicas: 2, pods: []*corev1.Pod{ func() *corev1.Pod { @@ -1068,14 +1052,11 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { }(), makePod(podName1, nil), }, - podRoles: map[string]string{ - podName1: "DRAINED", - }, wantAction: false, wantInProgress: true, wantNoDrains: true, }, - "multiple DRAINED pods with drain in progress drains none": { + "multiple pods with drain in progress drains none": { replicas: 3, pods: []*corev1.Pod{ makePod(podName0, map[string]string{ @@ -1084,10 +1065,6 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { makePod(podName1, nil), makePod(podName2, nil), }, - podRoles: map[string]string{ - podName1: "DRAINED", - podName2: "DRAINED", - }, wantAction: false, wantInProgress: true, wantNoDrains: true, @@ -1111,7 +1088,6 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { for testName, tc := range tests { t.Run(testName, func(t *testing.T) { shard := baseShard.DeepCopy() - shard.Status.PodRoles = tc.podRoles objects := make([]client.Object, 0, len(tc.pods)+1) objects = append(objects, shard) @@ -1141,14 +1117,6 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { poolSpec := multigresv1alpha1.PoolSpec{} - var drainedInTest int32 - for _, role := range tc.podRoles { - if role == "DRAINED" { - drainedInTest++ - } - } - effectiveReplicas := tc.replicas + drainedInTest - gotAction, gotInProgress, err := reconciler.handleScaleDown( context.Background(), shard, @@ -1156,7 +1124,7 @@ func TestHandleScaleDown_ConcurrentDrainPrevention(t *testing.T) { poolSpec, existingPods, tc.replicas, - effectiveReplicas, + tc.replicas, tc.actionTaken, ) if err != nil { @@ -1389,52 +1357,30 @@ func TestCleanupDrainedPod_PVCDeletion(t *testing.T) { podName string pvcName string siblingCount int - podRoles map[string]string policy *multigresv1alpha1.PVCDeletionPolicy wantPVC bool wantOrphan bool }{ - "DRAINED small pool (3) -> orphan": { + "idx orphan": { + "idx orphan": { podName: podName5, pvcName: pvcName5, siblingCount: 3, - podRoles: map[string]string{}, - policy: deletePolicy, - wantPVC: true, - wantOrphan: true, - }, - "DRAINED large pool (4, scaling to 3) -> orphan": { - podName: podName0, - pvcName: pvcName0, - siblingCount: 4, - podRoles: map[string]string{podName0: "DRAINED"}, policy: deletePolicy, wantPVC: true, wantOrphan: true, @@ -1443,7 +1389,6 @@ func TestCleanupDrainedPod_PVCDeletion(t *testing.T) { podName: podName5, pvcName: pvcName5, siblingCount: 4, - podRoles: map[string]string{}, policy: deletePolicy, wantPVC: true, wantOrphan: true, @@ -1453,12 +1398,8 @@ func TestCleanupDrainedPod_PVCDeletion(t *testing.T) { for tn, tc := range tests { t.Run(tn, func(t *testing.T) { shard := baseShard.DeepCopy() - shard.Status.PodRoles = tc.podRoles pod := makePod(tc.podName) - if tc.podRoles[tc.podName] == "DRAINED" { - pod.Labels[metadata.LabelPodRole] = "DRAINED" - } // Create the target PVC plus filler siblings so the pool+cell has // exactly siblingCount PVCs. @@ -2308,7 +2249,7 @@ func TestCreateMissingResources(t *testing.T) { desiredPod0, _ := BuildPoolPod(shard, poolName, cellName, poolSpec, 0, scheme) hash0 := ComputeSpecHash(desiredPod0) - // Pod 0: not ready, has matching spec hash (not drifted, not DRAINED) + // Pod 0: not ready, has matching spec hash (not drifted) pod0 := &corev1.Pod{ ObjectMeta: metav1.ObjectMeta{ Name: podName0, @@ -2779,9 +2720,9 @@ func TestIsPoolHealthy(t *testing.T) { } }) - t.Run("DRAINED pod that is not ready does not block health check", func(t *testing.T) { - drainedShard := shard.DeepCopy() - drainedShard.Status.PodRoles = map[string]string{"pod-0": "DRAINED"} + t.Run("QUARANTINED pod that is not ready does not block health check", func(t *testing.T) { + quarantinedShard := shard.DeepCopy() + quarantinedShard.Status.PodRoles = map[string]string{"pod-0": "QUARANTINED"} pods := map[string]*corev1.Pod{ "pod-0": { ObjectMeta: metav1.ObjectMeta{Name: "pod-0"}, @@ -2792,8 +2733,8 @@ func TestIsPoolHealthy(t *testing.T) { }, }, } - if !isPoolHealthy(pods, 1, drainedShard) { - t.Error("DRAINED pod should be excluded from health check") + if !isPoolHealthy(pods, 1, quarantinedShard) { + t.Error("QUARANTINED pod should be excluded from health check") } }) } @@ -5060,17 +5001,6 @@ func TestIsDrainStale(t *testing.T) { } }) - t.Run("DoesNotCancelDrainedPodDrain", func(t *testing.T) { - shardWithDrained := shard.DeepCopy() - pod := matchingPod(0, metadata.DrainStateRequested) - shardWithDrained.Status.PodRoles = map[string]string{ - pod.Name: "DRAINED", - } - if r.isDrainStale(shardWithDrained, pod, metadata.DrainStateRequested) { - t.Error("expected drain NOT to be stale (pod role is DRAINED)") - } - }) - t.Run("DoesNotCancelDrainOnDeletingPod", func(t *testing.T) { pod := matchingPod(4, metadata.DrainStateRequested) now := metav1.Now() @@ -6044,43 +5974,6 @@ func TestReconcilePoolPods_AdditionalErrorPaths(t *testing.T) { } poolSpec := shard.Spec.Pools["main"] - t.Run("syncDrainedLabels error", func(t *testing.T) { - pod := &corev1.Pod{ - ObjectMeta: metav1.ObjectMeta{ - Name: BuildPoolPodName(shard, "main", "z1", 0), - Namespace: "default", - Labels: map[string]string{ - metadata.LabelMultigresCluster: "test-cluster", - metadata.LabelMultigresDatabase: "db", - metadata.LabelMultigresTableGroup: "tg", - metadata.LabelMultigresShard: "test-shard", - metadata.LabelMultigresPool: "main", - metadata.LabelMultigresCell: "z1", - metadata.LabelPodRole: "DRAINED", - }, - }, - } - - baseClient := fake.NewClientBuilder().WithScheme(scheme).WithObjects(shard, pod).Build() - fails := testutil.NewFakeClientWithFailures(baseClient, &testutil.FailureConfig{ - OnPatch: testutil.FailOnObjectName(pod.Name, testutil.ErrPermissionError), - }) - - r := &ShardReconciler{Client: fails, Scheme: scheme, Recorder: record.NewFakeRecorder(10)} - - // Temporarily set the topology returned role to NOT DRAINED, which causes syncDrainedLabels to try and strip the label - // which triggers the patch failure - // Wait, the test uses resolvePodRole, which defaults to whatever unless mocked. - // Since topo is not mocked here, resolvePodRole returns "" -> not DRAINED -> tries to patch remove DRAINED -> fails - existingPods := map[string]*corev1.Pod{ - pod.Name: pod, - } - err := r.syncDrainedLabels(context.Background(), shard, existingPods) - if err == nil || !strings.Contains(err.Error(), "failed to remove DRAINED label") { - t.Fatalf("expected syncDrainedLabels error, got %v", err) - } - }) - t.Run("CreatePVC error", func(t *testing.T) { baseClient := fake.NewClientBuilder().WithScheme(scheme).WithObjects(shard).Build() pvcName := BuildPoolDataPVCName(shard, "main", "z1", 0) @@ -6103,10 +5996,11 @@ func TestReconcilePoolPods_AdditionalErrorPaths(t *testing.T) { }) t.Run("markPodPVCOrphan network error", func(t *testing.T) { - // Mock a DRAINED pod so markPodPVCOrphan is called during cleanupDrainedPod + // A scaled-down pod (index >= replicas) so cleanupDrainedPod orphans its + // PVC, exercising the markPodPVCOrphan patch-failure path. pod := &corev1.Pod{ ObjectMeta: metav1.ObjectMeta{ - Name: BuildPoolPodName(shard, "main", "z1", 0), + Name: BuildPoolPodName(shard, "main", "z1", 1), Namespace: "default", Labels: map[string]string{ metadata.LabelMultigresCluster: "test-cluster", @@ -6115,7 +6009,6 @@ func TestReconcilePoolPods_AdditionalErrorPaths(t *testing.T) { metadata.LabelMultigresShard: "test-shard", metadata.LabelMultigresPool: "main", metadata.LabelMultigresCell: "z1", - metadata.LabelPodRole: "DRAINED", }, Annotations: map[string]string{ metadata.AnnotationDrainState: metadata.DrainStateReadyForDeletion, @@ -6123,7 +6016,7 @@ func TestReconcilePoolPods_AdditionalErrorPaths(t *testing.T) { }, } - pvcName := BuildPoolDataPVCName(shard, "main", "z1", 0) + pvcName := BuildPoolDataPVCName(shard, "main", "z1", 1) pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{ Name: pvcName, diff --git a/pkg/resource-handler/controller/shard/shard_controller_test.go b/pkg/resource-handler/controller/shard/shard_controller_test.go index f07a4b3f..ea822b4a 100644 --- a/pkg/resource-handler/controller/shard/shard_controller_test.go +++ b/pkg/resource-handler/controller/shard/shard_controller_test.go @@ -1804,7 +1804,7 @@ func TestScaleDown_HealthGateBlocksDrain(t *testing.T) { context.Background(), shard, poolName, multigresv1alpha1.PoolSpec{}, existingPods, 2, // replicas: pod-2 is extra - 2, // effectiveReplicas + 2, // effectiveReplicas: no maintenance surge in this test false, ) if err != nil { @@ -1875,7 +1875,7 @@ func TestScaleDown_HealthGateBlocksDrain(t *testing.T) { context.Background(), shard, poolName, multigresv1alpha1.PoolSpec{}, existingPods, 2, // replicas: pod-2 is extra - 2, // effectiveReplicas + 2, // effectiveReplicas: no maintenance surge in this test false, ) if err != nil { @@ -1988,7 +1988,7 @@ func TestScaleDown_HealthGateBlocksDrain(t *testing.T) { context.Background(), shard, poolName, multigresv1alpha1.PoolSpec{}, existingPods, 2, // replicas: pod-2 is extra - 2, // effectiveReplicas + 2, // effectiveReplicas: no maintenance surge in this test false, ) if err != nil { @@ -2252,111 +2252,3 @@ func TestRollingUpdateWaitsForSiblingCell(t *testing.T) { ) } } - -func TestDrainedPodReplacement(t *testing.T) { - t.Parallel() - scheme := runtime.NewScheme() - _ = multigresv1alpha1.AddToScheme(scheme) - _ = corev1.AddToScheme(scheme) - - shardObj := &multigresv1alpha1.Shard{ - ObjectMeta: metav1.ObjectMeta{ - Name: "test-shard", Namespace: "default", - Labels: map[string]string{metadata.LabelMultigresCluster: "test-cluster"}, - }, - Spec: multigresv1alpha1.ShardSpec{ - DatabaseName: "db", - TableGroupName: "tg", - ShardName: "s1", - Pools: map[multigresv1alpha1.PoolName]multigresv1alpha1.PoolSpec{ - "primary": { - ReplicasPerCell: ptr.To(int32(1)), - Storage: multigresv1alpha1.StorageSpec{Size: "10Gi"}, - }, - }, - }, - } - - podName0 := BuildPoolPodName(shardObj, "primary", "zone1", 0) - shardObj.Status.PodRoles = map[string]string{ - podName0: "DRAINED", - } - - pod := &corev1.Pod{ - ObjectMeta: metav1.ObjectMeta{ - Name: podName0, - Namespace: "default", - Labels: map[string]string{ - "app.kubernetes.io/component": "shard-pool", - "app.kubernetes.io/instance": "test-cluster", - metadata.LabelMultigresCluster: "test-cluster", - metadata.LabelMultigresDatabase: "db", - metadata.LabelMultigresTableGroup: "tg", - metadata.LabelMultigresShard: "s1", - metadata.LabelMultigresPool: "primary", - metadata.LabelMultigresCell: "zone1", - metadata.LabelPodRole: "DRAINED", - }, - Annotations: map[string]string{}, - }, - Status: corev1.PodStatus{ - Conditions: []corev1.PodCondition{ - {Type: corev1.PodReady, Status: corev1.ConditionTrue}, - }, - }, - } - - // Pre-calculate hash so it doesn't get deleted as drifted - poolSpec := shardObj.Spec.Pools["primary"] - desiredPod, _ := BuildPoolPod(shardObj, "primary", "zone1", poolSpec, 0, scheme) - pod.Annotations[metadata.AnnotationSpecHash] = desiredPod.Annotations[metadata.AnnotationSpecHash] - - c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(shardObj).Build() - r := &ShardReconciler{Client: c, Scheme: scheme, Recorder: record.NewFakeRecorder(10)} - - if err := c.Create(context.Background(), pod); err != nil { - t.Fatalf("failed to create pod: %v", err) - } - - // Run reconcile loop - err := r.reconcilePoolPods( - context.Background(), - shardObj, - "primary", - "zone1", - poolSpec, - &shardRolloutTracker{}, - ) - if err != nil { - t.Fatalf("unexpected error: %v", err) - } - - // The DRAINED pod should NOT have a drain annotation (DRAINED pods are no longer auto-drained) - err = c.Get( - context.Background(), - types.NamespacedName{Name: podName0, Namespace: "default"}, - pod, - ) - if err != nil { - t.Fatalf("expected pod to exist, got %v", err) - } - - if pod.Annotations[metadata.AnnotationDrainState] != "" { - t.Errorf( - "Expected DRAINED pod to have no drain annotation, got %q", - pod.Annotations[metadata.AnnotationDrainState], - ) - } - - // A stand-in pod at index 1 should be created (replicas=1, effectiveReplicas=2) - podName1 := BuildPoolPodName(shardObj, "primary", "zone1", 1) - standInPod := &corev1.Pod{} - err = c.Get( - context.Background(), - types.NamespacedName{Name: podName1, Namespace: "default"}, - standInPod, - ) - if err != nil { - t.Fatalf("expected stand-in pod at index 1 to be created, got %v", err) - } -} diff --git a/pkg/util/metadata/labels.go b/pkg/util/metadata/labels.go index a741b17f..e352962d 100644 --- a/pkg/util/metadata/labels.go +++ b/pkg/util/metadata/labels.go @@ -160,12 +160,6 @@ const ( // cross-cell durability during a rollout or explicit maintenance request. AnnotationMaintenanceSurge = "maintenance.multigres.com/surge" - // LabelPodRole reflects the pod's topology role (e.g. "DRAINED"). - // Set by the resource-handler when the topology store reports a notable role. - // Used as the durable signal for DRAINED PVC cleanup, since PodRoles may be - // cleared by the data-handler during drain before cleanup runs. - LabelPodRole = "multigres.com/role" - // LabelOrphan marks a resource as orphaned (no longer referenced by an owner) // and eligible for garbage collection. Value is the UTC timestamp when the // resource became orphaned, formatted as OrphanTimestampFormat. From 0b64c8d41b95177154af341df6eeee9c3c61a4a2 Mon Sep 17 00:00:00 2001 From: Mats Kindahl Date: Fri, 25 Sep 2026 10:10:36 +0200 Subject: [PATCH 2/3] docs(shard): drop remaining DRAINED references from API and docs Follow-up to the DRAINED stand-in removal, addressing review feedback. - Update the Shard.Status.PodRoles godoc to list QUARANTINED instead of DRAINED and regenerate the CRD description. - Refresh the operator capability, pod-management-architecture, and pvc-lifecycle docs: replace the DRAINED stand-in / DRAINED-specific PVC handling with the quarantine-remediation and orphan-on-scale-down behavior that actually runs today. - Add a dated API-design changelog entry recording the removal. The multiorch-facing "DRAINED Pooler Replacement" narrative in pod-management-design.md is left in place but marked obsolete / needs rewrite, since rewriting the multiorch internals is out of scope here. Signed-off-by: Mats Kindahl Co-Authored-By: Claude Opus 4.8 --- api/v1alpha1/shard_types.go | 2 +- config/crd/bases/multigres.com_shards.yaml | 2 +- .../api-design/multigres-operator-api-v1alpha1-design.md | 1 + docs/development/pod-management-architecture.md | 8 ++++---- docs/development/pod-management-design.md | 3 +++ docs/development/pvc-lifecycle.md | 2 +- docs/operator-capability-levels.md | 9 +++++---- 7 files changed, 16 insertions(+), 11 deletions(-) diff --git a/api/v1alpha1/shard_types.go b/api/v1alpha1/shard_types.go index d0d4c7a6..c371ea90 100644 --- a/api/v1alpha1/shard_types.go +++ b/api/v1alpha1/shard_types.go @@ -326,7 +326,7 @@ type ShardStatus struct { // +kubebuilder:validation:Enum=full;diff;incr LastBackupType string `json:"lastBackupType,omitempty"` - // PodRoles maps pod names to their database roles (e.g. PRIMARY, REPLICA, DRAINED). + // PodRoles maps pod names to their database roles (e.g. PRIMARY, REPLICA, QUARANTINED). // +optional PodRoles map[string]string `json:"podRoles,omitempty"` diff --git a/config/crd/bases/multigres.com_shards.yaml b/config/crd/bases/multigres.com_shards.yaml index 082859d1..67e94a64 100644 --- a/config/crd/bases/multigres.com_shards.yaml +++ b/config/crd/bases/multigres.com_shards.yaml @@ -3042,7 +3042,7 @@ spec: additionalProperties: type: string description: PodRoles maps pod names to their database roles (e.g. - PRIMARY, REPLICA, DRAINED). + PRIMARY, REPLICA, QUARANTINED). type: object poolsReady: description: PoolsReady indicates if all data pools are ready. diff --git a/docs/development/api-design/multigres-operator-api-v1alpha1-design.md b/docs/development/api-design/multigres-operator-api-v1alpha1-design.md index e4ac0f87..baefca01 100644 --- a/docs/development/api-design/multigres-operator-api-v1alpha1-design.md +++ b/docs/development/api-design/multigres-operator-api-v1alpha1-design.md @@ -1540,6 +1540,7 @@ spec: * **2026-03-20:** Added `PostgresConfigRef` type (ConfigMap reference) at the shard level for custom `postgresql.conf` parameters. Operator mounts the user-provided ConfigMap and sets `POSTGRES_INITDB_EXTRA_CONF` on pgctld so the user-supplied lines are appended to pgctld's auto-tuned config. See design rationale below. * **2026-03-20:** Added `ExternalGatewayConfig` (`spec.externalGateway`) for external exposure of the global multigateway Service via `externalIPs` and user-provided annotations. Added `GatewayStatus` (`status.gateway.externalEndpoint`) and `GatewayExternalReady` condition. Added `InitdbArgs` field to `ShardTemplateSpec`, `ShardInlineSpec`, and `ShardOverrides` for passing extra arguments to `initdb` during PostgreSQL data directory initialization. Added `IPAddress` and `InitdbArgs` validated types to `common_types.go`. * **2026-03-20:** Added `ExternalAdminWebConfig` (`spec.externalAdminWeb`) for external exposure of the global multiadmin-web Service via `externalIPs` and user-provided annotations. Mirrors the `ExternalGatewayConfig` pattern. Added `AdminWebStatus` (`status.adminWeb.externalEndpoint`) and `AdminWebExternalReady` condition. Readiness is driven by the admin-web Deployment's `ReadyReplicas` (unlike the gateway which aggregates across Cell CRs). + * **2026-09-25:** Removed the dormant DRAINED stand-in machinery (superseding the 2026-03-17 entry). Multigres no longer emits the `DRAINED` topology role — poolers whose postgres cannot start are surfaced as `QUARANTINED` and remediated in place (delete pod + wipe data PVC + re-bootstrap from backup), so the stand-in-replica model and the `multigres.com/role=DRAINED` PVC handling were dead code. `PodRoles` values are now `PRIMARY`, `REPLICA`, or `QUARANTINED`. ## Design Rationale: PostgreSQL Configuration (`postgresConfigRef`) diff --git a/docs/development/pod-management-architecture.md b/docs/development/pod-management-architecture.md index 7a9c7e02..192a99e5 100644 --- a/docs/development/pod-management-architecture.md +++ b/docs/development/pod-management-architecture.md @@ -136,7 +136,7 @@ PVC lifecycle is managed through **conditional owner references** based on the ` - When `WhenDeleted` is `Delete`: PVCs are created with an ownerRef pointing to the Shard CR, enabling Kubernetes garbage collection to cascade-delete them when the Shard is removed. - When `WhenDeleted` is `Retain`: PVCs are created without ownerRefs, ensuring they persist after Shard deletion. - The shard controller's `reconcilePVCOwnerRefs` function ensures existing PVCs stay in sync with the current policy — adding or removing ownerRefs as the policy changes mid-lifecycle. -- During scale-down, `cleanupDrainedPod` checks `WhenScaled` and deletes data PVCs directly if the policy is `Delete` (the default). For DRAINED pods (identified by the `multigres.com/role=DRAINED` label), PVCs are always deleted regardless of the `WhenScaled` policy because DRAINED pod data is known-bad. +- During scale-down, `cleanupDrainedPod` checks `WhenScaled`: for a scaled-down pod (index >= replicas) under the `Delete` policy (the default) it orphans the data PVC for garbage collection, while a rolling-update pod (index < replicas) keeps its PVC. --- @@ -408,7 +408,7 @@ Metrics are emitted per pool via `monitoring.SetShardPoolReplicas()`. A `PoolEmp | **Etcd topology cleanup** | `UnregisterMultiPooler` called during drain flow; stale entries removed on pod termination | | **Topology registration & pruning** | Cell and database registration centralized in MultigresCluster controller; stale entries pruned when `topologyPruning.enabled` (default) | | **Backup health reporting** | Shard controller calls `GetBackups` RPC, sets `BackupHealthy` condition and `LastBackupTime` status | -| **DRAINED pod handling** | DRAINED pods (diverged data, pg_rewind failure) are kept alive for admin investigation. Stand-in replicas created at next index for availability. Admin discards via `kubectl delete pod`, triggering drain + PVC deletion | +| **QUARANTINED pod handling** | QUARANTINED pods (postgres cannot start) are remediated in place: the operator deletes the pod, wipes its data PVC, and re-bootstraps from backup at the same index | | **Shard-wide drain serialization** | Planned drains check all persisted drain/termination state across reconciles and wait for data-plane recovery before the next removal | | **Two-cell maintenance surge** | Creates and verifies temporary same-cell capacity before disrupting the final ready member of a two-cell cross-cell shard | | **Scale-down health gate** | Drains deferred when either the current pool or another shard pool/cell has non-ready pods | @@ -523,7 +523,7 @@ Pods reference this via `spec.subdomain`, combined with `spec.hostname` (set to ### Etcd Topology The operator reads from etcd topology (via the shard controller) to: -- Determine pod roles (`PRIMARY`, `REPLICA`, `DRAINED`) for scale-down and rolling-update decisions. +- Determine pod roles (`PRIMARY`, `REPLICA`, `QUARANTINED`) for scale-down and rolling-update decisions. - Clean up stale topology entries on permanent pod removal (`UnregisterMultiPooler`). The operator writes to etcd topology (via the MultigresCluster controller) to: @@ -559,7 +559,7 @@ Through the shard controller: | Scale-up (new replicas) | Create PVC + Pod, multigres handles the rest | | Scale-down | Drain state machine with standby removal + etcd unregistration | | Rolling update | Spec-hash detection, ordered recreation | -| DRAINED pod handling | Detects DRAINED role from etcd, keeps pod alive for investigation, creates stand-in replica | +| QUARANTINED pod handling | Detects QUARANTINED role from etcd, remediates in place (delete pod + wipe data PVC + re-bootstrap from backup) | | Backup health reporting | Calls `GetBackups` RPC, sets `BackupHealthy` condition | | PVC lifecycle | Direct creation/deletion per policy | | Certificate provisioning | `pkg/cert` for pgBackRest TLS | diff --git a/docs/development/pod-management-design.md b/docs/development/pod-management-design.md index 46cfe0bc..51bdda17 100644 --- a/docs/development/pod-management-design.md +++ b/docs/development/pod-management-design.md @@ -529,6 +529,9 @@ When scaling down, the operator must choose which pod to delete. The selection a ### DRAINED Pooler Replacement +> [!WARNING] +> **Obsolete — needs rewrite.** This section (and the other `DRAINED` references in this design doc) describes the retired "stand-in replica" model. Multigres no longer emits the `DRAINED` topology role (`PoolerType.DRAINED` is deprecated and never produced); poolers whose postgres cannot start are surfaced as `QUARANTINED` and remediated in place (delete pod + wipe data PVC + re-bootstrap from backup). The stand-in machinery has been removed from the operator. This narrative should be rewritten around quarantine remediation; it is left in place here pending that rewrite. + Multiorch may mark a pooler as `DRAINED` in etcd independently of the operator (e.g., during internal recovery or rebalancing). When this happens, the operator must: 1. **Detect the DRAINED pooler** by reading etcd topology during reconciliation and comparing `PoolerType` for each pod's `service-id` diff --git a/docs/development/pvc-lifecycle.md b/docs/development/pvc-lifecycle.md index 2b63e4c1..7c4e9df1 100644 --- a/docs/development/pvc-lifecycle.md +++ b/docs/development/pvc-lifecycle.md @@ -74,7 +74,7 @@ sts.Spec.PersistentVolumeClaimRetentionPolicy = pvc.BuildRetentionPolicy( ``` **Shard Pool Pods** (`pkg/resource-handler/controller/shard/reconcile_pool_pods.go`): -Pool pod PVCs are managed directly by the operator. During scale-down, the `cleanupDrainedPod` function checks `shard.Spec.PVCDeletionPolicy.WhenScaled` and deletes the data PVC if the policy is `Delete` (the default). For DRAINED pods (detected via the `multigres.com/role=DRAINED` label), the PVC is always deleted regardless of the `WhenScaled` policy because DRAINED pod data is known-bad. During shard/cluster deletion, PVCs are garbage-collected by Kubernetes via conditional owner references: when `WhenDeleted` is `Delete`, PVCs are created with an ownerRef to the Shard CR, enabling cascade deletion. When `WhenDeleted` is `Retain`, PVCs have no ownerRef and persist after deletion. The `reconcilePVCOwnerRefs` function ensures existing PVCs stay in sync with the current policy during mid-lifecycle changes. +Pool pod PVCs are managed directly by the operator. During scale-down, the `cleanupDrainedPod` function checks `shard.Spec.PVCDeletionPolicy.WhenScaled` and, for a scaled-down pod (index >= replicas) under the `Delete` policy (the default), orphans the data PVC for garbage collection; a rolling-update pod (index < replicas) keeps its PVC. During shard/cluster deletion, PVCs are garbage-collected by Kubernetes via conditional owner references: when `WhenDeleted` is `Delete`, PVCs are created with an ownerRef to the Shard CR, enabling cascade deletion. When `WhenDeleted` is `Retain`, PVCs have no ownerRef and persist after deletion. The `reconcilePVCOwnerRefs` function ensures existing PVCs stay in sync with the current policy during mid-lifecycle changes. **Utility Function** (`pkg/util/pvc/retention.go:BuildRetentionPolicy`): - Used by TopoServer StatefulSets to convert operator's `PVCDeletionPolicy` to Kubernetes `StatefulSetPersistentVolumeClaimRetentionPolicy` diff --git a/docs/operator-capability-levels.md b/docs/operator-capability-levels.md index 5686c45e..3fedc1cf 100644 --- a/docs/operator-capability-levels.md +++ b/docs/operator-capability-levels.md @@ -142,7 +142,7 @@ The operator continuously updates status on all CRs: - **MultigresCluster**: Phase (Healthy/Progressing/Error/Deleting), conditions (Available, Progressing), per-cell and per-database status summaries - **Shard**: Phase, conditions, per-cell pool status, PodRoles map - (PRIMARY/REPLICA/DRAINED), LastBackupTime, LastBackupType + (PRIMARY/REPLICA/QUARANTINED), LastBackupTime, LastBackupType - **Cell/TableGroup/TopoServer**: Phase, conditions, ObservedGeneration - All CRs expose Phase, Available, and Age via `kubectl get` print columns @@ -293,7 +293,7 @@ Two-dimensional PVC deletion policy: | Policy | Retain (default) | Delete | |:-------|:-----------------|:-------| | **WhenDeleted** | PVCs kept after Shard deletion | PVCs removed with Shard | -| **WhenScaled** | PVCs kept after scale-down (explicit `Retain`) | PVCs removed on scale-down (default); always deleted for DRAINED pods | +| **WhenScaled** | PVCs kept after scale-down (explicit `Retain`) | PVCs orphaned for garbage collection on scale-down (default) | --- @@ -411,8 +411,9 @@ The operator continuously reconciles desired state and recovers from failures: entries are missing - **PodRoles refresh**: Continuously updated from topology server to reflect actual database roles -- **DRAINED pod handling**: Detects DRAINED role from etcd, keeps pod alive - for admin investigation, creates stand-in replica for availability +- **QUARANTINED pod handling**: Detects QUARANTINED role from etcd (a pooler + whose postgres cannot start) and remediates in place — deletes the pod, wipes + its data PVC, and re-bootstraps from backup at the same index - **Scale-down safety**: Blocks scale-down when pool is already degraded ### Auto-healing (upstream Multigres) From 022e7a3d12f5d7e0df905c194de779da7f3fcb36 Mon Sep 17 00:00:00 2001 From: Mats Kindahl Date: Fri, 25 Sep 2026 10:10:36 +0200 Subject: [PATCH 3/3] fix(observer): classify QUARANTINED separately in replication checks The observer's replication checks special-cased the DRAINED role and lumped everything else into the replica bucket. Since the operator now surfaces unrecoverable poolers as QUARANTINED (never DRAINED), quarantined pods fell through as replicas: they inflated the expected replica count (producing false "primary has fewer standbys than expected" findings), got probed as replicas (spurious "no WAL receiver" errors), and were never surfaced as a distinct condition. Replace the dead DRAINED handling with QUARANTINED: count and report quarantined pods separately, exclude them from the replica set so expectedReplicas is correct, and update the finding wording to describe quarantine remediation. Add unit coverage for the classification and counting paths. Signed-off-by: Mats Kindahl Co-Authored-By: Claude Opus 4.8 --- tools/observer/pkg/observer/replication.go | 47 ++++---- .../observer/pkg/observer/replication_test.go | 100 ++++++++++++++++++ 2 files changed, 126 insertions(+), 21 deletions(-) create mode 100644 tools/observer/pkg/observer/replication_test.go diff --git a/tools/observer/pkg/observer/replication.go b/tools/observer/pkg/observer/replication.go index 3e01ce0d..eb3336a2 100644 --- a/tools/observer/pkg/observer/replication.go +++ b/tools/observer/pkg/observer/replication.go @@ -32,24 +32,24 @@ func (o *Observer) checkReplication(ctx context.Context) { } o.checkShardReplication(ctx, shard) - var primaryCount, replicaCount, drainedCount int + var primaryCount, replicaCount, quarantinedCount int for _, role := range shard.Status.PodRoles { - switch { - case role == "PRIMARY" || role == "primary": + switch role { + case "PRIMARY", "primary": primaryCount++ - case role == "DRAINED": - drainedCount++ + case "QUARANTINED": + quarantinedCount++ default: replicaCount++ } } replData = append(replData, map[string]any{ - "shard": shard.Name, - "namespace": shard.Namespace, - "primaryCount": primaryCount, - "replicaCount": replicaCount, - "drainedCount": drainedCount, - "podRoles": shard.Status.PodRoles, + "shard": shard.Name, + "namespace": shard.Namespace, + "primaryCount": primaryCount, + "replicaCount": replicaCount, + "quarantinedCount": quarantinedCount, + "podRoles": shard.Status.PodRoles, }) } @@ -66,32 +66,37 @@ func (o *Observer) checkShardReplication(ctx context.Context, shard *multigresv1 shardLabelValue := shard.Labels[common.LabelMultigresShard] password := o.fetchShardPassword(ctx, shard) - // Classify pods by role. + // Classify pods by role. QUARANTINED pods are kept out of the replica set: + // their postgres cannot start, so they neither stream from the primary nor + // count toward the expected replica connection count. Counting them as + // replicas would inflate expectedReplicas and produce false "primary has + // fewer standbys than expected" findings. They are surfaced separately and + // replaced by quarantine remediation (data PVC wipe + re-bootstrap). var primaryPodNames []string var replicaPodNames []string - var drainedPodNames []string + var quarantinedPodNames []string for podName, role := range shard.Status.PodRoles { - switch { - case role == "PRIMARY" || role == "primary": + switch role { + case "PRIMARY", "primary": primaryPodNames = append(primaryPodNames, podName) - case role == "DRAINED": - drainedPodNames = append(drainedPodNames, podName) + case "QUARANTINED": + quarantinedPodNames = append(quarantinedPodNames, podName) default: replicaPodNames = append(replicaPodNames, podName) } } - if len(drainedPodNames) > 0 { + if len(quarantinedPodNames) > 0 { o.reporter.Report(report.Finding{ Severity: report.SeverityWarn, Check: "replication", Component: comp, Message: fmt.Sprintf( - "Shard has %d DRAINED pod(s) awaiting admin intervention: %v", - len(drainedPodNames), drainedPodNames, + "Shard has %d QUARANTINED pod(s) awaiting remediation (data PVC wipe + re-bootstrap): %v", + len(quarantinedPodNames), quarantinedPodNames, ), Details: map[string]any{ - "drainedPods": drainedPodNames, + "quarantinedPods": quarantinedPodNames, }, }) } diff --git a/tools/observer/pkg/observer/replication_test.go b/tools/observer/pkg/observer/replication_test.go new file mode 100644 index 00000000..84565192 --- /dev/null +++ b/tools/observer/pkg/observer/replication_test.go @@ -0,0 +1,100 @@ +package observer + +import ( + "strings" + "testing" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + "github.com/multigres/multigres-operator/tools/observer/pkg/report" +) + +// TestCheckShardReplication_QuarantinedClassification verifies that a +// QUARANTINED pod is surfaced with its own finding and is not mistaken for a +// replica. Before this, QUARANTINED fell through the role switch into the +// replica bucket, which both suppressed the quarantine finding and inflated the +// expected replica count used by the primary replication probe. +func TestCheckShardReplication_QuarantinedClassification(t *testing.T) { + shard := &multigresv1alpha1.Shard{ + ObjectMeta: metav1.ObjectMeta{Name: "s1", Namespace: "test-ns"}, + Status: multigresv1alpha1.ShardStatus{ + PodRoles: map[string]string{ + "s1-primary-zone1-0": "PRIMARY", + "s1-primary-zone1-1": "REPLICA", + "s1-primary-zone1-2": "QUARANTINED", + }, + }, + } + + // No pod objects exist, so podIPs stays empty and the SQL probes are skipped; + // only the role-classification finding is exercised. + o := newTestObserver(shard) + o.checkShardReplication(t.Context(), shard) + + findings := collectFindings(o) + if len(findings) != 1 { + t.Fatalf("expected exactly 1 finding (the quarantine warning), got %d: %+v", len(findings), findings) + } + + f := findings[0] + if f.Severity != report.SeverityWarn { + t.Errorf("expected Warn severity, got %s", f.Severity) + } + if !strings.Contains(f.Message, "QUARANTINED") { + t.Errorf("expected message to mention QUARANTINED, got %q", f.Message) + } + if strings.Contains(f.Message, "DRAINED") { + t.Errorf("message should no longer reference DRAINED, got %q", f.Message) + } + if !strings.Contains(f.Message, "s1-primary-zone1-2") { + t.Errorf("expected message to name the quarantined pod, got %q", f.Message) + } + qp, ok := f.Details["quarantinedPods"].([]string) + if !ok { + t.Fatalf("expected Details[quarantinedPods] to be []string, got %T", f.Details["quarantinedPods"]) + } + if len(qp) != 1 || qp[0] != "s1-primary-zone1-2" { + t.Errorf("expected quarantinedPods=[s1-primary-zone1-2], got %v", qp) + } +} + +// TestCheckReplication_QuarantinedNotCountedAsReplica verifies the aggregate +// counting path: a QUARANTINED pod is tallied as quarantined, not as a replica. +func TestCheckReplication_QuarantinedNotCountedAsReplica(t *testing.T) { + shard := &multigresv1alpha1.Shard{ + ObjectMeta: metav1.ObjectMeta{Name: "s1", Namespace: "test-ns"}, + Status: multigresv1alpha1.ShardStatus{ + PodRoles: map[string]string{ + "s1-primary-zone1-0": "PRIMARY", + "s1-primary-zone1-1": "REPLICA", + "s1-primary-zone1-2": "QUARANTINED", + }, + }, + } + + o := newTestObserver(shard) + o.enableSQLProbe = true + o.probes = newProbeCollector() + + o.checkReplication(t.Context()) + + repl, ok := o.probes.Data()["replication"].(map[string]any) + if !ok { + t.Fatalf("expected replication probe data, got %T", o.probes.Data()["replication"]) + } + shardsData, ok := repl["shards"].([]map[string]any) + if !ok || len(shardsData) != 1 { + t.Fatalf("expected 1 shard entry, got %#v", repl["shards"]) + } + entry := shardsData[0] + if got := entry["quarantinedCount"]; got != 1 { + t.Errorf("expected quarantinedCount=1, got %v", got) + } + if got := entry["replicaCount"]; got != 1 { + t.Errorf("expected replicaCount=1 (quarantined pod excluded), got %v", got) + } + if got := entry["primaryCount"]; got != 1 { + t.Errorf("expected primaryCount=1, got %v", got) + } +}