From 0680092fa056f6895aabb632b30a6150b7129d87 Mon Sep 17 00:00:00 2001 From: Andrey Kolkov Date: Mon, 17 Aug 2026 21:53:46 +0400 Subject: [PATCH] =?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) + } +}