From 5c2b7cd9a9877692da0b58cd52826a72e73bf937 Mon Sep 17 00:00:00 2001 From: Jonathan Siegel <248302+usiegj00@users.noreply.github.com> Date: Tue, 29 Sep 2026 13:30:16 +0900 Subject: [PATCH] Adopt upstream's settled check in place of our replica quorum #24 added a ready-replica quorum gate because the readiness loop only iterates pods that are already Running, so a replica being recreated was invisible and the gate passed on an empty set. Upstream solved the same problem in redisPodsSettled, and solved it better: it counts pods against spec.redis.replicas and waits while any is terminating, rather than inferring from whichever replicas happen to answer. Carrying two implementations of one guard is how a fork drifts, so take theirs. Port util.PodIsReady and redisPodsSettled verbatim, drop our quorum gate, and replace our bespoke test with upstream's table test - adapted only where the master rows differ, since a stale master is now promoted away rather than deleted. What remains ours is the handover itself (#26), which is proposed upstream as Saremox/redis-operator#201. If that lands, this file converges. --- operator/redisfailover/checker.go | 81 +++++----- operator/redisfailover/checker_test.go | 206 ++++++++++--------------- operator/redisfailover/util/pod.go | 9 ++ 3 files changed, 130 insertions(+), 166 deletions(-) diff --git a/operator/redisfailover/checker.go b/operator/redisfailover/checker.go index 2a1c28c00..f16608546 100644 --- a/operator/redisfailover/checker.go +++ b/operator/redisfailover/checker.go @@ -3,15 +3,18 @@ package redisfailover import ( "context" "errors" + "fmt" "strconv" "time" "github.com/saremox/redis-operator/service/k8s" + appsv1 "k8s.io/api/apps/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" redisfailoverv1 "github.com/saremox/redis-operator/api/redisfailover/v1" "github.com/saremox/redis-operator/metrics" rfservice "github.com/saremox/redis-operator/operator/redisfailover/service" + "github.com/saremox/redis-operator/operator/redisfailover/util" "github.com/saremox/redis-operator/service/redis" ) @@ -28,7 +31,6 @@ func (r *RedisFailoverHandler) UpdateRedisesPods(rf *redisfailoverv1.RedisFailov r.logger.WithField("namespace", rf.Namespace).WithField("name", rf.Name).WithField("masterIP", masterIP).Debug("got master IP") } // No performed updates when nodes are syncing, still not connected, etc. - readyReplicas := int32(0) for _, rip := range redises { if rip != masterIP { ready, err := r.rfChecker.CheckRedisSlavesReady(rip, rf) @@ -39,43 +41,6 @@ func (r *RedisFailoverHandler) UpdateRedisesPods(rf *redisfailoverv1.RedisFailov if !ready { return nil } - readyReplicas++ - } - } - - // The loop above only sees pods that are already Running: GetRedisesIPs - // filters on Status.Phase, so a replica that is still terminating or being - // recreated is absent from the list entirely. "No replica reported unready" - // is therefore also true when there is no replica at all, and the gate - // passes on an empty set. The code then falls through to replacing the - // master while the replacement replica has not synced, leaving the failover - // with nothing to promote until the next reconcile elects one. - // - // In sentinel mode that is masked by the second gate further down - // (CheckSentinelSlavesNumberQuorumInMemory), which reads sentinel's own - // in-memory view of the replicas and does block. Operator-managed failover - // skips that gate, so it needs the count here instead. Scoped to - // OperatorManagedFailover so sentinel behaviour is unchanged. - // - // A quorum rather than the full expected count, mirroring the sentinel gate, - // so that one permanently unavailable replica (e.g. a PVC stuck in a dead - // zone) cannot block pod replacement forever while a safe failover is still - // available through the reachable majority. - if rf.OperatorManagedFailover() { - expectedReplicas := rf.Spec.Redis.Replicas - 1 - if rf.Bootstrapping() { - // Every redis pod replicates from the external bootstrap node, so - // none of them is the master and all count towards the quorum. - expectedReplicas = rf.Spec.Redis.Replicas - } - var quorum int32 - if expectedReplicas > 0 { - quorum = expectedReplicas/2 + 1 - } - if readyReplicas < quorum { - r.logger.WithField("namespace", rf.Namespace).WithField("name", rf.Name). - Infof("waiting for a quorum of ready replicas before replacing pods: have %d, need at least %d of %d expected", readyReplicas, quorum, expectedReplicas) - return nil } } @@ -119,6 +84,16 @@ func (r *RedisFailoverHandler) UpdateRedisesPods(rf *redisfailoverv1.RedisFailov return err } if masterRevision != ssUR { + // Upstream's settled check (Saremox/redis-operator): the readiness + // loop above only iterates pods that are already Running, so a + // replica that is terminating or being recreated is absent from it + // and "nothing reported unready" is trivially true. This counts the + // pods against the expected replica count instead, which is what + // actually keeps replica-first / master-last ordering intact. + if settled, err := r.redisPodsSettled(rf, ssUR); err != nil || !settled { + return err + } + // Deleting the master makes sentinel run a failover. Only do that once // every sentinel has a quorum (majority) of the freshly (re)started // slaves in memory - the redis-side readiness checked above is not @@ -804,3 +779,33 @@ func updateStatus(k8sservice k8s.Services, rf *redisfailoverv1.RedisFailover, ol } k8sservice.UpdateRedisFailoverStatus(context.Background(), rf.Namespace, rf, metav1.PatchOptions{}) } + +// redisPodsSettled reports whether the redis StatefulSet has finished the +// previous step of a rollout: every expected pod exists, none is terminating, +// and every pod already carrying the target revision is ready. +// +// Ported verbatim from upstream (Saremox/redis-operator) so this fork does not +// carry a second, divergent implementation of the same guard. +func (r *RedisFailoverHandler) redisPodsSettled(rf *redisfailoverv1.RedisFailover, updateRevision string) (bool, error) { + pods, err := r.k8sservice.GetStatefulSetPods(rf.Namespace, rfservice.GetRedisName(rf)) + if err != nil { + return false, err + } + wait := func(reason string) (bool, error) { + r.logger.WithField("namespace", rf.Namespace).WithField("name", rf.Name).Infof("redis rollout waits: %s", reason) + return false, nil + } + if len(pods.Items) < int(rf.Spec.Redis.Replicas) { + return wait(fmt.Sprintf("%d of %d pods exist", len(pods.Items), rf.Spec.Redis.Replicas)) + } + for i := range pods.Items { + pod := &pods.Items[i] + if pod.DeletionTimestamp != nil { + return wait("pod " + pod.Name + " is terminating") + } + if pod.Labels[appsv1.ControllerRevisionHashLabelKey] == updateRevision && !util.PodIsReady(pod) { + return wait("pod " + pod.Name + " is not ready") + } + } + return true, nil +} diff --git a/operator/redisfailover/checker_test.go b/operator/redisfailover/checker_test.go index 914658b7d..6df889e91 100644 --- a/operator/redisfailover/checker_test.go +++ b/operator/redisfailover/checker_test.go @@ -310,7 +310,7 @@ func TestCheckAndHeal(t *testing.T) { sentinel := "1.1.1.1" config := generateConfig() - mk := &mK8SService.Services{} + mk := settledK8sServices() // CheckAndHeal always defers updateStatus, on every return path. mk.On("UpdateRedisFailoverStatus", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return() mrfs := &mRFService.RedisFailoverClient{} @@ -767,7 +767,7 @@ func TestCheckAndHealOperatorManagedMode(t *testing.T) { rf := operatorManagedRF() config := generateConfig() - mk := &mK8SService.Services{} + mk := settledK8sServices() // CheckAndHeal always defers updateStatus, on every return path. mk.On("UpdateRedisFailoverStatus", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return() mrfs := &mRFService.RedisFailoverClient{} @@ -834,7 +834,7 @@ func TestUpdateStatusLastChanged(t *testing.T) { rf.Status.LastChanged = "2020-01-01T00:00:00Z" config := generateConfig() - mk := &mK8SService.Services{} + mk := settledK8sServices() mk.On("UpdateRedisFailoverStatus", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return() mrfs := &mRFService.RedisFailoverClient{} mrfc := &mRFService.RedisFailoverCheck{} @@ -862,7 +862,7 @@ func TestUpdateStatusLastChanged(t *testing.T) { rf.Status.LastChanged = previousLastChanged config := generateConfig() - mk := &mK8SService.Services{} + mk := settledK8sServices() mk.On("UpdateRedisFailoverStatus", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return() mrfs := &mRFService.RedisFailoverClient{} mrfc := &mRFService.RedisFailoverCheck{} @@ -1244,7 +1244,7 @@ func TestCheckAndHealPlainModeErrorBranches(t *testing.T) { } config := generateConfig() - mk := &mK8SService.Services{} + mk := settledK8sServices() // CheckAndHeal always defers updateStatus, on every return path. mk.On("UpdateRedisFailoverStatus", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return() mrfs := &mRFService.RedisFailoverClient{} @@ -1393,7 +1393,7 @@ func TestCheckAndHealBootstrapModeErrorBranches(t *testing.T) { rf.Spec.BootstrapNode.AllowSentinels = test.allowSentinels config := generateConfig() - mk := &mK8SService.Services{} + mk := settledK8sServices() // CheckAndHeal always defers updateStatus, on every return path. mk.On("UpdateRedisFailoverStatus", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return() mrfs := &mRFService.RedisFailoverClient{} @@ -1988,7 +1988,7 @@ func TestUpdate(t *testing.T) { } } - mk := &mK8SService.Services{} + mk := settledK8sServices() handler := rfOperator.NewRedisFailoverHandler(config, mrfs, mrfc, mrfh, mk, metrics.Dummy, log.Dummy) err := handler.UpdateRedisesPods(rf) @@ -2044,7 +2044,7 @@ func TestUpdateRedisesPodsOperatorManagedModeSkipsSentinelGate(t *testing.T) { mrfc.On("GetBestReplicaForPromotion", rf).Once().Return(&rfservice.ReplicaInfo{IP: "1.1.1.2"}, nil) mrfh.On("PromoteBestReplica", "1.1.1.2", rf).Once().Return(nil) - mk := &mK8SService.Services{} + mk := settledK8sServices() handler := rfOperator.NewRedisFailoverHandler(config, mrfs, mrfc, mrfh, mk, metrics.Dummy, log.Dummy) err := handler.UpdateRedisesPods(rf) @@ -2058,145 +2058,95 @@ func TestUpdateRedisesPodsOperatorManagedModeSkipsSentinelGate(t *testing.T) { mrfh.AssertExpectations(t) } -// TestUpdateRedisesPodsWaitsForReplicaQuorum covers #23: a rolling update must -// not replace the master while the replacement replica is still coming back. -// -// GetRedisesIPs only reports pods whose Status.Phase is Running, so a replica -// that is terminating or being recreated is missing from the list entirely. -// The readiness loop in UpdateRedisesPods then has nothing to iterate, and -// "no replica reported unready" is trivially true. Before the fix that empty -// set satisfied the gate and the master pod was deleted straight after, which -// left the RedisFailover with no master (and no synced replica to promote) -// until a later reconcile elected one - a short full outage on every routine -// pod-spec change. +// settledK8sServices returns a k8s service mock whose redis pods are all up, +// so the pod-replacement waits don't hold anything back. Ported from upstream. +func settledK8sServices() *mK8SService.Services { + mk := &mK8SService.Services{} + mk.On("GetStatefulSetPods", mock.Anything, mock.Anything).Maybe().Return(&corev1.PodList{Items: make([]corev1.Pod, 5)}, nil) + return mk +} + +func redisPod(revision string, ready, deleting bool) corev1.Pod { + pod := corev1.Pod{ObjectMeta: metav1.ObjectMeta{Labels: map[string]string{appsv1.ControllerRevisionHashLabelKey: revision}}} + if ready { + pod.Status.Conditions = []corev1.PodCondition{{Type: corev1.PodReady, Status: corev1.ConditionTrue}} + } + if deleting { + pod.DeletionTimestamp = &metav1.Time{Time: time.Now()} + } + return pod +} + +// TestUpdateRedisesPodsWaitsForTheLastReplacement is upstream's guard against +// replacing a pod before the previous replacement has settled, adopted here so +// this fork does not carry a second implementation of the same check. // -// Sentinel deployments never saw this: CheckSentinelSlavesNumberQuorumInMemory -// gates the master delete on sentinel's own view of the replicas, which lags -// pod creation and so does block. Operator-managed failover skips that gate, -// so the replica count has to be checked here instead. -func TestUpdateRedisesPodsWaitsForReplicaQuorum(t *testing.T) { - // rf asks for 3 replicas, so a healthy cluster is the master plus two - // replicas and the quorum required before replacing a pod is 2/2+1 = 2. +// The master rows differ from upstream: with the role handed over before the +// master pod is replaced, a stale master is promoted away rather than deleted. +func TestUpdateRedisesPodsWaitsForTheLastReplacement(t *testing.T) { tests := []struct { - name string - operatorMode bool - redisesIPs []string - readyReplicas map[string]bool - wantDelete bool // sentinel mode: the master pod is deleted so sentinel fails over - wantPromote bool // operator-managed: a replica is promoted, master pod left alone + name string + pods []corev1.Pod + podsErr error + wantAction bool }{ { - // The regression: both replicas absent from GetRedisesIPs because - // their pods are not Running yet. Nothing to iterate, so nothing - // reports unready - and the master must still be left alone. - name: "operator-managed: replicas missing entirely - master is not replaced", - operatorMode: true, - redisesIPs: []string{"10.0.0.1"}, - wantDelete: false, - }, - { - // One replica back but below quorum (1 of the 2 expected). - name: "operator-managed: replicas below quorum - master is not replaced", - operatorMode: true, - redisesIPs: []string{"10.0.0.1", "10.0.0.2"}, - readyReplicas: map[string]bool{"10.0.0.2": true}, - wantDelete: false, - }, - { - // With a quorum ready the rollout proceeds - but by handing the - // master role to a replica, not by deleting the live master. - name: "operator-managed: quorum of ready replicas - master is handed over, not deleted", - operatorMode: true, - redisesIPs: []string{"10.0.0.1", "10.0.0.2", "10.0.0.3"}, - readyReplicas: map[string]bool{"10.0.0.2": true, "10.0.0.3": true}, - wantPromote: true, - }, - { - // The no-op proof: identical fixture to the first case, but with - // Sentinel enabled. Sentinel mode keeps its existing behaviour and - // reaches its own sentinel-quorum gate, so the new count must not - // change anything for it. - name: "sentinel: replicas missing entirely - behaviour unchanged", - operatorMode: false, - redisesIPs: []string{"10.0.0.1"}, - wantDelete: true, + name: "all pods up", + pods: []corev1.Pod{redisPod("old", true, false), redisPod("old", true, false), redisPod("new", true, false)}, + wantAction: true, + }, + { + name: "a pod is being deleted", + pods: []corev1.Pod{redisPod("old", true, false), redisPod("old", true, false), redisPod("old", true, true)}, + }, + { + name: "the deleted pod is not recreated yet", + pods: []corev1.Pod{redisPod("old", true, false), redisPod("old", true, false)}, + }, + { + name: "the recreated pod is not ready yet", + pods: []corev1.Pod{redisPod("old", true, false), redisPod("old", true, false), redisPod("new", false, false)}, + }, + { + name: "an old pod that is not ready doesn't block", + pods: []corev1.Pod{redisPod("old", true, false), redisPod("old", false, false), redisPod("new", true, false)}, + wantAction: true, + }, + { + name: "listing pods fails", + podsErr: errors.New("list err"), }, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { - assertTest := assert.New(t) - - rf := generateRF(false, false) - if test.operatorMode { - rf.Spec.Sentinel.Enabled = ptr.To(false) - } - - config := generateConfig() - mrfs := &mRFService.RedisFailoverClient{} + rf := operatorManagedRF() + rf.Spec.Redis.Replicas = 3 mrfc := &mRFService.RedisFailoverCheck{} mrfh := &mRFService.RedisFailoverHeal{} - - mrfc.On("GetRedisesIPs", rf).Once().Return(test.redisesIPs, nil) + mk := &mK8SService.Services{} + mrfc.On("GetRedisesIPs", rf).Once().Return([]string{"10.0.0.1"}, nil) mrfc.On("GetMasterIP", rf).Once().Return("10.0.0.1", nil) - for ip, ready := range test.readyReplicas { - mrfc.On("CheckRedisSlavesReady", ip, rf).Once().Return(ready, nil) - } - - if test.wantDelete || test.wantPromote { - mrfc.On("GetStatefulSetUpdateRevision", rf).Once().Return("10", nil) - mrfc.On("GetRedisesSlavesPods", rf).Once().Return([]string{}, nil) - mrfc.On("GetRedisesMasterPod", rf).Once().Return("master", nil) - mrfc.On("GetRedisRevisionHash", "master", rf).Once().Return("9", nil) // stale - } - if test.wantDelete { - mrfh.On("DeletePod", "master", rf).Once().Return(nil) - if !test.operatorMode { - // Sentinel mode still consults its own quorum gate, then - // deletes the master so sentinel runs the failover. - mrfc.On("GetSentinelsIPs", rf).Once().Return([]string{"11.0.0.1"}, nil) - mrfc.On("CheckSentinelSlavesNumberQuorumInMemory", "11.0.0.1", rf).Once().Return(nil) - } - } - if test.wantPromote { + mrfc.On("GetStatefulSetUpdateRevision", rf).Once().Return("new", nil) + mrfc.On("GetRedisesSlavesPods", rf).Once().Return([]string{}, nil) + mrfc.On("GetRedisesMasterPod", rf).Once().Return("master", nil) + mrfc.On("GetRedisRevisionHash", "master", rf).Once().Return("old", nil) + mk.On("GetStatefulSetPods", rf.Namespace, rfservice.GetRedisName(rf)).Once().Return(&corev1.PodList{Items: test.pods}, test.podsErr) + if test.wantAction { mrfc.On("GetBestReplicaForPromotion", rf).Once().Return(&rfservice.ReplicaInfo{IP: "10.0.0.2"}, nil) mrfh.On("PromoteBestReplica", "10.0.0.2", rf).Once().Return(nil) } - if !test.wantDelete && !test.wantPromote { - // Permit, but do not require, everything past the gate. Without - // the fix the code runs straight on and deletes the master, and - // a permitted-but-unexpected call gives a readable assertion - // failure below instead of a mock panic that would abort the - // sibling subtests along with this one. - mrfc.On("GetStatefulSetUpdateRevision", rf).Maybe().Return("10", nil) - mrfc.On("GetRedisesSlavesPods", rf).Maybe().Return([]string{}, nil) - mrfc.On("GetRedisesMasterPod", rf).Maybe().Return("master", nil) - mrfc.On("GetRedisRevisionHash", "master", rf).Maybe().Return("9", nil) // stale - mrfh.On("DeletePod", "master", rf).Maybe().Return(nil) - } + // Permitted but never wanted: the master must not be deleted while serving. + mrfh.On("DeletePod", mock.Anything, mock.Anything).Maybe().Return(nil) - mk := &mK8SService.Services{} - handler := rfOperator.NewRedisFailoverHandler(config, mrfs, mrfc, mrfh, mk, metrics.Dummy, log.Dummy) + handler := rfOperator.NewRedisFailoverHandler(generateConfig(), &mRFService.RedisFailoverClient{}, mrfc, mrfh, mk, metrics.Dummy, log.Dummy) err := handler.UpdateRedisesPods(rf) - assertTest.NoError(err) - switch { - case test.wantDelete: - mrfh.AssertCalled(t, "DeletePod", "master", rf) - case test.wantPromote: - // The role moves first; the old master pod is left running and - // is replaced later as an ordinary stale replica. - mrfh.AssertCalled(t, "PromoteBestReplica", "10.0.0.2", rf) - mrfh.AssertNotCalled(t, "DeletePod", mock.Anything, mock.Anything) - default: - // The master survives, and the operator never even looks at the - // StatefulSet revision: it returned before getting that far. - mrfh.AssertNotCalled(t, "DeletePod", mock.Anything, mock.Anything) - mrfh.AssertNotCalled(t, "PromoteBestReplica", mock.Anything, mock.Anything) - mrfc.AssertNotCalled(t, "GetStatefulSetUpdateRevision", mock.Anything) - } + assert.Equal(t, test.podsErr, err) + mrfh.AssertNotCalled(t, "DeletePod", "master", mock.Anything) mrfc.AssertExpectations(t) mrfh.AssertExpectations(t) + mk.AssertExpectations(t) }) } } @@ -2305,7 +2255,7 @@ func TestUpdateRedisesPodsErrorBranches(t *testing.T) { rf := generateRF(false, false) config := generateConfig() - mk := &mK8SService.Services{} + mk := settledK8sServices() mrfs := &mRFService.RedisFailoverClient{} mrfc := &mRFService.RedisFailoverCheck{} mrfh := &mRFService.RedisFailoverHeal{} diff --git a/operator/redisfailover/util/pod.go b/operator/redisfailover/util/pod.go index d0e4b0137..d11fb1e13 100644 --- a/operator/redisfailover/util/pod.go +++ b/operator/redisfailover/util/pod.go @@ -9,3 +9,12 @@ func PodIsTerminal(pod *v1.Pod) bool { func PodIsScheduling(pod *v1.Pod) bool { return pod.DeletionTimestamp != nil || pod.Status.Phase == v1.PodPending } + +func PodIsReady(pod *v1.Pod) bool { + for _, c := range pod.Status.Conditions { + if c.Type == v1.PodReady { + return c.Status == v1.ConditionTrue + } + } + return false +}