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) 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. 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) + } +}