From dffee92126a79ab4110a3ecbf0788a93e59318ab Mon Sep 17 00:00:00 2001 From: Allen Ray Date: Thu, 9 Jul 2026 09:45:18 -0400 Subject: [PATCH 1/7] Defrag one member per sync cycle and move leadership before defrag Changes the defrag controller to defrag a single member per sync cycle instead of all members in one pass. This reduces the blast radius of defrag operations and avoids write blocks by preemptively transferring leadership away from the defrag target. After a successful leader transfer, the controller requeues with a short delay to allow etcd to settle before proceeding with defragmentation. Defrag failures no longer cause immediate requeue; instead the controller waits for the next periodic sync. The operator client informer is registered as a bare informer since the controller only uses it for cache warmup, not watch-triggered syncs. Co-Authored-By: Claude Opus 4.6 (1M context) --- pkg/etcdcli/etcdcli.go | 14 ++ pkg/etcdcli/helpers.go | 13 + pkg/etcdcli/interfaces.go | 5 + .../defragcontroller/defragcontroller.go | 224 +++++++++--------- .../defragcontroller/defragcontroller_test.go | 133 +++++++++-- pkg/operator/starter.go | 3 +- 6 files changed, 252 insertions(+), 140 deletions(-) diff --git a/pkg/etcdcli/etcdcli.go b/pkg/etcdcli/etcdcli.go index f067c02aae..3a9ffc62fb 100644 --- a/pkg/etcdcli/etcdcli.go +++ b/pkg/etcdcli/etcdcli.go @@ -246,6 +246,20 @@ func (g *etcdClientGetter) MemberUpdatePeerURL(ctx context.Context, id uint64, p return err } +func (g *etcdClientGetter) MoveLeader(ctx context.Context, toMember uint64) error { + cli, err := g.clientPool.Get() + if err != nil { + return err + } + + defer g.clientPool.Return(cli) + + ctx, cancel := context.WithTimeout(ctx, DefaultClientTimeout) + defer cancel() + _, err = cli.MoveLeader(ctx, toMember) + return err +} + func (g *etcdClientGetter) MemberRemove(ctx context.Context, memberID uint64) error { cli, err := g.clientPool.Get() if err != nil { diff --git a/pkg/etcdcli/helpers.go b/pkg/etcdcli/helpers.go index e16f068779..968a5b6b2f 100644 --- a/pkg/etcdcli/helpers.go +++ b/pkg/etcdcli/helpers.go @@ -22,6 +22,12 @@ func (f *fakeEtcdClient) Defragment(ctx context.Context, member *etcdserverpb.Me f.opts.defragErrors = f.opts.defragErrors[1:] return nil, err } + for _, status := range f.opts.status { + if status.Header.MemberId == member.ID { + status.DbSize = status.DbSizeInUse + break + } + } // dramatic simplification f.opts.dbSize = f.opts.dbSizeInUse return nil, nil @@ -82,6 +88,13 @@ func (f *fakeEtcdClient) VotingMemberList(ctx context.Context) ([]*etcdserverpb. return filterVotingMembers(members), nil } +func (f *fakeEtcdClient) MoveLeader(ctx context.Context, toMember uint64) error { + for _, status := range f.opts.status { + status.Leader = toMember + } + return nil +} + func (f *fakeEtcdClient) MemberRemove(ctx context.Context, memberID uint64) error { var memberExists bool for _, m := range f.members { diff --git a/pkg/etcdcli/interfaces.go b/pkg/etcdcli/interfaces.go index 6a31a9c620..ae6cf3e072 100644 --- a/pkg/etcdcli/interfaces.go +++ b/pkg/etcdcli/interfaces.go @@ -25,6 +25,7 @@ type EtcdClient interface { HealthyMemberLister UnhealthyMemberLister MemberStatusChecker + LeaderMover Status GetMember(ctx context.Context, name string) (*etcdserverpb.Member, error) @@ -64,6 +65,10 @@ type MemberRemover interface { MemberRemove(ctx context.Context, memberID uint64) error } +type LeaderMover interface { + MoveLeader(ctx context.Context, toMember uint64) error +} + type MemberLister interface { // MemberList lists all members in a cluster MemberList(ctx context.Context) ([]*etcdserverpb.Member, error) diff --git a/pkg/operator/defragcontroller/defragcontroller.go b/pkg/operator/defragcontroller/defragcontroller.go index 4438dd6fbc..17e6d03872 100644 --- a/pkg/operator/defragcontroller/defragcontroller.go +++ b/pkg/operator/defragcontroller/defragcontroller.go @@ -1,9 +1,11 @@ package defragcontroller import ( + "cmp" "context" "fmt" "math" + "slices" "time" configv1 "github.com/openshift/api/config/v1" @@ -16,7 +18,6 @@ import ( "go.etcd.io/etcd/api/v3/etcdserverpb" clientv3 "go.etcd.io/etcd/client/v3" k8serror "k8s.io/apimachinery/pkg/api/errors" - "k8s.io/apimachinery/pkg/util/wait" corev1listers "k8s.io/client-go/listers/core/v1" "k8s.io/klog/v2" @@ -27,12 +28,10 @@ import ( const ( minDefragBytes int64 = 100 * 1024 * 1024 // 100MB - minDefragWaitDuration = 36 * time.Second maxFragmentedPercentage float64 = 45 - pollWaitDuration = 2 * time.Second - pollTimeoutDuration = 60 * time.Second compactionInterval = 10 * time.Minute maxDefragFailuresBeforeDegrade = 3 + leaderTransferSettleTime = 5 * time.Second defragDisabledCondition = "DefragControllerDisabled" defragDisableConfigmapName = "etcd-disable-defrag" @@ -47,11 +46,11 @@ type DefragController struct { memberLister etcdcli.AllMemberLister defragClient etcdcli.Defragment statusClient etcdcli.Status + leaderMover etcdcli.LeaderMover infrastructureLister configv1listers.InfrastructureLister configmapLister corev1listers.ConfigMapLister - numDefragFailures int - defragWaitDuration time.Duration + numDefragFailures int } func NewDefragController( @@ -60,6 +59,7 @@ func NewDefragController( memberLister etcdcli.AllMemberLister, defragClient etcdcli.Defragment, statusClient etcdcli.Status, + leaderMover etcdcli.LeaderMover, infrastructureLister configv1listers.InfrastructureLister, eventRecorder events.Recorder, kubeInformers v1helpers.KubeInformersForNamespaces) factory.Controller { @@ -68,14 +68,14 @@ func NewDefragController( memberLister: memberLister, defragClient: defragClient, statusClient: statusClient, + leaderMover: leaderMover, infrastructureLister: infrastructureLister, configmapLister: kubeInformers.ConfigMapLister(), - defragWaitDuration: minDefragWaitDuration, } syncer := health.NewCheckingSyncWrapper(c.sync, 3*compactionInterval+1*time.Minute) livenessChecker.Add("DefragController", syncer) - return factory.New().ResyncEvery(compactionInterval+1*time.Minute).WithInformers( // attempt to sync outside of etcd compaction interval to ensure maximum gain by defragmentation. + return factory.New().ResyncEvery(compactionInterval+1*time.Minute).WithBareInformers( // attempt to sync outside of etcd compaction interval to ensure maximum gain by defragmentation. operatorClient.Informer(), ).WithSync(syncer.Sync).ToController("DefragController", eventRecorder.WithComponentSuffix("defrag-controller")) } @@ -90,7 +90,7 @@ func (c *DefragController) sync(ctx context.Context, syncCtx factory.SyncContext return nil } - return c.runDefrag(ctx, syncCtx.Recorder()) + return c.runDefrag(ctx, syncCtx) } func (c *DefragController) checkDefragEnabled(ctx context.Context, recorder events.Recorder) (bool, error) { @@ -124,7 +124,13 @@ func (c *DefragController) checkDefragEnabled(ctx context.Context, recorder even return true, nil } -func (c *DefragController) runDefrag(ctx context.Context, recorder events.Recorder) error { +type StatusMember struct { + Status *clientv3.StatusResponse + Member *etcdserverpb.Member +} + +func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncContext) error { + recorder := syncCtx.Recorder() // Do not defrag if any of the cluster members are unhealthy. memberHealth, err := c.memberLister.MemberHealth(ctx) if err != nil { @@ -139,120 +145,114 @@ func (c *DefragController) runDefrag(ctx context.Context, recorder events.Record return err } - // filter out learner members since they don't support the defragment API call - var etcdMembers []*etcdserverpb.Member - for _, m := range members { - if !m.IsLearner { - etcdMembers = append(etcdMembers, m) - } - } - - var endpointStatus []*clientv3.StatusResponse - var leader *clientv3.StatusResponse - for _, member := range etcdMembers { - if len(member.ClientURLs) == 0 { - // skip unstarted member + var ( + statusMembers = make(map[uint64]StatusMember) + defragTargets = make([]StatusMember, 0, len(members)) + ) + for _, member := range members { + // filter out learner members since they don't support the defragment API call + // and filter out unstarted members + if member.IsLearner || len(member.ClientURLs) == 0 { continue } + status, err := c.statusClient.Status(ctx, member.ClientURLs[0]) if err != nil { return err + } else if status == nil { + return fmt.Errorf("endpoint status returned nil for member %q (%s)", member.Name, member.ClientURLs[0]) } - if leader == nil && status.Leader == member.ID { - leader = status - continue + + sm := StatusMember{ + Status: status, + Member: member, + } + + statusMembers[member.ID] = sm + + if isEndpointBackendFragmented(member, status) { + defragTargets = append(defragTargets, sm) } - endpointStatus = append(endpointStatus, status) } - // Leader last if possible. - if leader != nil { - klog.V(4).Infof("Appending leader last, ID: %x", leader.Header.MemberId) - endpointStatus = append(endpointStatus, leader) + if len(defragTargets) == 0 { + recorder.Eventf("DefragControllerDefragmentSkipped", "No etcd members meet the conditions for defragmentation") + return nil } - successfulDefrags := 0 - var errors []error - for _, status := range endpointStatus { - member, err := getMemberFromStatus(etcdMembers, status) - if err != nil { - errors = append(errors, err) - continue - } + // Sort fragmented members so we defragment the most fragmented member first + slices.SortFunc(defragTargets, sortByMostFragmented) - // Check each member's status which includes the db size on disk "DbSize" and the db size in use "DbSizeInUse" - // compare the % difference and if that difference is over the max diff threshold and also above the minimum - // db size we defrag the members state file. In the case where this command only partially completed controller - // can clean that up on the next sync. Having the db sizes slightly different is not a problem in itself. - if isEndpointBackendFragmented(member, status) { - recorder.Eventf("DefragControllerDefragmentAttempt", "Attempting defrag on member: %s, memberID: %x, dbSize: %d, dbInUse: %d, leader ID: %d", member.Name, member.ID, status.DbSize, status.DbSizeInUse, status.Leader) - if _, err := c.defragClient.Defragment(ctx, member); err != nil { - // Defrag can timeout if defragmentation takes longer than etcdcli.DefragDialTimeout. - errMsg := fmt.Sprintf("failed defrag on member: %s, memberID: %x: %v", member.Name, member.ID, err) - recorder.Eventf("DefragControllerDefragmentFailed", errMsg) - errors = append(errors, fmt.Errorf("%s", errMsg)) + // If the leader is in the first slot, and we have other members to defrag, defrag the others first, + // so we limit the number of times we move the leader during defrag. + defragTargetStatus := defragTargets[0].Status + defragTargetMember := defragTargets[0].Member + if len(defragTargets) > 1 && defragTargetStatus.Leader == defragTargetMember.ID { + defragTargetStatus = defragTargets[1].Status + defragTargetMember = defragTargets[1].Member + } + + // Preemptively attempt to move the leadership away from the current defrag target to another valid follower. + // We try this to avoid the write block that occurs while the leader is being defragmented. + // We record any error that occurs while attempting this, but we do not halt defrag if the move fails; + // we just accept the write block. + if defragTargetStatus.Leader == defragTargetMember.ID && len(statusMembers) > 1 { + followers := make([]StatusMember, 0, len(statusMembers)) + for id, member := range statusMembers { + if defragTargetMember.ID == id { continue } + followers = append(followers, member) + } + + if len(followers) > 0 { + slices.SortFunc(followers, sortByLeastFragmented) - recorder.Eventf("DefragControllerDefragmentSuccess", "etcd member has been defragmented: %s, memberID: %d", member.Name, member.ID) - successfulDefrags++ - - // Give cluster time to recover before we move to the next member. - if err := wait.Poll( - pollWaitDuration, - pollTimeoutDuration, - func() (bool, error) { - // Ensure defragmentation attempts have clear observable signal. - klog.V(4).Infof("Sleeping to allow cluster to recover before defrag next member: %v", c.defragWaitDuration) - time.Sleep(c.defragWaitDuration) - - memberHealth, err := c.memberLister.MemberHealth(ctx) - if err != nil { - klog.Warningf("failed checking member health: %v", err) - return false, nil - } - if !etcdcli.IsClusterHealthy(memberHealth) { - klog.Warningf("cluster is unhealthy: %s", memberHealth.Status()) - return false, nil - } - return true, nil - }); err != nil { - errors = append(errors, fmt.Errorf("timeout waiting for cluster to stabilize after defrag: %w", err)) + newLeader := followers[0] + + err := c.leaderMover.MoveLeader(ctx, newLeader.Member.ID) + if err != nil { + recorder.Eventf("DefragControllerLeaderTransferFailed", "Failed to move leader away from member %s to member %s before defrag: %v", defragTargetMember.Name, newLeader.Member.Name, err) + } else { + recorder.Eventf("DefragControllerLeaderTransferred", "Moved leadership away from member %s (memberID: %x) to member %s (memberID: %x) before defrag, requeueing to allow etcd to settle", defragTargetMember.Name, defragTargetMember.ID, newLeader.Member.Name, newLeader.Member.ID) + syncCtx.Queue().AddAfter(syncCtx.QueueKey(), leaderTransferSettleTime) + return nil } - } else { - // no fragmentation needed is also a success - successfulDefrags++ } } - if successfulDefrags != len(endpointStatus) { + recorder.Eventf("DefragControllerDefragmentAttempt", "Attempting defrag on member: %s, memberID: %x, dbSize: %d, dbInUse: %d, leader ID: %d", defragTargetMember.Name, defragTargetMember.ID, defragTargetStatus.DbSize, defragTargetStatus.DbSizeInUse, defragTargetStatus.Leader) + if _, err := c.defragClient.Defragment(ctx, defragTargetMember); err != nil { + // Defrag can timeout if defragmentation takes longer than etcdcli.DefragDialTimeout. + errMsg := fmt.Sprintf("failed defrag on member: %s, memberID: %x: %v", defragTargetMember.Name, defragTargetMember.ID, err) + recorder.Eventf("DefragControllerDefragmentFailed", errMsg) + klog.Errorf("%s", errMsg) c.numDefragFailures++ - recorder.Eventf("DefragControllerDefragmentPartialFailure", - "only %d/%d members were successfully defragmented, %d tries left before controller degrades", - successfulDefrags, len(endpointStatus), maxDefragFailuresBeforeDegrade-c.numDefragFailures) - if c.numDefragFailures >= maxDefragFailuresBeforeDegrade { - _, _, updateErr := v1helpers.UpdateStatus(ctx, c.operatorClient, v1helpers.UpdateConditionFn(operatorv1.OperatorCondition{ - Type: defragDegradedCondition, - Status: operatorv1.ConditionTrue, - Reason: "Error", - Message: fmt.Sprintf("degraded after %d attempts at defragmenting all etcd members", c.numDefragFailures), - })) - if updateErr != nil { - recorder.Warning("DefragControllerUpdatingStatus", updateErr.Error()) - } + c.setDegraded(ctx, recorder) } - - // return all errors here for the sync loop to retry immediately - return v1helpers.NewMultiLineAggregate(errors) + return nil } - if len(errors) > 0 { - klog.Warningf("found errors even though all members have been successfully defragmented: %s", - v1helpers.NewMultiLineAggregate(errors).Error()) + recorder.Eventf("DefragControllerDefragmentSuccess", "etcd member has been defragmented: %s, memberID: %d", defragTargetMember.Name, defragTargetMember.ID) + c.numDefragFailures = 0 + c.clearDegraded(ctx, recorder) + return nil +} + +func (c *DefragController) setDegraded(ctx context.Context, recorder events.Recorder) { + _, _, updateErr := v1helpers.UpdateStatus(ctx, c.operatorClient, v1helpers.UpdateConditionFn(operatorv1.OperatorCondition{ + Type: defragDegradedCondition, + Status: operatorv1.ConditionTrue, + Reason: "Error", + Message: fmt.Sprintf("degraded after %d attempts at defragmenting etcd members", c.numDefragFailures), + })) + if updateErr != nil { + recorder.Warning("DefragControllerUpdatingStatus", updateErr.Error()) } +} - c.numDefragFailures = 0 +func (c *DefragController) clearDegraded(ctx context.Context, recorder events.Recorder) { _, _, updateErr := v1helpers.UpdateStatus(ctx, c.operatorClient, v1helpers.UpdateConditionFn(operatorv1.OperatorCondition{ Type: defragDegradedCondition, @@ -262,8 +262,6 @@ func (c *DefragController) runDefrag(ctx context.Context, recorder events.Record if updateErr != nil { recorder.Warning("DefragControllerUpdatingStatus", updateErr.Error()) } - - return updateErr } func (c *DefragController) ensureControllerDisabledCondition(ctx context.Context, desiredStatus operatorv1.ConditionStatus, recorder events.Recorder) error { @@ -293,10 +291,6 @@ func (c *DefragController) ensureControllerDisabledCondition(ctx context.Context // This can happen if the operator starts defrag of the cluster but then loses leader status and is rescheduled before // the operator can defrag all members. func isEndpointBackendFragmented(member *etcdserverpb.Member, endpointStatus *clientv3.StatusResponse) bool { - if endpointStatus == nil { - klog.Errorf("endpoint status validation failed: %v", endpointStatus) - return false - } fragmentedPercentage := checkFragmentationPercentage(endpointStatus.DbSize, endpointStatus.DbSizeInUse) if fragmentedPercentage > 0.00 { klog.Infof("etcd member %q backend store fragmented: %.2f %%, dbSize: %d", member.Name, fragmentedPercentage, endpointStatus.DbSize) @@ -310,14 +304,16 @@ func checkFragmentationPercentage(ondisk, inuse int64) float64 { return math.Round(fragmentedPercentage*100) / 100 } -func getMemberFromStatus(members []*etcdserverpb.Member, endpointStatus *clientv3.StatusResponse) (*etcdserverpb.Member, error) { - if endpointStatus == nil { - return nil, fmt.Errorf("endpoint status validation failed: %v", endpointStatus) - } - for _, member := range members { - if member.ID == endpointStatus.Header.MemberId { - return member, nil - } - } - return nil, fmt.Errorf("no member found in MemberList matching ID: %v", endpointStatus.Header.MemberId) +func sortByLeastFragmented(a, b StatusMember) int { + return cmp.Compare( + checkFragmentationPercentage(a.Status.DbSize, a.Status.DbSizeInUse), + checkFragmentationPercentage(b.Status.DbSize, b.Status.DbSizeInUse), + ) +} + +func sortByMostFragmented(a, b StatusMember) int { + return cmp.Compare( + checkFragmentationPercentage(b.Status.DbSize, b.Status.DbSizeInUse), + checkFragmentationPercentage(a.Status.DbSize, a.Status.DbSizeInUse), + ) } diff --git a/pkg/operator/defragcontroller/defragcontroller_test.go b/pkg/operator/defragcontroller/defragcontroller_test.go index c4b00314d1..be82946e68 100644 --- a/pkg/operator/defragcontroller/defragcontroller_test.go +++ b/pkg/operator/defragcontroller/defragcontroller_test.go @@ -4,7 +4,6 @@ import ( "context" "errors" "fmt" - "regexp" "strings" "testing" "time" @@ -69,6 +68,7 @@ func TestNewDefragController(t *testing.T) { staticPodStatus *operatorv1.StaticPodOperatorStatus objects []runtime.Object clusterSize int + syncLoops int memberHealth *etcdcli.FakeMemberHealth dbInUse int64 dbSize int64 @@ -84,6 +84,7 @@ func TestNewDefragController(t *testing.T) { dbInUse: minDefragBytes / 2, // 500MB defragSuccessEvents: 3, clusterSize: 3, + syncLoops: 4, // 2 non-leader defrags + 1 leader transfer + 1 former-leader defrag memberHealth: &etcdcli.FakeMemberHealth{Healthy: 3}, objects: []runtime.Object{ u.FakeInfrastructureTopology(configv1.HighlyAvailableTopologyMode), @@ -97,6 +98,7 @@ func TestNewDefragController(t *testing.T) { dbInUse: minDefragBytes / 2, // 500MB defragSuccessEvents: 2, clusterSize: 2, + syncLoops: 3, // 1 non-leader defrag + 1 leader transfer + 1 former-leader defrag memberHealth: &etcdcli.FakeMemberHealth{Healthy: 2}, objects: []runtime.Object{ u.FakeInfrastructureTopology(configv1.DualReplicaTopologyMode), @@ -110,6 +112,7 @@ func TestNewDefragController(t *testing.T) { dbInUse: minDefragBytes / 2, // 500MB defragSuccessEvents: 3, clusterSize: 3, + syncLoops: 4, // 2 non-leader defrags + 1 leader transfer + 1 former-leader defrag memberHealth: &etcdcli.FakeMemberHealth{Healthy: 3}, objects: []runtime.Object{ u.FakeInfrastructureTopology(configv1.HighlyAvailableArbiterMode), @@ -222,38 +225,35 @@ func TestNewDefragController(t *testing.T) { memberLister: fakeEtcdClient, defragClient: fakeEtcdClient, statusClient: fakeEtcdClient, + leaderMover: fakeEtcdClient, infrastructureLister: configv1listers.NewInfrastructureLister(indexer), configmapLister: corev1listers.NewConfigMapLister(indexer), - // to speed the tests up, in real life we use minDefragWaitDuration - defragWaitDuration: 1 * time.Second, } - err := controller.sync(context.TODO(), factory.NewSyncContext("defrag-controller", eventRecorder)) - if err != nil && !scenario.wantErr { - t.Fatalf("unexepected error %v", err) + syncLoops := scenario.syncLoops + if syncLoops == 0 { + syncLoops = 1 + } + var err error + for i := 0; i < syncLoops; i++ { + err = controller.sync(context.TODO(), factory.NewSyncContext("defrag-controller", eventRecorder)) + if err != nil && !scenario.wantErr { + t.Fatalf("unexpected error on sync %d: %v", i, err) + } } if err == nil && scenario.wantErr { t.Fatal("expected error got nil") } if err != nil && scenario.wantErr { if !strings.HasPrefix(err.Error(), scenario.wantErrMsg) { - t.Fatalf("unexepected error prefix want: %q got: %q", scenario.wantErrMsg, err.Error()) + t.Fatalf("unexpected error prefix want: %q got: %q", scenario.wantErrMsg, err.Error()) } } var defragSuccessEvents int - lastEvent := len(eventRecorder.Events()) - 1 - for i, event := range eventRecorder.Events() { + for _, event := range eventRecorder.Events() { if strings.HasPrefix(event.Message, "etcd member has been defragmented") { defragSuccessEvents++ } - // ensure the leader was defragged last - if i == lastEvent { - // last event will print leader ID - regex := regexp.MustCompile(fmt.Sprint(status[0].Leader)) - if len(regex.FindAll([]byte(event.Message), -1)) != 1 { - t.Fatalf("expected leader defrag event to be last got %q", event.Message) - } - } } if defragSuccessEvents != scenario.defragSuccessEvents { t.Fatalf("defragSuccessEvents invalid want %d got %d", scenario.defragSuccessEvents, defragSuccessEvents) @@ -304,13 +304,13 @@ func TestNewDefragControllerMultiSyncs(t *testing.T) { defragSuccessEvents: 0, clusterSize: 3, syncLoops: maxDefragFailuresBeforeDegrade, - errSyncLoops: maxDefragFailuresBeforeDegrade, + errSyncLoops: 0, memberHealth: &etcdcli.FakeMemberHealth{Healthy: 3}, objects: []runtime.Object{ u.FakeInfrastructureTopology(configv1.HighlyAvailableTopologyMode), }, fakeClientOpts: []etcdcli.FakeClientOption{ - etcdcli.WithFakeDefragErrors(generateErrors(maxDefragFailuresBeforeDegrade * 3)), + etcdcli.WithFakeDefragErrors(generateErrors(maxDefragFailuresBeforeDegrade)), }, wantDisabledCondition: operatorv1.ConditionFalse, wantDegradedCondition: operatorv1.ConditionTrue, @@ -322,16 +322,15 @@ func TestNewDefragControllerMultiSyncs(t *testing.T) { dbInUse: minDefragBytes / 2, defragSuccessEvents: 3, clusterSize: 3, - syncLoops: maxDefragFailuresBeforeDegrade + 1, - errSyncLoops: maxDefragFailuresBeforeDegrade, + syncLoops: maxDefragFailuresBeforeDegrade + 4, // +1 for leader transfer requeue + errSyncLoops: 0, memberHealth: &etcdcli.FakeMemberHealth{Healthy: 3}, objects: []runtime.Object{ u.FakeInfrastructureTopology(configv1.HighlyAvailableTopologyMode), }, fakeClientOpts: []etcdcli.FakeClientOption{ - etcdcli.WithFakeDefragErrors(generateErrors(maxDefragFailuresBeforeDegrade * 3)), + etcdcli.WithFakeDefragErrors(generateErrors(maxDefragFailuresBeforeDegrade)), }, - // ignoring errors here since the first maxDefragFailuresBeforeDegrade invocations will return an error, the ones after won't wantDisabledCondition: operatorv1.ConditionFalse, wantDegradedCondition: operatorv1.ConditionFalse, }, @@ -378,10 +377,9 @@ func TestNewDefragControllerMultiSyncs(t *testing.T) { memberLister: fakeEtcdClient, statusClient: fakeEtcdClient, defragClient: fakeEtcdClient, + leaderMover: fakeEtcdClient, infrastructureLister: configv1listers.NewInfrastructureLister(indexer), configmapLister: corev1listers.NewConfigMapLister(indexer), - // to speed the tests up, in real life we use minDefragWaitDuration - defragWaitDuration: 1 * time.Second, } numSyncErr := 0 @@ -413,6 +411,91 @@ func TestNewDefragControllerMultiSyncs(t *testing.T) { } } +func TestDefragMovesLeadershipBeforeDefrag(t *testing.T) { + fakeOperatorClient := v1helpers.NewFakeStaticPodOperatorClient( + &operatorv1.StaticPodOperatorSpec{ + OperatorSpec: operatorv1.OperatorSpec{ + ManagementState: operatorv1.Managed, + }, + }, + u.StaticPodOperatorStatus(), + nil, + nil, + ) + + integration.BeforeTestExternal(t) + testServer := integration.NewCluster(t, &integration.ClusterConfig{Size: 3}) + defer testServer.Terminate(t) + + etcdMembers := waitForMembersWithClientURLs(t, testServer) + + var status []*clientv3.StatusResponse + var leaderID uint64 + for _, member := range testServer.Members { + statusResp, err := testServer.Client(0).Status(context.TODO(), member.GRPCURL) + require.NoError(t, err) + if leaderID == 0 { + leaderID = statusResp.Leader + } + // Only the leader is fragmented, so it's the sole defrag target + // and must trigger a leader transfer. + statusResp.DbSize = minDefragBytes / 2 + statusResp.DbSizeInUse = minDefragBytes / 2 + for _, m := range etcdMembers { + if m.ID == leaderID && statusResp.Header.MemberId == leaderID { + statusResp.DbSize = minDefragBytes + statusResp.DbSizeInUse = minDefragBytes / 2 + } + } + status = append(status, statusResp) + } + + fakeEtcdClient, _ := etcdcli.NewFakeEtcdClient( + etcdMembers, + etcdcli.WithFakeClusterHealth(&etcdcli.FakeMemberHealth{Healthy: 3}), + etcdcli.WithFakeStatus(status), + ) + eventRecorder := events.NewInMemoryRecorder(t.Name(), clock.RealClock{}) + indexer := cache.NewIndexer(cache.MetaNamespaceKeyFunc, cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc}) + require.NoError(t, indexer.Add(u.FakeInfrastructureTopology(configv1.HighlyAvailableTopologyMode))) + + controller := &DefragController{ + operatorClient: fakeOperatorClient, + memberLister: fakeEtcdClient, + defragClient: fakeEtcdClient, + statusClient: fakeEtcdClient, + leaderMover: fakeEtcdClient, + infrastructureLister: configv1listers.NewInfrastructureLister(indexer), + configmapLister: corev1listers.NewConfigMapLister(indexer), + } + + syncCtx := factory.NewSyncContext("defrag-controller", eventRecorder) + + // First sync: leader is the most fragmented, so leadership is transferred and sync requeues. + err := controller.sync(context.TODO(), syncCtx) + require.NoError(t, err) + + var leaderTransferEvents int + for _, event := range eventRecorder.Events() { + if strings.HasPrefix(event.Message, "Moved leadership away from member") { + leaderTransferEvents++ + } + } + assert.Equal(t, 1, leaderTransferEvents, "expected exactly one leader transfer event") + + // Second sync: the former leader is now a follower and gets defragged. + err = controller.sync(context.TODO(), syncCtx) + require.NoError(t, err) + + var defragSuccessEvents int + for _, event := range eventRecorder.Events() { + if strings.HasPrefix(event.Message, "etcd member has been defragmented") { + defragSuccessEvents++ + } + } + assert.Equal(t, 1, defragSuccessEvents, "expected exactly one defrag success event") +} + func Test_isEndpointBackendFragmented(t *testing.T) { scenarios := []struct { name string diff --git a/pkg/operator/starter.go b/pkg/operator/starter.go index 026ec2c2ee..3e05dc7609 100644 --- a/pkg/operator/starter.go +++ b/pkg/operator/starter.go @@ -496,8 +496,9 @@ func RunOperator(ctx context.Context, controllerContext *controllercmd.Controlle AlivenessChecker, operatorClient, cachedMemberClient, // for cached List/Health calls - etcdClient, // for status calls etcdClient, // for defrag calls + etcdClient, // for status calls + etcdClient, // for leader transfer before defrag configInformers.Config().V1().Infrastructures().Lister(), controllerContext.EventRecorder, kubeInformersForNamespaces, From 13be04e4876e187be182ef3c4bd8be65af2de2f1 Mon Sep 17 00:00:00 2001 From: Allen Ray Date: Mon, 13 Jul 2026 11:03:13 -0400 Subject: [PATCH 2/7] Use non-cached etcd client and shorter settle time for defrag controller Replace the cached member client with the non-cached etcd client for the defrag controller's member health and list calls. The defrag controller syncs infrequently enough that caching provides no benefit, and the cached client's 1-minute TTL could mask unhealthy members between rapid defrag operations. Increase the leader transfer settle time from 5s to 10s (renamed from leaderTransferSettleTime to defragSettleTime) to allow sufficient time for the cluster to stabilize after a leadership change. After a successful defrag, if there are remaining defrag targets, requeue with the shorter defragSettleTime rather than waiting for the full ~11 minute compaction-aligned resync period. Non-leader defrags don't block cluster writes, and the live health check on each sync ensures the previously defragged member has recovered before proceeding. --- pkg/operator/defragcontroller/defragcontroller.go | 14 ++++++++++++-- pkg/operator/starter.go | 2 +- 2 files changed, 13 insertions(+), 3 deletions(-) diff --git a/pkg/operator/defragcontroller/defragcontroller.go b/pkg/operator/defragcontroller/defragcontroller.go index 17e6d03872..61ad8e0c63 100644 --- a/pkg/operator/defragcontroller/defragcontroller.go +++ b/pkg/operator/defragcontroller/defragcontroller.go @@ -31,7 +31,11 @@ const ( maxFragmentedPercentage float64 = 45 compactionInterval = 10 * time.Minute maxDefragFailuresBeforeDegrade = 3 - leaderTransferSettleTime = 5 * time.Second + + // defragSettleTime is the minimum time to wait between defrag operations + // (including leader transfers) to allow the affected member to recover + // and rejoin the cluster before the next operation. + defragSettleTime = 10 * time.Second defragDisabledCondition = "DefragControllerDisabled" defragDisableConfigmapName = "etcd-disable-defrag" @@ -215,7 +219,7 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo recorder.Eventf("DefragControllerLeaderTransferFailed", "Failed to move leader away from member %s to member %s before defrag: %v", defragTargetMember.Name, newLeader.Member.Name, err) } else { recorder.Eventf("DefragControllerLeaderTransferred", "Moved leadership away from member %s (memberID: %x) to member %s (memberID: %x) before defrag, requeueing to allow etcd to settle", defragTargetMember.Name, defragTargetMember.ID, newLeader.Member.Name, newLeader.Member.ID) - syncCtx.Queue().AddAfter(syncCtx.QueueKey(), leaderTransferSettleTime) + syncCtx.Queue().AddAfter(syncCtx.QueueKey(), defragSettleTime) return nil } } @@ -237,6 +241,12 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo recorder.Eventf("DefragControllerDefragmentSuccess", "etcd member has been defragmented: %s, memberID: %d", defragTargetMember.Name, defragTargetMember.ID) c.numDefragFailures = 0 c.clearDegraded(ctx, recorder) + + // If there are remaining defrag targets, requeue with a shorter interval + // rather than waiting for the full compaction-aligned resync period. + if len(defragTargets) > 1 { + syncCtx.Queue().AddAfter(syncCtx.QueueKey(), defragSettleTime) + } return nil } diff --git a/pkg/operator/starter.go b/pkg/operator/starter.go index 3e05dc7609..ffaca60093 100644 --- a/pkg/operator/starter.go +++ b/pkg/operator/starter.go @@ -495,7 +495,7 @@ func RunOperator(ctx context.Context, controllerContext *controllercmd.Controlle defragController := defragcontroller.NewDefragController( AlivenessChecker, operatorClient, - cachedMemberClient, // for cached List/Health calls + etcdClient, // for member list/health calls etcdClient, // for defrag calls etcdClient, // for status calls etcdClient, // for leader transfer before defrag From b4397f47ffb19227b0d415a71de162f48f94dd6f Mon Sep 17 00:00:00 2001 From: Allen Ray Date: Mon, 13 Jul 2026 11:15:13 -0400 Subject: [PATCH 3/7] Prevent leader starvation in high-churn defrag scenarios Track which members have been successfully defragged during the current defrag cycle via a defraggedMembers set. When picking the next defrag target, skip the leader only if there are non-leader targets that haven't been defragged yet this cycle. Once all non-leaders have been handled, allow the leader through for defrag (with leader transfer). The set is cleared when no members need defragmentation, starting a fresh cycle on the next compaction interval. This addresses the scenario where in a high-churn environment, all members stay above the fragmentation threshold across sync cycles. Without this tracking, len(defragTargets) > 1 would always be true and the leader would be perpetually skipped in favor of re-fragmented non-leaders. --- .../defragcontroller/defragcontroller.go | 56 ++++++++++-- .../defragcontroller/defragcontroller_test.go | 90 +++++++++++++++++++ 2 files changed, 138 insertions(+), 8 deletions(-) diff --git a/pkg/operator/defragcontroller/defragcontroller.go b/pkg/operator/defragcontroller/defragcontroller.go index 61ad8e0c63..f752a16d04 100644 --- a/pkg/operator/defragcontroller/defragcontroller.go +++ b/pkg/operator/defragcontroller/defragcontroller.go @@ -55,6 +55,13 @@ type DefragController struct { configmapLister corev1listers.ConfigMapLister numDefragFailures int + // defraggedMembers tracks which members have been successfully defragged + // during the current defrag cycle. This prevents leader starvation in + // high-churn environments where non-leaders could be repeatedly defragged + // while the leader is always skipped. Once all non-leader targets have + // been defragged, the leader is allowed through. The set is cleared when + // no members need defragmentation. + defraggedMembers map[uint64]struct{} } func NewDefragController( @@ -181,20 +188,21 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo if len(defragTargets) == 0 { recorder.Eventf("DefragControllerDefragmentSkipped", "No etcd members meet the conditions for defragmentation") + // Clear the defragged members set when there are no more targets, + // so the next defrag cycle starts fresh. + c.defraggedMembers = nil return nil } // Sort fragmented members so we defragment the most fragmented member first slices.SortFunc(defragTargets, sortByMostFragmented) - // If the leader is in the first slot, and we have other members to defrag, defrag the others first, - // so we limit the number of times we move the leader during defrag. - defragTargetStatus := defragTargets[0].Status - defragTargetMember := defragTargets[0].Member - if len(defragTargets) > 1 && defragTargetStatus.Leader == defragTargetMember.ID { - defragTargetStatus = defragTargets[1].Status - defragTargetMember = defragTargets[1].Member - } + // Pick the defrag target. We prefer non-leader members first to avoid + // unnecessary leader transfers. However, if all non-leader targets have + // already been defragged in this cycle (tracked via defraggedMembers), + // we allow the leader through to prevent leader starvation in high-churn + // environments. + defragTargetStatus, defragTargetMember := c.pickDefragTarget(defragTargets) // Preemptively attempt to move the leadership away from the current defrag target to another valid follower. // We try this to avoid the write block that occurs while the leader is being defragmented. @@ -242,6 +250,11 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo c.numDefragFailures = 0 c.clearDegraded(ctx, recorder) + if c.defraggedMembers == nil { + c.defraggedMembers = make(map[uint64]struct{}) + } + c.defraggedMembers[defragTargetMember.ID] = struct{}{} + // If there are remaining defrag targets, requeue with a shorter interval // rather than waiting for the full compaction-aligned resync period. if len(defragTargets) > 1 { @@ -297,6 +310,33 @@ func (c *DefragController) ensureControllerDisabledCondition(ctx context.Context return nil } +// pickDefragTarget selects the next member to defragment from the sorted list +// of defrag targets. It prefers non-leader members to avoid unnecessary leader +// transfers, but will select the leader if all non-leader targets have already +// been defragged in this cycle (preventing leader starvation in high-churn +// environments). +func (c *DefragController) pickDefragTarget(defragTargets []StatusMember) (*clientv3.StatusResponse, *etcdserverpb.Member) { + leaderIdx := -1 + for i, target := range defragTargets { + if target.Status.Leader == target.Member.ID { + leaderIdx = i + continue + } + // Pick the first non-leader target that hasn't been defragged yet this cycle. + if _, alreadyDefragged := c.defraggedMembers[target.Member.ID]; !alreadyDefragged { + return target.Status, target.Member + } + } + + // All non-leader targets have been defragged (or there are none); allow the + // leader through if it needs defrag, otherwise fall back to the most + // fragmented member. + if leaderIdx >= 0 { + return defragTargets[leaderIdx].Status, defragTargets[leaderIdx].Member + } + return defragTargets[0].Status, defragTargets[0].Member +} + // isEndpointBackendFragmented checks the status of all cluster members to ensure that no members have a fragmented store. // This can happen if the operator starts defrag of the cluster but then loses leader status and is rescheduled before // the operator can defrag all members. diff --git a/pkg/operator/defragcontroller/defragcontroller_test.go b/pkg/operator/defragcontroller/defragcontroller_test.go index be82946e68..d94b20d97f 100644 --- a/pkg/operator/defragcontroller/defragcontroller_test.go +++ b/pkg/operator/defragcontroller/defragcontroller_test.go @@ -496,6 +496,96 @@ func TestDefragMovesLeadershipBeforeDefrag(t *testing.T) { assert.Equal(t, 1, defragSuccessEvents, "expected exactly one defrag success event") } +// TestDefragLeaderNotStarvedInHighChurn verifies that in a high-churn environment +// where non-leader members re-fragment between sync cycles, the leader is eventually +// defragged rather than being perpetually skipped. +func TestDefragLeaderNotStarvedInHighChurn(t *testing.T) { + fakeOperatorClient := v1helpers.NewFakeStaticPodOperatorClient( + &operatorv1.StaticPodOperatorSpec{ + OperatorSpec: operatorv1.OperatorSpec{ + ManagementState: operatorv1.Managed, + }, + }, + u.StaticPodOperatorStatus(), + nil, + nil, + ) + + integration.BeforeTestExternal(t) + testServer := integration.NewCluster(t, &integration.ClusterConfig{Size: 3}) + defer testServer.Terminate(t) + + etcdMembers := waitForMembersWithClientURLs(t, testServer) + + var status []*clientv3.StatusResponse + for _, member := range testServer.Members { + statusResp, err := testServer.Client(0).Status(context.TODO(), member.GRPCURL) + require.NoError(t, err) + // All members are fragmented. + statusResp.DbSize = minDefragBytes + statusResp.DbSizeInUse = minDefragBytes / 2 + status = append(status, statusResp) + } + + fakeEtcdClient, _ := etcdcli.NewFakeEtcdClient( + etcdMembers, + etcdcli.WithFakeClusterHealth(&etcdcli.FakeMemberHealth{Healthy: 3}), + etcdcli.WithFakeStatus(status), + ) + + eventRecorder := events.NewInMemoryRecorder(t.Name(), clock.RealClock{}) + indexer := cache.NewIndexer(cache.MetaNamespaceKeyFunc, cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc}) + require.NoError(t, indexer.Add(u.FakeInfrastructureTopology(configv1.HighlyAvailableTopologyMode))) + + controller := &DefragController{ + operatorClient: fakeOperatorClient, + memberLister: fakeEtcdClient, + defragClient: fakeEtcdClient, + statusClient: fakeEtcdClient, + leaderMover: fakeEtcdClient, + infrastructureLister: configv1listers.NewInfrastructureLister(indexer), + configmapLister: corev1listers.NewConfigMapLister(indexer), + } + + syncCtx := factory.NewSyncContext("defrag-controller", eventRecorder) + + // Simulate high churn: after each sync, re-fragment all members. + // Without the defraggedMembers tracking, the leader would be starved + // because len(defragTargets) > 1 is always true. + // + // Expected sequence for a 3-node cluster: + // sync 1: defrag non-leader A + // sync 2: defrag non-leader B (A re-fragmented but already tracked) + // sync 3: leader transfer (all non-leaders defragged this cycle) + // sync 4: defrag former leader + maxSyncs := 10 + var leaderDefragged bool + for i := 0; i < maxSyncs; i++ { + // Re-fragment all members to simulate high churn. + for _, s := range status { + s.DbSize = minDefragBytes + s.DbSizeInUse = minDefragBytes / 2 + } + + err := controller.sync(context.TODO(), syncCtx) + require.NoError(t, err) + + // Check if the leader was defragged by looking for a defrag attempt + // on the leader member. + for _, event := range eventRecorder.Events() { + if strings.HasPrefix(event.Message, "Moved leadership away from member") { + leaderDefragged = true + break + } + } + if leaderDefragged { + break + } + } + + assert.True(t, leaderDefragged, "leader should have been selected for defrag (via leader transfer) within %d syncs", maxSyncs) +} + func Test_isEndpointBackendFragmented(t *testing.T) { scenarios := []struct { name string From f4bf8bbff5985b6100c81ab97760798a24fe171c Mon Sep 17 00:00:00 2001 From: Allen Ray Date: Tue, 14 Jul 2026 12:56:26 -0400 Subject: [PATCH 4/7] Address comments --- .../defragcontroller/defragcontroller.go | 106 +++++++----------- pkg/operator/starter.go | 8 +- 2 files changed, 43 insertions(+), 71 deletions(-) diff --git a/pkg/operator/defragcontroller/defragcontroller.go b/pkg/operator/defragcontroller/defragcontroller.go index f752a16d04..023ca397d7 100644 --- a/pkg/operator/defragcontroller/defragcontroller.go +++ b/pkg/operator/defragcontroller/defragcontroller.go @@ -55,13 +55,8 @@ type DefragController struct { configmapLister corev1listers.ConfigMapLister numDefragFailures int - // defraggedMembers tracks which members have been successfully defragged - // during the current defrag cycle. This prevents leader starvation in - // high-churn environments where non-leaders could be repeatedly defragged - // while the leader is always skipped. Once all non-leader targets have - // been defragged, the leader is allowed through. The set is cleared when - // no members need defragmentation. - defraggedMembers map[uint64]struct{} + // defragTargets tracks the ids of members that need to be defragged during the current cycle + defragTargets []uint64 } func NewDefragController( @@ -140,6 +135,10 @@ type StatusMember struct { Member *etcdserverpb.Member } +func (sm StatusMember) IsLeader() bool { + return sm.Status.Leader == sm.Member.ID +} + func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncContext) error { recorder := syncCtx.Recorder() // Do not defrag if any of the cluster members are unhealthy. @@ -157,8 +156,8 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo } var ( + isNewCycle = len(c.defragTargets) == 0 statusMembers = make(map[uint64]StatusMember) - defragTargets = make([]StatusMember, 0, len(members)) ) for _, member := range members { // filter out learner members since they don't support the defragment API call @@ -181,34 +180,35 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo statusMembers[member.ID] = sm - if isEndpointBackendFragmented(member, status) { - defragTargets = append(defragTargets, sm) + if isNewCycle && isEndpointBackendFragmented(member, status) { + c.defragTargets = append(c.defragTargets, member.ID) } } - if len(defragTargets) == 0 { + if len(c.defragTargets) == 0 { recorder.Eventf("DefragControllerDefragmentSkipped", "No etcd members meet the conditions for defragmentation") - // Clear the defragged members set when there are no more targets, - // so the next defrag cycle starts fresh. - c.defraggedMembers = nil return nil } - // Sort fragmented members so we defragment the most fragmented member first - slices.SortFunc(defragTargets, sortByMostFragmented) + // Sort fragmented members so we defragment the most fragmented member first while defragging the leader last + slices.SortFunc(c.defragTargets, func(a, b uint64) int { + aIsLeader, bIsLeader := statusMembers[a].IsLeader(), statusMembers[b].IsLeader() + if aIsLeader && !bIsLeader { + return 1 + } else if !aIsLeader && bIsLeader { + return -1 + } + return sortByMostFragmented(statusMembers, a, b) + }) - // Pick the defrag target. We prefer non-leader members first to avoid - // unnecessary leader transfers. However, if all non-leader targets have - // already been defragged in this cycle (tracked via defraggedMembers), - // we allow the leader through to prevent leader starvation in high-churn - // environments. - defragTargetStatus, defragTargetMember := c.pickDefragTarget(defragTargets) + defragTarget := statusMembers[c.defragTargets[0]] + defragTargetStatus, defragTargetMember := defragTarget.Status, defragTarget.Member // Preemptively attempt to move the leadership away from the current defrag target to another valid follower. // We try this to avoid the write block that occurs while the leader is being defragmented. // We record any error that occurs while attempting this, but we do not halt defrag if the move fails; // we just accept the write block. - if defragTargetStatus.Leader == defragTargetMember.ID && len(statusMembers) > 1 { + if defragTarget.IsLeader() && len(statusMembers) > 1 { followers := make([]StatusMember, 0, len(statusMembers)) for id, member := range statusMembers { if defragTargetMember.ID == id { @@ -217,20 +217,20 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo followers = append(followers, member) } - if len(followers) > 0 { - slices.SortFunc(followers, sortByLeastFragmented) - - newLeader := followers[0] + slices.SortFunc(followers, sortByLeastFragmented) + for _, newLeader := range followers { err := c.leaderMover.MoveLeader(ctx, newLeader.Member.ID) if err != nil { - recorder.Eventf("DefragControllerLeaderTransferFailed", "Failed to move leader away from member %s to member %s before defrag: %v", defragTargetMember.Name, newLeader.Member.Name, err) - } else { - recorder.Eventf("DefragControllerLeaderTransferred", "Moved leadership away from member %s (memberID: %x) to member %s (memberID: %x) before defrag, requeueing to allow etcd to settle", defragTargetMember.Name, defragTargetMember.ID, newLeader.Member.Name, newLeader.Member.ID) - syncCtx.Queue().AddAfter(syncCtx.QueueKey(), defragSettleTime) - return nil + recorder.Warningf("DefragControllerLeaderTransferAttemptFailed", "Failed to move leader away from member %s to member %s before defrag: %v", defragTargetMember.Name, newLeader.Member.Name, err) + continue } + + recorder.Eventf("DefragControllerLeaderTransferSuccess", "Moved leadership away from member %s (memberID: %x) to member %s (memberID: %x) before defrag, requeueing to allow etcd to settle", defragTargetMember.Name, defragTargetMember.ID, newLeader.Member.Name, newLeader.Member.ID) + syncCtx.Queue().AddAfter(syncCtx.QueueKey(), defragSettleTime) + return nil } + recorder.Warning("DefragControllerLeaderTransferFailed", "Failed to move leader away from member %s, continuing with blocking leader defrag") } recorder.Eventf("DefragControllerDefragmentAttempt", "Attempting defrag on member: %s, memberID: %x, dbSize: %d, dbInUse: %d, leader ID: %d", defragTargetMember.Name, defragTargetMember.ID, defragTargetStatus.DbSize, defragTargetStatus.DbSizeInUse, defragTargetStatus.Leader) @@ -243,6 +243,7 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo if c.numDefragFailures >= maxDefragFailuresBeforeDegrade { c.setDegraded(ctx, recorder) } + syncCtx.Queue().AddAfter(syncCtx.QueueKey(), defragSettleTime) return nil } @@ -250,14 +251,11 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo c.numDefragFailures = 0 c.clearDegraded(ctx, recorder) - if c.defraggedMembers == nil { - c.defraggedMembers = make(map[uint64]struct{}) - } - c.defraggedMembers[defragTargetMember.ID] = struct{}{} + c.defragTargets = c.defragTargets[1:] // If there are remaining defrag targets, requeue with a shorter interval // rather than waiting for the full compaction-aligned resync period. - if len(defragTargets) > 1 { + if len(c.defragTargets) > 0 { syncCtx.Queue().AddAfter(syncCtx.QueueKey(), defragSettleTime) } return nil @@ -310,33 +308,6 @@ func (c *DefragController) ensureControllerDisabledCondition(ctx context.Context return nil } -// pickDefragTarget selects the next member to defragment from the sorted list -// of defrag targets. It prefers non-leader members to avoid unnecessary leader -// transfers, but will select the leader if all non-leader targets have already -// been defragged in this cycle (preventing leader starvation in high-churn -// environments). -func (c *DefragController) pickDefragTarget(defragTargets []StatusMember) (*clientv3.StatusResponse, *etcdserverpb.Member) { - leaderIdx := -1 - for i, target := range defragTargets { - if target.Status.Leader == target.Member.ID { - leaderIdx = i - continue - } - // Pick the first non-leader target that hasn't been defragged yet this cycle. - if _, alreadyDefragged := c.defraggedMembers[target.Member.ID]; !alreadyDefragged { - return target.Status, target.Member - } - } - - // All non-leader targets have been defragged (or there are none); allow the - // leader through if it needs defrag, otherwise fall back to the most - // fragmented member. - if leaderIdx >= 0 { - return defragTargets[leaderIdx].Status, defragTargets[leaderIdx].Member - } - return defragTargets[0].Status, defragTargets[0].Member -} - // isEndpointBackendFragmented checks the status of all cluster members to ensure that no members have a fragmented store. // This can happen if the operator starts defrag of the cluster but then loses leader status and is rescheduled before // the operator can defrag all members. @@ -361,9 +332,10 @@ func sortByLeastFragmented(a, b StatusMember) int { ) } -func sortByMostFragmented(a, b StatusMember) int { +func sortByMostFragmented(statusMembers map[uint64]StatusMember, a, b uint64) int { + aStatus, bStatus := statusMembers[a].Status, statusMembers[b].Status return cmp.Compare( - checkFragmentationPercentage(b.Status.DbSize, b.Status.DbSizeInUse), - checkFragmentationPercentage(a.Status.DbSize, a.Status.DbSizeInUse), + checkFragmentationPercentage(bStatus.DbSize, bStatus.DbSizeInUse), + checkFragmentationPercentage(aStatus.DbSize, aStatus.DbSizeInUse), ) } diff --git a/pkg/operator/starter.go b/pkg/operator/starter.go index ffaca60093..7ab623affc 100644 --- a/pkg/operator/starter.go +++ b/pkg/operator/starter.go @@ -495,10 +495,10 @@ func RunOperator(ctx context.Context, controllerContext *controllercmd.Controlle defragController := defragcontroller.NewDefragController( AlivenessChecker, operatorClient, - etcdClient, // for member list/health calls - etcdClient, // for defrag calls - etcdClient, // for status calls - etcdClient, // for leader transfer before defrag + etcdClient, // for member list/health calls + etcdClient, // for defrag calls + etcdClient, // for status calls + etcdClient, // for leader transfer before defrag configInformers.Config().V1().Infrastructures().Lister(), controllerContext.EventRecorder, kubeInformersForNamespaces, From 8c0304ff85109ad607dbd80fc2462579161b85ff Mon Sep 17 00:00:00 2001 From: Allen Ray Date: Wed, 15 Jul 2026 12:03:14 -0400 Subject: [PATCH 5/7] Fix MoveLeader routing and enable StopGRPCServiceOnDefrag MoveLeader must be sent to the current leader, but the client pool may route to any member, causing consistent "etcdserver: not leader" errors and falling through to blocking leader defrag every time. Fix by accepting the leader member at the call site and creating a dedicated client connected directly to the leader endpoint, following the same pattern used by Defragment(). Also enable the StopGRPCServiceOnDefrag feature gate so etcd reports NOT_SERVING during defrag, allowing load-balancer-aware clients to route away from the member being defragged. Additionally fix the DefragControllerLeaderTransferFailed event which used Warning instead of Warningf, leaving %s uninterpolated. --- bindata/etcd/pod.gotpl.yaml | 2 +- pkg/etcdcli/etcdcli.go | 17 ++++++++++++----- pkg/etcdcli/helpers.go | 2 +- pkg/etcdcli/interfaces.go | 2 +- .../defragcontroller/defragcontroller.go | 4 ++-- 5 files changed, 17 insertions(+), 10 deletions(-) diff --git a/bindata/etcd/pod.gotpl.yaml b/bindata/etcd/pod.gotpl.yaml index 2dbf1a7ddc..c56bf8cdef 100644 --- a/bindata/etcd/pod.gotpl.yaml +++ b/bindata/etcd/pod.gotpl.yaml @@ -170,7 +170,7 @@ spec: exec nice -n -19 ionice -c2 -n0 etcd \ --logger=zap \ --log-level={{.LogLevel}} \ - --feature-gates=InitialCorruptCheck=true \ + --feature-gates=InitialCorruptCheck=true,StopGRPCServiceOnDefrag=true \ --initial-advertise-peer-urls=https://${NODE_NODE_ENVVAR_NAME_IP}:2380 \ --cert-file=/etc/kubernetes/static-pod-certs/secrets/etcd-all-certs/etcd-serving-NODE_NAME.crt \ --key-file=/etc/kubernetes/static-pod-certs/secrets/etcd-all-certs/etcd-serving-NODE_NAME.key \ diff --git a/pkg/etcdcli/etcdcli.go b/pkg/etcdcli/etcdcli.go index 3a9ffc62fb..ad752f9044 100644 --- a/pkg/etcdcli/etcdcli.go +++ b/pkg/etcdcli/etcdcli.go @@ -246,13 +246,20 @@ func (g *etcdClientGetter) MemberUpdatePeerURL(ctx context.Context, id uint64, p return err } -func (g *etcdClientGetter) MoveLeader(ctx context.Context, toMember uint64) error { - cli, err := g.clientPool.Get() +// MoveLeader creates a new client connected directly to the given leader member +// and issues the MoveLeader RPC. The MoveLeader API requires the request to be +// sent to the current leader; using the client pool can route to a follower and +// return "etcdserver: not leader". +func (g *etcdClientGetter) MoveLeader(ctx context.Context, leader *etcdserverpb.Member, toMember uint64) error { + cli, err := newEtcdClientWithClientOpts([]string{leader.ClientURLs[0]}, false) if err != nil { - return err + return fmt.Errorf("failed to create client to leader %s for MoveLeader: %w", leader.ClientURLs[0], err) } - - defer g.clientPool.Return(cli) + defer func() { + if err := cli.Close(); err != nil { + klog.Errorf("error closing leader client for MoveLeader: %v", err) + } + }() ctx, cancel := context.WithTimeout(ctx, DefaultClientTimeout) defer cancel() diff --git a/pkg/etcdcli/helpers.go b/pkg/etcdcli/helpers.go index 968a5b6b2f..f20888366a 100644 --- a/pkg/etcdcli/helpers.go +++ b/pkg/etcdcli/helpers.go @@ -88,7 +88,7 @@ func (f *fakeEtcdClient) VotingMemberList(ctx context.Context) ([]*etcdserverpb. return filterVotingMembers(members), nil } -func (f *fakeEtcdClient) MoveLeader(ctx context.Context, toMember uint64) error { +func (f *fakeEtcdClient) MoveLeader(ctx context.Context, leader *etcdserverpb.Member, toMember uint64) error { for _, status := range f.opts.status { status.Leader = toMember } diff --git a/pkg/etcdcli/interfaces.go b/pkg/etcdcli/interfaces.go index ae6cf3e072..3cfe7c896b 100644 --- a/pkg/etcdcli/interfaces.go +++ b/pkg/etcdcli/interfaces.go @@ -66,7 +66,7 @@ type MemberRemover interface { } type LeaderMover interface { - MoveLeader(ctx context.Context, toMember uint64) error + MoveLeader(ctx context.Context, leader *etcdserverpb.Member, toMember uint64) error } type MemberLister interface { diff --git a/pkg/operator/defragcontroller/defragcontroller.go b/pkg/operator/defragcontroller/defragcontroller.go index 023ca397d7..14f5619019 100644 --- a/pkg/operator/defragcontroller/defragcontroller.go +++ b/pkg/operator/defragcontroller/defragcontroller.go @@ -220,7 +220,7 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo slices.SortFunc(followers, sortByLeastFragmented) for _, newLeader := range followers { - err := c.leaderMover.MoveLeader(ctx, newLeader.Member.ID) + err := c.leaderMover.MoveLeader(ctx, defragTargetMember, newLeader.Member.ID) if err != nil { recorder.Warningf("DefragControllerLeaderTransferAttemptFailed", "Failed to move leader away from member %s to member %s before defrag: %v", defragTargetMember.Name, newLeader.Member.Name, err) continue @@ -230,7 +230,7 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo syncCtx.Queue().AddAfter(syncCtx.QueueKey(), defragSettleTime) return nil } - recorder.Warning("DefragControllerLeaderTransferFailed", "Failed to move leader away from member %s, continuing with blocking leader defrag") + recorder.Warningf("DefragControllerLeaderTransferFailed", "Failed to move leader away from member %s, continuing with blocking leader defrag", defragTargetMember.Name) } recorder.Eventf("DefragControllerDefragmentAttempt", "Attempting defrag on member: %s, memberID: %x, dbSize: %d, dbInUse: %d, leader ID: %d", defragTargetMember.Name, defragTargetMember.ID, defragTargetStatus.DbSize, defragTargetStatus.DbSizeInUse, defragTargetStatus.Leader) From 680d8a1aa15b587a9143296920e616470cdd77d0 Mon Sep 17 00:00:00 2001 From: Allen Ray Date: Fri, 17 Jul 2026 12:28:09 -0400 Subject: [PATCH 6/7] Enable NonBlockingDefrag feature gate and defrag one member per sync Enable the NonBlockingDefrag etcd feature gate so that defragmentation does not hold exclusive locks for the entire copy operation, avoiding write stalls during defrag. Also change the defrag controller to only defrag one member per sync cycle rather than immediately requeueing for the next member. This lets the normal 11-minute resync handle subsequent members, giving the cluster a full compaction interval to recover between defrags. --- bindata/etcd/pod.gotpl.yaml | 2 +- .../defragcontroller/defragcontroller.go | 19 ++++++------------- 2 files changed, 7 insertions(+), 14 deletions(-) diff --git a/bindata/etcd/pod.gotpl.yaml b/bindata/etcd/pod.gotpl.yaml index c56bf8cdef..f5955660d7 100644 --- a/bindata/etcd/pod.gotpl.yaml +++ b/bindata/etcd/pod.gotpl.yaml @@ -170,7 +170,7 @@ spec: exec nice -n -19 ionice -c2 -n0 etcd \ --logger=zap \ --log-level={{.LogLevel}} \ - --feature-gates=InitialCorruptCheck=true,StopGRPCServiceOnDefrag=true \ + --feature-gates=InitialCorruptCheck=true,StopGRPCServiceOnDefrag=true,NonBlockingDefrag=true \ --initial-advertise-peer-urls=https://${NODE_NODE_ENVVAR_NAME_IP}:2380 \ --cert-file=/etc/kubernetes/static-pod-certs/secrets/etcd-all-certs/etcd-serving-NODE_NAME.crt \ --key-file=/etc/kubernetes/static-pod-certs/secrets/etcd-all-certs/etcd-serving-NODE_NAME.key \ diff --git a/pkg/operator/defragcontroller/defragcontroller.go b/pkg/operator/defragcontroller/defragcontroller.go index 14f5619019..1509d47f6f 100644 --- a/pkg/operator/defragcontroller/defragcontroller.go +++ b/pkg/operator/defragcontroller/defragcontroller.go @@ -32,9 +32,8 @@ const ( compactionInterval = 10 * time.Minute maxDefragFailuresBeforeDegrade = 3 - // defragSettleTime is the minimum time to wait between defrag operations - // (including leader transfers) to allow the affected member to recover - // and rejoin the cluster before the next operation. + // defragSettleTime is the minimum time to wait after a leader transfer + // to allow the cluster to stabilize before proceeding with defrag. defragSettleTime = 10 * time.Second defragDisabledCondition = "DefragControllerDisabled" @@ -205,9 +204,9 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo defragTargetStatus, defragTargetMember := defragTarget.Status, defragTarget.Member // Preemptively attempt to move the leadership away from the current defrag target to another valid follower. - // We try this to avoid the write block that occurs while the leader is being defragmented. - // We record any error that occurs while attempting this, but we do not halt defrag if the move fails; - // we just accept the write block. + // We try this to avoid multiple leader elections in the case where defragging the leader causes leadership + // to move to a member we've yet to defrag, which could in turn lose leadership, etc. causing a lot of churn. + // We record any error that occurs while attempting this, but we do not halt defrag if the move fails. if defragTarget.IsLeader() && len(statusMembers) > 1 { followers := make([]StatusMember, 0, len(statusMembers)) for id, member := range statusMembers { @@ -235,7 +234,6 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo recorder.Eventf("DefragControllerDefragmentAttempt", "Attempting defrag on member: %s, memberID: %x, dbSize: %d, dbInUse: %d, leader ID: %d", defragTargetMember.Name, defragTargetMember.ID, defragTargetStatus.DbSize, defragTargetStatus.DbSizeInUse, defragTargetStatus.Leader) if _, err := c.defragClient.Defragment(ctx, defragTargetMember); err != nil { - // Defrag can timeout if defragmentation takes longer than etcdcli.DefragDialTimeout. errMsg := fmt.Sprintf("failed defrag on member: %s, memberID: %x: %v", defragTargetMember.Name, defragTargetMember.ID, err) recorder.Eventf("DefragControllerDefragmentFailed", errMsg) klog.Errorf("%s", errMsg) @@ -252,12 +250,7 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo c.clearDegraded(ctx, recorder) c.defragTargets = c.defragTargets[1:] - - // If there are remaining defrag targets, requeue with a shorter interval - // rather than waiting for the full compaction-aligned resync period. - if len(c.defragTargets) > 0 { - syncCtx.Queue().AddAfter(syncCtx.QueueKey(), defragSettleTime) - } + syncCtx.Queue().AddAfter(syncCtx.QueueKey(), defragSettleTime) return nil } From ccc4a52a02e65420c4530511b3a895785b31ea65 Mon Sep 17 00:00:00 2001 From: Allen Ray Date: Thu, 23 Jul 2026 08:59:46 -0400 Subject: [PATCH 7/7] Address comments --- pkg/operator/defragcontroller/defragcontroller.go | 15 +++++++++++++-- 1 file changed, 13 insertions(+), 2 deletions(-) diff --git a/pkg/operator/defragcontroller/defragcontroller.go b/pkg/operator/defragcontroller/defragcontroller.go index 1509d47f6f..9ff301f2cb 100644 --- a/pkg/operator/defragcontroller/defragcontroller.go +++ b/pkg/operator/defragcontroller/defragcontroller.go @@ -184,8 +184,19 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo } } + c.defragTargets = slices.DeleteFunc(c.defragTargets, func(i uint64) bool { + statusMember, has := statusMembers[i] + if !has { + // Remove any defrag targets that we don't have a status for. + return true + } + + // Remove any members that no longer meet the conditions for defrag. + return !isEndpointBackendFragmented(statusMember.Member, statusMember.Status) + }) + if len(c.defragTargets) == 0 { - recorder.Eventf("DefragControllerDefragmentSkipped", "No etcd members meet the conditions for defragmentation") + klog.V(4).Info("Defrag skipped: no etcd members meet the conditions for defragmentation") return nil } @@ -235,7 +246,7 @@ func (c *DefragController) runDefrag(ctx context.Context, syncCtx factory.SyncCo recorder.Eventf("DefragControllerDefragmentAttempt", "Attempting defrag on member: %s, memberID: %x, dbSize: %d, dbInUse: %d, leader ID: %d", defragTargetMember.Name, defragTargetMember.ID, defragTargetStatus.DbSize, defragTargetStatus.DbSizeInUse, defragTargetStatus.Leader) if _, err := c.defragClient.Defragment(ctx, defragTargetMember); err != nil { errMsg := fmt.Sprintf("failed defrag on member: %s, memberID: %x: %v", defragTargetMember.Name, defragTargetMember.ID, err) - recorder.Eventf("DefragControllerDefragmentFailed", errMsg) + recorder.Warningf("DefragControllerDefragmentFailed", errMsg) klog.Errorf("%s", errMsg) c.numDefragFailures++ if c.numDefragFailures >= maxDefragFailuresBeforeDegrade {