From 58c9fe0257883d08b0a8021a0d4ddf573a7c94ef Mon Sep 17 00:00:00 2001 From: Andrey Kolkov Date: Mon, 17 Aug 2026 21:53:46 +0400 Subject: [PATCH 1/4] =?UTF-8?q?feat(defrag):=20EtcdDefrag=20controller=20?= =?UTF-8?q?=E2=80=94=20operator-driven=20defragmentation?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Reworks the defragmentation implementation from the rejected EtcdCluster.spec.defrag shape (#361 review) into a controller for the dedicated EtcdDefrag API (proposal/etcd-defrag-api, which this is stacked on). The controller reconciles an EtcdDefrag as a one-shot, run-to-completion sweep: resolve the cluster; serialize per cluster (oldest non-terminal run acts, the rest wait Pending); gate on real health (every desired member present, reachable, alarm-free, agreeing on a leader — not just "Status answered"); then defragment members one at a time, followers before the leader, one per reconcile pass with the health re-checked between. A due defrag on an unhealthy cluster is deferred (Pending + DefragChecked=False/ClusterNotHealthy + a DefragDeferred event), never forced. Per-member outcomes/sizes land in status.members; phase moves Pending -> Running -> Complete|Failed; an overall active-deadline bounds a stuck run; ttlSecondsAfterFinished GCs a finished record. The rule matches the EtcdDefrag API: rule.all is unconditional; otherwise the reclaimable floor (freeSpaceAbove, default 200Mi) is always applied and the quota arm only fires with at least minReclaim to reclaim — so a full-but-unfragmented backend is never defragmented for nothing. Adds Defragment to the etcd client interface, wires the controller (RBAC + main.go), and covers it with unit + controller-integration tests (defrag when needed, skip below threshold, defer-not-force without quorum, failed RPC, per-cluster serialization) and an e2e retargeted to create EtcdDefrag objects. Capacity metrics/alerts are intentionally out of scope here (tracked in #357); no changes to EtcdClusterSpec. Refs #221, #357; supersedes the #361 spec.defrag approach. Assisted-By: Claude Opus 4.8 (1M context) Signed-off-by: Andrey Kolkov --- README.md | 2 +- api/v1alpha2/etcddefrag_types.go | 4 - ...tcd-operator.cozystack.io_etcddefrags.yaml | 4 - .../files/manager-role-rules.yaml | 19 + controllers/etcd_client.go | 16 +- controllers/etcddefrag_controller.go | 537 ++++++++++++++++++ controllers/etcddefrag_controller_test.go | 311 ++++++++++ controllers/testing_helpers_test.go | 33 +- docs/etcd-defrag.md | 15 +- main.go | 8 + test/e2e/defrag_test.go | 211 +++++++ 11 files changed, 1134 insertions(+), 26 deletions(-) create mode 100644 controllers/etcddefrag_controller.go create mode 100644 controllers/etcddefrag_controller_test.go create mode 100644 test/e2e/defrag_test.go diff --git a/README.md b/README.md index 4c246728..70184d73 100644 --- a/README.md +++ b/README.md @@ -76,7 +76,7 @@ For step-by-step setup, RBAC, image versions, and teardown see [docs/installatio - **[Installation](docs/installation.md)** — deploy the operator, create your first cluster, networking pitfalls, upgrades. - **[Concepts](docs/concepts.md)** — design rationale: locking pattern, single-seed bootstrap, GenerateName naming, scale-to-zero mechanics, conditions reference. - **[Operations](docs/operations.md)** — runbook for day-2: scaling, pausing/resuming, decoding conditions, escalating stuck reconciles, broken-member recovery. -- **[Defragmentation](docs/etcd-defrag.md)** — the `EtcdDefrag` resource: reclaiming etcd backend disk and its safety model (API type; reconciling controller is a follow-up). +- **[Defragmentation](docs/etcd-defrag.md)** — the `EtcdDefrag` resource: reclaiming etcd backend disk, one-shot and driven from outside on a schedule, and its safety model. - **[Migration](docs/migration.md)** — moving onto this operator from the legacy aenix operator; tracks behavioural changes that need an explicit migration step — currently the BYO root-credentials requirement when enabling auth. ## Testing diff --git a/api/v1alpha2/etcddefrag_types.go b/api/v1alpha2/etcddefrag_types.go index 7b11db9f..4e2bfa44 100644 --- a/api/v1alpha2/etcddefrag_types.go +++ b/api/v1alpha2/etcddefrag_types.go @@ -201,10 +201,6 @@ type MemberDefragStatus struct { // run-to-completion defragmentation of an EtcdCluster's members. Like // EtcdSnapshot it is a record: the operator drives it through status.phase and // it never re-runs. -// -// NOTE: this ships the API type ahead of its reconciling controller. Until that -// controller lands, an EtcdDefrag is inert — creating one records intent but -// nothing acts on it (no sweep runs, status stays empty, TTL does not fire). type EtcdDefrag struct { metav1.TypeMeta `json:",inline"` metav1.ObjectMeta `json:"metadata,omitempty"` diff --git a/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml b/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml index 18934ae2..febca3ac 100644 --- a/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml +++ b/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml @@ -35,10 +35,6 @@ spec: run-to-completion defragmentation of an EtcdCluster's members. Like EtcdSnapshot it is a record: the operator drives it through status.phase and it never re-runs. - - NOTE: this ships the API type ahead of its reconciling controller. Until that - controller lands, an EtcdDefrag is inert — creating one records intent but - nothing acts on it (no sweep runs, status stays empty, TTL does not fire). properties: apiVersion: description: |- diff --git a/charts/etcd-operator/files/manager-role-rules.yaml b/charts/etcd-operator/files/manager-role-rules.yaml index 4ad4b945..fe21bed2 100644 --- a/charts/etcd-operator/files/manager-role-rules.yaml +++ b/charts/etcd-operator/files/manager-role-rules.yaml @@ -1,3 +1,10 @@ +- apiGroups: + - "" + resources: + - events + verbs: + - create + - patch - apiGroups: - "" resources: @@ -87,12 +94,24 @@ - etcd-operator.cozystack.io resources: - etcdclusters/status + - etcddefrags/status - etcdmembers/status - etcdsnapshots/status verbs: - get - patch - update +- apiGroups: + - etcd-operator.cozystack.io + resources: + - etcddefrags + verbs: + - delete + - get + - list + - patch + - update + - watch - apiGroups: - etcd-operator.cozystack.io resources: diff --git a/controllers/etcd_client.go b/controllers/etcd_client.go index e00d8fc0..c4de7e05 100644 --- a/controllers/etcd_client.go +++ b/controllers/etcd_client.go @@ -35,13 +35,21 @@ type EtcdClusterClient interface { MemberPromote(ctx context.Context, id uint64) (*clientv3.MemberPromoteResponse, error) MemberRemove(ctx context.Context, id uint64) (*clientv3.MemberRemoveResponse, error) - // Status returns a single endpoint's server status, including the etcd - // version it is actually running (StatusResponse.Version). Used by the - // member controller to observe the running version into - // EtcdMember.status.version. *clientv3.Client satisfies this via its + // Status returns a single endpoint's server status: the etcd version it is + // actually running (StatusResponse.Version, observed into + // EtcdMember.status.version), its backend sizes (DbSize / DbSizeInUse), the + // leader it sees and any alarms — all read by the EtcdDefrag controller to + // gate and decide a defragmentation. *clientv3.Client satisfies this via its // embedded Maintenance interface. Status(ctx context.Context, endpoint string) (*clientv3.StatusResponse, error) + // Defragment releases a single endpoint's reclaimable backend space to the + // filesystem. It is per-endpoint (defrag is a member-local operation) and + // briefly blocks that member, so the EtcdDefrag controller calls it one + // member at a time on a healthy cluster. *clientv3.Client satisfies this via + // its embedded Maintenance interface. + Defragment(ctx context.Context, endpoint string) (*clientv3.DefragmentResponse, error) + // Auth surface — used by reconcileAuth to provision the single root // user/role and turn on authentication. The "root" role is built into // etcd, so a RoleAdd is not needed: UserAdd("root", …) + diff --git a/controllers/etcddefrag_controller.go b/controllers/etcddefrag_controller.go new file mode 100644 index 00000000..0a4dc98f --- /dev/null +++ b/controllers/etcddefrag_controller.go @@ -0,0 +1,537 @@ +/* +Copyright 2023 Timofey Larkin. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 +*/ + +package controllers + +import ( + "context" + "fmt" + "sort" + "strconv" + "strings" + "time" + + clientv3 "go.etcd.io/etcd/client/v3" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/tools/record" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/log" + + lll "github.com/cozystack/etcd-operator/api/v1alpha2" +) + +const ( + // Rule defaults (see api DefragRule). A defrag can only reclaim + // DbSize-DbSizeInUse, so freeSpaceAbove is the always-applied floor. + defaultDefragFreeSpace = int64(200 << 20) // 200Mi + defaultDefragMinReclaim = int64(32 << 20) // 32Mi + // defaultEtcdQuotaBytes mirrors etcd's built-in --quota-backend-bytes + // default (2Gi), used when the cluster leaves spec.options.quotaBackendBytes + // unset. + defaultEtcdQuotaBytes = int64(2 << 30) + + // defragStatusTimeout bounds a per-member Maintenance Status probe; + // defragRPCTimeout bounds a single stop-the-world Defragment call. + defragStatusTimeout = 5 * time.Second + defragRPCTimeout = 5 * time.Minute + // defragRequeueAfter paces the one-member-per-pass sweep and the retry while + // deferred on an unhealthy cluster. + defragRequeueAfter = 10 * time.Second + // defragActiveDeadline bounds a whole run (Running + waiting-while-Pending), + // so a run stuck on an unhealthy cluster can't hold the per-cluster slot + // forever. + defragActiveDeadline = 30 * time.Minute +) + +// EtcdDefragReconciler drives an EtcdDefrag: a one-shot, run-to-completion +// defragmentation of an EtcdCluster's members. It defragments members one at a +// time (followers before the leader) on a healthy cluster, deferring rather +// than forcing while the cluster is degraded, and records per-member outcomes +// in status. +type EtcdDefragReconciler struct { + client.Client + Scheme *runtime.Scheme + + // EtcdClientFactory builds an etcd client; tests inject a fake. + EtcdClientFactory EtcdClientFactory + + // Recorder emits events (e.g. when a due defrag is deferred). Tests may + // leave it nil. + Recorder record.EventRecorder +} + +//+kubebuilder:rbac:groups=etcd-operator.cozystack.io,resources=etcddefrags,verbs=get;list;watch;update;patch;delete +//+kubebuilder:rbac:groups=etcd-operator.cozystack.io,resources=etcddefrags/status,verbs=get;update;patch +//+kubebuilder:rbac:groups=etcd-operator.cozystack.io,resources=etcdclusters,verbs=get;list;watch +//+kubebuilder:rbac:groups=etcd-operator.cozystack.io,resources=etcdmembers,verbs=get;list;watch +//+kubebuilder:rbac:groups="",resources=events,verbs=create;patch + +func (r *EtcdDefragReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) { + logger := log.FromContext(ctx) + + df := &lll.EtcdDefrag{} + if err := r.Get(ctx, req.NamespacedName, df); err != nil { + return ctrl.Result{}, client.IgnoreNotFound(err) + } + + // Terminal: only TTL GC remains. + if df.Status.Phase == lll.EtcdDefragPhaseComplete || df.Status.Phase == lll.EtcdDefragPhaseFailed { + return r.handleTTL(ctx, df) + } + + cluster := &lll.EtcdCluster{} + if err := r.Get(ctx, types.NamespacedName{Namespace: df.Namespace, Name: df.Spec.ClusterRef.Name}, cluster); err != nil { + if apierrors.IsNotFound(err) { + return r.fail(ctx, df, "ClusterNotFound", + fmt.Sprintf("EtcdCluster %q not found in namespace %q", df.Spec.ClusterRef.Name, df.Namespace)) + } + return ctrl.Result{}, err + } + + // Serialize per cluster: only the oldest non-terminal EtcdDefrag targeting + // this cluster may act; the rest wait in Pending. + if active, err := r.oldestActive(ctx, df.Namespace, cluster.Name); err != nil { + return ctrl.Result{}, err + } else if active != "" && active != df.Name { + if setDefragCondition(df, metav1.ConditionFalse, "Queued", + fmt.Sprintf("waiting for EtcdDefrag %q to finish", active)) || df.Status.Phase != lll.EtcdDefragPhasePending { + df.Status.Phase = lll.EtcdDefragPhasePending + if err := r.Status().Update(ctx, df); err != nil { + return ctrl.Result{}, err + } + } + return ctrl.Result{RequeueAfter: defragRequeueAfter}, nil + } + + // Overall deadline: a run that can't make progress must fail, not linger. + // (CreationTimestamp is zero before the apiserver stamps it — e.g. in unit + // tests — so only enforce against a non-zero reference time.) + if started := df.Status.StartedAt; started != nil && time.Since(started.Time) > defragActiveDeadline { + return r.fail(ctx, df, "DeadlineExceeded", + fmt.Sprintf("defragmentation did not complete within %s", defragActiveDeadline)) + } else if started == nil && !df.CreationTimestamp.IsZero() && time.Since(df.CreationTimestamp.Time) > defragActiveDeadline { + return r.fail(ctx, df, "DeadlineExceeded", + fmt.Sprintf("defragmentation could not start within %s (cluster never became healthy)", defragActiveDeadline)) + } + + var memberList lll.EtcdMemberList + if err := r.List(ctx, &memberList, client.InNamespace(df.Namespace), + client.MatchingLabels{LabelCluster: cluster.Name}); err != nil { + return ctrl.Result{}, err + } + running := filterRunningMembers(memberList.Items) + + c, backends, err := r.dialAndProbe(ctx, cluster, running) + if err != nil { + logger.Error(err, "defrag: cannot dial/probe cluster; retrying") + return ctrl.Result{RequeueAfter: defragRequeueAfter}, nil + } + defer c.Close() + + // Health gate: defrag is only safe when every desired member is present, + // reachable, alarm-free and agrees on a leader. + if reason, ok := clusterDefragHealthy(cluster, running, backends); !ok { + msg := "a defragmentation is due but the cluster is not fully healthy; deferring to protect quorum: " + reason + if setDefragCondition(df, metav1.ConditionFalse, "ClusterNotHealthy", msg) { + r.event(df, corev1.EventTypeWarning, "DefragDeferred", msg) + } + df.Status.Phase = lll.EtcdDefragPhasePending + if err := r.Status().Update(ctx, df); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: defragRequeueAfter}, nil + } + + // Initialize the per-member work list once (followers first, leader last). + if len(df.Status.Members) == 0 { + df.Status.Members = plannedMembers(backends) + } + + next := firstPendingMember(df) + if next == nil { + return r.finalize(ctx, df) + } + + b := backendByName(backends, next.Name) + if b == nil { + // A planned member vanished between passes despite the health gate. + markMember(df, next.Name, lll.DefragOutcomeFailed, "MemberGone", nil) + return r.persistAndRequeue(ctx, df) + } + + if df.Status.Phase != lll.EtcdDefragPhaseRunning { + df.Status.Phase = lll.EtcdDefragPhaseRunning + now := metav1.Now() + df.Status.StartedAt = &now + setDefragCondition(df, metav1.ConditionTrue, "Running", "defragmenting members") + } + + quota := effectiveQuotaBytes(cluster) + if trig, reason := defragRuleTriggered(df.Spec.Rule, b.status.DbSize, b.status.DbSizeInUse, quota); !trig { + markMember(df, b.member.Name, lll.DefragOutcomeSkipped, reason, b) + return r.persistAndRequeue(ctx, df) + } + + dctx, cancel := context.WithTimeout(ctx, defragRPCTimeout) + defer cancel() + if _, derr := c.Defragment(dctx, b.endpoint); derr != nil { + msg := fmt.Sprintf("defragmentation of member %s failed: %v", b.member.Name, derr) + markMember(df, b.member.Name, lll.DefragOutcomeFailed, "RPCError", b) + r.event(df, corev1.EventTypeWarning, "DefragFailed", msg) + logger.Error(derr, "defrag: member Defragment failed", "member", b.member.Name) + return r.persistAndRequeue(ctx, df) + } + + // Read the post-defrag size for the record (best-effort). + if after, serr := statusWithTimeout(ctx, c, b.endpoint); serr == nil { + b.after = after.DbSize + } + markMember(df, b.member.Name, lll.DefragOutcomeDefragmented, "", b) + df.Status.Defragmented++ + logger.Info("defragmented member", "member", b.member.Name, "reclaimed", b.status.DbSize-b.after) + return r.persistAndRequeue(ctx, df) +} + +// memberBackend pairs a running member with its live Maintenance Status. +type memberBackend struct { + member *lll.EtcdMember + endpoint string + status *clientv3.StatusResponse + after int64 // DbSize after a successful defrag +} + +// dialAndProbe opens one client and reads every running member's Status. The +// caller closes the client. +func (r *EtcdDefragReconciler) dialAndProbe(ctx context.Context, cluster *lll.EtcdCluster, running []lll.EtcdMember) (EtcdClusterClient, []memberBackend, error) { + tlsCfg, err := buildOperatorTLSConfig(ctx, r.Client, cluster) + if err != nil { + return nil, nil, err + } + user, pass, _, err := resolveEtcdCredentials(ctx, r.Client, cluster) + if err != nil { + return nil, nil, err + } + scheme := clusterClientScheme(cluster) + endpoints := make([]string, len(running)) + for i := range running { + endpoints[i] = clientURL(scheme, running[i].Name, memberServiceName(&running[i]), cluster.Namespace) + } + c, err := r.EtcdClientFactory(ctx, endpoints, tlsCfg, user, pass) + if err != nil { + return nil, nil, err + } + backends := make([]memberBackend, 0, len(running)) + for i := range running { + resp, serr := statusWithTimeout(ctx, c, endpoints[i]) + if serr != nil { + // Unreachable member: record a nil-status backend so the health gate + // sees the gap. + backends = append(backends, memberBackend{member: &running[i], endpoint: endpoints[i]}) + continue + } + backends = append(backends, memberBackend{member: &running[i], endpoint: endpoints[i], status: resp, after: resp.DbSize}) + } + return c, backends, nil +} + +func statusWithTimeout(ctx context.Context, c EtcdClusterClient, endpoint string) (*clientv3.StatusResponse, error) { + sctx, cancel := context.WithTimeout(ctx, defragStatusTimeout) + defer cancel() + return c.Status(sctx, endpoint) +} + +// clusterDefragHealthy reports whether the cluster is safe to defragment: every +// desired member present and reachable, alarm-free, and agreeing on a single +// non-zero leader. Status is a local read — a member answers it while +// partitioned or alarmed — so agreement and Errors are checked, not just that +// the RPC returned. +func clusterDefragHealthy(cluster *lll.EtcdCluster, running []lll.EtcdMember, backends []memberBackend) (string, bool) { + desired := 0 + if cluster.Status.Observed != nil { + desired = int(cluster.Status.Observed.Replicas) + } + if desired == 0 || len(running) != desired { + return fmt.Sprintf("have %d running members, want %d", len(running), desired), false + } + var leader uint64 + for i := range backends { + b := &backends[i] + if b.status == nil { + return fmt.Sprintf("member %s is unreachable", b.member.Name), false + } + if len(b.status.Errors) > 0 { + return fmt.Sprintf("member %s reports alarms: %s", b.member.Name, strings.Join(b.status.Errors, ",")), false + } + if b.status.Leader == 0 { + return fmt.Sprintf("member %s reports no leader", b.member.Name), false + } + if leader == 0 { + leader = b.status.Leader + } else if b.status.Leader != leader { + return "members disagree on the leader", false + } + } + return "", true +} + +// plannedMembers builds the ordered work list: followers first, the leader +// last (its defrag is the most disruptive and is done only after followers +// prove defrag is healthy on this cluster). +func plannedMembers(backends []memberBackend) []lll.MemberDefragStatus { + followers := make([]lll.MemberDefragStatus, 0, len(backends)) + var leader []lll.MemberDefragStatus + for i := range backends { + b := &backends[i] + row := lll.MemberDefragStatus{Name: b.member.Name, Outcome: lll.DefragOutcomePending} + if isLeaderStatus(b.status) { + row.Role = lll.MemberRoleLeader + leader = append(leader, row) + } else { + row.Role = lll.MemberRoleFollower + followers = append(followers, row) + } + } + return append(followers, leader...) +} + +func isLeaderStatus(s *clientv3.StatusResponse) bool { + return s != nil && s.Header != nil && s.Leader != 0 && s.Leader == s.Header.MemberId +} + +func firstPendingMember(df *lll.EtcdDefrag) *lll.MemberDefragStatus { + for i := range df.Status.Members { + if df.Status.Members[i].Outcome == lll.DefragOutcomePending { + return &df.Status.Members[i] + } + } + return nil +} + +func backendByName(backends []memberBackend, name string) *memberBackend { + for i := range backends { + if backends[i].member.Name == name { + return &backends[i] + } + } + return nil +} + +func markMember(df *lll.EtcdDefrag, name string, outcome lll.DefragOutcome, reason string, b *memberBackend) { + for i := range df.Status.Members { + m := &df.Status.Members[i] + if m.Name != name { + continue + } + m.Outcome = outcome + m.Reason = reason + now := metav1.Now() + m.FinishedAt = &now + if b != nil && b.status != nil { + m.DBSizeBefore = b.status.DbSize + m.DBSizeAfter = b.after + if outcome == lll.DefragOutcomeDefragmented && b.status.DbSize >= b.after { + m.ReclaimedBytes = b.status.DbSize - b.after + } + } + return + } +} + +// effectiveQuotaBytes is the backend quota the cluster's members run with: the +// latched spec.options.quotaBackendBytes, or etcd's 2Gi default. +func effectiveQuotaBytes(cluster *lll.EtcdCluster) int64 { + if o := cluster.Status.Observed; o != nil && o.Options != nil && + o.Options.QuotaBackendBytes != nil && *o.Options.QuotaBackendBytes > 0 { + return *o.Options.QuotaBackendBytes + } + return defaultEtcdQuotaBytes +} + +// defragRuleTriggered reports whether a member's backend meets the rule, and a +// reason when it does not. rule.All is unconditional. Otherwise the free-space +// arm (reclaimable > freeSpaceAbove, default 200Mi) is always applied; the +// quota arm additionally fires under quota pressure but only when there is at +// least MinReclaim to reclaim — so a full-but-unfragmented backend is never +// defragmented for nothing. +func defragRuleTriggered(rule *lll.DefragRule, dbSize, dbSizeInUse, quota int64) (bool, string) { + if rule != nil && rule.All { + return true, "" + } + + free := defaultDefragFreeSpace + minReclaim := defaultDefragMinReclaim + quotaUsage := 0.0 + quotaSet := false + if rule != nil { + if rule.FreeSpaceAbove != nil { + free = rule.FreeSpaceAbove.Value() + } + if rule.MinReclaim != nil { + minReclaim = rule.MinReclaim.Value() + } + if p, ok := parsePercent(rule.QuotaUsageAbove); ok { + quotaUsage = p + quotaSet = true + } + } + + reclaimable := dbSize - dbSizeInUse + if reclaimable > free { + return true, "" + } + if quotaSet && quota > 0 && float64(dbSize) > quotaUsage*float64(quota) && reclaimable >= minReclaim { + return true, "" + } + return false, "BelowThreshold" +} + +// parsePercent parses "80%" into 0.80. Returns ok=false for anything outside +// 1–100 or non-numeric; the CRD pattern rejects such values at admission, so +// this is a defensive fallback. +func parsePercent(s string) (float64, bool) { + n, err := strconv.Atoi(strings.TrimSuffix(s, "%")) + if err != nil || n <= 0 || n > 100 { + return 0, false + } + return float64(n) / 100, true +} + +// oldestActive returns the name of the oldest non-terminal EtcdDefrag targeting +// clusterName in namespace (creationTimestamp, name tiebreak), or "" if none. +func (r *EtcdDefragReconciler) oldestActive(ctx context.Context, namespace, clusterName string) (string, error) { + var list lll.EtcdDefragList + if err := r.List(ctx, &list, client.InNamespace(namespace)); err != nil { + return "", err + } + active := make([]lll.EtcdDefrag, 0, len(list.Items)) + for _, d := range list.Items { + if d.Spec.ClusterRef.Name != clusterName { + continue + } + if d.Status.Phase == lll.EtcdDefragPhaseComplete || d.Status.Phase == lll.EtcdDefragPhaseFailed { + continue + } + active = append(active, d) + } + if len(active) == 0 { + return "", nil + } + sort.Slice(active, func(i, j int) bool { + if !active[i].CreationTimestamp.Equal(&active[j].CreationTimestamp) { + return active[i].CreationTimestamp.Before(&active[j].CreationTimestamp) + } + return active[i].Name < active[j].Name + }) + return active[0].Name, nil +} + +func (r *EtcdDefragReconciler) finalize(ctx context.Context, df *lll.EtcdDefrag) (ctrl.Result, error) { + phase := lll.EtcdDefragPhaseComplete + reason := "Complete" + for _, m := range df.Status.Members { + if m.Outcome == lll.DefragOutcomeFailed { + phase = lll.EtcdDefragPhaseFailed + reason = "MemberFailed" + break + } + } + df.Status.Phase = phase + now := metav1.Now() + df.Status.CompletedAt = &now + status := metav1.ConditionTrue + if phase == lll.EtcdDefragPhaseFailed { + status = metav1.ConditionFalse + } + setDefragCondition(df, status, reason, + fmt.Sprintf("%d/%d members defragmented", df.Status.Defragmented, len(df.Status.Members))) + if err := r.Status().Update(ctx, df); err != nil { + return ctrl.Result{}, err + } + return r.handleTTL(ctx, df) +} + +func (r *EtcdDefragReconciler) fail(ctx context.Context, df *lll.EtcdDefrag, reason, msg string) (ctrl.Result, error) { + df.Status.Phase = lll.EtcdDefragPhaseFailed + now := metav1.Now() + df.Status.CompletedAt = &now + if setDefragCondition(df, metav1.ConditionFalse, reason, msg) { + r.event(df, corev1.EventTypeWarning, reason, msg) + } + if err := r.Status().Update(ctx, df); err != nil { + return ctrl.Result{}, err + } + return r.handleTTL(ctx, df) +} + +func (r *EtcdDefragReconciler) persistAndRequeue(ctx context.Context, df *lll.EtcdDefrag) (ctrl.Result, error) { + if err := r.Status().Update(ctx, df); err != nil { + return ctrl.Result{}, err + } + return ctrl.Result{RequeueAfter: defragRequeueAfter}, nil +} + +// handleTTL deletes a finished EtcdDefrag once TTLSecondsAfterFinished has +// elapsed, or requeues to delete it later. No TTL means keep it as history. +func (r *EtcdDefragReconciler) handleTTL(ctx context.Context, df *lll.EtcdDefrag) (ctrl.Result, error) { + if df.Spec.TTLSecondsAfterFinished == nil || df.Status.CompletedAt == nil { + return ctrl.Result{}, nil + } + expiry := df.Status.CompletedAt.Add(time.Duration(*df.Spec.TTLSecondsAfterFinished) * time.Second) + if now := time.Now(); now.Before(expiry) { + return ctrl.Result{RequeueAfter: expiry.Sub(now)}, nil + } + if err := r.Delete(ctx, df); err != nil { + return ctrl.Result{}, client.IgnoreNotFound(err) + } + return ctrl.Result{}, nil +} + +func (r *EtcdDefragReconciler) event(obj client.Object, eventType, reason, msg string) { + if r.Recorder != nil { + r.Recorder.Event(obj, eventType, reason, msg) + } +} + +// setDefragCondition upserts the DefragChecked condition and reports whether it +// changed — so callers emit an event only on a real transition, not every pass. +func setDefragCondition(df *lll.EtcdDefrag, status metav1.ConditionStatus, reason, msg string) bool { + want := metav1.Condition{ + Type: "DefragChecked", + Status: status, + Reason: reason, + Message: msg, + ObservedGeneration: df.Generation, + } + for _, existing := range df.Status.Conditions { + if existing.Type == want.Type { + if existing.Status == want.Status && existing.Reason == want.Reason && + existing.Message == want.Message && existing.ObservedGeneration == want.ObservedGeneration { + return false + } + break + } + } + setCondition(&df.Status.Conditions, want) + return true +} + +func (r *EtcdDefragReconciler) SetupWithManager(mgr ctrl.Manager) error { + if r.EtcdClientFactory == nil { + r.EtcdClientFactory = DefaultEtcdClientFactory + } + return ctrl.NewControllerManagedBy(mgr). + For(&lll.EtcdDefrag{}). + Complete(r) +} diff --git a/controllers/etcddefrag_controller_test.go b/controllers/etcddefrag_controller_test.go new file mode 100644 index 00000000..6fb90084 --- /dev/null +++ b/controllers/etcddefrag_controller_test.go @@ -0,0 +1,311 @@ +/* +Copyright 2023 Timofey Larkin. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 +*/ + +package controllers + +import ( + "context" + "errors" + "fmt" + "strings" + "testing" + + etcdserverpb "go.etcd.io/etcd/api/v3/etcdserverpb" + clientv3 "go.etcd.io/etcd/client/v3" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/tools/record" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + + lll "github.com/cozystack/etcd-operator/api/v1alpha2" +) + +var gib = int64(1 << 30) + +func TestDefragRuleTriggered(t *testing.T) { + q := func(s string) *resource.Quantity { v := resource.MustParse(s); return &v } + quota := 2 * gib + cases := []struct { + name string + rule *lll.DefragRule + dbSize, dbSizeInUse int64 + want bool + }{ + {"nil rule, no fragmentation", nil, 100 << 20, 100 << 20, false}, + {"nil rule, fragmentation over 200Mi default", nil, 500 << 20, 100 << 20, true}, + {"all: unconditional even with nothing to reclaim", &lll.DefragRule{All: true}, 10 << 20, 10 << 20, true}, + {"freeSpaceAbove not met", &lll.DefragRule{FreeSpaceAbove: q("1Gi")}, 500 << 20, 100 << 20, false}, + {"freeSpaceAbove met", &lll.DefragRule{FreeSpaceAbove: q("200Mi")}, 500 << 20, 100 << 20, true}, + {"quota arm: full but unfragmented never fires", &lll.DefragRule{QuotaUsageAbove: "80%"}, int64(1.9 * float64(gib)), int64(1.9 * float64(gib)), false}, + {"quota arm: under pressure with reclaimable fires", &lll.DefragRule{QuotaUsageAbove: "80%", MinReclaim: q("32Mi")}, int64(1.9 * float64(gib)), int64(1.9*float64(gib)) - (64 << 20), true}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + got, _ := defragRuleTriggered(tc.rule, tc.dbSize, tc.dbSizeInUse, quota) + if got != tc.want { + t.Errorf("defragRuleTriggered = %v, want %v", got, tc.want) + } + }) + } +} + +// ── controller integration (fake etcd + fake kube client) ─────────────────── + +func defragEndpoint(name string) string { return fmt.Sprintf("http://%s.c1.ns.svc:2379", name) } + +func defragCluster3() (*lll.EtcdCluster, []lll.EtcdMember) { + cluster := &lll.EtcdCluster{ + ObjectMeta: metav1.ObjectMeta{Name: "c1", Namespace: "ns"}, + Spec: lll.EtcdClusterSpec{Replicas: ptrInt32(3), Version: "3.6.11", Storage: lll.StorageSpec{Size: resource.MustParse("1Gi")}}, + Status: lll.EtcdClusterStatus{ClusterID: "abc", Observed: &lll.ObservedClusterSpec{Replicas: 3, Version: "3.6.11", Storage: lll.StorageSpec{Size: resource.MustParse("1Gi")}}}, + } + var members []lll.EtcdMember + for i := 0; i < 3; i++ { + name := fmt.Sprintf("c1-%d", i) + members = append(members, lll.EtcdMember{ + ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "ns", Labels: memberLabels("c1", name)}, + Spec: lll.EtcdMemberSpec{ClusterName: "c1", Version: "3.6.11", Storage: lll.StorageSpec{Size: resource.MustParse("1Gi")}, InitialCluster: "x", ClusterToken: "t"}, + Status: lll.EtcdMemberStatus{PodName: name}, + }) + } + return cluster, members +} + +func status(memberID, leader uint64, dbSize, inUse int64) *clientv3.StatusResponse { + return &clientv3.StatusResponse{Header: &etcdserverpb.ResponseHeader{MemberId: memberID}, Leader: leader, DbSize: dbSize, DbSizeInUse: inUse} +} + +func objs(cluster *lll.EtcdCluster, members []lll.EtcdMember, dfs ...*lll.EtcdDefrag) []client.Object { + out := []client.Object{cluster} + for i := range members { + out = append(out, &members[i]) + } + for _, d := range dfs { + out = append(out, d) + } + return out +} + +// driveDefrag reconciles d1 until it reaches a terminal phase or the cap. +func driveDefrag(t *testing.T, r *EtcdDefragReconciler, c client.Client, name string) *lll.EtcdDefrag { + t.Helper() + ctx := context.Background() + req := ctrl.Request{NamespacedName: nn(name, "ns")} + for i := 0; i < 20; i++ { + if _, err := r.Reconcile(ctx, req); err != nil { + t.Fatalf("reconcile: %v", err) + } + df := &lll.EtcdDefrag{} + if err := c.Get(ctx, nn(name, "ns"), df); err != nil { + t.Fatalf("get %s: %v", name, err) + } + if df.Status.Phase == lll.EtcdDefragPhaseComplete || df.Status.Phase == lll.EtcdDefragPhaseFailed { + return df + } + } + df := &lll.EtcdDefrag{} + _ = c.Get(ctx, nn(name, "ns"), df) + return df +} + +func nn(name, ns string) client.ObjectKey { return client.ObjectKey{Name: name, Namespace: ns} } + +// A fragmented follower on a healthy cluster is defragmented; the leader and an +// unfragmented follower are skipped; the follower is done before the leader. +func TestEtcdDefrag_DefragmentsFragmentedMember(t *testing.T) { + cluster, members := defragCluster3() + df := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}}} + c, s := newTestClient(t, objs(cluster, members, df)...) + + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): status(10, 10, 100<<20, 100<<20), // leader, clean + defragEndpoint("c1-1"): status(11, 10, 500<<20, 100<<20), // follower, fragmented + defragEndpoint("c1-2"): status(12, 10, 100<<20, 100<<20), // follower, clean + } + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: record.NewFakeRecorder(20)} + + got := driveDefrag(t, r, c, "d1") + if got.Status.Phase != lll.EtcdDefragPhaseComplete { + t.Fatalf("phase = %q, want Complete", got.Status.Phase) + } + if len(fe.defragCalls) != 1 || fe.defragCalls[0] != defragEndpoint("c1-1") { + t.Fatalf("defragCalls = %v, want exactly [c1-1]", fe.defragCalls) + } + if got.Status.Defragmented != 1 { + t.Errorf("Defragmented = %d, want 1", got.Status.Defragmented) + } + // Members recorded; the fragmented follower Defragmented, leader last. + outcome := map[string]lll.DefragOutcome{} + for _, m := range got.Status.Members { + outcome[m.Name] = m.Outcome + } + if outcome["c1-1"] != lll.DefragOutcomeDefragmented { + t.Errorf("c1-1 outcome = %q, want Defragmented", outcome["c1-1"]) + } + if got.Status.Members[len(got.Status.Members)-1].Role != lll.MemberRoleLeader { + t.Errorf("leader not last in the plan: %+v", got.Status.Members) + } +} + +// rule.all defragments every member unconditionally. +func TestEtcdDefrag_RuleAllDefragmentsEveryone(t *testing.T) { + cluster, members := defragCluster3() + df := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}, Rule: &lll.DefragRule{All: true}}} + c, s := newTestClient(t, objs(cluster, members, df)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): status(10, 10, 50<<20, 50<<20), + defragEndpoint("c1-1"): status(11, 10, 50<<20, 50<<20), + defragEndpoint("c1-2"): status(12, 10, 50<<20, 50<<20), + } + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: record.NewFakeRecorder(20)} + + got := driveDefrag(t, r, c, "d1") + if got.Status.Phase != lll.EtcdDefragPhaseComplete || got.Status.Defragmented != 3 { + t.Fatalf("phase=%q defragmented=%d, want Complete/3", got.Status.Phase, got.Status.Defragmented) + } + if len(fe.defragCalls) != 3 { + t.Fatalf("defragCalls = %v, want all 3", fe.defragCalls) + } + // Leader defragmented last. + if fe.defragCalls[2] != defragEndpoint("c1-0") { + t.Errorf("leader c1-0 not defragmented last: %v", fe.defragCalls) + } +} + +// Quorum lost (two members unreachable) → no defrag, phase Pending, +// DefragChecked=False/ClusterNotHealthy, a DefragDeferred event. +func TestEtcdDefrag_DeferredWhenUnhealthy(t *testing.T) { + cluster, members := defragCluster3() + df := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}}} + c, s := newTestClient(t, objs(cluster, members, df)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): status(10, 10, 500<<20, 100<<20), + } + fe.statusErrByEndpoint = map[string]error{ + defragEndpoint("c1-1"): errors.New("context deadline exceeded"), + defragEndpoint("c1-2"): errors.New("context deadline exceeded"), + } + rec := record.NewFakeRecorder(20) + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: rec} + + ctx := context.Background() + res, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: nn("d1", "ns")}) + if err != nil { + t.Fatalf("reconcile: %v", err) + } + if res.RequeueAfter == 0 { + t.Errorf("expected a requeue while deferred, got %+v", res) + } + if len(fe.defragCalls) != 0 { + t.Fatalf("defrag ran on an unhealthy cluster: %v", fe.defragCalls) + } + got := &lll.EtcdDefrag{} + if err := c.Get(ctx, nn("d1", "ns"), got); err != nil { + t.Fatal(err) + } + if got.Status.Phase != lll.EtcdDefragPhasePending { + t.Errorf("phase = %q, want Pending", got.Status.Phase) + } + cond := findDefragCond(got) + if cond == nil || cond.Status != metav1.ConditionFalse || cond.Reason != "ClusterNotHealthy" { + t.Errorf("DefragChecked = %+v, want False/ClusterNotHealthy", cond) + } + assertDefragEvent(t, rec, "DefragDeferred") +} + +// A failed Defragment RPC lands the run in Failed with the member marked Failed. +func TestEtcdDefrag_FailedRPC(t *testing.T) { + cluster, members := defragCluster3() + df := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}, Rule: &lll.DefragRule{All: true}}} + c, s := newTestClient(t, objs(cluster, members, df)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.defragErr = errors.New("boom") + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): status(10, 10, 50<<20, 50<<20), + defragEndpoint("c1-1"): status(11, 10, 50<<20, 50<<20), + defragEndpoint("c1-2"): status(12, 10, 50<<20, 50<<20), + } + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: record.NewFakeRecorder(20)} + + got := driveDefrag(t, r, c, "d1") + if got.Status.Phase != lll.EtcdDefragPhaseFailed { + t.Fatalf("phase = %q, want Failed", got.Status.Phase) + } + failed := false + for _, m := range got.Status.Members { + if m.Outcome == lll.DefragOutcomeFailed { + failed = true + } + } + if !failed { + t.Errorf("no member marked Failed: %+v", got.Status.Members) + } +} + +// Two EtcdDefrags for one cluster are serialized: only the oldest acts; the +// newer stays Pending and performs no defrag until the first finishes. +func TestEtcdDefrag_SerializedPerCluster(t *testing.T) { + cluster, members := defragCluster3() + older := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d-old", Namespace: "ns", CreationTimestamp: metav1.Unix(100, 0)}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}, Rule: &lll.DefragRule{All: true}}} + newer := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d-new", Namespace: "ns", CreationTimestamp: metav1.Unix(200, 0)}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}, Rule: &lll.DefragRule{All: true}}} + c, s := newTestClient(t, objs(cluster, members, older, newer)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): status(10, 10, 50<<20, 50<<20), + defragEndpoint("c1-1"): status(11, 10, 50<<20, 50<<20), + defragEndpoint("c1-2"): status(12, 10, 50<<20, 50<<20), + } + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: record.NewFakeRecorder(20)} + + // Reconcile the NEWER one: it must stay Pending and not defrag. + ctx := context.Background() + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: nn("d-new", "ns")}); err != nil { + t.Fatalf("reconcile d-new: %v", err) + } + if len(fe.defragCalls) != 0 { + t.Fatalf("newer EtcdDefrag defragged while an older one is active: %v", fe.defragCalls) + } + got := &lll.EtcdDefrag{} + _ = c.Get(ctx, nn("d-new", "ns"), got) + if got.Status.Phase != lll.EtcdDefragPhasePending { + t.Errorf("d-new phase = %q, want Pending (queued)", got.Status.Phase) + } +} + +func findDefragCond(df *lll.EtcdDefrag) *metav1.Condition { + for i := range df.Status.Conditions { + if df.Status.Conditions[i].Type == "DefragChecked" { + return &df.Status.Conditions[i] + } + } + return nil +} + +func assertDefragEvent(t *testing.T, rec *record.FakeRecorder, wantReason string) { + t.Helper() + select { + case ev := <-rec.Events: + if !strings.Contains(ev, wantReason) { + t.Errorf("event = %q, want one mentioning %q", ev, wantReason) + } + default: + t.Errorf("no event emitted, want one mentioning %q", wantReason) + } +} diff --git a/controllers/testing_helpers_test.go b/controllers/testing_helpers_test.go index 5c63020a..b1c621db 100644 --- a/controllers/testing_helpers_test.go +++ b/controllers/testing_helpers_test.go @@ -48,6 +48,17 @@ type fakeEtcd struct { statusVersion string statusErr error + // Defrag surface. statusByEndpoint overrides Status per endpoint (backend + // sizes, leader); statusErrByEndpoint fails Status for a specific endpoint + // (an unhealthy member); leader is the default StatusResponse.Leader. + // defragCalls records each Defragment endpoint in call order; defragErr, + // when set, fails Defragment. + statusByEndpoint map[string]*clientv3.StatusResponse + statusErrByEndpoint map[string]error + leader uint64 + defragCalls []string + defragErr error + addCalls []string addLearnerCalls []string promoteCalls []uint64 @@ -181,16 +192,34 @@ func (f *fakeEtcd) UserGrantRole(_ context.Context, user, role string) (*clientv return &clientv3.AuthUserGrantRoleResponse{Header: &etcdserverpb.ResponseHeader{ClusterId: f.clusterID}}, nil } -func (f *fakeEtcd) Status(_ context.Context, _ string) (*clientv3.StatusResponse, error) { +func (f *fakeEtcd) Status(_ context.Context, endpoint string) (*clientv3.StatusResponse, error) { if f.statusErr != nil { return nil, f.statusErr } + if err := f.statusErrByEndpoint[endpoint]; err != nil { + return nil, err + } + if resp, ok := f.statusByEndpoint[endpoint]; ok { + if resp.Header == nil { + resp.Header = &etcdserverpb.ResponseHeader{ClusterId: f.clusterID} + } + return resp, nil + } return &clientv3.StatusResponse{ Header: &etcdserverpb.ResponseHeader{ClusterId: f.clusterID}, Version: f.statusVersion, + Leader: f.leader, }, nil } +func (f *fakeEtcd) Defragment(_ context.Context, endpoint string) (*clientv3.DefragmentResponse, error) { + f.defragCalls = append(f.defragCalls, endpoint) + if f.defragErr != nil { + return nil, f.defragErr + } + return &clientv3.DefragmentResponse{Header: &etcdserverpb.ResponseHeader{ClusterId: f.clusterID}}, nil +} + func (f *fakeEtcd) Close() error { f.closed = true; return nil } func factoryReturning(c EtcdClusterClient) EtcdClientFactory { @@ -251,7 +280,7 @@ func newTestClient(t *testing.T, objs ...client.Object) (client.Client, *runtime c := fake.NewClientBuilder(). WithScheme(s). WithObjects(objs...). - WithStatusSubresource(&lll.EtcdCluster{}, &lll.EtcdMember{}, &lll.EtcdSnapshot{}). + WithStatusSubresource(&lll.EtcdCluster{}, &lll.EtcdMember{}, &lll.EtcdSnapshot{}, &lll.EtcdDefrag{}). Build() return c, s } diff --git a/docs/etcd-defrag.md b/docs/etcd-defrag.md index 60cd5a50..a976e125 100644 --- a/docs/etcd-defrag.md +++ b/docs/etcd-defrag.md @@ -1,12 +1,5 @@ # Defragmentation (`EtcdDefrag`) -> **Status:** this ships the `EtcdDefrag` **API type** ahead of its reconciling -> controller. Until that controller lands the resource is **inert** — creating -> one records intent but nothing acts on it (no sweep runs, `status` stays -> empty, `ttlSecondsAfterFinished` does not fire). The "Safety model", -> "Timeouts and retries" and `status` sections below describe the **controller -> contract the follow-up implements**, not behaviour that exists today. - etcd never reclaims backend disk on its own: compaction frees pages logically, but the file — and the space counted against `--quota-backend-bytes` — stays allocated until a **defragment** returns it. `EtcdDefrag` is how you ask the @@ -66,7 +59,7 @@ spec: minReclaim: 32Mi # … but never a no-op defrag ``` -Inspect progress and history (once the controller exists): +Inspect progress and history: ```sh kubectl get etcddefrag.etcd-operator.cozystack.io -n team-a @@ -101,7 +94,7 @@ out. > `autoCompactionRetention`) has `DbSizeInUse ≈ DbSize` and little to reclaim. > Set auto-compaction if you rely on defrag to hold the backend down. -## Safety model (planned controller behaviour) +## Safety model - **One member at a time, followers before the leader**, only while the whole cluster is healthy. A defrag due on a not-fully-healthy cluster is **deferred** @@ -113,7 +106,7 @@ out. local status read while partitioned, alarmed (`NOSPACE`/`CORRUPT`), or behind in raft, so those are checked before acting. -## Status (planned controller behaviour) +## Status `status.phase` moves `Pending → Running → Complete | Failed`; a `Pending` run waiting on cluster health carries a condition saying so. `status.members[]` @@ -121,7 +114,7 @@ records, per member (keyed by name), the role at processing time, the outcome (`Skipped` / `Defragmented` / `Failed`), the before/after `DbSize`, and the bytes reclaimed — the run's full history, not a single rolled-up condition. -## Timeouts and retries (planned controller behaviour) +## Timeouts and retries Following [`EtcdSnapshot`](concepts.md#snapshots--restore) — where the Job's deadlines are controller constants and terminal phases are sticky — this needs diff --git a/main.go b/main.go index 6d961ce3..8a1d18a4 100644 --- a/main.go +++ b/main.go @@ -252,6 +252,14 @@ func main() { setupLog.Error(err, "unable to create controller", "controller", "EtcdSnapshot") os.Exit(1) } + if err = (&controllers.EtcdDefragReconciler{ + Client: mgr.GetClient(), + Scheme: mgr.GetScheme(), + Recorder: mgr.GetEventRecorderFor("etcd-operator"), + }).SetupWithManager(mgr); err != nil { + setupLog.Error(err, "unable to create controller", "controller", "EtcdDefrag") + os.Exit(1) + } //+kubebuilder:scaffold:builder if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil { diff --git a/test/e2e/defrag_test.go b/test/e2e/defrag_test.go new file mode 100644 index 00000000..cef5ae4d --- /dev/null +++ b/test/e2e/defrag_test.go @@ -0,0 +1,211 @@ +//go:build e2e + +package e2e + +import ( + "context" + "encoding/json" + "fmt" + "strings" + "testing" + "time" + + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" + + etcdv1alpha2 "github.com/cozystack/etcd-operator/api/v1alpha2" +) + +// The cluster name is shared; each test gets its OWN namespace so one test's +// namespace teardown (which is asynchronous — the namespace lingers in +// Terminating) can't block the next test from creating content in it. +const defragCluster = "etcd" + +// TestEtcdDefragReclaimsSpace proves the EtcdDefrag controller end to end on a +// real cluster: a member accrues reclaimable free space (write a few MB, delete +// it, compact — which frees pages logically but leaves the file allocated), an +// EtcdDefrag is created, and the controller defragments it so the physical +// DbSize shrinks and the run reaches phase Complete. +func TestEtcdDefragReclaimsSpace(t *testing.T) { + ctx := context.Background() + ns := "defrag-reclaim-e2e" + createDefragNamespace(ctx, t, ns) + + if err := kube.Create(ctx, defragClusterObject(ns)); err != nil { + t.Fatalf("create EtcdCluster: %v", err) + } + waitFor(ctx, t, 5*time.Minute, "cluster Available", etcdClusterAvailable(ns, defragCluster)) + waitFor(ctx, t, 2*time.Minute, "3 members ready", readyMembersIsIn(ns, defragCluster, 3)) + + pod := aReadyMemberPod(ctx, t, ns) + fragmentEtcd(ctx, t, ns, pod) + frag := endpointDBSize(ctx, t, ns, pod) + t.Logf("db after fragmenting: size=%d inUse=%d free=%d", frag.dbSize, frag.dbSizeInUse, frag.dbSize-frag.dbSizeInUse) + if frag.dbSize-frag.dbSizeInUse < 1<<20 { + t.Fatalf("expected >1Mi reclaimable free space after fragmenting, got %d", frag.dbSize-frag.dbSizeInUse) + } + + // Ask for an unconditional defrag now. + createEtcdDefrag(ctx, t, ns, "defrag-now", &etcdv1alpha2.DefragRule{All: true}) + + waitFor(ctx, t, 3*time.Minute, "EtcdDefrag Complete", etcdDefragPhaseIs(ns, "defrag-now", etcdv1alpha2.EtcdDefragPhaseComplete)) + waitFor(ctx, t, 2*time.Minute, "physical DbSize reclaimed", func(ctx context.Context) error { + now := endpointDBSize(ctx, t, ns, pod) + if now.dbSize >= frag.dbSize { + return fmt.Errorf("dbSize not reclaimed: was %d, still %d", frag.dbSize, now.dbSize) + } + return nil + }) +} + +// The "defer-not-force while the cluster is unhealthy" invariant is covered +// deterministically by the controller unit test +// TestEtcdDefrag_DeferredWhenUnhealthy — it asserts no Defragment call plus the +// DefragDeferred event and DefragChecked=False/ClusterNotHealthy. It has no +// reliable e2e counterpart: deleting member Pods to break quorum races the +// operator, which recreates the PVC-backed Pods and heals faster than a test +// window can observe, so an e2e negative-window assertion is inherently flaky. + +// ── helpers ───────────────────────────────────────────────────────────────── + +func createDefragNamespace(ctx context.Context, t *testing.T, ns string) { + t.Helper() + nsObj := &corev1.Namespace{ + TypeMeta: metav1.TypeMeta{APIVersion: "v1", Kind: "Namespace"}, + ObjectMeta: metav1.ObjectMeta{Name: ns}, + } + if err := kube.Patch(ctx, nsObj, client.Apply, fieldOwner, client.ForceOwnership); err != nil { + t.Fatalf("create namespace %s: %v", ns, err) + } + t.Cleanup(func() { + _ = kube.Delete(context.Background(), &corev1.Namespace{ObjectMeta: metav1.ObjectMeta{Name: ns}}) + }) +} + +func defragClusterObject(ns string) *etcdv1alpha2.EtcdCluster { + three := int32(3) + return &etcdv1alpha2.EtcdCluster{ + ObjectMeta: metav1.ObjectMeta{Name: defragCluster, Namespace: ns}, + Spec: etcdv1alpha2.EtcdClusterSpec{ + Replicas: &three, + Version: "3.6.11", + Storage: etcdv1alpha2.StorageSpec{Size: resource.MustParse("1Gi")}, + }, + } +} + +func createEtcdDefrag(ctx context.Context, t *testing.T, ns, name string, rule *etcdv1alpha2.DefragRule) { + t.Helper() + d := &etcdv1alpha2.EtcdDefrag{ + ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: ns}, + Spec: etcdv1alpha2.EtcdDefragSpec{ + ClusterRef: corev1.LocalObjectReference{Name: defragCluster}, + Rule: rule, + }, + } + if err := kube.Create(ctx, d); err != nil { + t.Fatalf("create EtcdDefrag %s: %v", name, err) + } +} + +func etcdDefragPhaseIs(ns, name string, phase etcdv1alpha2.EtcdDefragPhase) func(context.Context) error { + return func(ctx context.Context) error { + d := &etcdv1alpha2.EtcdDefrag{} + if err := kube.Get(ctx, client.ObjectKey{Namespace: ns, Name: name}, d); err != nil { + return err + } + if d.Status.Phase != phase { + return fmt.Errorf("EtcdDefrag %s phase=%q, want %q", name, d.Status.Phase, phase) + } + return nil + } +} + +func defragMemberNames(ctx context.Context, t *testing.T, ns string) []string { + t.Helper() + list := &etcdv1alpha2.EtcdMemberList{} + if err := kube.List(ctx, list, client.InNamespace(ns), + client.MatchingLabels{"etcd-operator.cozystack.io/cluster": defragCluster}); err != nil { + t.Fatalf("list members: %v", err) + } + names := make([]string, 0, len(list.Items)) + for i := range list.Items { + names = append(names, list.Items[i].Name) + } + return names +} + +func aReadyMemberPod(ctx context.Context, t *testing.T, ns string) string { + t.Helper() + for _, name := range defragMemberNames(ctx, t, ns) { + p, err := clientset.CoreV1().Pods(ns).Get(ctx, name, metav1.GetOptions{}) + if err != nil { + continue + } + for _, cs := range p.Status.ContainerStatuses { + if cs.Name == "etcd" && cs.Ready { + return name + } + } + } + t.Fatal("no ready member pod found") + return "" +} + +type dbStat struct { + dbSize int64 + dbSizeInUse int64 + revision int64 +} + +// endpointDBSize reads the member's own backend sizes via `etcdctl endpoint +// status -w json`. +func endpointDBSize(ctx context.Context, t *testing.T, ns, pod string) dbStat { + t.Helper() + out, stderr, err := podExec(ctx, ns, pod, "etcd", + []string{"etcdctl", "endpoint", "status", "-w", "json"}) + if err != nil { + t.Fatalf("endpoint status on %s: %v (stderr: %s)", pod, err, stderr) + } + var rows []struct { + Status struct { + DbSize int64 `json:"dbSize"` + DbSizeInUse int64 `json:"dbSizeInUse"` + Header struct { + Revision int64 `json:"revision"` + } `json:"header"` + } `json:"Status"` + } + if err := json.Unmarshal([]byte(out), &rows); err != nil || len(rows) == 0 { + t.Fatalf("parse endpoint status %q: %v", out, err) + } + return dbStat{rows[0].Status.DbSize, rows[0].Status.DbSizeInUse, rows[0].Status.Header.Revision} +} + +// fragmentEtcd creates reclaimable free space: write a few MB across many keys, +// delete them all, then compact (which frees the pages logically but leaves the +// file allocated — exactly what defrag reclaims). etcdctl runs one arg-only +// command per call (the etcd image is distroless, no shell), so the write loop +// is driven from Go. +func fragmentEtcd(ctx context.Context, t *testing.T, ns, pod string) { + t.Helper() + value := strings.Repeat("x", 8<<10) // 8Ki per key + const keys = 512 // ~4Mi of data + for i := 0; i < keys; i++ { + if _, stderr, err := podExec(ctx, ns, pod, "etcd", + []string{"etcdctl", "put", fmt.Sprintf("frag/%04d", i), value}); err != nil { + t.Fatalf("etcdctl put %d: %v (stderr: %s)", i, err, stderr) + } + } + if _, stderr, err := podExec(ctx, ns, pod, "etcd", + []string{"etcdctl", "del", "frag/", "--prefix"}); err != nil { + t.Fatalf("etcdctl del: %v (stderr: %s)", err, stderr) + } + rev := endpointDBSize(ctx, t, ns, pod).revision + if _, stderr, err := podExec(ctx, ns, pod, "etcd", + []string{"etcdctl", "compact", fmt.Sprintf("%d", rev), "--physical"}); err != nil { + t.Fatalf("etcdctl compact: %v (stderr: %s)", err, stderr) + } +} From c1c0402e3d92235e5bf0178e73fa15fa8e292f20 Mon Sep 17 00:00:00 2001 From: Andrey Kolkov Date: Thu, 20 Aug 2026 11:45:30 +0400 Subject: [PATCH 2/4] fix(defrag): admit NOSPACE, hold the run deadline, record reclaim honestly MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses the review on the EtcdDefrag controller. Blocking: - A NOSPACE alarm no longer blocks the run. The health gate refused any member reporting a non-empty status.Errors, which is where etcd reports active alarms — so a cluster that had hit its backend quota, the one case this feature exists to resolve, was the one cluster the controller would not touch. The gate now discriminates: NOSPACE is admitted, CORRUPT (and any other health error) still defers. After a sweep that reclaimed space, the controller disarms the NOSPACE alarm (AlarmList/AlarmDisarm added to the client interface) so the cluster leaves read-only. - status.startedAt is stamped exactly once, on the first Running transition, instead of on every Pending->Running edge. The active-deadline is measured from it; re-stamping on a cluster that flapped between health-gate blips reset the deadline every pass, so a stuck run could hold the per-cluster serialization slot forever. Non-blocking: - A failed post-defrag Status read leaves dbSizeAfter/reclaimedBytes unset (reason AfterSizeUnavailable) rather than pre-seeding the after-size to the before-size and reporting a real defrag as reclaiming zero. - A run keeps phase Running across a mid-sweep health-gate flap (the condition carries the reason) rather than flipping back to Pending with partial results. - MaxConcurrentReconciles set to 4 so one wedged member's Defragment RPC no longer stalls every other cluster; per-cluster serialization is still enforced by oldestActive. - docs: drop the "defrag that doesn't shrink DbSize is backed off" line — the one-shot controller does not do this. Tests: NOSPACE-admits-and-disarms, CORRUPT-blocks, deadline-not-reset-on-flap, and after-size-unknown, all against the fake-etcd harness. Assisted-By: Claude Opus 4.8 (1M context) Signed-off-by: Andrey Kolkov --- controllers/etcd_client.go | 9 ++ controllers/etcddefrag_controller.go | 137 ++++++++++++++++++---- controllers/etcddefrag_controller_test.go | 131 +++++++++++++++++++++ controllers/testing_helpers_test.go | 21 ++++ docs/etcd-defrag.md | 3 +- 5 files changed, 275 insertions(+), 26 deletions(-) diff --git a/controllers/etcd_client.go b/controllers/etcd_client.go index c4de7e05..1ddd9cff 100644 --- a/controllers/etcd_client.go +++ b/controllers/etcd_client.go @@ -50,6 +50,15 @@ type EtcdClusterClient interface { // its embedded Maintenance interface. Defragment(ctx context.Context, endpoint string) (*clientv3.DefragmentResponse, error) + // AlarmList and AlarmDisarm let the EtcdDefrag controller close the NOSPACE + // loop: a cluster at its backend quota raises NOSPACE and goes read-only, + // which is the case defrag exists to relieve, so the health gate permits the + // run — but the alarm stays armed after the space is reclaimed until it is + // explicitly disarmed. A CORRUPT alarm, by contrast, blocks the run. + // *clientv3.Client satisfies both via its embedded Maintenance interface. + AlarmList(ctx context.Context) (*clientv3.AlarmResponse, error) + AlarmDisarm(ctx context.Context, m *clientv3.AlarmMember) (*clientv3.AlarmResponse, error) + // Auth surface — used by reconcileAuth to provision the single root // user/role and turn on authentication. The "root" role is built into // etcd, so a RoleAdd is not needed: UserAdd("root", …) + diff --git a/controllers/etcddefrag_controller.go b/controllers/etcddefrag_controller.go index 0a4dc98f..acf56ee4 100644 --- a/controllers/etcddefrag_controller.go +++ b/controllers/etcddefrag_controller.go @@ -18,6 +18,7 @@ import ( "strings" "time" + etcdserverpb "go.etcd.io/etcd/api/v3/etcdserverpb" clientv3 "go.etcd.io/etcd/client/v3" corev1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" @@ -27,6 +28,7 @@ import ( "k8s.io/client-go/tools/record" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/controller" "sigs.k8s.io/controller-runtime/pkg/log" lll "github.com/cozystack/etcd-operator/api/v1alpha2" @@ -53,6 +55,10 @@ const ( // so a run stuck on an unhealthy cluster can't hold the per-cluster slot // forever. defragActiveDeadline = 30 * time.Minute + + // defragMaxConcurrentReconciles lets distinct clusters' runs proceed in + // parallel (per-cluster serialization is enforced separately by oldestActive). + defragMaxConcurrentReconciles = 4 ) // EtcdDefragReconciler drives an EtcdDefrag: a one-shot, run-to-completion @@ -141,13 +147,20 @@ func (r *EtcdDefragReconciler) Reconcile(ctx context.Context, req ctrl.Request) defer c.Close() // Health gate: defrag is only safe when every desired member is present, - // reachable, alarm-free and agrees on a leader. + // reachable, free of blocking alarms and agreeing on a leader. if reason, ok := clusterDefragHealthy(cluster, running, backends); !ok { msg := "a defragmentation is due but the cluster is not fully healthy; deferring to protect quorum: " + reason if setDefragCondition(df, metav1.ConditionFalse, "ClusterNotHealthy", msg) { r.event(df, corev1.EventTypeWarning, "DefragDeferred", msg) } - df.Status.Phase = lll.EtcdDefragPhasePending + // Keep a run that has already started (partial results in + // status.members) in Running and let the condition carry the reason; + // only a run that never started falls back to Pending. StartedAt is what + // the active-deadline is measured against, so it must not be re-stamped + // by such a flap — see the Running transition below. + if df.Status.StartedAt == nil { + df.Status.Phase = lll.EtcdDefragPhasePending + } if err := r.Status().Update(ctx, df); err != nil { return ctrl.Result{}, err } @@ -161,7 +174,7 @@ func (r *EtcdDefragReconciler) Reconcile(ctx context.Context, req ctrl.Request) next := firstPendingMember(df) if next == nil { - return r.finalize(ctx, df) + return r.finalize(ctx, df, c) } b := backendByName(backends, next.Name) @@ -171,12 +184,19 @@ func (r *EtcdDefragReconciler) Reconcile(ctx context.Context, req ctrl.Request) return r.persistAndRequeue(ctx, df) } - if df.Status.Phase != lll.EtcdDefragPhaseRunning { - df.Status.Phase = lll.EtcdDefragPhaseRunning + df.Status.Phase = lll.EtcdDefragPhaseRunning + // Stamp StartedAt exactly once, on the first Running transition. The + // active-deadline is measured from it, so re-stamping on a Pending->Running + // flap (a cluster that recovers between health-gate blips) would keep + // resetting the deadline and a stuck run could hold the per-cluster slot + // forever. + if df.Status.StartedAt == nil { now := metav1.Now() df.Status.StartedAt = &now - setDefragCondition(df, metav1.ConditionTrue, "Running", "defragmenting members") } + // Refresh the condition each active pass so a run that resumes after a + // health-gate flap does not stay reading ClusterNotHealthy. + setDefragCondition(df, metav1.ConditionTrue, "Running", "defragmenting members") quota := effectiveQuotaBytes(cluster) if trig, reason := defragRuleTriggered(df.Spec.Rule, b.status.DbSize, b.status.DbSizeInUse, quota); !trig { @@ -194,22 +214,32 @@ func (r *EtcdDefragReconciler) Reconcile(ctx context.Context, req ctrl.Request) return r.persistAndRequeue(ctx, df) } - // Read the post-defrag size for the record (best-effort). + // Read the post-defrag size for the record. When the read fails, leave the + // after-size unset rather than defaulting it to the before-size — reporting + // a successful defrag as reclaiming zero would be silently wrong in the very + // field the per-member status exists to provide. + outcomeReason := "" if after, serr := statusWithTimeout(ctx, c, b.endpoint); serr == nil { b.after = after.DbSize + b.afterKnown = true + } else { + outcomeReason = "AfterSizeUnavailable" + logger.Info("defrag: post-defrag Status read failed; reclaimed size unrecorded", + "member", b.member.Name, "err", serr) } - markMember(df, b.member.Name, lll.DefragOutcomeDefragmented, "", b) + markMember(df, b.member.Name, lll.DefragOutcomeDefragmented, outcomeReason, b) df.Status.Defragmented++ - logger.Info("defragmented member", "member", b.member.Name, "reclaimed", b.status.DbSize-b.after) + logger.Info("defragmented member", "member", b.member.Name) return r.persistAndRequeue(ctx, df) } // memberBackend pairs a running member with its live Maintenance Status. type memberBackend struct { - member *lll.EtcdMember - endpoint string - status *clientv3.StatusResponse - after int64 // DbSize after a successful defrag + member *lll.EtcdMember + endpoint string + status *clientv3.StatusResponse + after int64 // DbSize after a successful defrag; valid only if afterKnown + afterKnown bool // whether the post-defrag Status read succeeded } // dialAndProbe opens one client and reads every running member's Status. The @@ -241,7 +271,7 @@ func (r *EtcdDefragReconciler) dialAndProbe(ctx context.Context, cluster *lll.Et backends = append(backends, memberBackend{member: &running[i], endpoint: endpoints[i]}) continue } - backends = append(backends, memberBackend{member: &running[i], endpoint: endpoints[i], status: resp, after: resp.DbSize}) + backends = append(backends, memberBackend{member: &running[i], endpoint: endpoints[i], status: resp}) } return c, backends, nil } @@ -252,11 +282,41 @@ func statusWithTimeout(ctx context.Context, c EtcdClusterClient, endpoint string return c.Status(sctx, endpoint) } +// disarmNoSpaceAlarms clears any armed NOSPACE alarm. The health gate lets a +// NOSPACE cluster through so the run can reclaim its backend space; etcd keeps +// the alarm armed — and the cluster read-only — until it is explicitly disarmed. +// Best-effort: etcd re-arms on the next write if the space was not actually +// freed, so a failure here is logged by the caller rather than failing the run. +func (r *EtcdDefragReconciler) disarmNoSpaceAlarms(ctx context.Context, c EtcdClusterClient) error { + lctx, cancel := context.WithTimeout(ctx, defragStatusTimeout) + defer cancel() + resp, err := c.AlarmList(lctx) + if err != nil { + return err + } + for _, a := range resp.Alarms { + if a == nil || a.Alarm != etcdserverpb.AlarmType_NOSPACE { + continue + } + dctx, dcancel := context.WithTimeout(ctx, defragStatusTimeout) + _, derr := c.AlarmDisarm(dctx, (*clientv3.AlarmMember)(a)) + dcancel() + if derr != nil { + return derr + } + } + return nil +} + // clusterDefragHealthy reports whether the cluster is safe to defragment: every -// desired member present and reachable, alarm-free, and agreeing on a single -// non-zero leader. Status is a local read — a member answers it while -// partitioned or alarmed — so agreement and Errors are checked, not just that -// the RPC returned. +// desired member present and reachable, free of blocking alarms, and agreeing on +// a single non-zero leader. Status is a local read — a member answers it while +// partitioned or alarmed — so agreement and reported alarms are checked, not just +// that the RPC returned. +// +// A NOSPACE alarm does not block: a backend at its quota is exactly what a defrag +// is meant to relieve, and refusing it would leave the one cluster this feature +// exists for read-only forever. Every other alarm (notably CORRUPT) blocks. func clusterDefragHealthy(cluster *lll.EtcdCluster, running []lll.EtcdMember, backends []memberBackend) (string, bool) { desired := 0 if cluster.Status.Observed != nil { @@ -271,8 +331,8 @@ func clusterDefragHealthy(cluster *lll.EtcdCluster, running []lll.EtcdMember, ba if b.status == nil { return fmt.Sprintf("member %s is unreachable", b.member.Name), false } - if len(b.status.Errors) > 0 { - return fmt.Sprintf("member %s reports alarms: %s", b.member.Name, strings.Join(b.status.Errors, ",")), false + if e, blocking := blockingStatusError(b.status.Errors); blocking { + return fmt.Sprintf("member %s reports: %s", b.member.Name, e), false } if b.status.Leader == 0 { return fmt.Sprintf("member %s reports no leader", b.member.Name), false @@ -286,6 +346,21 @@ func clusterDefragHealthy(cluster *lll.EtcdCluster, running []lll.EtcdMember, ba return "", true } +// blockingStatusError returns the first StatusResponse.Errors entry that should +// block a defragmentation, and whether one exists. etcd reports active alarms +// there as "memberID: alarm:"; a NOSPACE line is the reclaimable-space +// case defrag relieves and does not block, while anything else (a CORRUPT alarm +// or any other health string) does. +func blockingStatusError(errs []string) (string, bool) { + for _, e := range errs { + if strings.Contains(e, "alarm:"+etcdserverpb.AlarmType_NOSPACE.String()) { + continue + } + return e, true + } + return "", false +} + // plannedMembers builds the ordered work list: followers first, the leader // last (its defrag is the most disruptive and is done only after followers // prove defrag is healthy on this cluster). @@ -340,9 +415,11 @@ func markMember(df *lll.EtcdDefrag, name string, outcome lll.DefragOutcome, reas m.FinishedAt = &now if b != nil && b.status != nil { m.DBSizeBefore = b.status.DbSize - m.DBSizeAfter = b.after - if outcome == lll.DefragOutcomeDefragmented && b.status.DbSize >= b.after { - m.ReclaimedBytes = b.status.DbSize - b.after + if b.afterKnown { + m.DBSizeAfter = b.after + if outcome == lll.DefragOutcomeDefragmented && b.status.DbSize >= b.after { + m.ReclaimedBytes = b.status.DbSize - b.after + } } } return @@ -437,7 +514,7 @@ func (r *EtcdDefragReconciler) oldestActive(ctx context.Context, namespace, clus return active[0].Name, nil } -func (r *EtcdDefragReconciler) finalize(ctx context.Context, df *lll.EtcdDefrag) (ctrl.Result, error) { +func (r *EtcdDefragReconciler) finalize(ctx context.Context, df *lll.EtcdDefrag, c EtcdClusterClient) (ctrl.Result, error) { phase := lll.EtcdDefragPhaseComplete reason := "Complete" for _, m := range df.Status.Members { @@ -447,6 +524,13 @@ func (r *EtcdDefragReconciler) finalize(ctx context.Context, df *lll.EtcdDefrag) break } } + // A cluster admitted with a NOSPACE alarm stays read-only until the alarm is + // disarmed; do it once the sweep has actually reclaimed space. + if phase == lll.EtcdDefragPhaseComplete && df.Status.Defragmented > 0 { + if err := r.disarmNoSpaceAlarms(ctx, c); err != nil { + log.FromContext(ctx).Error(err, "defrag: could not disarm NOSPACE alarm after sweep") + } + } df.Status.Phase = phase now := metav1.Now() df.Status.CompletedAt = &now @@ -533,5 +617,10 @@ func (r *EtcdDefragReconciler) SetupWithManager(mgr ctrl.Manager) error { } return ctrl.NewControllerManagedBy(mgr). For(&lll.EtcdDefrag{}). + // A wedged member holds a Defragment RPC for up to defragRPCTimeout; + // with a single worker that would stall every other cluster's run too. + // oldestActive already serializes runs per cluster, so distinct clusters + // are safe to reconcile in parallel. + WithOptions(controller.Options{MaxConcurrentReconciles: defragMaxConcurrentReconciles}). Complete(r) } diff --git a/controllers/etcddefrag_controller_test.go b/controllers/etcddefrag_controller_test.go index 6fb90084..506a0ccc 100644 --- a/controllers/etcddefrag_controller_test.go +++ b/controllers/etcddefrag_controller_test.go @@ -16,6 +16,7 @@ import ( "fmt" "strings" "testing" + "time" etcdserverpb "go.etcd.io/etcd/api/v3/etcdserverpb" clientv3 "go.etcd.io/etcd/client/v3" @@ -84,6 +85,16 @@ func status(memberID, leader uint64, dbSize, inUse int64) *clientv3.StatusRespon return &clientv3.StatusResponse{Header: &etcdserverpb.ResponseHeader{MemberId: memberID}, Leader: leader, DbSize: dbSize, DbSizeInUse: inUse} } +// statusAlarm is status() with the given active-alarm lines, formatted the way +// etcd's Status RPC reports them in StatusResponse.Errors. +func statusAlarm(memberID, leader uint64, dbSize, inUse int64, alarms ...etcdserverpb.AlarmType) *clientv3.StatusResponse { + s := status(memberID, leader, dbSize, inUse) + for _, a := range alarms { + s.Errors = append(s.Errors, fmt.Sprintf("memberID:%d alarm:%s", memberID, a.String())) + } + return s +} + func objs(cluster *lll.EtcdCluster, members []lll.EtcdMember, dfs ...*lll.EtcdDefrag) []client.Object { out := []client.Object{cluster} for i := range members { @@ -289,6 +300,126 @@ func TestEtcdDefrag_SerializedPerCluster(t *testing.T) { } } +// A NOSPACE alarm is the case defrag exists to relieve: the health gate admits +// the run rather than refusing the one cluster that needs it, and the alarm is +// disarmed once space has been reclaimed. +func TestEtcdDefrag_NoSpaceAlarmPermitsRunAndDisarms(t *testing.T) { + cluster, members := defragCluster3() + df := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}, Rule: &lll.DefragRule{All: true}}} + c, s := newTestClient(t, objs(cluster, members, df)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): statusAlarm(10, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_NOSPACE), + defragEndpoint("c1-1"): statusAlarm(11, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_NOSPACE), + defragEndpoint("c1-2"): statusAlarm(12, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_NOSPACE), + } + fe.alarms = []*etcdserverpb.AlarmMember{{MemberID: 10, Alarm: etcdserverpb.AlarmType_NOSPACE}} + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: record.NewFakeRecorder(20)} + + got := driveDefrag(t, r, c, "d1") + if got.Status.Phase != lll.EtcdDefragPhaseComplete { + t.Fatalf("phase = %q, want Complete (NOSPACE must not block)", got.Status.Phase) + } + if len(fe.defragCalls) != 3 { + t.Fatalf("defragCalls = %v, want all 3 under NOSPACE", fe.defragCalls) + } + if len(fe.disarmCalls) != 1 || fe.disarmCalls[0].Alarm != etcdserverpb.AlarmType_NOSPACE { + t.Fatalf("disarmCalls = %+v, want one NOSPACE disarm after the sweep", fe.disarmCalls) + } +} + +// A CORRUPT alarm blocks the run: it is deferred, not forced. +func TestEtcdDefrag_CorruptAlarmBlocksRun(t *testing.T) { + cluster, members := defragCluster3() + df := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}, Rule: &lll.DefragRule{All: true}}} + c, s := newTestClient(t, objs(cluster, members, df)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): statusAlarm(10, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_CORRUPT), + defragEndpoint("c1-1"): status(11, 10, 500<<20, 100<<20), + defragEndpoint("c1-2"): status(12, 10, 500<<20, 100<<20), + } + rec := record.NewFakeRecorder(20) + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: rec} + + res, err := r.Reconcile(context.Background(), ctrl.Request{NamespacedName: nn("d1", "ns")}) + if err != nil { + t.Fatalf("reconcile: %v", err) + } + if res.RequeueAfter == 0 { + t.Errorf("expected a requeue while deferred, got %+v", res) + } + if len(fe.defragCalls) != 0 { + t.Fatalf("defrag ran despite a CORRUPT alarm: %v", fe.defragCalls) + } + got := mustGet(t, c, "d1", "ns", &lll.EtcdDefrag{}) + if got.Status.Phase != lll.EtcdDefragPhasePending { + t.Errorf("phase = %q, want Pending", got.Status.Phase) + } + if cond := findDefragCond(got); cond == nil || cond.Reason != "ClusterNotHealthy" { + t.Errorf("DefragChecked = %+v, want ClusterNotHealthy", cond) + } +} + +// A run that flaps Pending->Running->Pending must not keep re-stamping StartedAt: +// the active-deadline is measured from it, and re-stamping would let a stuck run +// hold the per-cluster slot forever. +func TestEtcdDefrag_StartedAtNotResetAcrossFlap(t *testing.T) { + cluster, members := defragCluster3() + seeded := metav1.NewTime(time.Now().Add(-25 * time.Minute)) + df := &lll.EtcdDefrag{ + ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, + Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}}, + Status: lll.EtcdDefragStatus{Phase: lll.EtcdDefragPhasePending, StartedAt: &seeded}, + } + c, s := newTestClient(t, objs(cluster, members, df)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): status(10, 10, 50<<20, 50<<20), + defragEndpoint("c1-1"): status(11, 10, 50<<20, 50<<20), + defragEndpoint("c1-2"): status(12, 10, 50<<20, 50<<20), + } + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: record.NewFakeRecorder(20)} + + if _, err := r.Reconcile(context.Background(), ctrl.Request{NamespacedName: nn("d1", "ns")}); err != nil { + t.Fatalf("reconcile: %v", err) + } + got := mustGet(t, c, "d1", "ns", &lll.EtcdDefrag{}) + if got.Status.StartedAt == nil { + t.Fatal("StartedAt was cleared") + } + if age := time.Since(got.Status.StartedAt.Time); age < 20*time.Minute { + t.Fatalf("StartedAt age = %s, want ~25m (it was reset on the Pending->Running transition)", age) + } +} + +// A defrag whose post-defrag Status read fails records the after-size and +// reclaimed bytes as unset rather than reporting a real defrag as reclaiming 0. +func TestMarkMember_AfterSizeUnknownLeavesReclaimedUnset(t *testing.T) { + newDF := func() *lll.EtcdDefrag { + return &lll.EtcdDefrag{Status: lll.EtcdDefragStatus{Members: []lll.MemberDefragStatus{{Name: "m", Outcome: lll.DefragOutcomePending}}}} + } + + unknown := newDF() + markMember(unknown, "m", lll.DefragOutcomeDefragmented, "AfterSizeUnavailable", + &memberBackend{status: status(1, 1, 500<<20, 100<<20)}) + if m := unknown.Status.Members[0]; m.DBSizeBefore != 500<<20 || m.DBSizeAfter != 0 || m.ReclaimedBytes != 0 { + t.Fatalf("after-size unknown: got before=%d after=%d reclaimed=%d, want before=%d after/reclaimed=0", + m.DBSizeBefore, m.DBSizeAfter, m.ReclaimedBytes, 500<<20) + } + + known := newDF() + markMember(known, "m", lll.DefragOutcomeDefragmented, "", + &memberBackend{status: status(1, 1, 500<<20, 100<<20), after: 120 << 20, afterKnown: true}) + if m := known.Status.Members[0]; m.DBSizeAfter != 120<<20 || m.ReclaimedBytes != 380<<20 { + t.Fatalf("after-size known: got after=%d reclaimed=%d, want after=%d reclaimed=%d", + m.DBSizeAfter, m.ReclaimedBytes, 120<<20, 380<<20) + } +} + func findDefragCond(df *lll.EtcdDefrag) *metav1.Condition { for i := range df.Status.Conditions { if df.Status.Conditions[i].Type == "DefragChecked" { diff --git a/controllers/testing_helpers_test.go b/controllers/testing_helpers_test.go index b1c621db..fd4353c8 100644 --- a/controllers/testing_helpers_test.go +++ b/controllers/testing_helpers_test.go @@ -59,6 +59,12 @@ type fakeEtcd struct { defragCalls []string defragErr error + // Alarm surface. alarms is what AlarmList returns; alarmListErr fails it; + // disarmCalls records each AlarmDisarm target in call order. + alarms []*etcdserverpb.AlarmMember + alarmListErr error + disarmCalls []*etcdserverpb.AlarmMember + addCalls []string addLearnerCalls []string promoteCalls []uint64 @@ -220,6 +226,21 @@ func (f *fakeEtcd) Defragment(_ context.Context, endpoint string) (*clientv3.Def return &clientv3.DefragmentResponse{Header: &etcdserverpb.ResponseHeader{ClusterId: f.clusterID}}, nil } +func (f *fakeEtcd) AlarmList(_ context.Context) (*clientv3.AlarmResponse, error) { + if f.alarmListErr != nil { + return nil, f.alarmListErr + } + return &clientv3.AlarmResponse{ + Header: &etcdserverpb.ResponseHeader{ClusterId: f.clusterID}, + Alarms: f.alarms, + }, nil +} + +func (f *fakeEtcd) AlarmDisarm(_ context.Context, m *clientv3.AlarmMember) (*clientv3.AlarmResponse, error) { + f.disarmCalls = append(f.disarmCalls, (*etcdserverpb.AlarmMember)(m)) + return &clientv3.AlarmResponse{Header: &etcdserverpb.ResponseHeader{ClusterId: f.clusterID}}, nil +} + func (f *fakeEtcd) Close() error { f.closed = true; return nil } func factoryReturning(c EtcdClusterClient) EtcdClientFactory { diff --git a/docs/etcd-defrag.md b/docs/etcd-defrag.md index a976e125..03a552bb 100644 --- a/docs/etcd-defrag.md +++ b/docs/etcd-defrag.md @@ -129,8 +129,7 @@ no `spec` knobs: the next one forever. - **Retry within a run:** a deferred `Pending` re-checks cluster health with backoff up to the deadline; a failed per-member RPC is retried a bounded number - of times then marked `Failed` (a failing leader fails the run); and a defrag - that doesn't shrink `DbSize` is backed off rather than repeated. + of times then marked `Failed` (a failing leader fails the run). - **Retry across runs:** terminal phases (`Complete`/`Failed`) are sticky — an `EtcdDefrag` never re-runs itself. A retry is a *new* `EtcdDefrag`: the external scheduler's next tick for periodic use, or a re-create for a one-shot. Each From 4c2a28aa21e9f3482e7679c13a892b96b3221b25 Mon Sep 17 00:00:00 2001 From: Andrey Kolkov Date: Thu, 20 Aug 2026 20:23:08 +0400 Subject: [PATCH 3/4] fix(defrag): close NOSPACE disarm holes, harden deadline, honest docs Addresses the review blockers on the EtcdDefrag controller. - Disarm NOSPACE whenever the sweep reclaimed space, not only on a wholly clean run. A partial sweep (one member's Defragment fails, or the run hits the deadline mid-way) still relieved the read-only wedge on the members it compacted, so gating the disarm on phase==Complete stranded the exact cluster this feature rescues. finalize and the deadline path now both disarm when defragmented > 0. - disarmNoSpaceAlarms is now truly best-effort: it continues past a transient disarm failure on one member, joins the errors naming each, and surfaces an AlarmDisarmFailed event so a still-read-only cluster after a Complete run is visible in kubectl describe. - The active deadline is derived from the worst-case serial sweep for the largest supported cluster (defragMaxSupportedMembers) instead of a bare 30m literal that a 7-member big-backend run overran, killing a healthy sweep mid-flight. A constant-relation test pins it. - StartedAt/Running is stamped before the backend lookup so a run whose first planned member has vanished is still measured from work, not creation. - Docs: the safety model no longer claims raft-lag is gated or that NOSPACE blocks; the retry section states the shipped behaviour (a failed per-member RPC fails the run, no backoff) instead of promising retries that do not exist; the scheduling note drops the "this API PR" framing. - CRD: drop the stale "(not-yet-implemented)" note from ttlSecondsAfterFinished now that handleTTL acts on it; regenerated. - Tests: partial-failure disarm, disarm-continues-after-failure, TTL GC, DeadlineExceeded, ClusterNotFound, and the deadline/worst-case relation. Signed-off-by: Andrey Kolkov Co-Authored-By: Claude Opus 4.8 (1M context) --- api/v1alpha2/etcddefrag_types.go | 5 +- ...tcd-operator.cozystack.io_etcddefrags.yaml | 5 +- controllers/etcddefrag_controller.go | 122 +++++++++---- controllers/etcddefrag_controller_test.go | 164 ++++++++++++++++++ controllers/testing_helpers_test.go | 21 ++- docs/etcd-defrag.md | 20 ++- 6 files changed, 285 insertions(+), 52 deletions(-) diff --git a/api/v1alpha2/etcddefrag_types.go b/api/v1alpha2/etcddefrag_types.go index 4e2bfa44..f101a08b 100644 --- a/api/v1alpha2/etcddefrag_types.go +++ b/api/v1alpha2/etcddefrag_types.go @@ -38,9 +38,8 @@ type EtcdDefragSpec struct { // TTLSecondsAfterFinished records how long after a terminal phase this // object should be garbage-collected — meaningful for objects a scheduler - // stamps out. NOTE: acted on by the (not-yet-implemented) reconciling - // controller; the API server does not garbage-collect custom resources on - // its own. Absent means the record is kept. + // stamps out. Acted on by the reconciling controller; the API server does not + // garbage-collect custom resources on its own. Absent means the record is kept. // +kubebuilder:validation:Minimum=0 // +optional TTLSecondsAfterFinished *int32 `json:"ttlSecondsAfterFinished,omitempty"` diff --git a/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml b/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml index febca3ac..ae38991d 100644 --- a/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml +++ b/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml @@ -132,9 +132,8 @@ spec: description: |- TTLSecondsAfterFinished records how long after a terminal phase this object should be garbage-collected — meaningful for objects a scheduler - stamps out. NOTE: acted on by the (not-yet-implemented) reconciling - controller; the API server does not garbage-collect custom resources on - its own. Absent means the record is kept. + stamps out. Acted on by the reconciling controller; the API server does not + garbage-collect custom resources on its own. Absent means the record is kept. format: int32 minimum: 0 type: integer diff --git a/controllers/etcddefrag_controller.go b/controllers/etcddefrag_controller.go index acf56ee4..eba5f3bb 100644 --- a/controllers/etcddefrag_controller.go +++ b/controllers/etcddefrag_controller.go @@ -12,6 +12,7 @@ package controllers import ( "context" + "errors" "fmt" "sort" "strconv" @@ -51,10 +52,17 @@ const ( // defragRequeueAfter paces the one-member-per-pass sweep and the retry while // deferred on an unhealthy cluster. defragRequeueAfter = 10 * time.Second + + // defragMaxSupportedMembers bounds the worst-case serial sweep the active + // deadline must outlast. etcd clusters are odd-sized and rarely exceed 7. + defragMaxSupportedMembers = 7 // defragActiveDeadline bounds a whole run (Running + waiting-while-Pending), // so a run stuck on an unhealthy cluster can't hold the per-cluster slot - // forever. - defragActiveDeadline = 30 * time.Minute + // forever. It must outlast the worst-case serial sweep — one member per pass, + // each pass probing every member then a stop-the-world Defragment and a + // requeue gap — or a healthy-but-slow big-backend cluster is killed mid-sweep; + // derived from the constants above so it stays consistent if any is retuned. + defragActiveDeadline = defragMaxSupportedMembers*(defragRPCTimeout+defragRequeueAfter+defragMaxSupportedMembers*defragStatusTimeout) + 5*time.Minute // defragMaxConcurrentReconciles lets distinct clusters' runs proceed in // parallel (per-cluster serialization is enforced separately by oldestActive). @@ -121,17 +129,6 @@ func (r *EtcdDefragReconciler) Reconcile(ctx context.Context, req ctrl.Request) return ctrl.Result{RequeueAfter: defragRequeueAfter}, nil } - // Overall deadline: a run that can't make progress must fail, not linger. - // (CreationTimestamp is zero before the apiserver stamps it — e.g. in unit - // tests — so only enforce against a non-zero reference time.) - if started := df.Status.StartedAt; started != nil && time.Since(started.Time) > defragActiveDeadline { - return r.fail(ctx, df, "DeadlineExceeded", - fmt.Sprintf("defragmentation did not complete within %s", defragActiveDeadline)) - } else if started == nil && !df.CreationTimestamp.IsZero() && time.Since(df.CreationTimestamp.Time) > defragActiveDeadline { - return r.fail(ctx, df, "DeadlineExceeded", - fmt.Sprintf("defragmentation could not start within %s (cluster never became healthy)", defragActiveDeadline)) - } - var memberList lll.EtcdMemberList if err := r.List(ctx, &memberList, client.InNamespace(df.Namespace), client.MatchingLabels{LabelCluster: cluster.Name}); err != nil { @@ -141,11 +138,25 @@ func (r *EtcdDefragReconciler) Reconcile(ctx context.Context, req ctrl.Request) c, backends, err := r.dialAndProbe(ctx, cluster, running) if err != nil { + // Without a client we can neither defragment nor disarm; if the run has + // outlived its deadline while stuck here, fail it rather than requeue + // forever. + if exceededDeadline(df) { + return r.fail(ctx, df, "DeadlineExceeded", deadlineMsg(df)) + } logger.Error(err, "defrag: cannot dial/probe cluster; retrying") return ctrl.Result{RequeueAfter: defragRequeueAfter}, nil } defer c.Close() + // Overall deadline: a run that can't make progress must fail, not linger. A + // run that already reclaimed space disarms NOSPACE on the way out (failRun), + // so a deadline-terminated partial sweep still lifts the wedge that admitted + // it. Checked after dialing so the disarm has a client. + if exceededDeadline(df) { + return r.failRun(ctx, df, c, "DeadlineExceeded", deadlineMsg(df)) + } + // Health gate: defrag is only safe when every desired member is present, // reachable, free of blocking alarms and agreeing on a leader. if reason, ok := clusterDefragHealthy(cluster, running, backends); !ok { @@ -177,19 +188,14 @@ func (r *EtcdDefragReconciler) Reconcile(ctx context.Context, req ctrl.Request) return r.finalize(ctx, df, c) } - b := backendByName(backends, next.Name) - if b == nil { - // A planned member vanished between passes despite the health gate. - markMember(df, next.Name, lll.DefragOutcomeFailed, "MemberGone", nil) - return r.persistAndRequeue(ctx, df) - } - df.Status.Phase = lll.EtcdDefragPhaseRunning // Stamp StartedAt exactly once, on the first Running transition. The // active-deadline is measured from it, so re-stamping on a Pending->Running // flap (a cluster that recovers between health-gate blips) would keep // resetting the deadline and a stuck run could hold the per-cluster slot - // forever. + // forever. Done before the backend lookup so a run whose first planned member + // has vanished is still stamped Running and its deadline measured from work, + // not creation. if df.Status.StartedAt == nil { now := metav1.Now() df.Status.StartedAt = &now @@ -198,6 +204,13 @@ func (r *EtcdDefragReconciler) Reconcile(ctx context.Context, req ctrl.Request) // health-gate flap does not stay reading ClusterNotHealthy. setDefragCondition(df, metav1.ConditionTrue, "Running", "defragmenting members") + b := backendByName(backends, next.Name) + if b == nil { + // A planned member vanished between passes despite the health gate. + markMember(df, next.Name, lll.DefragOutcomeFailed, "MemberGone", nil) + return r.persistAndRequeue(ctx, df) + } + quota := effectiveQuotaBytes(cluster) if trig, reason := defragRuleTriggered(df.Spec.Rule, b.status.DbSize, b.status.DbSizeInUse, quota); !trig { markMember(df, b.member.Name, lll.DefragOutcomeSkipped, reason, b) @@ -282,11 +295,13 @@ func statusWithTimeout(ctx context.Context, c EtcdClusterClient, endpoint string return c.Status(sctx, endpoint) } -// disarmNoSpaceAlarms clears any armed NOSPACE alarm. The health gate lets a +// disarmNoSpaceAlarms clears every armed NOSPACE alarm. The health gate lets a // NOSPACE cluster through so the run can reclaim its backend space; etcd keeps -// the alarm armed — and the cluster read-only — until it is explicitly disarmed. -// Best-effort: etcd re-arms on the next write if the space was not actually -// freed, so a failure here is logged by the caller rather than failing the run. +// each member's alarm armed — and the cluster read-only — until it is explicitly +// disarmed. AlarmList returns one entry per member that raised NOSPACE, so a +// transient failure on one must not abandon the rest: the loop continues and +// joins the errors, naming every member it could not disarm. Best-effort in that +// etcd re-arms on the next write if the space was not actually freed. func (r *EtcdDefragReconciler) disarmNoSpaceAlarms(ctx context.Context, c EtcdClusterClient) error { lctx, cancel := context.WithTimeout(ctx, defragStatusTimeout) defer cancel() @@ -294,6 +309,7 @@ func (r *EtcdDefragReconciler) disarmNoSpaceAlarms(ctx context.Context, c EtcdCl if err != nil { return err } + var errs error for _, a := range resp.Alarms { if a == nil || a.Alarm != etcdserverpb.AlarmType_NOSPACE { continue @@ -302,10 +318,27 @@ func (r *EtcdDefragReconciler) disarmNoSpaceAlarms(ctx context.Context, c EtcdCl _, derr := c.AlarmDisarm(dctx, (*clientv3.AlarmMember)(a)) dcancel() if derr != nil { - return derr + errs = errors.Join(errs, fmt.Errorf("member %d: %w", a.MemberID, derr)) } } - return nil + return errs +} + +// maybeDisarm clears any armed NOSPACE alarm when the run reclaimed backend +// space, on any terminal path — a partial or deadline-terminated sweep still +// relieves the wedge that admitted it, and etcd holds the cluster read-only +// until the alarm is disarmed. Best-effort: a failure is logged and surfaced as +// an event rather than propagated, since etcd re-arms on the next write if the +// space was not actually freed. +func (r *EtcdDefragReconciler) maybeDisarm(ctx context.Context, df *lll.EtcdDefrag, c EtcdClusterClient) { + if c == nil || df.Status.Defragmented == 0 { + return + } + if err := r.disarmNoSpaceAlarms(ctx, c); err != nil { + log.FromContext(ctx).Error(err, "defrag: could not disarm NOSPACE alarm after sweep") + r.event(df, corev1.EventTypeWarning, "AlarmDisarmFailed", + fmt.Sprintf("reclaimed backend space but could not disarm NOSPACE alarm; cluster may stay read-only: %v", err)) + } } // clusterDefragHealthy reports whether the cluster is safe to defragment: every @@ -525,12 +558,10 @@ func (r *EtcdDefragReconciler) finalize(ctx context.Context, df *lll.EtcdDefrag, } } // A cluster admitted with a NOSPACE alarm stays read-only until the alarm is - // disarmed; do it once the sweep has actually reclaimed space. - if phase == lll.EtcdDefragPhaseComplete && df.Status.Defragmented > 0 { - if err := r.disarmNoSpaceAlarms(ctx, c); err != nil { - log.FromContext(ctx).Error(err, "defrag: could not disarm NOSPACE alarm after sweep") - } - } + // disarmed; disarm whenever the sweep reclaimed space, even if a later member + // failed — the reclaimed space is what lifts the wedge, and gating this on a + // wholly-clean run would leave the one cluster this feature rescues read-only. + r.maybeDisarm(ctx, df, c) df.Status.Phase = phase now := metav1.Now() df.Status.CompletedAt = &now @@ -546,6 +577,31 @@ func (r *EtcdDefragReconciler) finalize(ctx context.Context, df *lll.EtcdDefrag, return r.handleTTL(ctx, df) } +// failRun is fail() for a terminal path that holds a client: it disarms any +// NOSPACE alarm the run relieved before recording the failure. +func (r *EtcdDefragReconciler) failRun(ctx context.Context, df *lll.EtcdDefrag, c EtcdClusterClient, reason, msg string) (ctrl.Result, error) { + r.maybeDisarm(ctx, df, c) + return r.fail(ctx, df, reason, msg) +} + +// exceededDeadline reports whether the run has outlived defragActiveDeadline, +// measured from StartedAt once work began, else from creation. CreationTimestamp +// is zero before the apiserver stamps it (e.g. in unit tests), so a zero +// reference never trips the deadline. +func exceededDeadline(df *lll.EtcdDefrag) bool { + if s := df.Status.StartedAt; s != nil { + return time.Since(s.Time) > defragActiveDeadline + } + return !df.CreationTimestamp.IsZero() && time.Since(df.CreationTimestamp.Time) > defragActiveDeadline +} + +func deadlineMsg(df *lll.EtcdDefrag) string { + if df.Status.StartedAt != nil { + return fmt.Sprintf("defragmentation did not complete within %s", defragActiveDeadline) + } + return fmt.Sprintf("defragmentation could not start within %s (cluster never became healthy)", defragActiveDeadline) +} + func (r *EtcdDefragReconciler) fail(ctx context.Context, df *lll.EtcdDefrag, reason, msg string) (ctrl.Result, error) { df.Status.Phase = lll.EtcdDefragPhaseFailed now := metav1.Now() diff --git a/controllers/etcddefrag_controller_test.go b/controllers/etcddefrag_controller_test.go index 506a0ccc..499e9995 100644 --- a/controllers/etcddefrag_controller_test.go +++ b/controllers/etcddefrag_controller_test.go @@ -420,6 +420,170 @@ func TestMarkMember_AfterSizeUnknownLeavesReclaimedUnset(t *testing.T) { } } +// A NOSPACE sweep where one member's Defragment fails still disarms the alarm: +// the space reclaimed on the members that succeeded is what lifts the read-only +// wedge, so gating the disarm on a wholly-clean run would strand the cluster. +func TestEtcdDefrag_PartialFailureStillDisarmsNoSpace(t *testing.T) { + cluster, members := defragCluster3() + df := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}, Rule: &lll.DefragRule{All: true}}} + c, s := newTestClient(t, objs(cluster, members, df)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): statusAlarm(10, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_NOSPACE), + defragEndpoint("c1-1"): statusAlarm(11, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_NOSPACE), + defragEndpoint("c1-2"): statusAlarm(12, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_NOSPACE), + } + fe.alarms = []*etcdserverpb.AlarmMember{{MemberID: 10, Alarm: etcdserverpb.AlarmType_NOSPACE}} + // One follower's Defragment fails; the other follower and the leader succeed. + fe.defragErrByEndpoint = map[string]error{defragEndpoint("c1-1"): errors.New("boom")} + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: record.NewFakeRecorder(20)} + + got := driveDefrag(t, r, c, "d1") + if got.Status.Phase != lll.EtcdDefragPhaseFailed { + t.Fatalf("phase = %q, want Failed (one member's Defragment failed)", got.Status.Phase) + } + if got.Status.Defragmented != 2 { + t.Fatalf("Defragmented = %d, want 2", got.Status.Defragmented) + } + if len(fe.disarmCalls) != 1 { + t.Fatalf("disarmCalls = %+v, want one NOSPACE disarm despite the failed run", fe.disarmCalls) + } +} + +// AlarmList returns one entry per member that raised NOSPACE; a transient disarm +// failure on one must not abandon the rest. +func TestEtcdDefrag_DisarmContinuesAfterFailure(t *testing.T) { + cluster, members := defragCluster3() + df := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}, Rule: &lll.DefragRule{All: true}}} + c, s := newTestClient(t, objs(cluster, members, df)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusByEndpoint = map[string]*clientv3.StatusResponse{ + defragEndpoint("c1-0"): statusAlarm(10, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_NOSPACE), + defragEndpoint("c1-1"): statusAlarm(11, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_NOSPACE), + defragEndpoint("c1-2"): statusAlarm(12, 10, 500<<20, 100<<20, etcdserverpb.AlarmType_NOSPACE), + } + fe.alarms = []*etcdserverpb.AlarmMember{ + {MemberID: 10, Alarm: etcdserverpb.AlarmType_NOSPACE}, + {MemberID: 11, Alarm: etcdserverpb.AlarmType_NOSPACE}, + {MemberID: 12, Alarm: etcdserverpb.AlarmType_NOSPACE}, + } + fe.disarmErrByMember = map[uint64]error{10: errors.New("transient")} + rec := record.NewFakeRecorder(20) + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: rec} + + got := driveDefrag(t, r, c, "d1") + if got.Status.Phase != lll.EtcdDefragPhaseComplete { + t.Fatalf("phase = %q, want Complete", got.Status.Phase) + } + if len(fe.disarmCalls) != 3 { + t.Fatalf("disarmCalls = %d, want 3 (loop must not abort on the first failure)", len(fe.disarmCalls)) + } + assertDefragEvent(t, rec, "AlarmDisarmFailed") +} + +// ttlSecondsAfterFinished GCs a finished record once it expires, and requeues +// (not deletes) one that has not. +func TestEtcdDefrag_TTLGarbageCollects(t *testing.T) { + ctx := context.Background() + newFinished := func(name string, completedAt metav1.Time) *lll.EtcdDefrag { + return &lll.EtcdDefrag{ + ObjectMeta: metav1.ObjectMeta{Name: name, Namespace: "ns"}, + Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}, TTLSecondsAfterFinished: ptrInt32(3600)}, + Status: lll.EtcdDefragStatus{Phase: lll.EtcdDefragPhaseComplete, CompletedAt: &completedAt}, + } + } + + expired := newFinished("d-expired", metav1.NewTime(time.Now().Add(-2*time.Hour))) + fresh := newFinished("d-fresh", metav1.NewTime(time.Now().Add(-1*time.Second))) + c, s := newTestClient(t, expired, fresh) + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(newFakeEtcd(0xabc)), Recorder: record.NewFakeRecorder(20)} + + if _, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: nn("d-expired", "ns")}); err != nil { + t.Fatalf("reconcile d-expired: %v", err) + } + if err := c.Get(ctx, nn("d-expired", "ns"), &lll.EtcdDefrag{}); err == nil { + t.Fatal("expired EtcdDefrag was not garbage-collected") + } else if client.IgnoreNotFound(err) != nil { + t.Fatalf("get d-expired: %v", err) + } + + res, err := r.Reconcile(ctx, ctrl.Request{NamespacedName: nn("d-fresh", "ns")}) + if err != nil { + t.Fatalf("reconcile d-fresh: %v", err) + } + if res.RequeueAfter <= 0 || res.RequeueAfter > 3600*time.Second { + t.Fatalf("RequeueAfter = %s, want 0 < requeue <= 3600s", res.RequeueAfter) + } + if err := c.Get(ctx, nn("d-fresh", "ns"), &lll.EtcdDefrag{}); err != nil { + t.Fatalf("fresh EtcdDefrag was deleted before expiry: %v", err) + } +} + +// A run that outlives the active deadline fails with DeadlineExceeded rather than +// lingering and holding the per-cluster slot. +func TestEtcdDefrag_DeadlineExceededFails(t *testing.T) { + cluster, members := defragCluster3() + started := metav1.NewTime(time.Now().Add(-defragActiveDeadline - time.Minute)) + df := &lll.EtcdDefrag{ + ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, + Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "c1"}}, + Status: lll.EtcdDefragStatus{Phase: lll.EtcdDefragPhaseRunning, StartedAt: &started}, + } + c, s := newTestClient(t, objs(cluster, members, df)...) + fe := newFakeEtcd(0xabc) + fe.leader = 10 + fe.statusErr = errors.New("context deadline exceeded") // cluster unreachable + rec := record.NewFakeRecorder(20) + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(fe), Recorder: rec} + + if _, err := r.Reconcile(context.Background(), ctrl.Request{NamespacedName: nn("d1", "ns")}); err != nil { + t.Fatalf("reconcile: %v", err) + } + got := mustGet(t, c, "d1", "ns", &lll.EtcdDefrag{}) + if got.Status.Phase != lll.EtcdDefragPhaseFailed { + t.Fatalf("phase = %q, want Failed", got.Status.Phase) + } + if cond := findDefragCond(got); cond == nil || cond.Reason != "DeadlineExceeded" { + t.Fatalf("DefragChecked = %+v, want DeadlineExceeded", cond) + } + assertDefragEvent(t, rec, "DeadlineExceeded") +} + +// A run whose clusterRef names no EtcdCluster fails with ClusterNotFound. +func TestEtcdDefrag_ClusterNotFoundFails(t *testing.T) { + df := &lll.EtcdDefrag{ObjectMeta: metav1.ObjectMeta{Name: "d1", Namespace: "ns"}, Spec: lll.EtcdDefragSpec{ClusterRef: corev1.LocalObjectReference{Name: "nope"}}} + c, s := newTestClient(t, df) + rec := record.NewFakeRecorder(20) + r := &EtcdDefragReconciler{Client: c, Scheme: s, EtcdClientFactory: factoryReturning(newFakeEtcd(0xabc)), Recorder: rec} + + if _, err := r.Reconcile(context.Background(), ctrl.Request{NamespacedName: nn("d1", "ns")}); err != nil { + t.Fatalf("reconcile: %v", err) + } + got := mustGet(t, c, "d1", "ns", &lll.EtcdDefrag{}) + if got.Status.Phase != lll.EtcdDefragPhaseFailed { + t.Fatalf("phase = %q, want Failed", got.Status.Phase) + } + if cond := findDefragCond(got); cond == nil || cond.Reason != "ClusterNotFound" { + t.Fatalf("DefragChecked = %+v, want ClusterNotFound", cond) + } + assertDefragEvent(t, rec, "ClusterNotFound") +} + +// The active deadline must outlast the worst-case serial sweep for the largest +// supported cluster: one member per pass, each probing every member then a +// stop-the-world Defragment and a requeue gap. This reads the constants the +// controller actually uses, so retuning any of them without widening the deadline +// fails the build. +func TestDefragActiveDeadlineCoversWorstCaseSweep(t *testing.T) { + worst := defragMaxSupportedMembers * (defragRPCTimeout + defragRequeueAfter + defragMaxSupportedMembers*defragStatusTimeout) + if defragActiveDeadline < worst { + t.Fatalf("defragActiveDeadline %s < worst-case sweep %s for %d members", + defragActiveDeadline, worst, defragMaxSupportedMembers) + } +} + func findDefragCond(df *lll.EtcdDefrag) *metav1.Condition { for i := range df.Status.Conditions { if df.Status.Conditions[i].Type == "DefragChecked" { diff --git a/controllers/testing_helpers_test.go b/controllers/testing_helpers_test.go index fd4353c8..2ab616e8 100644 --- a/controllers/testing_helpers_test.go +++ b/controllers/testing_helpers_test.go @@ -51,19 +51,22 @@ type fakeEtcd struct { // Defrag surface. statusByEndpoint overrides Status per endpoint (backend // sizes, leader); statusErrByEndpoint fails Status for a specific endpoint // (an unhealthy member); leader is the default StatusResponse.Leader. - // defragCalls records each Defragment endpoint in call order; defragErr, - // when set, fails Defragment. + // defragCalls records each Defragment endpoint in call order; defragErr fails + // every Defragment, defragErrByEndpoint fails it for a specific endpoint. statusByEndpoint map[string]*clientv3.StatusResponse statusErrByEndpoint map[string]error leader uint64 defragCalls []string defragErr error + defragErrByEndpoint map[string]error // Alarm surface. alarms is what AlarmList returns; alarmListErr fails it; - // disarmCalls records each AlarmDisarm target in call order. - alarms []*etcdserverpb.AlarmMember - alarmListErr error - disarmCalls []*etcdserverpb.AlarmMember + // disarmCalls records each AlarmDisarm target in call order; + // disarmErrByMember fails AlarmDisarm for a specific member. + alarms []*etcdserverpb.AlarmMember + alarmListErr error + disarmCalls []*etcdserverpb.AlarmMember + disarmErrByMember map[uint64]error addCalls []string addLearnerCalls []string @@ -223,6 +226,9 @@ func (f *fakeEtcd) Defragment(_ context.Context, endpoint string) (*clientv3.Def if f.defragErr != nil { return nil, f.defragErr } + if err := f.defragErrByEndpoint[endpoint]; err != nil { + return nil, err + } return &clientv3.DefragmentResponse{Header: &etcdserverpb.ResponseHeader{ClusterId: f.clusterID}}, nil } @@ -238,6 +244,9 @@ func (f *fakeEtcd) AlarmList(_ context.Context) (*clientv3.AlarmResponse, error) func (f *fakeEtcd) AlarmDisarm(_ context.Context, m *clientv3.AlarmMember) (*clientv3.AlarmResponse, error) { f.disarmCalls = append(f.disarmCalls, (*etcdserverpb.AlarmMember)(m)) + if err := f.disarmErrByMember[m.MemberID]; err != nil { + return nil, err + } return &clientv3.AlarmResponse{Header: &etcdserverpb.ResponseHeader{ClusterId: f.clusterID}}, nil } diff --git a/docs/etcd-defrag.md b/docs/etcd-defrag.md index 03a552bb..12624f9c 100644 --- a/docs/etcd-defrag.md +++ b/docs/etcd-defrag.md @@ -12,8 +12,8 @@ the operator drives it through `status.phase` and it never re-runs. Today, recurring defragmentation is driven by creating `EtcdDefrag` objects from outside (a `CronJob`, a GitOps cron). A companion `EtcdDefragPolicy` kind — a cadence (`schedule`) and/or a condition (`when`) that stamps out `EtcdDefrag` -runs — is planned so the operator absorbs that scheduling itself; it is not part -of this API PR. +runs — is planned so the operator absorbs that scheduling itself; it is not +implemented yet. ## Why in the operator (not a bare CronJob) @@ -103,8 +103,12 @@ out. - **Serialized per cluster:** at most one `EtcdDefrag` runs against a given `EtcdCluster` at a time; others wait in `Pending`. - Health is judged from more than "the member answered": a member replies to a - local status read while partitioned, alarmed (`NOSPACE`/`CORRUPT`), or behind - in raft, so those are checked before acting. + local status read while partitioned or alarmed, so the gate checks that every + desired member is present and reachable, that they agree on a single non-zero + leader, and that no member reports a blocking alarm. A `CORRUPT` alarm blocks; + a `NOSPACE` alarm does **not** — a backend at its quota is exactly what a defrag + relieves, so the run is admitted and the alarm is disarmed once space has been + reclaimed. Raft lag is not yet part of the gate. ## Status @@ -127,9 +131,11 @@ no `spec` knobs: together; on expiry the run is `Failed`. This also protects the per-cluster serialization slot — a run stuck waiting on an unhealthy cluster can't block the next one forever. -- **Retry within a run:** a deferred `Pending` re-checks cluster health with - backoff up to the deadline; a failed per-member RPC is retried a bounded number - of times then marked `Failed` (a failing leader fails the run). +- **Retry within a run:** a deferred `Pending` re-checks cluster health each pass + up to the deadline. A failed per-member `Defragment` RPC marks that member + `Failed` immediately, and any failed member fails the run — a partial sweep that + reclaimed space still disarms `NOSPACE` on the way out. (Per-member RPC retry is + a possible follow-up, not shipped here.) - **Retry across runs:** terminal phases (`Complete`/`Failed`) are sticky — an `EtcdDefrag` never re-runs itself. A retry is a *new* `EtcdDefrag`: the external scheduler's next tick for periodic use, or a re-create for a one-shot. Each From a18378e5c26e963562a9f7e087927bfcd6e9ae52 Mon Sep 17 00:00:00 2001 From: Andrey Kolkov Date: Thu, 20 Aug 2026 20:34:13 +0400 Subject: [PATCH 4/4] fix(defrag): reject inert quotaUsageAbove 100%, retry status write on conflict MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two non-blocking items from the review. - quotaUsageAbove: "100%" passed the CRD pattern but could never fire — a backend never exceeds its quota (etcd raises NOSPACE first), so the quota arm's dbSize > 1.0*quota test is always false. Tighten the pattern to 1..99 and the parsePercent fallback to match, so the dead value is rejected at admission instead of silently doing nothing. - persistAndRequeue now re-fetches and re-applies the computed status on a conflict (retry.RetryOnConflict) instead of returning the error. A member's Defragment RPC can block for defragRPCTimeout; an unrelated metadata write in that window would otherwise discard the recorded outcome and the run would defragment the member a second time next pass. The defrag controller owns these status fields, so re-applying them is safe. Signed-off-by: Andrey Kolkov Co-Authored-By: Claude Opus 4.8 (1M context) --- api/v1alpha2/etcddefrag_types.go | 7 +++--- ...tcd-operator.cozystack.io_etcddefrags.yaml | 7 +++--- controllers/etcddefrag_controller.go | 25 ++++++++++++++++--- controllers/etcddefrag_controller_test.go | 1 + 4 files changed, 30 insertions(+), 10 deletions(-) diff --git a/api/v1alpha2/etcddefrag_types.go b/api/v1alpha2/etcddefrag_types.go index f101a08b..e69fea64 100644 --- a/api/v1alpha2/etcddefrag_types.go +++ b/api/v1alpha2/etcddefrag_types.go @@ -70,9 +70,10 @@ type DefragRule struct { // QuotaUsageAbove: when DbSize exceeds this fraction of the backend quota // (approaching NOSPACE), lower the reclaimable floor to MinReclaim so small // wins are taken under pressure. A member is never defragmented when its - // reclaimable space is below MinReclaim. Integer percent 1..100 with a "%" - // suffix, e.g. "80%". - // +kubebuilder:validation:Pattern=`^([1-9][0-9]?|100)%$` + // reclaimable space is below MinReclaim. Integer percent 1..99 with a "%" + // suffix, e.g. "80%"; 100% is rejected because a backend never exceeds its + // quota (etcd raises NOSPACE first), so the arm could never fire. + // +kubebuilder:validation:Pattern=`^[1-9][0-9]?%$` // +optional QuotaUsageAbove string `json:"quotaUsageAbove,omitempty"` diff --git a/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml b/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml index ae38991d..abc07b7e 100644 --- a/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml +++ b/charts/etcd-operator/crd-bases/etcd-operator.cozystack.io_etcddefrags.yaml @@ -110,9 +110,10 @@ spec: QuotaUsageAbove: when DbSize exceeds this fraction of the backend quota (approaching NOSPACE), lower the reclaimable floor to MinReclaim so small wins are taken under pressure. A member is never defragmented when its - reclaimable space is below MinReclaim. Integer percent 1..100 with a "%" - suffix, e.g. "80%". - pattern: ^([1-9][0-9]?|100)%$ + reclaimable space is below MinReclaim. Integer percent 1..99 with a "%" + suffix, e.g. "80%"; 100% is rejected because a backend never exceeds its + quota (etcd raises NOSPACE first), so the arm could never fire. + pattern: ^[1-9][0-9]?%$ type: string type: object x-kubernetes-validations: diff --git a/controllers/etcddefrag_controller.go b/controllers/etcddefrag_controller.go index eba5f3bb..ff049685 100644 --- a/controllers/etcddefrag_controller.go +++ b/controllers/etcddefrag_controller.go @@ -27,6 +27,7 @@ import ( "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/types" "k8s.io/client-go/tools/record" + "k8s.io/client-go/util/retry" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller" @@ -508,11 +509,12 @@ func defragRuleTriggered(rule *lll.DefragRule, dbSize, dbSizeInUse, quota int64) } // parsePercent parses "80%" into 0.80. Returns ok=false for anything outside -// 1–100 or non-numeric; the CRD pattern rejects such values at admission, so -// this is a defensive fallback. +// 1–99 or non-numeric; the CRD pattern rejects such values at admission, so this +// is a defensive fallback. 100% is excluded on purpose: a backend never exceeds +// its quota (etcd raises NOSPACE first), so the quota arm could never fire. func parsePercent(s string) (float64, bool) { n, err := strconv.Atoi(strings.TrimSuffix(s, "%")) - if err != nil || n <= 0 || n > 100 { + if err != nil || n <= 0 || n >= 100 { return 0, false } return float64(n) / 100, true @@ -616,7 +618,22 @@ func (r *EtcdDefragReconciler) fail(ctx context.Context, df *lll.EtcdDefrag, rea } func (r *EtcdDefragReconciler) persistAndRequeue(ctx context.Context, df *lll.EtcdDefrag) (ctrl.Result, error) { - if err := r.Status().Update(ctx, df); err != nil { + // A Defragment RPC can block for defragRPCTimeout; an unrelated metadata write + // (a kubectl label/annotate) in that window would make this status write + // conflict and discard the recorded outcome, so the member would be + // defragmented again next pass. The defrag controller owns these status + // fields, so on conflict re-fetch and re-apply the computed status rather than + // dropping it. + desired := df.Status + err := retry.RetryOnConflict(retry.DefaultRetry, func() error { + latest := &lll.EtcdDefrag{} + if err := r.Get(ctx, client.ObjectKeyFromObject(df), latest); err != nil { + return err + } + latest.Status = desired + return r.Status().Update(ctx, latest) + }) + if err != nil { return ctrl.Result{}, err } return ctrl.Result{RequeueAfter: defragRequeueAfter}, nil diff --git a/controllers/etcddefrag_controller_test.go b/controllers/etcddefrag_controller_test.go index 499e9995..256dc227 100644 --- a/controllers/etcddefrag_controller_test.go +++ b/controllers/etcddefrag_controller_test.go @@ -48,6 +48,7 @@ func TestDefragRuleTriggered(t *testing.T) { {"freeSpaceAbove met", &lll.DefragRule{FreeSpaceAbove: q("200Mi")}, 500 << 20, 100 << 20, true}, {"quota arm: full but unfragmented never fires", &lll.DefragRule{QuotaUsageAbove: "80%"}, int64(1.9 * float64(gib)), int64(1.9 * float64(gib)), false}, {"quota arm: under pressure with reclaimable fires", &lll.DefragRule{QuotaUsageAbove: "80%", MinReclaim: q("32Mi")}, int64(1.9 * float64(gib)), int64(1.9*float64(gib)) - (64 << 20), true}, + {"quota arm: 100% is inert, only the free-space floor applies", &lll.DefragRule{QuotaUsageAbove: "100%", MinReclaim: q("32Mi")}, int64(1.9 * float64(gib)), int64(1.9*float64(gib)) - (64 << 20), false}, } for _, tc := range cases { t.Run(tc.name, func(t *testing.T) {