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..e69fea64 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"` @@ -71,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"` @@ -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..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 @@ -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: |- @@ -114,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: @@ -136,9 +133,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/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..1ddd9cff 100644 --- a/controllers/etcd_client.go +++ b/controllers/etcd_client.go @@ -35,13 +35,30 @@ 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) + + // 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 new file mode 100644 index 00000000..ff049685 --- /dev/null +++ b/controllers/etcddefrag_controller.go @@ -0,0 +1,699 @@ +/* +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" + "sort" + "strconv" + "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" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "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" + "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 + + // 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. 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). + defragMaxConcurrentReconciles = 4 +) + +// 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 + } + + 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 { + // 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 { + 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) + } + // 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 + } + 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, c) + } + + 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. 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 + } + // 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") + + 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) + 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. 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, outcomeReason, b) + df.Status.Defragmented++ + 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; 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 +// 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}) + } + 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) +} + +// disarmNoSpaceAlarms clears every armed NOSPACE alarm. The health gate lets a +// NOSPACE cluster through so the run can reclaim its backend space; etcd keeps +// 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() + resp, err := c.AlarmList(lctx) + if err != nil { + return err + } + var errs error + 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 { + errs = errors.Join(errs, fmt.Errorf("member %d: %w", a.MemberID, derr)) + } + } + 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 +// 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 { + 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 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 + } + if leader == 0 { + leader = b.status.Leader + } else if b.status.Leader != leader { + return "members disagree on the leader", false + } + } + 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). +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 + if b.afterKnown { + 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–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 { + 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, c EtcdClusterClient) (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 + } + } + // A cluster admitted with a NOSPACE alarm stays read-only until the alarm is + // 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 + 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) +} + +// 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() + 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) { + // 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 +} + +// 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{}). + // 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 new file mode 100644 index 00000000..256dc227 --- /dev/null +++ b/controllers/etcddefrag_controller_test.go @@ -0,0 +1,607 @@ +/* +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" + "time" + + 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}, + {"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) { + 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} +} + +// 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 { + 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) + } +} + +// 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) + } +} + +// 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" { + 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..2ab616e8 100644 --- a/controllers/testing_helpers_test.go +++ b/controllers/testing_helpers_test.go @@ -48,6 +48,26 @@ 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 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; + // disarmErrByMember fails AlarmDisarm for a specific member. + alarms []*etcdserverpb.AlarmMember + alarmListErr error + disarmCalls []*etcdserverpb.AlarmMember + disarmErrByMember map[uint64]error + addCalls []string addLearnerCalls []string promoteCalls []uint64 @@ -181,16 +201,55 @@ 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 + } + if err := f.defragErrByEndpoint[endpoint]; err != nil { + return nil, err + } + 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)) + if err := f.disarmErrByMember[m.MemberID]; err != nil { + return nil, err + } + 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 { @@ -251,7 +310,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..12624f9c 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 @@ -19,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) @@ -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** @@ -110,10 +103,14 @@ 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 (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 +118,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 @@ -134,10 +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); and a defrag - that doesn't shrink `DbSize` is backed off rather than repeated. +- **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 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) + } +}