From 50f8345618c8d912bc5824f6608d03429a1b949b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ver=C3=B3nica=20L=C3=B3pez?= Date: Fri, 25 Sep 2026 00:45:06 -0600 Subject: [PATCH 1/3] fix(topology): skip unchanged writes and compact etcd history MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Compare owned topology fields before updating cell and database records. - Configure MVCC compaction with one-hour retention and an explicit backend quota. - Add opt-in defragmentation with health checks and persisted maintenance reservations. Signed-off-by: Verónica López --- api/v1alpha1/etcd_maintenance.go | 53 +++ api/v1alpha1/etcd_maintenance_test.go | 31 ++ api/v1alpha1/toposerver_types.go | 56 +++ api/v1alpha1/zz_generated.deepcopy.go | 51 +++ config/crd/bases/multigres.com_cells.yaml | 47 +++ .../bases/multigres.com_celltemplates.yaml | 47 +++ .../bases/multigres.com_coretemplates.yaml | 47 +++ .../multigres.com_multigresclusters.yaml | 95 +++++ .../crd/bases/multigres.com_toposervers.yaml | 69 ++++ docs/README.md | 1 + docs/topology-maintenance.md | 112 ++++++ go.mod | 4 +- .../multigrescluster/builders_global.go | 11 +- .../multigrescluster_controller.go | 6 +- .../multigrescluster/reconcile_cells.go | 9 +- .../reconcile_topology_test.go | 58 +++ pkg/data-handler/topo/cell.go | 7 +- pkg/data-handler/topo/database.go | 7 +- pkg/data-handler/topo/idempotency_test.go | 119 +++++++ pkg/data-handler/topo/metadata.go | 28 ++ pkg/data-handler/topo/topology.go | 9 +- pkg/resolver/cell.go | 3 + pkg/resolver/cluster.go | 23 ++ pkg/resolver/etcd_maintenance_test.go | 35 ++ .../controller/toposerver/container_env.go | 18 + .../toposerver/container_env_test.go | 31 ++ .../controller/toposerver/integration_test.go | 6 + .../controller/toposerver/maintenance.go | 337 ++++++++++++++++++ .../toposerver/maintenance_client.go | 135 +++++++ .../toposerver/maintenance_client_test.go | 84 +++++ .../toposerver/maintenance_etcd_test.go | 271 ++++++++++++++ .../maintenance_status_integration_test.go | 72 ++++ .../controller/toposerver/maintenance_test.go | 283 +++++++++++++++ .../controller/toposerver/statefulset.go | 12 +- .../controller/toposerver/statefulset_test.go | 6 + .../toposerver/toposerver_controller.go | 26 +- .../toposerver_controller_internal_test.go | 30 ++ .../etcd_maintenance_validation_test.go | 76 ++++ 38 files changed, 2289 insertions(+), 26 deletions(-) create mode 100644 api/v1alpha1/etcd_maintenance.go create mode 100644 api/v1alpha1/etcd_maintenance_test.go create mode 100644 docs/topology-maintenance.md create mode 100644 pkg/data-handler/topo/idempotency_test.go create mode 100644 pkg/data-handler/topo/metadata.go create mode 100644 pkg/resolver/etcd_maintenance_test.go create mode 100644 pkg/resource-handler/controller/toposerver/maintenance.go create mode 100644 pkg/resource-handler/controller/toposerver/maintenance_client.go create mode 100644 pkg/resource-handler/controller/toposerver/maintenance_client_test.go create mode 100644 pkg/resource-handler/controller/toposerver/maintenance_etcd_test.go create mode 100644 pkg/resource-handler/controller/toposerver/maintenance_status_integration_test.go create mode 100644 pkg/resource-handler/controller/toposerver/maintenance_test.go create mode 100644 pkg/webhook/etcd_maintenance_validation_test.go diff --git a/api/v1alpha1/etcd_maintenance.go b/api/v1alpha1/etcd_maintenance.go new file mode 100644 index 000000000..3e1f9c886 --- /dev/null +++ b/api/v1alpha1/etcd_maintenance.go @@ -0,0 +1,53 @@ +package v1alpha1 + +import ( + "fmt" + "strconv" + "time" +) + +// EffectiveCompaction returns etcd's configured mode and retention, including +// defaults. Validate here as well as in the CRD for callers without admission. +func (c *EtcdMaintenanceConfig) EffectiveCompaction() (mode, retention string, err error) { + mode = "periodic" + if c != nil && c.AutoCompactionMode != "" { + mode = c.AutoCompactionMode + } + switch mode { + case "periodic": + retention = "1h" + case "revision": + retention = "10000" + default: + return "", "", fmt.Errorf("invalid etcd compaction mode %q", mode) + } + if c != nil && c.AutoCompactionRetention != "" { + retention = c.AutoCompactionRetention + } + if mode == "periodic" { + d, parseErr := time.ParseDuration(retention) + if parseErr != nil || d <= 0 { + return "", "", fmt.Errorf("invalid periodic etcd compaction retention %q", retention) + } + } else { + n, parseErr := strconv.ParseInt(retention, 10, 64) + if parseErr != nil || n <= 0 { + return "", "", fmt.Errorf("invalid revision etcd compaction retention %q", retention) + } + } + return mode, retention, nil +} + +// EffectiveQuotaBackendBytes preserves the pre-existing etcd quota unless an +// administrator explicitly supplies one. Quota cannot predict restore memory. +func (c *EtcdMaintenanceConfig) EffectiveQuotaBackendBytes() int64 { + if c == nil || c.QuotaBackendBytes == nil { + return 2 * 1024 * 1024 * 1024 + } + return *c.QuotaBackendBytes +} + +// DefragmentationIsEnabled reports whether automatic defragmentation is opted in. +func (c *EtcdMaintenanceConfig) DefragmentationIsEnabled() bool { + return c != nil && c.DefragmentationEnabled != nil && *c.DefragmentationEnabled +} diff --git a/api/v1alpha1/etcd_maintenance_test.go b/api/v1alpha1/etcd_maintenance_test.go new file mode 100644 index 000000000..f4ddd9858 --- /dev/null +++ b/api/v1alpha1/etcd_maintenance_test.go @@ -0,0 +1,31 @@ +package v1alpha1 + +import "testing" + +func TestEffectiveCompaction(t *testing.T) { + for _, tc := range []struct { + name string + config *EtcdMaintenanceConfig + mode, retention string + wantErr bool + }{ + {"omitted", nil, "periodic", "1h", false}, + {"empty", &EtcdMaintenanceConfig{}, "periodic", "1h", false}, + {"periodic", &EtcdMaintenanceConfig{AutoCompactionRetention: "30m"}, "periodic", "30m", false}, + {"revision default", &EtcdMaintenanceConfig{AutoCompactionMode: "revision"}, "revision", "10000", false}, + {"revision", &EtcdMaintenanceConfig{AutoCompactionMode: "revision", AutoCompactionRetention: "5000"}, "revision", "5000", false}, + {"unknown mode", &EtcdMaintenanceConfig{AutoCompactionMode: "bad"}, "", "", true}, + {"disabled retention", &EtcdMaintenanceConfig{AutoCompactionRetention: "0h"}, "", "", true}, + {"negative retention", &EtcdMaintenanceConfig{AutoCompactionRetention: "-1h"}, "", "", true}, + {"invalid duration", &EtcdMaintenanceConfig{AutoCompactionRetention: "forever"}, "", "", true}, + {"revision duration", &EtcdMaintenanceConfig{AutoCompactionMode: "revision", AutoCompactionRetention: "1h"}, "", "", true}, + {"zero revisions", &EtcdMaintenanceConfig{AutoCompactionMode: "revision", AutoCompactionRetention: "0"}, "", "", true}, + } { + t.Run(tc.name, func(t *testing.T) { + mode, retention, err := tc.config.EffectiveCompaction() + if (err != nil) != tc.wantErr || mode != tc.mode || retention != tc.retention { + t.Fatalf("got (%q,%q,%v)", mode, retention, err) + } + }) + } +} diff --git a/api/v1alpha1/toposerver_types.go b/api/v1alpha1/toposerver_types.go index d6ecd4173..52c9aa8ed 100644 --- a/api/v1alpha1/toposerver_types.go +++ b/api/v1alpha1/toposerver_types.go @@ -96,6 +96,22 @@ type TopoServerStatus struct { // PeerService is the name of the service for peers. // +optional PeerService string `json:"peerService,omitempty"` + + // EtcdMaintenance records a maintenance reservation before contacting etcd, + // preventing overlapping defragmentation across reconciles and restarts. + // +optional + EtcdMaintenance *EtcdMaintenanceStatus `json:"etcdMaintenance,omitempty"` +} + +// EtcdMaintenanceStatus tracks the most recent automatic defragmentation. +type EtcdMaintenanceStatus struct { + // LastAttemptTime starts the one-hour minimum interval between members. + LastAttemptTime metav1.Time `json:"lastAttemptTime"` + // Endpoint identifies the member reserved for defragmentation. + Endpoint string `json:"endpoint"` + // InProgress remains true after an interrupted or uncertain operation until + // all members pass health checks again. While true, pod rollouts are paused. + InProgress bool `json:"inProgress"` } // ============================================================================ @@ -127,6 +143,12 @@ type EtcdSpec struct { // +optional Resources corev1.ResourceRequirements `json:"resources,omitempty"` + // Maintenance configures MVCC history retention, backend quota, and optional + // defragmentation. Applies only to operator-managed etcd. Omitted fields use + // the defaults documented on EtcdMaintenanceConfig. + // +optional + Maintenance *EtcdMaintenanceConfig `json:"maintenance,omitempty"` + // RootPath is the etcd prefix for this cluster. // +optional // +kubebuilder:validation:MinLength=1 @@ -139,6 +161,40 @@ type EtcdSpec struct { PVCDeletionPolicy *PVCDeletionPolicy `json:"pvcDeletionPolicy,omitempty"` } +// EtcdMaintenanceConfig controls storage maintenance independently of logical +// topology pruning. Defaults are applied at rendering time so template values +// remain inheritable. Changing compaction or quota rolls the etcd pods. +// +kubebuilder:validation:XValidation:rule="!has(self.autoCompactionRetention) || ((has(self.autoCompactionMode) && self.autoCompactionMode == 'revision') ? self.autoCompactionRetention.matches('^[1-9][0-9]{0,9}$') : self.autoCompactionRetention.matches('^([0-9]{1,6}h)?([0-9]{1,6}m)?([0-9]{1,6}s)?$'))",message="compaction retention must be a positive revision count in revision mode or a duration using h, m, s in periodic mode" +type EtcdMaintenanceConfig struct { + // AutoCompactionMode defaults to periodic (time-based retention). + // +kubebuilder:validation:Enum=periodic;revision + // +optional + AutoCompactionMode string `json:"autoCompactionMode,omitempty"` + + // AutoCompactionRetention defaults to 1h in periodic mode or 10000 in + // revision mode. Periodic values must be positive durations (e.g. 30m, 1h). + // +kubebuilder:validation:MinLength=1 + // +kubebuilder:validation:MaxLength=32 + // +kubebuilder:validation:XValidation:rule="self.matches('[1-9]')",message="compaction retention must be positive" + // +optional + AutoCompactionRetention string `json:"autoCompactionRetention,omitempty"` + + // QuotaBackendBytes defaults to 2147483648 (2 GiB), preserving etcd's + // existing default. This is a storage quota, not a memory limit: provision + // memory for measured restore-time usage and disk for the backend and WAL. + // Do not lower the quota below an existing backend's size. + // +kubebuilder:validation:Minimum=1048576 + // +kubebuilder:validation:Maximum=8589934592 + // +optional + QuotaBackendBytes *int64 `json:"quotaBackendBytes,omitempty"` + + // DefragmentationEnabled opts into hourly, health-gated maintenance of at + // most one member. Requires at least three healthy members, no rollout, + // and at least 100 MiB and 30% reclaimable space. Disabled by default. + // +optional + DefragmentationEnabled *bool `json:"defragmentationEnabled,omitempty"` +} + // GlobalTopoServerSpec defines the configuration for the global topology server. // It can be either an inline Etcd spec, an External reference, or a Template reference. // +kubebuilder:validation:XValidation:rule="[has(self.etcd), has(self.external), has(self.templateRef)].filter(x, x).size() == 1",message="must specify exactly one of 'etcd', 'external', or 'templateRef'" diff --git a/api/v1alpha1/zz_generated.deepcopy.go b/api/v1alpha1/zz_generated.deepcopy.go index 3b3481986..dc88876c7 100644 --- a/api/v1alpha1/zz_generated.deepcopy.go +++ b/api/v1alpha1/zz_generated.deepcopy.go @@ -611,6 +611,47 @@ func (in *DatabaseStatusSummary) DeepCopy() *DatabaseStatusSummary { return out } +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *EtcdMaintenanceConfig) DeepCopyInto(out *EtcdMaintenanceConfig) { + *out = *in + if in.QuotaBackendBytes != nil { + in, out := &in.QuotaBackendBytes, &out.QuotaBackendBytes + *out = new(int64) + **out = **in + } + if in.DefragmentationEnabled != nil { + in, out := &in.DefragmentationEnabled, &out.DefragmentationEnabled + *out = new(bool) + **out = **in + } +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new EtcdMaintenanceConfig. +func (in *EtcdMaintenanceConfig) DeepCopy() *EtcdMaintenanceConfig { + if in == nil { + return nil + } + out := new(EtcdMaintenanceConfig) + in.DeepCopyInto(out) + return out +} + +// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. +func (in *EtcdMaintenanceStatus) DeepCopyInto(out *EtcdMaintenanceStatus) { + *out = *in + in.LastAttemptTime.DeepCopyInto(&out.LastAttemptTime) +} + +// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new EtcdMaintenanceStatus. +func (in *EtcdMaintenanceStatus) DeepCopy() *EtcdMaintenanceStatus { + if in == nil { + return nil + } + out := new(EtcdMaintenanceStatus) + in.DeepCopyInto(out) + return out +} + // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *EtcdSpec) DeepCopyInto(out *EtcdSpec) { *out = *in @@ -621,6 +662,11 @@ func (in *EtcdSpec) DeepCopyInto(out *EtcdSpec) { } in.Storage.DeepCopyInto(&out.Storage) in.Resources.DeepCopyInto(&out.Resources) + if in.Maintenance != nil { + in, out := &in.Maintenance, &out.Maintenance + *out = new(EtcdMaintenanceConfig) + (*in).DeepCopyInto(*out) + } if in.PVCDeletionPolicy != nil { in, out := &in.PVCDeletionPolicy, &out.PVCDeletionPolicy *out = new(PVCDeletionPolicy) @@ -2380,6 +2426,11 @@ func (in *TopoServerStatus) DeepCopyInto(out *TopoServerStatus) { (*in)[i].DeepCopyInto(&(*out)[i]) } } + if in.EtcdMaintenance != nil { + in, out := &in.EtcdMaintenance, &out.EtcdMaintenance + *out = new(EtcdMaintenanceStatus) + (*in).DeepCopyInto(*out) + } } // DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new TopoServerStatus. diff --git a/config/crd/bases/multigres.com_cells.yaml b/config/crd/bases/multigres.com_cells.yaml index 050cd7b86..ef2d1ed14 100644 --- a/config/crd/bases/multigres.com_cells.yaml +++ b/config/crd/bases/multigres.com_cells.yaml @@ -1453,6 +1453,53 @@ spec: maxLength: 512 minLength: 1 type: string + maintenance: + description: |- + Maintenance configures MVCC history retention, backend quota, and optional + defragmentation. Applies only to operator-managed etcd. Omitted fields use + the defaults documented on EtcdMaintenanceConfig. + properties: + autoCompactionMode: + description: AutoCompactionMode defaults to periodic (time-based + retention). + enum: + - periodic + - revision + type: string + autoCompactionRetention: + description: |- + AutoCompactionRetention defaults to 1h in periodic mode or 10000 in + revision mode. Periodic values must be positive durations (e.g. 30m, 1h). + maxLength: 32 + minLength: 1 + type: string + x-kubernetes-validations: + - message: compaction retention must be positive + rule: self.matches('[1-9]') + defragmentationEnabled: + description: |- + DefragmentationEnabled opts into hourly, health-gated maintenance of at + most one member. Requires at least three healthy members, no rollout, + and at least 100 MiB and 30% reclaimable space. Disabled by default. + type: boolean + quotaBackendBytes: + description: |- + QuotaBackendBytes defaults to 2147483648 (2 GiB), preserving etcd's + existing default. This is a storage quota, not a memory limit: provision + memory for measured restore-time usage and disk for the backend and WAL. + Do not lower the quota below an existing backend's size. + format: int64 + maximum: 8589934592 + minimum: 1048576 + type: integer + type: object + x-kubernetes-validations: + - message: compaction retention must be a positive revision + count in revision mode or a duration using h, m, s in + periodic mode + rule: '!has(self.autoCompactionRetention) || ((has(self.autoCompactionMode) + && self.autoCompactionMode == ''revision'') ? self.autoCompactionRetention.matches(''^[1-9][0-9]{0,9}$'') + : self.autoCompactionRetention.matches(''^([0-9]{1,6}h)?([0-9]{1,6}m)?([0-9]{1,6}s)?$''))' pvcDeletionPolicy: description: |- PVCDeletionPolicy controls PVC lifecycle for etcd volumes. diff --git a/config/crd/bases/multigres.com_celltemplates.yaml b/config/crd/bases/multigres.com_celltemplates.yaml index 7f2639c21..b2f949163 100644 --- a/config/crd/bases/multigres.com_celltemplates.yaml +++ b/config/crd/bases/multigres.com_celltemplates.yaml @@ -54,6 +54,53 @@ spec: maxLength: 512 minLength: 1 type: string + maintenance: + description: |- + Maintenance configures MVCC history retention, backend quota, and optional + defragmentation. Applies only to operator-managed etcd. Omitted fields use + the defaults documented on EtcdMaintenanceConfig. + properties: + autoCompactionMode: + description: AutoCompactionMode defaults to periodic (time-based + retention). + enum: + - periodic + - revision + type: string + autoCompactionRetention: + description: |- + AutoCompactionRetention defaults to 1h in periodic mode or 10000 in + revision mode. Periodic values must be positive durations (e.g. 30m, 1h). + maxLength: 32 + minLength: 1 + type: string + x-kubernetes-validations: + - message: compaction retention must be positive + rule: self.matches('[1-9]') + defragmentationEnabled: + description: |- + DefragmentationEnabled opts into hourly, health-gated maintenance of at + most one member. Requires at least three healthy members, no rollout, + and at least 100 MiB and 30% reclaimable space. Disabled by default. + type: boolean + quotaBackendBytes: + description: |- + QuotaBackendBytes defaults to 2147483648 (2 GiB), preserving etcd's + existing default. This is a storage quota, not a memory limit: provision + memory for measured restore-time usage and disk for the backend and WAL. + Do not lower the quota below an existing backend's size. + format: int64 + maximum: 8589934592 + minimum: 1048576 + type: integer + type: object + x-kubernetes-validations: + - message: compaction retention must be a positive revision + count in revision mode or a duration using h, m, s in + periodic mode + rule: '!has(self.autoCompactionRetention) || ((has(self.autoCompactionMode) + && self.autoCompactionMode == ''revision'') ? self.autoCompactionRetention.matches(''^[1-9][0-9]{0,9}$'') + : self.autoCompactionRetention.matches(''^([0-9]{1,6}h)?([0-9]{1,6}m)?([0-9]{1,6}s)?$''))' pvcDeletionPolicy: description: |- PVCDeletionPolicy controls PVC lifecycle for etcd volumes. diff --git a/config/crd/bases/multigres.com_coretemplates.yaml b/config/crd/bases/multigres.com_coretemplates.yaml index c27c18456..237ec4af5 100644 --- a/config/crd/bases/multigres.com_coretemplates.yaml +++ b/config/crd/bases/multigres.com_coretemplates.yaml @@ -56,6 +56,53 @@ spec: maxLength: 512 minLength: 1 type: string + maintenance: + description: |- + Maintenance configures MVCC history retention, backend quota, and optional + defragmentation. Applies only to operator-managed etcd. Omitted fields use + the defaults documented on EtcdMaintenanceConfig. + properties: + autoCompactionMode: + description: AutoCompactionMode defaults to periodic (time-based + retention). + enum: + - periodic + - revision + type: string + autoCompactionRetention: + description: |- + AutoCompactionRetention defaults to 1h in periodic mode or 10000 in + revision mode. Periodic values must be positive durations (e.g. 30m, 1h). + maxLength: 32 + minLength: 1 + type: string + x-kubernetes-validations: + - message: compaction retention must be positive + rule: self.matches('[1-9]') + defragmentationEnabled: + description: |- + DefragmentationEnabled opts into hourly, health-gated maintenance of at + most one member. Requires at least three healthy members, no rollout, + and at least 100 MiB and 30% reclaimable space. Disabled by default. + type: boolean + quotaBackendBytes: + description: |- + QuotaBackendBytes defaults to 2147483648 (2 GiB), preserving etcd's + existing default. This is a storage quota, not a memory limit: provision + memory for measured restore-time usage and disk for the backend and WAL. + Do not lower the quota below an existing backend's size. + format: int64 + maximum: 8589934592 + minimum: 1048576 + type: integer + type: object + x-kubernetes-validations: + - message: compaction retention must be a positive revision + count in revision mode or a duration using h, m, s in + periodic mode + rule: '!has(self.autoCompactionRetention) || ((has(self.autoCompactionMode) + && self.autoCompactionMode == ''revision'') ? self.autoCompactionRetention.matches(''^[1-9][0-9]{0,9}$'') + : self.autoCompactionRetention.matches(''^([0-9]{1,6}h)?([0-9]{1,6}m)?([0-9]{1,6}s)?$''))' pvcDeletionPolicy: description: |- PVCDeletionPolicy controls PVC lifecycle for etcd volumes. diff --git a/config/crd/bases/multigres.com_multigresclusters.yaml b/config/crd/bases/multigres.com_multigresclusters.yaml index 1120e4683..866611310 100644 --- a/config/crd/bases/multigres.com_multigresclusters.yaml +++ b/config/crd/bases/multigres.com_multigresclusters.yaml @@ -1370,6 +1370,54 @@ spec: maxLength: 512 minLength: 1 type: string + maintenance: + description: |- + Maintenance configures MVCC history retention, backend quota, and optional + defragmentation. Applies only to operator-managed etcd. Omitted fields use + the defaults documented on EtcdMaintenanceConfig. + properties: + autoCompactionMode: + description: AutoCompactionMode defaults to + periodic (time-based retention). + enum: + - periodic + - revision + type: string + autoCompactionRetention: + description: |- + AutoCompactionRetention defaults to 1h in periodic mode or 10000 in + revision mode. Periodic values must be positive durations (e.g. 30m, 1h). + maxLength: 32 + minLength: 1 + type: string + x-kubernetes-validations: + - message: compaction retention must be positive + rule: self.matches('[1-9]') + defragmentationEnabled: + description: |- + DefragmentationEnabled opts into hourly, health-gated maintenance of at + most one member. Requires at least three healthy members, no rollout, + and at least 100 MiB and 30% reclaimable space. Disabled by default. + type: boolean + quotaBackendBytes: + description: |- + QuotaBackendBytes defaults to 2147483648 (2 GiB), preserving etcd's + existing default. This is a storage quota, not a memory limit: provision + memory for measured restore-time usage and disk for the backend and WAL. + Do not lower the quota below an existing backend's size. + format: int64 + maximum: 8589934592 + minimum: 1048576 + type: integer + type: object + x-kubernetes-validations: + - message: compaction retention must be a positive + revision count in revision mode or a duration + using h, m, s in periodic mode + rule: '!has(self.autoCompactionRetention) || ((has(self.autoCompactionMode) + && self.autoCompactionMode == ''revision'') + ? self.autoCompactionRetention.matches(''^[1-9][0-9]{0,9}$'') + : self.autoCompactionRetention.matches(''^([0-9]{1,6}h)?([0-9]{1,6}m)?([0-9]{1,6}s)?$''))' pvcDeletionPolicy: description: |- PVCDeletionPolicy controls PVC lifecycle for etcd volumes. @@ -8352,6 +8400,53 @@ spec: maxLength: 512 minLength: 1 type: string + maintenance: + description: |- + Maintenance configures MVCC history retention, backend quota, and optional + defragmentation. Applies only to operator-managed etcd. Omitted fields use + the defaults documented on EtcdMaintenanceConfig. + properties: + autoCompactionMode: + description: AutoCompactionMode defaults to periodic (time-based + retention). + enum: + - periodic + - revision + type: string + autoCompactionRetention: + description: |- + AutoCompactionRetention defaults to 1h in periodic mode or 10000 in + revision mode. Periodic values must be positive durations (e.g. 30m, 1h). + maxLength: 32 + minLength: 1 + type: string + x-kubernetes-validations: + - message: compaction retention must be positive + rule: self.matches('[1-9]') + defragmentationEnabled: + description: |- + DefragmentationEnabled opts into hourly, health-gated maintenance of at + most one member. Requires at least three healthy members, no rollout, + and at least 100 MiB and 30% reclaimable space. Disabled by default. + type: boolean + quotaBackendBytes: + description: |- + QuotaBackendBytes defaults to 2147483648 (2 GiB), preserving etcd's + existing default. This is a storage quota, not a memory limit: provision + memory for measured restore-time usage and disk for the backend and WAL. + Do not lower the quota below an existing backend's size. + format: int64 + maximum: 8589934592 + minimum: 1048576 + type: integer + type: object + x-kubernetes-validations: + - message: compaction retention must be a positive revision + count in revision mode or a duration using h, m, s in + periodic mode + rule: '!has(self.autoCompactionRetention) || ((has(self.autoCompactionMode) + && self.autoCompactionMode == ''revision'') ? self.autoCompactionRetention.matches(''^[1-9][0-9]{0,9}$'') + : self.autoCompactionRetention.matches(''^([0-9]{1,6}h)?([0-9]{1,6}m)?([0-9]{1,6}s)?$''))' pvcDeletionPolicy: description: |- PVCDeletionPolicy controls PVC lifecycle for etcd volumes. diff --git a/config/crd/bases/multigres.com_toposervers.yaml b/config/crd/bases/multigres.com_toposervers.yaml index 1c8c0dd38..d71d98929 100644 --- a/config/crd/bases/multigres.com_toposervers.yaml +++ b/config/crd/bases/multigres.com_toposervers.yaml @@ -58,6 +58,52 @@ spec: maxLength: 512 minLength: 1 type: string + maintenance: + description: |- + Maintenance configures MVCC history retention, backend quota, and optional + defragmentation. Applies only to operator-managed etcd. Omitted fields use + the defaults documented on EtcdMaintenanceConfig. + properties: + autoCompactionMode: + description: AutoCompactionMode defaults to periodic (time-based + retention). + enum: + - periodic + - revision + type: string + autoCompactionRetention: + description: |- + AutoCompactionRetention defaults to 1h in periodic mode or 10000 in + revision mode. Periodic values must be positive durations (e.g. 30m, 1h). + maxLength: 32 + minLength: 1 + type: string + x-kubernetes-validations: + - message: compaction retention must be positive + rule: self.matches('[1-9]') + defragmentationEnabled: + description: |- + DefragmentationEnabled opts into hourly, health-gated maintenance of at + most one member. Requires at least three healthy members, no rollout, + and at least 100 MiB and 30% reclaimable space. Disabled by default. + type: boolean + quotaBackendBytes: + description: |- + QuotaBackendBytes defaults to 2147483648 (2 GiB), preserving etcd's + existing default. This is a storage quota, not a memory limit: provision + memory for measured restore-time usage and disk for the backend and WAL. + Do not lower the quota below an existing backend's size. + format: int64 + maximum: 8589934592 + minimum: 1048576 + type: integer + type: object + x-kubernetes-validations: + - message: compaction retention must be a positive revision count + in revision mode or a duration using h, m, s in periodic mode + rule: '!has(self.autoCompactionRetention) || ((has(self.autoCompactionMode) + && self.autoCompactionMode == ''revision'') ? self.autoCompactionRetention.matches(''^[1-9][0-9]{0,9}$'') + : self.autoCompactionRetention.matches(''^([0-9]{1,6}h)?([0-9]{1,6}m)?([0-9]{1,6}s)?$''))' pvcDeletionPolicy: description: |- PVCDeletionPolicy controls PVC lifecycle for etcd volumes. @@ -1451,6 +1497,29 @@ spec: x-kubernetes-list-map-keys: - type x-kubernetes-list-type: map + etcdMaintenance: + description: |- + EtcdMaintenance records a maintenance reservation before contacting etcd, + preventing overlapping defragmentation across reconciles and restarts. + properties: + endpoint: + description: Endpoint identifies the member reserved for defragmentation. + type: string + inProgress: + description: |- + InProgress remains true after an interrupted or uncertain operation until + all members pass health checks again. While true, pod rollouts are paused. + type: boolean + lastAttemptTime: + description: LastAttemptTime starts the one-hour minimum interval + between members. + format: date-time + type: string + required: + - endpoint + - inProgress + - lastAttemptTime + type: object message: description: Message provides details about the current phase. type: string diff --git a/docs/README.md b/docs/README.md index 1e41e6605..744a2d923 100644 --- a/docs/README.md +++ b/docs/README.md @@ -10,6 +10,7 @@ | [Runtime Log Schema](observability/log-schema.md) | Field definitions for enriched Multigres runtime logs | | [Durability Policy](durability-policy.md) | Durability policies, cross-AZ quorum, per-database overrides | | [Storage Management](storage.md) | PVC deletion policies and volume expansion | +| [Topology Maintenance](topology-maintenance.md) | etcd compaction, backend quota, and guarded defragmentation | | [External Gateway](external-gateway.md) | External multigateway exposure, DNS wiring, GatewayExternalReady condition | | [External Admin Web](external-admin-web.md) | External multiadmin-web exposure, AdminWebExternalReady condition | | [PostgreSQL Initialization](postgresql-initialization.md) | Custom initdb arguments (locale, encoding, WAL segment size) | diff --git a/docs/topology-maintenance.md b/docs/topology-maintenance.md new file mode 100644 index 000000000..799f80063 --- /dev/null +++ b/docs/topology-maintenance.md @@ -0,0 +1,112 @@ +# Topology storage maintenance + +Managed global and cell-local topology servers enable etcd MVCC auto-compaction +with a one-hour retention window. The operator compares existing cell and +database metadata before updating it, so an unchanged registration does not +consume a new etcd revision. Changes to owned fields and deleted records are +still reconciled; other writers' fields are preserved. + +`spec.topologyPruning` deletes logical topology records that are no longer +desired. It does not compact their revision history or return unused backend +pages to the filesystem. These are separate operations: + +| Operation | Effect | +| --- | --- | +| Topology pruning | Removes stale logical records | +| MVCC compaction | Discards superseded history and makes pages reusable | +| Defragmentation | Returns unused backend pages to the filesystem | + +Compaction bounds retained history for a bounded workload. The latest revision +number still increases with writes; live data can still grow. Use backend bytes, +bytes in use, memory utilization, and write rate to assess storage health. +[etcd maintenance reference](https://etcd.io/docs/v3.6/op-guide/maintenance/) + +## Configuration + +Set `maintenance` on the managed `etcd` spec, including a `CoreTemplate`, a +`CellTemplate`'s `localTopoServer`, or a standalone `TopoServer`: + +```yaml +spec: + globalTopoServer: + etcd: + maintenance: + autoCompactionMode: periodic + autoCompactionRetention: 1h + quotaBackendBytes: 2147483648 + defragmentationEnabled: false +``` + +This fragment shows defaults, not a production resource-sizing recommendation. +For cell-local topology, place the same block under +`spec.cells[].spec.localTopoServer.etcd`. External topology is not modified. + +| Field | Default | Meaning | +| --- | --- | --- | +| `autoCompactionMode` | `periodic` | `periodic` or `revision` | +| `autoCompactionRetention` | `1h` / `10000` | Positive h/m/s duration in periodic mode; positive revision count in revision mode | +| `quotaBackendBytes` | `2147483648` | Backend quota in bytes, from 1 MiB through 8 GiB | +| `defragmentationEnabled` | `false` | Opt into maintenance of healthy clusters with at least three members | + +Global topology template overrides merge explicitly supplied fields. Set both +mode and retention when switching a template's compaction mode or overriding a +revision count. Explicit `false` +disables defragmentation inherited from a template. Local topology follows the +existing whole-block replacement behavior for an inline `localTopoServer`. + +## Upgrade and sizing + +Applying compaction or quota settings changes the StatefulSet pod template and +causes a rolling update. Apply to healthy topology clusters first. This change +does not implement recovery of an already stalled, fully-down StatefulSet. + +The explicit default quota preserves etcd's existing 2 GiB default. It is not +lowered automatically based on memory limits or PVC size: lowering a quota below +an existing database can cause `NOSPACE`. Check current backend size before +choosing a smaller quota. Compaction can reuse pages without reducing the +physical file; defragmentation is needed to reduce that file. + +The existing 512 MiB memory limit and 1 GiB resolver storage default are **not a +capacity guarantee for a 2 GiB quota**. Provision disk for backend, WAL, and +maintenance headroom. Size memory using restore-time measurements at the target +backend size; quota is not a reliable memory sizing formula. This change leaves +resource sizing to the separate defaults work. + +Watchers resuming before the compaction boundary must re-list current state and +restart their watches. Choose a retention window that accommodates expected +disconnects, and validate custom topology consumers before shortening it. The +Multigres topology watch retry helpers obtain a fresh snapshot on reconnect. + +## Automatic defragmentation + +When enabled, the TopoServer controller checks for maintenance at least hourly +and defragments at most one member per hour. A candidate must have both at least +100 MiB and 30% reclaimable space. It requires: + +- At least three expected voting members, all reporting the same cluster and + leader, with successful linearizable reads through every member. +- Every owned pod Ready for at least one minute, with no terminating pod or + pending StatefulSet rollout. +- No member reporting an etcd error, including a backend quota alarm. + +The controller transfers leadership before defragmenting a leader and verifies +health again before and after maintenance. It does not delete pods, alter etcd +membership, or delete PVCs. Single-member topology requires planned manual +maintenance because defragmentation blocks the member being processed. + +`status.etcdMaintenance` reserves the target using a Kubernetes resource-version +precondition before any maintenance mutation. This reservation and the minimum +interval survive operator restarts. An interrupted or timed-out request remains +in progress: pod rollouts pause until the operation deadline has passed and all +members pass direct health checks. An `EtcdDefragmented` event records success; +`EtcdMaintenanceFailed` records a failed maintenance attempt. + +If maintenance cannot recover its reservation, inspect etcd member health and +the operator error. After the two-minute operation deadline, explicitly setting +`defragmentationEnabled: false` releases an abandoned reservation and permits a +corrected workload configuration to roll out. This does not cancel a server-side +defragmentation already underway; verify member health before manual disruption. + +Managed topology TLS is supported using the TopoServer's existing certificate +with client-auth usage and its CA. The operator needs network access to each +member's client endpoint. It does not run maintenance against external topology. diff --git a/go.mod b/go.mod index faa883817..c2bf7b1d3 100644 --- a/go.mod +++ b/go.mod @@ -9,6 +9,8 @@ require ( github.com/prometheus/client_golang v1.24.1 github.com/prometheus/client_model v0.6.3 github.com/stretchr/testify v1.12.1 + go.etcd.io/etcd/api/v3 v3.7.1 + go.etcd.io/etcd/client/v3 v3.7.1 go.opentelemetry.io/contrib/exporters/autoexport v0.71.0 go.opentelemetry.io/otel v1.46.0 go.opentelemetry.io/otel/sdk v1.46.0 @@ -102,9 +104,7 @@ require ( github.com/vladimirvivien/gexe v0.5.0 // indirect github.com/x448/float16 v0.8.4 // indirect github.com/yusufpapurcu/wmi v1.2.4 // indirect - go.etcd.io/etcd/api/v3 v3.7.1 // indirect go.etcd.io/etcd/client/pkg/v3 v3.7.1 // indirect - go.etcd.io/etcd/client/v3 v3.7.1 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect go.opentelemetry.io/contrib/bridges/otelslog v0.20.1 // indirect go.opentelemetry.io/contrib/bridges/prometheus v0.71.0 // indirect diff --git a/pkg/cluster-handler/controller/multigrescluster/builders_global.go b/pkg/cluster-handler/controller/multigrescluster/builders_global.go index 32410fcd1..7f9c24888 100644 --- a/pkg/cluster-handler/controller/multigrescluster/builders_global.go +++ b/pkg/cluster-handler/controller/multigrescluster/builders_global.go @@ -80,11 +80,12 @@ func BuildGlobalTopoServer( }, Spec: multigresv1alpha1.TopoServerSpec{ Etcd: &multigresv1alpha1.EtcdSpec{ - Image: spec.Etcd.Image, - Replicas: spec.Etcd.Replicas, - Storage: spec.Etcd.Storage, - Resources: spec.Etcd.Resources, - RootPath: spec.Etcd.RootPath, + Image: spec.Etcd.Image, + Replicas: spec.Etcd.Replicas, + Storage: spec.Etcd.Storage, + Resources: spec.Etcd.Resources, + RootPath: spec.Etcd.RootPath, + Maintenance: spec.Etcd.Maintenance.DeepCopy(), }, PVCDeletionPolicy: finalPolicy, TLS: cluster.Spec.TopoTLS.DeepCopy(), diff --git a/pkg/cluster-handler/controller/multigrescluster/multigrescluster_controller.go b/pkg/cluster-handler/controller/multigrescluster/multigrescluster_controller.go index cfa85bac7..bb4336806 100644 --- a/pkg/cluster-handler/controller/multigrescluster/multigrescluster_controller.go +++ b/pkg/cluster-handler/controller/multigrescluster/multigrescluster_controller.go @@ -167,8 +167,10 @@ func (r *MultigresClusterReconciler) Reconcile( } childSpan.End() - for _, decision := range decisions { - r.Recorder.Event(cluster, "Normal", "ImplicitDefault", decision) + if cluster.Status.ObservedGeneration != cluster.Generation || cluster.Status.Phase == "" { + for _, decision := range decisions { + r.Recorder.Event(cluster, "Normal", "ImplicitDefault", decision) + } } } diff --git a/pkg/cluster-handler/controller/multigrescluster/reconcile_cells.go b/pkg/cluster-handler/controller/multigrescluster/reconcile_cells.go index e73d89441..57c0b1aa9 100644 --- a/pkg/cluster-handler/controller/multigrescluster/reconcile_cells.go +++ b/pkg/cluster-handler/controller/multigrescluster/reconcile_cells.go @@ -35,6 +35,10 @@ func (r *MultigresClusterReconciler) reconcileCells( } activeCellNames := make(map[multigresv1alpha1.CellName]bool, len(cluster.Spec.Cells)) + existingVersions := make(map[string]string, len(existingCells.Items)) + for i := range existingCells.Items { + existingVersions[existingCells.Items[i].Name] = existingCells.Items[i].ResourceVersion + } allCellNames := []multigresv1alpha1.CellName{} for _, cellCfg := range cluster.Spec.Cells { @@ -78,7 +82,10 @@ func (r *MultigresClusterReconciler) reconcileCells( ); err != nil { return false, fmt.Errorf("failed to apply cell '%s': %w", cellCfg.Name, err) } - r.Recorder.Eventf(cluster, "Normal", "Applied", "Applied Cell %s", desired.Name) + if before, found := existingVersions[desired.Name]; !found || + before != desired.ResourceVersion { + r.Recorder.Eventf(cluster, "Normal", "Applied", "Applied Cell %s", desired.Name) + } } var pendingDeletion bool diff --git a/pkg/cluster-handler/controller/multigrescluster/reconcile_topology_test.go b/pkg/cluster-handler/controller/multigrescluster/reconcile_topology_test.go index 66ca65090..51cdbc4c8 100644 --- a/pkg/cluster-handler/controller/multigrescluster/reconcile_topology_test.go +++ b/pkg/cluster-handler/controller/multigrescluster/reconcile_topology_test.go @@ -2,6 +2,7 @@ package multigrescluster import ( "context" + "path" "reflect" "testing" @@ -19,6 +20,63 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) +func TestConvergedTopologyReconcileDoesNotWrite(t *testing.T) { + t.Parallel() + ctx := t.Context() + cluster := &multigresv1alpha1.MultigresCluster{ + ObjectMeta: metav1.ObjectMeta{Name: "cluster", Namespace: "default"}, + Spec: multigresv1alpha1.MultigresClusterSpec{ + GlobalTopoServer: &multigresv1alpha1.GlobalTopoServerSpec{ + External: &multigresv1alpha1.ExternalTopoServerSpec{ + Endpoints: []multigresv1alpha1.EndpointUrl{"http://topo:2379"}, + }, + }, + Cells: []multigresv1alpha1.CellConfig{ + {Name: "cell1"}, + }, + Databases: []multigresv1alpha1.DatabaseConfig{{Name: "db"}}, + }, + } + store := newClusterTopologyMemoryStore(t) + c := fake.NewClientBuilder().WithScheme(setupScheme()).WithObjects(cluster).Build() + r := &MultigresClusterReconciler{ + Client: c, + Scheme: setupScheme(), + Recorder: record.NewFakeRecorder(100), + CreateTopoStore: func(multigresv1alpha1.GlobalTopoServerRef) (topoclient.Store, error) { return noCloseStore{store}, nil }, + } + res := resolver.NewResolver(c, cluster.Namespace) + if _, err := r.reconcileTopology(ctx, cluster, res); err != nil { + t.Fatal(err) + } + conn, err := store.ConnForCell(ctx, topoclient.GlobalCell) + if err != nil { + t.Fatal(err) + } + versions := map[string]string{} + for _, file := range []string{path.Join(topoclient.CellsPath, "cell1", topoclient.CellFile), path.Join(topoclient.DatabasesPath, "db", topoclient.DatabaseFile)} { + _, v, err := conn.Get(ctx, file) + if err != nil { + t.Fatal(err) + } + versions[file] = v.String() + } + for range 5 { + if _, err := r.reconcileTopology(ctx, cluster, res); err != nil { + t.Fatal(err) + } + for file, want := range versions { + _, v, err := conn.Get(ctx, file) + if err != nil { + t.Fatal(err) + } + if v.String() != want { + t.Fatalf("converged reconcile rewrote %s", file) + } + } + } +} + func TestReconcileTopologySharedTopo(t *testing.T) { t.Parallel() diff --git a/pkg/data-handler/topo/cell.go b/pkg/data-handler/topo/cell.go index f351d6f0e..fffe4cb6e 100644 --- a/pkg/data-handler/topo/cell.go +++ b/pkg/data-handler/topo/cell.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "slices" "github.com/multigres/multigres/go/common/topoclient" "github.com/multigres/multigres/go/pb/clustermetadata" @@ -63,7 +64,7 @@ func RegisterCell( } if !created { - logger.V(1).Info("Updated existing cell in topology", "cellName", cellName) + logger.V(1).Info("Reconciled cell in topology", "cellName", cellName) return nil } @@ -132,6 +133,10 @@ func createOrUpdateCell( ctx, cellName, func(existing *clustermetadata.Cell) error { + if existing.Name == cellMetadata.Name && existing.Root == cellMetadata.Root && + existing.Metadata == cellMetadata.Metadata && slices.Equal(existing.ServerAddresses, cellMetadata.ServerAddresses) { + return &topoclient.TopoError{Code: topoclient.NoUpdateNeeded} + } existing.Name = cellMetadata.Name existing.ServerAddresses = cellMetadata.ServerAddresses existing.Root = cellMetadata.Root diff --git a/pkg/data-handler/topo/database.go b/pkg/data-handler/topo/database.go index fa36ff566..7eab2eac8 100644 --- a/pkg/data-handler/topo/database.go +++ b/pkg/data-handler/topo/database.go @@ -52,15 +52,12 @@ func RegisterDatabase( ctx, dbName, func(existing *clustermetadatapb.Database) error { - existing.BackupLocation = dbMetadata.BackupLocation - existing.Cells = dbMetadata.Cells - existing.BootstrapDurabilityPolicy = dbMetadata.BootstrapDurabilityPolicy - return nil + return updateDatabaseMetadata(existing, dbMetadata) }, ); err != nil { return fmt.Errorf("updating existing database %s in topology: %w", dbName, err) } - logger.V(1).Info("Updated existing database in topology", "database", dbName) + logger.V(1).Info("Reconciled database in topology", "database", dbName) return nil } recorder.Eventf( diff --git a/pkg/data-handler/topo/idempotency_test.go b/pkg/data-handler/topo/idempotency_test.go new file mode 100644 index 000000000..0c91b3645 --- /dev/null +++ b/pkg/data-handler/topo/idempotency_test.go @@ -0,0 +1,119 @@ +package topo_test + +import ( + "path" + "testing" + + "github.com/multigres/multigres/go/common/topoclient" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/tools/record" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + "github.com/multigres/multigres-operator/pkg/data-handler/topo" +) + +// Inspect the backing record version, not just the returned value: rewriting +// identical protobuf bytes still consumes a new etcd revision. +func TestRegistrationDoesNotRewriteUnchangedRecords(t *testing.T) { + t.Parallel() + for _, kind := range []string{"cell", "database", "shard-database"} { + t.Run(kind, func(t *testing.T) { + t.Parallel() + ctx := t.Context() + store := newMemoryStore(t) + recorder := record.NewFakeRecorder(100) + owner := &multigresv1alpha1.MultigresCluster{ + ObjectMeta: metav1.ObjectMeta{Name: "cluster", Namespace: "default"}, + } + cell := multigresv1alpha1.CellConfig{Name: "cell1"} + backup := &multigresv1alpha1.BackupConfig{ + Type: multigresv1alpha1.BackupTypeFilesystem, + Filesystem: &multigresv1alpha1.FilesystemBackupConfig{Path: "/backups"}, + } + shard := newTestShard("shard") + register := func() error { + switch kind { + case "cell": + return topo.RegisterCellFromSpec( + ctx, + store, + recorder, + owner, + cell, + nil, + multigresv1alpha1.GlobalTopoServerRef{ + Address: "topo:2379", + RootPath: "/global", + }, + ) + case "database": + return topo.RegisterDatabaseFromSpec( + ctx, + store, + recorder, + owner, + multigresv1alpha1.DatabaseConfig{ + Name: "test-db", + }, + []string{"cell1"}, + backup, + "", + ) + default: + return topo.RegisterDatabase(ctx, store, recorder, shard) + } + } + file := path.Join(topoclient.DatabasesPath, "test-db", topoclient.DatabaseFile) + if kind == "cell" { + file = path.Join(topoclient.CellsPath, "cell1", topoclient.CellFile) + } + conn, err := store.ConnForCell(ctx, topoclient.GlobalCell) + if err != nil { + t.Fatal(err) + } + version := func() string { + t.Helper() + _, v, err := conn.Get(ctx, file) + if err != nil { + t.Fatal(err) + } + return v.String() + } + if err := register(); err != nil { + t.Fatal(err) + } + initial := version() + for range 5 { + if err := register(); err != nil { + t.Fatal(err) + } + if got := version(); got != initial { + t.Fatalf("unchanged registration rewrote record: %s -> %s", initial, got) + } + } + switch kind { + case "cell": + cell.Metadata = `{"region":"new"}` + case "database": + backup.Filesystem.Path = "/new-backups" + default: + shard.Spec.Pools["pool2"] = multigresv1alpha1.PoolSpec{ + Cells: []multigresv1alpha1.CellName{"cell2"}, + } + } + if err := register(); err != nil { + t.Fatal(err) + } + changed := version() + if changed == initial { + t.Fatal("changed registration did not update record") + } + if err := register(); err != nil { + t.Fatal(err) + } + if version() != changed { + t.Fatal("registration did not converge after change") + } + }) + } +} diff --git a/pkg/data-handler/topo/metadata.go b/pkg/data-handler/topo/metadata.go new file mode 100644 index 000000000..ac7bd76e8 --- /dev/null +++ b/pkg/data-handler/topo/metadata.go @@ -0,0 +1,28 @@ +package topo + +import ( + "slices" + + "github.com/multigres/multigres/go/common/topoclient" + clustermetadatapb "github.com/multigres/multigres/go/pb/clustermetadata" + "google.golang.org/protobuf/proto" +) + +// updateDatabaseMetadata compares only fields managed by the operator, inside +// the store's optimistic concurrency retry. Comparing before UpdateDatabaseFields +// would allow an intervening write to be missed. Preserve all other fields. +func updateDatabaseMetadata(existing, desired *clustermetadatapb.Database) error { + existingCells := slices.Clone(existing.Cells) + desiredCells := slices.Clone(desired.Cells) + slices.Sort(existingCells) + slices.Sort(desiredCells) + if slices.Equal(existingCells, desiredCells) && + proto.Equal(existing.BackupLocation, desired.BackupLocation) && + proto.Equal(existing.BootstrapDurabilityPolicy, desired.BootstrapDurabilityPolicy) { + return &topoclient.TopoError{Code: topoclient.NoUpdateNeeded} + } + existing.Cells = desiredCells + existing.BackupLocation = desired.BackupLocation + existing.BootstrapDurabilityPolicy = desired.BootstrapDurabilityPolicy + return nil +} diff --git a/pkg/data-handler/topo/topology.go b/pkg/data-handler/topo/topology.go index 819b5725e..e8c89cc7b 100644 --- a/pkg/data-handler/topo/topology.go +++ b/pkg/data-handler/topo/topology.go @@ -86,15 +86,12 @@ func RegisterDatabaseFromSpec( ctx, dbName, func(existing *clustermetadatapb.Database) error { - existing.BackupLocation = dbMetadata.BackupLocation - existing.Cells = dbMetadata.Cells - existing.BootstrapDurabilityPolicy = dbMetadata.BootstrapDurabilityPolicy - return nil + return updateDatabaseMetadata(existing, dbMetadata) }, ); err != nil { return fmt.Errorf("updating existing database %s in topology: %w", dbName, err) } - logger.V(1).Info("Updated existing database in topology", "database", dbName) + logger.V(1).Info("Reconciled database in topology", "database", dbName) return nil } return fmt.Errorf("failed to create database in topology: %w", err) @@ -157,7 +154,7 @@ func RegisterCellFromSpec( } if !created { - logger.Info("Updated existing cell in topology", "cellName", cellName) + logger.V(1).Info("Reconciled cell in topology", "cellName", cellName) return nil } diff --git a/pkg/resolver/cell.go b/pkg/resolver/cell.go index 6874c0453..d389a806e 100644 --- a/pkg/resolver/cell.go +++ b/pkg/resolver/cell.go @@ -53,6 +53,9 @@ func (r *Resolver) ResolveCell( } if localTopo.Etcd != nil { defaultEtcdSpec(localTopo.Etcd, defaultRootPath) + if _, _, err := localTopo.Etcd.Maintenance.EffectiveCompaction(); err != nil { + return nil, nil, nil, fmt.Errorf("local topology maintenance: %w", err) + } } if localTopo.External != nil { defaultExternalTopoSpec(localTopo.External, defaultRootPath) diff --git a/pkg/resolver/cluster.go b/pkg/resolver/cluster.go index de3f6bece..4ca16fd7f 100644 --- a/pkg/resolver/cluster.go +++ b/pkg/resolver/cluster.go @@ -219,6 +219,9 @@ func (r *Resolver) ResolveGlobalTopo( if finalSpec.Etcd != nil { defaultEtcdSpec(finalSpec.Etcd, roots.Global()) + if _, _, err := finalSpec.Etcd.Maintenance.EffectiveCompaction(); err != nil { + return nil, fmt.Errorf("global topology maintenance: %w", err) + } } if finalSpec.External != nil { defaultExternalTopoSpec(finalSpec.External, roots.Global()) @@ -381,4 +384,24 @@ func mergeEtcdSpec(base *multigresv1alpha1.EtcdSpec, override *multigresv1alpha1 if override.PVCDeletionPolicy != nil { base.PVCDeletionPolicy = override.PVCDeletionPolicy } + if override.Maintenance != nil { + if base.Maintenance == nil { + base.Maintenance = &multigresv1alpha1.EtcdMaintenanceConfig{} + } + m := override.Maintenance + if m.AutoCompactionMode != "" { + base.Maintenance.AutoCompactionMode = m.AutoCompactionMode + } + if m.AutoCompactionRetention != "" { + base.Maintenance.AutoCompactionRetention = m.AutoCompactionRetention + } + if m.QuotaBackendBytes != nil { + v := *m.QuotaBackendBytes + base.Maintenance.QuotaBackendBytes = &v + } + if m.DefragmentationEnabled != nil { + v := *m.DefragmentationEnabled + base.Maintenance.DefragmentationEnabled = &v + } + } } diff --git a/pkg/resolver/etcd_maintenance_test.go b/pkg/resolver/etcd_maintenance_test.go new file mode 100644 index 000000000..f7b9164b0 --- /dev/null +++ b/pkg/resolver/etcd_maintenance_test.go @@ -0,0 +1,35 @@ +package resolver + +import ( + "testing" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + "k8s.io/utils/ptr" +) + +func TestMergeEtcdMaintenance(t *testing.T) { + base := &multigresv1alpha1.EtcdSpec{Maintenance: &multigresv1alpha1.EtcdMaintenanceConfig{ + AutoCompactionMode: "periodic", + AutoCompactionRetention: "6h", + QuotaBackendBytes: ptr.To(int64(1 << 30)), + DefragmentationEnabled: ptr.To(true), + }} + override := &multigresv1alpha1.EtcdSpec{ + Maintenance: &multigresv1alpha1.EtcdMaintenanceConfig{ + DefragmentationEnabled: ptr.To(false), + QuotaBackendBytes: ptr.To(int64(512 << 20)), + }, + } + mergeEtcdSpec(base, override) + if base.Maintenance.AutoCompactionRetention != "6h" || + base.Maintenance.DefragmentationIsEnabled() || + base.Maintenance.EffectiveQuotaBackendBytes() != 512<<20 { + t.Fatalf("incorrect merge: %+v", base.Maintenance) + } + *override.Maintenance.DefragmentationEnabled = true + *override.Maintenance.QuotaBackendBytes = 2 << 30 + if base.Maintenance.DefragmentationIsEnabled() || + base.Maintenance.EffectiveQuotaBackendBytes() != 512<<20 { + t.Fatal("merged maintenance aliases override") + } +} diff --git a/pkg/resource-handler/controller/toposerver/container_env.go b/pkg/resource-handler/controller/toposerver/container_env.go index d42e8262a..51e041888 100644 --- a/pkg/resource-handler/controller/toposerver/container_env.go +++ b/pkg/resource-handler/controller/toposerver/container_env.go @@ -2,11 +2,29 @@ package toposerver import ( "fmt" + "strconv" "strings" + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" corev1 "k8s.io/api/core/v1" ) +func buildMaintenanceEnv(config *multigresv1alpha1.EtcdMaintenanceConfig) ([]corev1.EnvVar, error) { + mode, retention, err := config.EffectiveCompaction() + if err != nil { + return nil, err + } + quota := config.EffectiveQuotaBackendBytes() + if quota < 1024*1024 || quota > 8*1024*1024*1024 { + return nil, fmt.Errorf("etcd backend quota must be between 1 MiB and 8 GiB") + } + return []corev1.EnvVar{ + {Name: "ETCD_AUTO_COMPACTION_MODE", Value: mode}, + {Name: "ETCD_AUTO_COMPACTION_RETENTION", Value: retention}, + {Name: "ETCD_QUOTA_BACKEND_BYTES", Value: strconv.FormatInt(quota, 10)}, + }, nil +} + // buildContainerEnv constructs all environment variables for etcd clustering in // StatefulSets. This combines pod identity, etcd config, and cluster peer // discovery details. When tlsEnabled is set, etcd serves its client and peer diff --git a/pkg/resource-handler/controller/toposerver/container_env_test.go b/pkg/resource-handler/controller/toposerver/container_env_test.go index 1fd5842f6..f9351a66a 100644 --- a/pkg/resource-handler/controller/toposerver/container_env_test.go +++ b/pkg/resource-handler/controller/toposerver/container_env_test.go @@ -4,9 +4,40 @@ import ( "testing" "github.com/google/go-cmp/cmp" + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" corev1 "k8s.io/api/core/v1" + "k8s.io/utils/ptr" ) +func TestMaintenanceEnvironment(t *testing.T) { + ts := certTestTopoServer(nil) + ts.Spec.Etcd.Maintenance = &multigresv1alpha1.EtcdMaintenanceConfig{ + AutoCompactionMode: "revision", + AutoCompactionRetention: "20000", + QuotaBackendBytes: ptr.To(int64(512 << 20)), + } + sts, err := BuildStatefulSet(ts, certScheme()) + if err != nil { + t.Fatal(err) + } + env := etcdEnvMap(t, sts) + for key, want := range map[string]string{"ETCD_AUTO_COMPACTION_MODE": "revision", "ETCD_AUTO_COMPACTION_RETENTION": "20000", "ETCD_QUOTA_BACKEND_BYTES": "536870912"} { + if env[key] != want { + t.Errorf("%s=%q, want %q", key, env[key], want) + } + } + for _, config := range []*multigresv1alpha1.EtcdMaintenanceConfig{ + {AutoCompactionRetention: "0h"}, + {QuotaBackendBytes: ptr.To(int64(0))}, + {QuotaBackendBytes: ptr.To(int64(9 << 30))}, + } { + ts.Spec.Etcd.Maintenance = config + if _, err := BuildStatefulSet(ts, certScheme()); err == nil { + t.Errorf("accepted invalid maintenance: %+v", config) + } + } +} + func TestBuildPodIdentityEnv(t *testing.T) { got := buildPodIdentityEnv() diff --git a/pkg/resource-handler/controller/toposerver/integration_test.go b/pkg/resource-handler/controller/toposerver/integration_test.go index 8b3ece46b..96fa332ce 100644 --- a/pkg/resource-handler/controller/toposerver/integration_test.go +++ b/pkg/resource-handler/controller/toposerver/integration_test.go @@ -145,6 +145,9 @@ func TestTopoServerReconciliation(t *testing.T) { {Name: "ETCD_INITIAL_CLUSTER_STATE", Value: "new"}, {Name: "ETCD_INITIAL_CLUSTER_TOKEN", Value: "test-toposerver"}, {Name: "ETCD_INITIAL_CLUSTER", Value: "test-toposerver-0=http://test-toposerver-0.test-toposerver-headless.default.svc.cluster.local:2380,test-toposerver-1=http://test-toposerver-1.test-toposerver-headless.default.svc.cluster.local:2380,test-toposerver-2=http://test-toposerver-2.test-toposerver-headless.default.svc.cluster.local:2380"}, + {Name: "ETCD_AUTO_COMPACTION_MODE", Value: "periodic"}, + {Name: "ETCD_AUTO_COMPACTION_RETENTION", Value: "1h"}, + {Name: "ETCD_QUOTA_BACKEND_BYTES", Value: "2147483648"}, }, VolumeMounts: []corev1.VolumeMount{ {Name: "data", MountPath: "/var/lib/etcd"}, @@ -327,6 +330,9 @@ func TestTopoServerReconciliation(t *testing.T) { {Name: "ETCD_INITIAL_CLUSTER_STATE", Value: "new"}, {Name: "ETCD_INITIAL_CLUSTER_TOKEN", Value: "delete-policy-topo"}, {Name: "ETCD_INITIAL_CLUSTER", Value: "delete-policy-topo-0=http://delete-policy-topo-0.delete-policy-topo-headless.default.svc.cluster.local:2380,delete-policy-topo-1=http://delete-policy-topo-1.delete-policy-topo-headless.default.svc.cluster.local:2380,delete-policy-topo-2=http://delete-policy-topo-2.delete-policy-topo-headless.default.svc.cluster.local:2380"}, + {Name: "ETCD_AUTO_COMPACTION_MODE", Value: "periodic"}, + {Name: "ETCD_AUTO_COMPACTION_RETENTION", Value: "1h"}, + {Name: "ETCD_QUOTA_BACKEND_BYTES", Value: "2147483648"}, }, VolumeMounts: []corev1.VolumeMount{ {Name: "data", MountPath: "/var/lib/etcd"}, diff --git a/pkg/resource-handler/controller/toposerver/maintenance.go b/pkg/resource-handler/controller/toposerver/maintenance.go new file mode 100644 index 000000000..fd9bd4cce --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/maintenance.go @@ -0,0 +1,337 @@ +package toposerver + +import ( + "context" + "fmt" + "time" + + clientv3 "go.etcd.io/etcd/client/v3" + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" +) + +const ( + maintenanceInterval = time.Hour + maintenanceTimeout = 2 * time.Minute + maintenanceHealthTimeout = 15 * time.Second + maintenanceStablePeriod = time.Minute + minimumReclaimBytes int64 = 100 * 1024 * 1024 + maintenanceFieldOwner = "multigres-etcd-maintenance" +) + +func maintenanceEnabled(ts *multigresv1alpha1.TopoServer) bool { + return ts.Spec.Etcd != nil && ts.Spec.Etcd.Maintenance.DefragmentationIsEnabled() +} + +func (r *TopoServerReconciler) maintenanceReader() client.Reader { + if r.APIReader != nil { + return r.APIReader + } + return r.Client +} + +func maintenanceEndpoints(ts *multigresv1alpha1.TopoServer) []string { + replicas := DefaultReplicas + if ts.Spec.Etcd != nil && ts.Spec.Etcd.Replicas != nil { + replicas = *ts.Spec.Etcd.Replicas + } + endpoints := make([]string, 0, replicas) + for i := int32(0); i < replicas; i++ { + endpoints = append( + endpoints, + fmt.Sprintf( + "%s://%s-%d.%s-headless.%s.svc.cluster.local:2379", + clientScheme(ts.Spec.TLS.IsEnabled()), + ts.Name, + i, + ts.Name, + ts.Namespace, + ), + ) + } + return endpoints +} + +// healthyMembers requires all expected voting members, a common leader, and a +// successful linearizable read through each member. Pod readiness alone cannot +// establish quorum health. Status also rejects members reporting etcd alarms. +func healthyMembers( + ctx context.Context, + c etcdMaintenanceClient, + endpoints []string, +) (map[string]*clientv3.StatusResponse, error) { + ctx, cancel := context.WithTimeout(ctx, maintenanceHealthTimeout) + defer cancel() + members, err := c.Members(ctx) + if err != nil { + return nil, err + } + if len(members.Members) != len(endpoints) { + return nil, fmt.Errorf("etcd membership differs from expected replicas") + } + ids := make(map[uint64]bool, len(endpoints)) + for _, m := range members.Members { + if m.IsLearner || m.ID == 0 { + return nil, fmt.Errorf("etcd membership includes an unready voting member") + } + ids[m.ID] = true + } + statuses := make(map[string]*clientv3.StatusResponse, len(endpoints)) + var leader, clusterID uint64 + for _, ep := range endpoints { + s, err := c.Status(ctx, ep) + if err != nil { + return nil, fmt.Errorf("etcd member %s status: %w", ep, err) + } + if s == nil || s.Header == nil || s.Header.ClusterId == 0 || !ids[s.Header.MemberId] || + s.Leader == 0 || + s.IsLearner || + len(s.Errors) != 0 { + return nil, fmt.Errorf("etcd member %s is not healthy", ep) + } + if leader == 0 { + leader, clusterID = s.Leader, s.Header.ClusterId + } + if s.Leader != leader || s.Header.ClusterId != clusterID { + return nil, fmt.Errorf("etcd members disagree on leader or cluster identity") + } + delete(ids, s.Header.MemberId) + if err := c.Health(ctx, ep); err != nil { + return nil, fmt.Errorf("etcd member %s linearizable read: %w", ep, err) + } + statuses[ep] = s + } + foundLeader := false + for _, s := range statuses { + foundLeader = foundLeader || s.Header.MemberId == leader + } + if !foundLeader { + return nil, fmt.Errorf("etcd leader is not an expected member") + } + return statuses, nil +} + +// maintenanceWorkloadReady uses uncached observations, checks ownership, and +// excludes partially applied templates, terminating pods, and recent restarts. +func (r *TopoServerReconciler) maintenanceWorkloadReady( + ctx context.Context, + ts *multigresv1alpha1.TopoServer, +) (bool, error) { + sts := &appsv1.StatefulSet{} + if err := r.maintenanceReader().Get(ctx, client.ObjectKeyFromObject(ts), sts); err != nil { + return false, err + } + owner := metav1.GetControllerOf(sts) + if owner == nil || owner.UID != ts.UID || !sts.DeletionTimestamp.IsZero() || + sts.Spec.Replicas == nil { + return false, nil + } + n := *sts.Spec.Replicas + if n < 3 || int64(n) != int64(len(maintenanceEndpoints(ts))) || + sts.Status.ObservedGeneration != sts.Generation || + sts.Status.ReadyReplicas != n || + sts.Status.UpdatedReplicas != n || + sts.Status.CurrentRevision == "" || + sts.Status.CurrentRevision != sts.Status.UpdateRevision { + return false, nil + } + for i := int32(0); i < n; i++ { + pod := &corev1.Pod{} + if err := r.maintenanceReader(). + Get(ctx, client.ObjectKey{Namespace: ts.Namespace, Name: fmt.Sprintf("%s-%d", ts.Name, i)}, pod); err != nil { + return false, err + } + owner := metav1.GetControllerOf(pod) + if owner == nil || owner.UID != sts.UID || !pod.DeletionTimestamp.IsZero() || + pod.Labels[appsv1.StatefulSetRevisionLabel] != sts.Status.UpdateRevision { + return false, nil + } + ready := false + for _, condition := range pod.Status.Conditions { + if condition.Type == corev1.PodReady && condition.Status == corev1.ConditionTrue && + !condition.LastTransitionTime.IsZero() && + time.Since(condition.LastTransitionTime.Time) >= maintenanceStablePeriod { + ready = true + } + } + if !ready { + return false, nil + } + } + return true, nil +} + +func (r *TopoServerReconciler) saveMaintenance( + ctx context.Context, + ts *multigresv1alpha1.TopoServer, + state *multigresv1alpha1.EtcdMaintenanceStatus, +) error { + before := ts.DeepCopy() + ts.Status.EtcdMaintenance = state + // The resourceVersion precondition is the lock: only one controller can + // reserve this generation. A conflict never authorizes a maintenance RPC. + return r.Status(). + Patch(ctx, ts, client.MergeFromWithOptions(before, client.MergeFromWithOptimisticLock{}), client.FieldOwner(maintenanceFieldOwner)) +} + +func (r *TopoServerReconciler) resumeMaintenance( + ctx context.Context, + ts *multigresv1alpha1.TopoServer, +) (bool, error) { + if !maintenanceEnabled(ts) && ts.Status.EtcdMaintenance == nil { + return false, nil + } + fresh := &multigresv1alpha1.TopoServer{} + if err := r.maintenanceReader().Get(ctx, client.ObjectKeyFromObject(ts), fresh); err != nil { + return false, err + } + *ts = *fresh + state := ts.Status.EtcdMaintenance + if state == nil || !state.InProgress { + return false, nil + } + if time.Since(state.LastAttemptTime.Time) < maintenanceTimeout { + return true, nil + } + // Disabling maintenance explicitly releases an abandoned reservation after + // its RPC deadline, allowing an administrator to repair an unhealthy member. + if maintenanceEnabled(ts) { + c, err := r.maintenanceClient(ctx, ts) + if err != nil { + return true, err + } + defer c.Close() + if _, err := healthyMembers(ctx, c, maintenanceEndpoints(ts)); err != nil { + return true, err + } + } + state = state.DeepCopy() + state.InProgress = false + return false, r.saveMaintenance(ctx, ts, state) +} + +func (r *TopoServerReconciler) reconcileMaintenance( + ctx context.Context, + ts *multigresv1alpha1.TopoServer, +) error { + if !maintenanceEnabled(ts) { + return nil + } + if state := ts.Status.EtcdMaintenance; state != nil && + (state.InProgress || time.Since(state.LastAttemptTime.Time) < maintenanceInterval) { + return nil + } + ready, err := r.maintenanceWorkloadReady(ctx, ts) + if err != nil || !ready { + return err + } + c, err := r.maintenanceClient(ctx, ts) + if err != nil { + return err + } + defer c.Close() + endpoints := maintenanceEndpoints(ts) + statuses, err := healthyMembers(ctx, c, endpoints) + if err != nil { + return err + } + var target string + var reclaim int64 + for _, ep := range endpoints { + s := statuses[ep] + free := s.DbSize - s.DbSizeInUse + if s.DbSizeInUse <= 0 || free < minimumReclaimBytes || + float64(free)/float64(s.DbSize) < 0.3 { + continue + } + if free > reclaim { + target, reclaim = ep, free + } + } + if target == "" { + return nil + } + + // Re-read the CR immediately before the CAS. A changed spec is handled by + // the next reconcile, never with the old endpoint or maintenance settings. + fresh := &multigresv1alpha1.TopoServer{} + if err := r.maintenanceReader().Get(ctx, client.ObjectKeyFromObject(ts), fresh); err != nil { + return err + } + if fresh.UID != ts.UID || fresh.Generation != ts.Generation || + !fresh.DeletionTimestamp.IsZero() { + return nil + } + if state := fresh.Status.EtcdMaintenance; state != nil && + (state.InProgress || time.Since(state.LastAttemptTime.Time) < maintenanceInterval) { + return nil + } + state := &multigresv1alpha1.EtcdMaintenanceStatus{ + LastAttemptTime: metav1.Now(), + Endpoint: target, + InProgress: true, + } + if err := r.saveMaintenance(ctx, fresh, state); err != nil { + return err + } + ts.Status.EtcdMaintenance = state.DeepCopy() + + operationCtx, cancel := context.WithTimeout(ctx, maintenanceTimeout) + defer cancel() + // Once reserved, leave InProgress set on any uncertainty. A later reconcile + // must verify health before releasing it; it cannot jump to another member. + ready, err = r.maintenanceWorkloadReady(operationCtx, fresh) + if err != nil { + return err + } + if !ready { + return fmt.Errorf("etcd workload changed before defragmentation") + } + statuses, err = healthyMembers(operationCtx, c, endpoints) + if err != nil { + return err + } + if s := statuses[target]; s.Header.MemberId == s.Leader { + for _, ep := range endpoints { + if ep == target { + continue + } + if err := c.MoveLeader(operationCtx, target, statuses[ep].Header.MemberId); err != nil { + return err + } + break + } + statuses, err = healthyMembers(operationCtx, c, endpoints) + if err != nil { + return err + } + if statuses[target].Header.MemberId == statuses[target].Leader { + return fmt.Errorf("etcd leadership transfer has not completed") + } + } + if err := c.Defragment(operationCtx, target); err != nil { + return fmt.Errorf("defragmenting %s: %w", target, err) + } + if _, err := healthyMembers(operationCtx, c, endpoints); err != nil { + return err + } + state = state.DeepCopy() + state.InProgress = false + if err := r.saveMaintenance(ctx, fresh, state); err != nil { + return err + } + ts.Status.EtcdMaintenance = state + r.Recorder.Eventf( + ts, + "Normal", + "EtcdDefragmented", + "Defragmented member %s (estimated reclaimable bytes: %d)", + target, + reclaim, + ) + return nil +} diff --git a/pkg/resource-handler/controller/toposerver/maintenance_client.go b/pkg/resource-handler/controller/toposerver/maintenance_client.go new file mode 100644 index 000000000..674ecdc9d --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/maintenance_client.go @@ -0,0 +1,135 @@ +package toposerver + +import ( + "context" + "crypto/tls" + "crypto/x509" + "fmt" + "time" + + clientv3 "go.etcd.io/etcd/client/v3" + corev1 "k8s.io/api/core/v1" + "sigs.k8s.io/controller-runtime/pkg/client" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" +) + +type etcdMaintenanceClient interface { + Status(context.Context, string) (*clientv3.StatusResponse, error) + Members(context.Context) (*clientv3.MemberListResponse, error) + Health(context.Context, string) error + MoveLeader(context.Context, string, uint64) error + Defragment(context.Context, string) error + Close() +} + +// Each endpoint has a separate client. A linearizable health read must reach +// that member, not silently fail over to a healthy endpoint in a client pool. +type memberClients struct { + clients map[string]*clientv3.Client + first *clientv3.Client +} + +func (r *TopoServerReconciler) maintenanceClient( + ctx context.Context, + ts *multigresv1alpha1.TopoServer, +) (etcdMaintenanceClient, error) { + if r.newMaintenanceClient != nil { + return r.newMaintenanceClient(ctx, ts) + } + tlsConfig, err := r.maintenanceTLSConfig(ctx, ts) + if err != nil { + return nil, err + } + result := &memberClients{clients: make(map[string]*clientv3.Client)} + for _, endpoint := range maintenanceEndpoints(ts) { + c, err := clientv3.New( + clientv3.Config{ + Endpoints: []string{endpoint}, + TLS: tlsConfig, + DialTimeout: 5 * time.Second, + Context: ctx, + }, + ) + if err != nil { + result.Close() + return nil, err + } + result.clients[endpoint] = c + if result.first == nil { + result.first = c + } + } + return result, nil +} + +func (c *memberClients) Status( + ctx context.Context, + endpoint string, +) (*clientv3.StatusResponse, error) { + return c.clients[endpoint].Status(ctx, endpoint) +} + +func (c *memberClients) Members(ctx context.Context) (*clientv3.MemberListResponse, error) { + return c.first.MemberList(ctx) +} + +func (c *memberClients) Health(ctx context.Context, endpoint string) error { + // Get is linearizable by default. This probe never writes a health key. + _, err := c.clients[endpoint].Get( + ctx, + "/multigres-operator/maintenance-health", + clientv3.WithCountOnly(), + ) + return err +} + +func (c *memberClients) MoveLeader(ctx context.Context, endpoint string, id uint64) error { + _, err := c.clients[endpoint].MoveLeader(ctx, id) + return err +} + +func (c *memberClients) Defragment(ctx context.Context, endpoint string) error { + _, err := c.clients[endpoint].Defragment(ctx, endpoint) + return err +} + +func (c *memberClients) Close() { + for _, c := range c.clients { + _ = c.Close() + } +} + +func (r *TopoServerReconciler) maintenanceTLSConfig( + ctx context.Context, + ts *multigresv1alpha1.TopoServer, +) (*tls.Config, error) { + if !ts.Spec.TLS.IsEnabled() { + return nil, nil + } + secret := &corev1.Secret{} + key := client.ObjectKey{ + Namespace: ts.Namespace, + Name: multigresv1alpha1.TopoServerCertSecretName(ts.Name), + } + if err := r.maintenanceReader().Get(ctx, key, secret); err != nil { + return nil, fmt.Errorf("reading etcd maintenance credential: %w", err) + } + // The managed serving identity also has client-auth usage for peer TLS. + cert, err := tls.X509KeyPair( + secret.Data[corev1.TLSCertKey], + secret.Data[corev1.TLSPrivateKeyKey], + ) + if err != nil { + return nil, fmt.Errorf("loading etcd maintenance certificate: %w", err) + } + roots := x509.NewCertPool() + if !roots.AppendCertsFromPEM(secret.Data["ca.crt"]) { + return nil, fmt.Errorf("etcd maintenance Secret %q has no valid ca.crt", key.Name) + } + return &tls.Config{ + MinVersion: tls.VersionTLS12, + RootCAs: roots, + Certificates: []tls.Certificate{cert}, + }, nil +} diff --git a/pkg/resource-handler/controller/toposerver/maintenance_client_test.go b/pkg/resource-handler/controller/toposerver/maintenance_client_test.go new file mode 100644 index 000000000..595b57f83 --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/maintenance_client_test.go @@ -0,0 +1,84 @@ +package toposerver + +import ( + "crypto/ed25519" + "crypto/rand" + "crypto/tls" + "crypto/x509" + "encoding/pem" + "math/big" + "testing" + "time" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/utils/ptr" + "sigs.k8s.io/controller-runtime/pkg/client/fake" +) + +func TestMaintenanceTLSUsesUncachedCredentialAndFailsClosed(t *testing.T) { + ts := certTestTopoServer(&multigresv1alpha1.TopoTLSConfig{Enabled: ptr.To(true)}) + pub, key, err := ed25519.GenerateKey(rand.Reader) + if err != nil { + t.Fatal(err) + } + template := &x509.Certificate{ + SerialNumber: big.NewInt(1), + NotBefore: time.Now().Add(-time.Hour), + NotAfter: time.Now().Add(time.Hour), + IsCA: true, + BasicConstraintsValid: true, + KeyUsage: x509.KeyUsageCertSign | x509.KeyUsageDigitalSignature, + ExtKeyUsage: []x509.ExtKeyUsage{ + x509.ExtKeyUsageClientAuth, + x509.ExtKeyUsageServerAuth, + }, + } + der, err := x509.CreateCertificate(rand.Reader, template, template, pub, key) + if err != nil { + t.Fatal(err) + } + keyDER, err := x509.MarshalPKCS8PrivateKey(key) + if err != nil { + t.Fatal(err) + } + certPEM := pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der}) + secret := &corev1.Secret{ + ObjectMeta: metav1.ObjectMeta{ + Name: multigresv1alpha1.TopoServerCertSecretName(ts.Name), + Namespace: ts.Namespace, + }, + Data: map[string][]byte{ + "tls.crt": certPEM, + "tls.key": pem.EncodeToMemory(&pem.Block{Type: "PRIVATE KEY", Bytes: keyDER}), + "ca.crt": certPEM, + }, + } + for _, test := range []struct { + name string + change func(*corev1.Secret) + wantErr bool + }{ + {"valid", func(*corev1.Secret) {}, false}, + {"missing key", func(s *corev1.Secret) { delete(s.Data, "tls.key") }, true}, + {"invalid CA", func(s *corev1.Secret) { s.Data["ca.crt"] = []byte("invalid") }, true}, + } { + t.Run(test.name, func(t *testing.T) { + s := secret.DeepCopy() + test.change(s) + r := &TopoServerReconciler{ + Client: fake.NewClientBuilder().WithScheme(certScheme()).Build(), + APIReader: fake.NewClientBuilder().WithScheme(certScheme()).WithObjects(s).Build(), + } + cfg, err := r.maintenanceTLSConfig(t.Context(), ts) + if (err != nil) != test.wantErr { + t.Fatalf("error=%v", err) + } + if !test.wantErr && + (cfg.MinVersion < tls.VersionTLS12 || cfg.InsecureSkipVerify || cfg.RootCAs == nil || len(cfg.Certificates) != 1) { + t.Fatal("invalid TLS configuration") + } + }) + } +} diff --git a/pkg/resource-handler/controller/toposerver/maintenance_etcd_test.go b/pkg/resource-handler/controller/toposerver/maintenance_etcd_test.go new file mode 100644 index 000000000..3316d8e77 --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/maintenance_etcd_test.go @@ -0,0 +1,271 @@ +//go:build integration + +package toposerver + +import ( + "context" + "errors" + "fmt" + "net" + "os" + "os/exec" + "path/filepath" + "strings" + "testing" + "time" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + "github.com/multigres/multigres-operator/pkg/data-handler/topo" + "github.com/multigres/multigres/go/common/topoclient" + _ "github.com/multigres/multigres/go/common/topoclient/etcdtopo" + "go.etcd.io/etcd/api/v3/v3rpc/rpctypes" + clientv3 "go.etcd.io/etcd/client/v3" + "k8s.io/client-go/tools/record" +) + +// These tests launch disposable local etcd processes, never a configured +// Kubernetes cluster. Prefer the same binary as envtest when available. +func startMaintenanceEtcd(t *testing.T) (*memberClients, []string) { + t.Helper() + binary := filepath.Join(os.Getenv("KUBEBUILDER_ASSETS"), "etcd") + if _, err := os.Stat(binary); err != nil { + var lookupErr error + binary, lookupErr = exec.LookPath("etcd") + if lookupErr != nil { + t.Fatal("integration test requires etcd or KUBEBUILDER_ASSETS") + } + } + allocateURL := func() string { + t.Helper() + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + url := "http://" + listener.Addr().String() + if err := listener.Close(); err != nil { + t.Fatal(err) + } + return url + } + endpoints, peers, cluster := make([]string, 3), make([]string, 3), make([]string, 3) + for i := range endpoints { + endpoints[i] = allocateURL() + peers[i] = allocateURL() + cluster[i] = fmt.Sprintf("member-%d=%s", i, peers[i]) + } + for i := range endpoints { + dir := t.TempDir() + logFile, err := os.Create(filepath.Join(dir, "etcd.log")) + if err != nil { + t.Fatal(err) + } + cmd := exec.Command( + binary, + "--name", + fmt.Sprintf("member-%d", i), + "--data-dir", + filepath.Join(dir, "data"), + "--listen-client-urls", + endpoints[i], + "--advertise-client-urls", + endpoints[i], + "--listen-peer-urls", + peers[i], + "--initial-advertise-peer-urls", + peers[i], + "--initial-cluster", + strings.Join(cluster, ","), + "--initial-cluster-token", + "maintenance-test", + "--auto-compaction-mode", + "periodic", + "--auto-compaction-retention", + "2s", + "--quota-backend-bytes", + "536870912", + "--log-level", + "error", + ) + cmd.Stdout, cmd.Stderr = logFile, logFile + if err := cmd.Start(); err != nil { + _ = logFile.Close() + t.Fatal(err) + } + t.Cleanup(func() { + _ = cmd.Process.Kill() + _ = cmd.Wait() + _ = logFile.Close() + if t.Failed() { + data, _ := os.ReadFile(filepath.Join(dir, "etcd.log")) + t.Logf("etcd log: %s", data) + } + }) + } + c := &memberClients{clients: map[string]*clientv3.Client{}} + for _, ep := range endpoints { + cl, err := clientv3.New(clientv3.Config{Endpoints: []string{ep}, DialTimeout: time.Second}) + if err != nil { + t.Fatal(err) + } + c.clients[ep] = cl + if c.first == nil { + c.first = cl + } + } + t.Cleanup(c.Close) + deadline := time.Now().Add(30 * time.Second) + for { + if _, err := healthyMembers(t.Context(), c, endpoints); err == nil { + break + } else if time.Now().After(deadline) { + t.Fatalf("etcd did not become healthy: %v", err) + } + time.Sleep(100 * time.Millisecond) + } + return c, endpoints +} + +func TestLiveEtcdCompactionAndMaintenance(t *testing.T) { + c, endpoints := startMaintenanceEtcd(t) + ctx, cancel := context.WithTimeout(t.Context(), 90*time.Second) + defer cancel() + value := strings.Repeat("x", 128*1024) + var firstSize int64 + for cycle := range 5 { + first, err := c.first.Put(ctx, "/maintenance-test/data", value) + if err != nil { + t.Fatal(err) + } + for range 63 { + if _, err := c.first.Put(ctx, "/maintenance-test/data", value); err != nil { + t.Fatal(err) + } + } + deadline := time.Now().Add(15 * time.Second) + for { + _, err := c.first.Get( + ctx, + "/maintenance-test/data", + clientv3.WithRev(first.Header.Revision), + ) + if errors.Is(err, rpctypes.ErrCompacted) { + break + } + if err != nil || time.Now().After(deadline) { + t.Fatalf("history was not compacted: %v", err) + } + time.Sleep(200 * time.Millisecond) + } + statuses, err := healthyMembers(ctx, c, endpoints) + if err != nil { + t.Fatal(err) + } + s := statuses[endpoints[0]] + if cycle == 0 { + firstSize = s.DbSize + } + t.Logf( + "cycle %d: backend=%d in-use=%d revision=%d", + cycle, + s.DbSize, + s.DbSizeInUse, + s.Header.Revision, + ) + if cycle == 4 && s.DbSize > firstSize+16*1024*1024 { + t.Fatalf( + "backend kept growing across compaction cycles: initial=%d final=%d", + firstSize, + s.DbSize, + ) + } + } + + // Maintenance probes and converged registration must not manufacture MVCC writes. + store, err := topoclient.OpenServer( + "etcd", + "/operator-test", + endpoints, + topoclient.NewDefaultTopoConfig(), + ) + if err != nil { + t.Fatal(err) + } + defer func() { _ = store.Close() }() + owner := &multigresv1alpha1.MultigresCluster{} + register := func() { + t.Helper() + if err := topo.RegisterDatabaseFromSpec( + ctx, + store, + record.NewFakeRecorder(10), + owner, + multigresv1alpha1.DatabaseConfig{Name: "db"}, + []string{"cell1"}, + nil, + "", + ); err != nil { + t.Fatal(err) + } + } + register() + before, err := c.first.Get(ctx, "/operator-test", clientv3.WithPrefix()) + if err != nil { + t.Fatal(err) + } + for range 5 { + register() + if _, err := healthyMembers(ctx, c, endpoints); err != nil { + t.Fatal(err) + } + } + after, err := c.first.Get(ctx, "/operator-test", clientv3.WithPrefix()) + if err != nil { + t.Fatal(err) + } + if before.Header.Revision != after.Header.Revision { + t.Fatalf( + "no-op reconciles advanced etcd revision: %d -> %d", + before.Header.Revision, + after.Header.Revision, + ) + } + + // Reclaim each member independently, transferring leadership first where + // needed and requiring healthy, linearizable reads between every operation. + for _, ep := range endpoints { + statuses, err := healthyMembers(ctx, c, endpoints) + if err != nil { + t.Fatal(err) + } + s := statuses[ep] + if s.Header.MemberId == s.Leader { + for _, other := range endpoints { + if other != ep { + if err := c.MoveLeader(ctx, ep, statuses[other].Header.MemberId); err != nil { + t.Fatal(err) + } + break + } + } + } + if err := c.Defragment(ctx, ep); err != nil { + t.Fatal(err) + } + statuses, err = healthyMembers(ctx, c, endpoints) + if err != nil { + t.Fatal(err) + } + if statuses[ep].DbSize >= s.DbSize { + t.Fatalf( + "defragmentation did not shrink %s: %d -> %d", + ep, + s.DbSize, + statuses[ep].DbSize, + ) + } + } + got, err := c.first.Get(ctx, "/maintenance-test/data") + if err != nil || len(got.Kvs) != 1 || string(got.Kvs[0].Value) != value { + t.Fatalf("live data lost after maintenance: %v", err) + } +} diff --git a/pkg/resource-handler/controller/toposerver/maintenance_status_integration_test.go b/pkg/resource-handler/controller/toposerver/maintenance_status_integration_test.go new file mode 100644 index 000000000..d959add02 --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/maintenance_status_integration_test.go @@ -0,0 +1,72 @@ +//go:build integration + +package toposerver + +import ( + "path/filepath" + "testing" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + "github.com/multigres/multigres-operator/pkg/testutil" + appsv1 "k8s.io/api/apps/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/tools/record" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +func TestMaintenanceReservationSurvivesStatusApply(t *testing.T) { + scheme := certScheme() + _ = appsv1.AddToScheme(scheme) + cfg := testutil.SetUpEnvtest( + t, + testutil.WithCRDPaths(filepath.Join("../../../..", "config", "crd", "bases")), + ) + c, err := client.New(cfg, client.Options{Scheme: scheme}) + if err != nil { + t.Fatal(err) + } + ts := certTestTopoServer(nil) + ts.Namespace = "default" + ts.UID = "" + if err := c.Create(t.Context(), ts); err != nil { + t.Fatal(err) + } + sts, err := BuildStatefulSet(ts, scheme) + if err != nil { + t.Fatal(err) + } + if err := c.Create(t.Context(), sts); err != nil { + t.Fatal(err) + } + r := &TopoServerReconciler{ + Client: c, + APIReader: c, + Scheme: scheme, + Recorder: record.NewFakeRecorder(100), + } + stale := ts.DeepCopy() + state := &multigresv1alpha1.EtcdMaintenanceStatus{ + LastAttemptTime: metav1.Now(), + Endpoint: maintenanceEndpoints(ts)[0], + InProgress: true, + } + if err := r.saveMaintenance(t.Context(), ts, state); err != nil { + t.Fatal(err) + } + if err := r.saveMaintenance(t.Context(), stale, state); !apierrors.IsConflict(err) { + t.Fatalf("stale reservation should conflict, got %v", err) + } + // The ordinary status writer intentionally omits maintenance fields; SSA + // must preserve the independently owned reservation, including after restart. + if err := r.updateStatus(t.Context(), stale); err != nil { + t.Fatal(err) + } + fresh := &multigresv1alpha1.TopoServer{} + if err := c.Get(t.Context(), client.ObjectKeyFromObject(ts), fresh); err != nil { + t.Fatal(err) + } + if fresh.Status.EtcdMaintenance == nil || !fresh.Status.EtcdMaintenance.InProgress { + t.Fatal("status apply removed active maintenance reservation") + } +} diff --git a/pkg/resource-handler/controller/toposerver/maintenance_test.go b/pkg/resource-handler/controller/toposerver/maintenance_test.go new file mode 100644 index 000000000..714da5f5e --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/maintenance_test.go @@ -0,0 +1,283 @@ +package toposerver + +import ( + "context" + "errors" + "fmt" + "testing" + "time" + + pb "go.etcd.io/etcd/api/v3/etcdserverpb" + clientv3 "go.etcd.io/etcd/client/v3" + appsv1 "k8s.io/api/apps/v1" + 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/client-go/tools/record" + "k8s.io/utils/ptr" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + "sigs.k8s.io/controller-runtime/pkg/client/interceptor" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" +) + +type fakeEtcdMaintenance struct { + statuses map[string]*clientv3.StatusResponse + defragged []string + moved []uint64 + healthErr error + defragErr error + afterDefrag func() +} + +func (f *fakeEtcdMaintenance) Status( + _ context.Context, + ep string, +) (*clientv3.StatusResponse, error) { + return f.statuses[ep], nil +} + +func (f *fakeEtcdMaintenance) Members(context.Context) (*clientv3.MemberListResponse, error) { + return &clientv3.MemberListResponse{Members: []*pb.Member{{ID: 1}, {ID: 2}, {ID: 3}}}, nil +} +func (f *fakeEtcdMaintenance) Health(context.Context, string) error { return f.healthErr } +func (f *fakeEtcdMaintenance) Defragment(_ context.Context, ep string) error { + f.defragged = append(f.defragged, ep) + if f.afterDefrag != nil { + f.afterDefrag() + } + return f.defragErr +} + +func (f *fakeEtcdMaintenance) MoveLeader(_ context.Context, _ string, id uint64) error { + f.moved = append(f.moved, id) + for _, s := range f.statuses { + s.Leader = id + } + return nil +} +func (f *fakeEtcdMaintenance) Close() {} + +func maintenanceFixture( + t *testing.T, +) (*TopoServerReconciler, *multigresv1alpha1.TopoServer, *fakeEtcdMaintenance) { + t.Helper() + scheme := runtime.NewScheme() + _ = multigresv1alpha1.AddToScheme(scheme) + _ = appsv1.AddToScheme(scheme) + _ = corev1.AddToScheme(scheme) + ts := certTestTopoServer(nil) + ts.Generation = 1 + ts.Spec.Etcd.Maintenance = &multigresv1alpha1.EtcdMaintenanceConfig{ + DefragmentationEnabled: ptr.To(true), + } + sts, err := BuildStatefulSet(ts, scheme) + if err != nil { + t.Fatal(err) + } + sts.UID, sts.Generation = "sts-uid", 1 + sts.Status = appsv1.StatefulSetStatus{ + ObservedGeneration: 1, + Replicas: 3, + ReadyReplicas: 3, + UpdatedReplicas: 3, + CurrentRevision: "rev1", + UpdateRevision: "rev1", + } + objects := []client.Object{ts, sts} + f := &fakeEtcdMaintenance{statuses: map[string]*clientv3.StatusResponse{}} + for i, ep := range maintenanceEndpoints(ts) { + f.statuses[ep] = &clientv3.StatusResponse{ + Header: &pb.ResponseHeader{ClusterId: 123, MemberId: uint64(i + 1)}, + Leader: 1, + DbSize: 400 << 20, + DbSizeInUse: 100 << 20, + } + objects = append(objects, &corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: fmt.Sprintf("%s-%d", ts.Name, i), + Namespace: ts.Namespace, + Labels: map[string]string{appsv1.StatefulSetRevisionLabel: "rev1"}, + OwnerReferences: []metav1.OwnerReference{ + *metav1.NewControllerRef(sts, appsv1.SchemeGroupVersion.WithKind("StatefulSet")), + }, + }, + Status: corev1.PodStatus{ + Conditions: []corev1.PodCondition{ + { + Type: corev1.PodReady, + Status: corev1.ConditionTrue, + LastTransitionTime: metav1.NewTime(time.Now().Add(-time.Hour)), + }, + }, + }, + }) + } + c := fake.NewClientBuilder(). + WithScheme(scheme). + WithStatusSubresource(&multigresv1alpha1.TopoServer{}, &appsv1.StatefulSet{}, &corev1.Pod{}). + WithObjects(objects...). + Build() + if err := c.Get(t.Context(), client.ObjectKeyFromObject(ts), ts); err != nil { + t.Fatal(err) + } + r := &TopoServerReconciler{ + Client: c, + APIReader: c, + Scheme: scheme, + Recorder: record.NewFakeRecorder(100), + newMaintenanceClient: func(context.Context, *multigresv1alpha1.TopoServer) (etcdMaintenanceClient, error) { return f, nil }, + } + return r, ts, f +} + +func TestEtcdMaintenanceSerializesMembersAndRestarts(t *testing.T) { + r, ts, f := maintenanceFixture(t) + if err := r.reconcileMaintenance(t.Context(), ts); err != nil { + t.Fatal(err) + } + if len(f.defragged) != 1 || len(f.moved) != 1 { + t.Fatalf("defrags=%v leader transfers=%v", f.defragged, f.moved) + } + fresh := &multigresv1alpha1.TopoServer{} + if err := r.Get(t.Context(), client.ObjectKeyFromObject(ts), fresh); err != nil { + t.Fatal(err) + } + if fresh.Status.EtcdMaintenance == nil || fresh.Status.EtcdMaintenance.InProgress { + t.Fatal("completed reservation not persisted") + } + // A new controller instance, or a stale reconcile, must obey the persisted interval. + r2 := *r + if err := r2.reconcileMaintenance(t.Context(), fresh); err != nil { + t.Fatal(err) + } + if len(f.defragged) != 1 { + t.Fatal("maintenance repeated inside the interval") + } +} + +func TestEtcdMaintenanceHealthGates(t *testing.T) { + for _, test := range []struct { + name string + mutate func(*TopoServerReconciler, *multigresv1alpha1.TopoServer, *fakeEtcdMaintenance) + wantErr bool + }{ + {"disabled", func(_ *TopoServerReconciler, ts *multigresv1alpha1.TopoServer, _ *fakeEtcdMaintenance) { + ts.Spec.Etcd.Maintenance = nil + }, false}, + {"single member", func(_ *TopoServerReconciler, ts *multigresv1alpha1.TopoServer, _ *fakeEtcdMaintenance) { + ts.Spec.Etcd.Replicas = ptr.To(int32(1)) + }, false}, + {"no fragmentation", func(_ *TopoServerReconciler, _ *multigresv1alpha1.TopoServer, f *fakeEtcdMaintenance) { + for _, s := range f.statuses { + s.DbSizeInUse = s.DbSize + } + }, false}, + {"small fragmentation", func(_ *TopoServerReconciler, _ *multigresv1alpha1.TopoServer, f *fakeEtcdMaintenance) { + for _, s := range f.statuses { + s.DbSize = 100 << 20 + s.DbSizeInUse = 10 << 20 + } + }, false}, + {"low fragmentation ratio", func(_ *TopoServerReconciler, _ *multigresv1alpha1.TopoServer, f *fakeEtcdMaintenance) { + for _, s := range f.statuses { + s.DbSize = 1 << 30 + s.DbSizeInUse = 900 << 20 + } + }, false}, + {"linearizable read fails", func(_ *TopoServerReconciler, _ *multigresv1alpha1.TopoServer, f *fakeEtcdMaintenance) { + f.healthErr = errors.New("no quorum") + }, true}, + {"leader disagreement", func(_ *TopoServerReconciler, ts *multigresv1alpha1.TopoServer, f *fakeEtcdMaintenance) { + f.statuses[maintenanceEndpoints(ts)[1]].Leader = 2 + }, true}, + {"member alarm", func(_ *TopoServerReconciler, ts *multigresv1alpha1.TopoServer, f *fakeEtcdMaintenance) { + f.statuses[maintenanceEndpoints(ts)[1]].Errors = []string{"NOSPACE"} + }, true}, + {"foreign cluster", func(_ *TopoServerReconciler, ts *multigresv1alpha1.TopoServer, f *fakeEtcdMaintenance) { + f.statuses[maintenanceEndpoints(ts)[1]].Header.ClusterId = 999 + }, true}, + {"duplicate member", func(_ *TopoServerReconciler, ts *multigresv1alpha1.TopoServer, f *fakeEtcdMaintenance) { + f.statuses[maintenanceEndpoints(ts)[1]].Header.MemberId = 1 + }, true}, + {"rolling update", func(r *TopoServerReconciler, ts *multigresv1alpha1.TopoServer, _ *fakeEtcdMaintenance) { + sts := &appsv1.StatefulSet{} + _ = r.Get(t.Context(), client.ObjectKeyFromObject(ts), sts) + sts.Status.UpdateRevision = "rev2" + if err := r.Status().Update(t.Context(), sts); err != nil { + t.Fatal(err) + } + }, false}, + {"reservation conflict", func(r *TopoServerReconciler, _ *multigresv1alpha1.TopoServer, _ *fakeEtcdMaintenance) { + r.Client = interceptor.NewClient(r.Client.(client.WithWatch), interceptor.Funcs{SubResourcePatch: func(context.Context, client.Client, string, client.Object, client.Patch, ...client.SubResourcePatchOption) error { + return apierrors.NewConflict(multigresv1alpha1.GroupVersion.WithResource("toposervers").GroupResource(), "topo", errors.New("conflict")) + }}) + }, true}, + } { + t.Run(test.name, func(t *testing.T) { + r, ts, f := maintenanceFixture(t) + test.mutate(r, ts, f) + err := r.reconcileMaintenance(t.Context(), ts) + if (err != nil) != test.wantErr { + t.Fatalf("error=%v", err) + } + if len(f.defragged) != 0 { + t.Fatalf("unsafe defragmentation: %v", f.defragged) + } + }) + } +} + +func TestEtcdMaintenanceInterruptedOperation(t *testing.T) { + r, ts, f := maintenanceFixture(t) + f.defragErr = context.DeadlineExceeded + if err := r.reconcileMaintenance(t.Context(), ts); err == nil { + t.Fatal("expected timeout") + } + fresh := &multigresv1alpha1.TopoServer{} + if err := r.Get(t.Context(), client.ObjectKeyFromObject(ts), fresh); err != nil { + t.Fatal(err) + } + if !fresh.Status.EtcdMaintenance.InProgress { + t.Fatal("uncertain operation released its reservation") + } + waiting, err := r.resumeMaintenance(t.Context(), fresh) + if err != nil || !waiting { + t.Fatalf("resume immediately: waiting=%v err=%v", waiting, err) + } + fresh.Status.EtcdMaintenance.LastAttemptTime = metav1.NewTime(time.Now().Add(-3 * time.Minute)) + if err := r.Status().Update(t.Context(), fresh); err != nil { + t.Fatal(err) + } + f.healthErr = errors.New("previous member still unavailable") + if waiting, err = r.resumeMaintenance(t.Context(), fresh); !waiting || err == nil { + t.Fatalf("unhealthy resume: waiting=%v err=%v", waiting, err) + } + f.healthErr = nil + if waiting, err = r.resumeMaintenance(t.Context(), fresh); waiting || err != nil { + t.Fatalf("healthy resume: waiting=%v err=%v", waiting, err) + } + if err := r.reconcileMaintenance(t.Context(), fresh); err != nil { + t.Fatal(err) + } + if len(f.defragged) != 1 { + t.Fatal("interrupted operation started another defrag") + } +} + +func TestEtcdMaintenancePostHealthFailureKeepsReservation(t *testing.T) { + r, ts, f := maintenanceFixture(t) + f.afterDefrag = func() { f.healthErr = errors.New("member unhealthy after defrag") } + if err := r.reconcileMaintenance(t.Context(), ts); err == nil { + t.Fatal("expected post-defrag health failure") + } + fresh := &multigresv1alpha1.TopoServer{} + if err := r.Get(t.Context(), client.ObjectKeyFromObject(ts), fresh); err != nil { + t.Fatal(err) + } + if !fresh.Status.EtcdMaintenance.InProgress { + t.Fatal("failed post-check released reservation") + } +} diff --git a/pkg/resource-handler/controller/toposerver/statefulset.go b/pkg/resource-handler/controller/toposerver/statefulset.go index 16ab203f8..ca6e0a2a2 100644 --- a/pkg/resource-handler/controller/toposerver/statefulset.go +++ b/pkg/resource-handler/controller/toposerver/statefulset.go @@ -88,6 +88,14 @@ func BuildStatefulSet( } tlsEnabled := toposerver.Spec.TLS.IsEnabled() + var maintenance *multigresv1alpha1.EtcdMaintenanceConfig + if toposerver.Spec.Etcd != nil { + maintenance = toposerver.Spec.Etcd.Maintenance + } + maintenanceEnv, err := buildMaintenanceEnv(maintenance) + if err != nil { + return nil, err + } volumeMounts := []corev1.VolumeMount{ { @@ -143,13 +151,13 @@ func BuildStatefulSet( Name: "etcd", Image: image, Resources: resources, - Env: buildContainerEnv( + Env: append(buildContainerEnv( toposerver.Name, toposerver.Namespace, replicas, headlessServiceName, tlsEnabled, - ), + ), maintenanceEnv...), Ports: buildContainerPorts(toposerver), VolumeMounts: volumeMounts, StartupProbe: &corev1.Probe{ diff --git a/pkg/resource-handler/controller/toposerver/statefulset_test.go b/pkg/resource-handler/controller/toposerver/statefulset_test.go index 3cee6d2da..2b084b7b3 100644 --- a/pkg/resource-handler/controller/toposerver/statefulset_test.go +++ b/pkg/resource-handler/controller/toposerver/statefulset_test.go @@ -640,6 +640,12 @@ func TestBuildStatefulSet(t *testing.T) { if tc.wantErr { return } + tc.want.Spec.Template.Spec.Containers[0].Env = append( + tc.want.Spec.Template.Spec.Containers[0].Env, + corev1.EnvVar{Name: "ETCD_AUTO_COMPACTION_MODE", Value: "periodic"}, + corev1.EnvVar{Name: "ETCD_AUTO_COMPACTION_RETENTION", Value: "1h"}, + corev1.EnvVar{Name: "ETCD_QUOTA_BACKEND_BYTES", Value: "2147483648"}, + ) if diff := cmp.Diff(tc.want, got); diff != "" { t.Errorf("BuildStatefulSet() mismatch (-want +got):\n%s", diff) diff --git a/pkg/resource-handler/controller/toposerver/toposerver_controller.go b/pkg/resource-handler/controller/toposerver/toposerver_controller.go index bc906658c..f3a04f3f1 100644 --- a/pkg/resource-handler/controller/toposerver/toposerver_controller.go +++ b/pkg/resource-handler/controller/toposerver/toposerver_controller.go @@ -37,8 +37,11 @@ const ( // TopoServerReconciler reconciles a TopoServer object. type TopoServerReconciler struct { client.Client - Scheme *runtime.Scheme - Recorder record.EventRecorder + APIReader client.Reader + Scheme *runtime.Scheme + Recorder record.EventRecorder + // newMaintenanceClient supplies direct per-member clients; overridden in tests. + newMaintenanceClient func(context.Context, *multigresv1alpha1.TopoServer) (etcdMaintenanceClient, error) } // Reconcile manages the etcd StatefulSet, headless service, and client service for a TopoServer. @@ -78,6 +81,12 @@ func (r *TopoServerReconciler) Reconcile( return ctrl.Result{}, nil } + // An interrupted defrag may still be running on the server. Recover its + // reservation before allowing a pod rollout to take another member down. + if waiting, err := r.resumeMaintenance(ctx, toposerver); err != nil || waiting { + return ctrl.Result{RequeueAfter: statusRecheckDelay}, err + } + // Validate StorageClass dependency before StatefulSet apply if err := r.validateEtcdStorageClassDependency(ctx, toposerver); err != nil { if isStorageClassDependencyError(err) { @@ -210,6 +219,16 @@ func (r *TopoServerReconciler) Reconcile( if toposerver.Status.Phase != multigresv1alpha1.PhaseHealthy { return ctrl.Result{RequeueAfter: statusRecheckDelay}, nil } + if maintenanceEnabled(toposerver) { + if err := r.reconcileMaintenance(ctx, toposerver); err != nil { + logger.Error(err, "Etcd maintenance deferred") + r.Recorder.Eventf(toposerver, "Warning", "EtcdMaintenanceFailed", "%v", err) + } + if state := toposerver.Status.EtcdMaintenance; state != nil && state.InProgress { + return ctrl.Result{RequeueAfter: statusRecheckDelay}, nil + } + return ctrl.Result{RequeueAfter: maintenanceInterval}, nil + } return ctrl.Result{}, nil } @@ -521,6 +540,9 @@ func (r *TopoServerReconciler) SetupWithManagerReconciler( reconciler reconcile.Reconciler, opts ...controller.Options, ) error { + if r.APIReader == nil { + r.APIReader = mgr.GetAPIReader() + } controllerOpts := controller.Options{ MaxConcurrentReconciles: 20, } diff --git a/pkg/resource-handler/controller/toposerver/toposerver_controller_internal_test.go b/pkg/resource-handler/controller/toposerver/toposerver_controller_internal_test.go index abdeb503c..c5d0a41c0 100644 --- a/pkg/resource-handler/controller/toposerver/toposerver_controller_internal_test.go +++ b/pkg/resource-handler/controller/toposerver/toposerver_controller_internal_test.go @@ -251,6 +251,9 @@ func TestSetupWithManager(t *testing.T) { if err := r.SetupWithManager(mgr); err != nil { t.Errorf("SetupWithManager() error = %v", err) } + if r.APIReader != mgr.GetAPIReader() { + t.Error("maintenance must use the manager's uncached API reader") + } }) t.Run("with options", func(t *testing.T) { @@ -267,6 +270,33 @@ func TestSetupWithManager(t *testing.T) { t.Errorf("SetupWithManager() with opts error = %v", err) } }) + + for name, reader := range map[string]client.Reader{ + "shared setup defaults to uncached reader": nil, + "shared setup preserves injected reader": fake.NewClientBuilder().WithScheme(scheme).Build(), + } { + t.Run(name, func(t *testing.T) { + mgr := createMgr() + r := &TopoServerReconciler{ + Client: mgr.GetClient(), + APIReader: reader, + Scheme: scheme, + Recorder: record.NewFakeRecorder(100), + } + if err := r.SetupWithManagerReconciler(mgr, r, controller.Options{ + SkipNameValidation: ptr.To(true), + }); err != nil { + t.Fatalf("SetupWithManagerReconciler() error = %v", err) + } + want := reader + if want == nil { + want = mgr.GetAPIReader() + } + if r.APIReader != want { + t.Error("shared setup selected the wrong maintenance reader") + } + }) + } } func TestUpdateStatus_DegradedOnCrashLoop(t *testing.T) { diff --git a/pkg/webhook/etcd_maintenance_validation_test.go b/pkg/webhook/etcd_maintenance_validation_test.go new file mode 100644 index 000000000..1701d7619 --- /dev/null +++ b/pkg/webhook/etcd_maintenance_validation_test.go @@ -0,0 +1,76 @@ +//go:build integration + +package webhook_test + +import ( + "strings" + "testing" + + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +func TestCEL_EtcdMaintenance(t *testing.T) { + c := getPrivilegedClient(t) + for _, tc := range []struct { + name string + config map[string]any + wantError string + }{ + {"default", map[string]any{}, ""}, + {"periodic", map[string]any{"autoCompactionRetention": "30m"}, ""}, + {"revision", map[string]any{"autoCompactionMode": "revision", "autoCompactionRetention": "5000"}, ""}, + {"quota", map[string]any{"quotaBackendBytes": int64(512 << 20), "defragmentationEnabled": true}, ""}, + {"zero-retention", map[string]any{"autoCompactionRetention": "0h"}, "positive"}, + {"bad-duration", map[string]any{"autoCompactionRetention": "forever"}, "compaction retention"}, + {"bad-revision", map[string]any{"autoCompactionMode": "revision", "autoCompactionRetention": "1h"}, "compaction retention"}, + {"negative-revision", map[string]any{"autoCompactionMode": "revision", "autoCompactionRetention": "-1"}, "compaction retention"}, + {"bad-mode", map[string]any{"autoCompactionMode": "disabled"}, "Unsupported value"}, + {"zero-quota", map[string]any{"quotaBackendBytes": int64(0)}, "greater than or equal"}, + {"large-quota", map[string]any{"quotaBackendBytes": int64(9 << 30)}, "less than or equal"}, + } { + t.Run(tc.name, func(t *testing.T) { + obj := &unstructured.Unstructured{ + Object: map[string]any{ + "apiVersion": "multigres.com/v1alpha1", + "kind": "TopoServer", + "metadata": map[string]any{ + "name": "maintenance-" + tc.name, + "namespace": testNamespace, + }, + "spec": map[string]any{ + "etcd": map[string]any{"replicas": int64(3), "maintenance": tc.config}, + }, + }, + } + err := c.Create(t.Context(), obj) + if tc.wantError != "" { + if err == nil || !strings.Contains(err.Error(), tc.wantError) { + t.Fatalf("expected %q, got %v", tc.wantError, err) + } + return + } + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = c.Delete(t.Context(), obj) }) + stored := &unstructured.Unstructured{} + stored.SetGroupVersionKind(obj.GroupVersionKind()) + if err := c.Get(t.Context(), client.ObjectKeyFromObject(obj), stored); err != nil { + t.Fatal(err) + } + for key, want := range tc.config { + got, found, err := unstructured.NestedFieldNoCopy( + stored.Object, + "spec", + "etcd", + "maintenance", + key, + ) + if err != nil || !found || got != want { + t.Errorf("%s: got %v, want %v (err %v)", key, got, want, err) + } + } + }) + } +} From b29513a1caa5edd269ef4a7aee622b75896503fc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ver=C3=B3nica=20L=C3=B3pez?= Date: Fri, 25 Sep 2026 02:08:44 -0600 Subject: [PATCH 2/3] fix(toposerver): restore services before maintenance recovery MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Reconcile Services and the serving Certificate before maintenance health checks. - Keep StatefulSet updates blocked while maintenance recovery is pending. - Test dependency repair and rollout resumption after recovery. Signed-off-by: Verónica López --- .../toposerver/maintenance_reconcile_test.go | 169 ++++++++++++++++++ .../toposerver/toposerver_controller.go | 87 ++++----- 2 files changed, 213 insertions(+), 43 deletions(-) create mode 100644 pkg/resource-handler/controller/toposerver/maintenance_reconcile_test.go diff --git a/pkg/resource-handler/controller/toposerver/maintenance_reconcile_test.go b/pkg/resource-handler/controller/toposerver/maintenance_reconcile_test.go new file mode 100644 index 000000000..424b4fc58 --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/maintenance_reconcile_test.go @@ -0,0 +1,169 @@ +package toposerver + +import ( + "context" + "errors" + "fmt" + "testing" + "time" + + "github.com/google/go-cmp/cmp" + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + policyv1 "k8s.io/api/policy/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/utils/ptr" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + "github.com/multigres/multigres-operator/pkg/util/certs" +) + +func TestReconcileRepairsMaintenanceDependencies(t *testing.T) { + for _, tc := range []struct { + name string + age time.Duration + tls bool + unhealthy bool + wantActive bool + }{ + {name: "operation still running", age: time.Minute, wantActive: true}, + {name: "TLS operation still running", age: time.Minute, tls: true, wantActive: true}, + {name: "member still unhealthy", age: 3 * time.Minute, unhealthy: true, wantActive: true}, + {name: "member recovered", age: 3 * time.Minute}, + } { + t.Run(tc.name, func(t *testing.T) { + r, ts, etcd := maintenanceFixture(t) + if err := policyv1.AddToScheme(r.Scheme); err != nil { + t.Fatal(err) + } + key := client.ObjectKeyFromObject(ts) + before := &appsv1.StatefulSet{} + if err := r.Get(t.Context(), key, before); err != nil { + t.Fatal(err) + } + if tc.tls { + ts.Spec.TLS = &multigresv1alpha1.TopoTLSConfig{ + Enabled: ptr.To(true), IssuerName: "topology-issuer", + } + if err := r.Update(t.Context(), ts); err != nil { + t.Fatal(err) + } + desired, err := BuildStatefulSet(ts, r.Scheme) + if err != nil { + t.Fatal(err) + } + before.Spec = desired.Spec + if err := r.Update(t.Context(), before); err != nil { + t.Fatal(err) + } + } + ts.Status.EtcdMaintenance = &multigresv1alpha1.EtcdMaintenanceStatus{ + LastAttemptTime: metav1.NewTime(time.Now().Add(-tc.age)), + Endpoint: maintenanceEndpoints(ts)[0], + InProgress: true, + } + if err := r.Status().Update(t.Context(), ts); err != nil { + t.Fatal(err) + } + ts.Spec.Etcd.Image = "etcd:pending-rollout" + if err := r.Update(t.Context(), ts); err != nil { + t.Fatal(err) + } + + healthErr := errors.New("member is still unavailable") + if tc.unhealthy { + etcd.healthErr = healthErr + } + connectionAttempts := 0 + r.newMaintenanceClient = func(ctx context.Context, current *multigresv1alpha1.TopoServer) (etcdMaintenanceClient, error) { + connectionAttempts++ + svc := &corev1.Service{} + serviceKey := client.ObjectKey{ + Namespace: current.Namespace, + Name: current.Name + "-headless", + } + if err := r.APIReader.Get(ctx, serviceKey, svc); err != nil { + return nil, fmt.Errorf("member DNS requires the headless Service: %w", err) + } + return etcd, nil + } + + result, err := r.Reconcile(t.Context(), ctrl.Request{NamespacedName: key}) + if tc.unhealthy { + if !errors.Is(err, healthErr) { + t.Errorf("expected member health error, got %v", err) + } + } else if err != nil { + t.Errorf("Reconcile() error = %v", err) + } + if tc.wantActive && result.RequeueAfter != statusRecheckDelay { + t.Errorf("requeue = %v, want %v", result.RequeueAfter, statusRecheckDelay) + } + for _, name := range []string{ts.Name + "-headless", ts.Name} { + svc := &corev1.Service{} + if err := r.Get( + t.Context(), + client.ObjectKey{Namespace: ts.Namespace, Name: name}, + svc, + ); err != nil { + t.Errorf("maintenance blocked Service repair: %v", err) + continue + } + if !metav1.IsControlledBy(svc, ts) { + t.Errorf("Service %s is not owned by the TopoServer", name) + } + if name == ts.Name+"-headless" && + (svc.Spec.ClusterIP != corev1.ClusterIPNone || !svc.Spec.PublishNotReadyAddresses) { + t.Error("headless Service does not publish member DNS during recovery") + } + } + if tc.tls { + cert, err := certs.Get( + t.Context(), + r.Client, + ts.Namespace, + multigresv1alpha1.TopoServerCertName(ts.Name), + ) + if err != nil || cert == nil { + t.Errorf( + "maintenance blocked Certificate repair: certificate=%v error=%v", + cert, + err, + ) + } + } + + fresh := &multigresv1alpha1.TopoServer{} + if err := r.Get(t.Context(), key, fresh); err != nil { + t.Fatal(err) + } + if fresh.Status.EtcdMaintenance == nil || + fresh.Status.EtcdMaintenance.InProgress != tc.wantActive { + t.Errorf( + "maintenance state = %+v, want active=%v", + fresh.Status.EtcdMaintenance, + tc.wantActive, + ) + } + after := &appsv1.StatefulSet{} + if err := r.Get(t.Context(), key, after); err != nil { + t.Fatal(err) + } + if tc.wantActive { + if diff := cmp.Diff(before.Spec, after.Spec); diff != "" { + t.Errorf("StatefulSet changed during maintenance (-before +after):\n%s", diff) + } + } else if after.Spec.Template.Spec.Containers[0].Image != string(ts.Spec.Etcd.Image) { + t.Error("StatefulSet update did not resume after maintenance recovery") + } + if (connectionAttempts > 0) != (tc.age >= maintenanceTimeout) { + t.Errorf("maintenance connection attempts = %d", connectionAttempts) + } + if len(etcd.defragged) != 0 { + t.Error("recovery started another defragmentation") + } + }) + } +} diff --git a/pkg/resource-handler/controller/toposerver/toposerver_controller.go b/pkg/resource-handler/controller/toposerver/toposerver_controller.go index f3a04f3f1..7a8a5ef83 100644 --- a/pkg/resource-handler/controller/toposerver/toposerver_controller.go +++ b/pkg/resource-handler/controller/toposerver/toposerver_controller.go @@ -81,36 +81,19 @@ func (r *TopoServerReconciler) Reconcile( return ctrl.Result{}, nil } - // An interrupted defrag may still be running on the server. Recover its - // reservation before allowing a pod rollout to take another member down. - if waiting, err := r.resumeMaintenance(ctx, toposerver); err != nil || waiting { - return ctrl.Result{RequeueAfter: statusRecheckDelay}, err - } - - // Validate StorageClass dependency before StatefulSet apply - if err := r.validateEtcdStorageClassDependency(ctx, toposerver); err != nil { - if isStorageClassDependencyError(err) { - logger.Info("StorageClass dependency is not ready; requeueing", - "after", storageClassDependencyRequeue) - return ctrl.Result{RequeueAfter: storageClassDependencyRequeue}, nil - } - monitoring.RecordSpanError(span, err) - logger.Error(err, "Failed to validate etcd StorageClass") - return ctrl.Result{}, err - } - - // Reconcile StatefulSet + // Restore Services and the serving Certificate before checking maintenance health. + // Reconcile headless Service { - ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcileStatefulSet") - if err := r.reconcileStatefulSet(ctx, toposerver); err != nil { + ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcileHeadlessService") + if err := r.reconcileHeadlessService(ctx, toposerver); err != nil { monitoring.RecordSpanError(childSpan, err) childSpan.End() - logger.Error(err, "Failed to reconcile StatefulSet") + logger.Error(err, "Failed to reconcile headless Service") r.Recorder.Eventf( toposerver, "Warning", "FailedApply", - "Failed to reconcile StatefulSet: %v", + "Failed to reconcile headless Service: %v", err, ) return ctrl.Result{}, err @@ -118,18 +101,18 @@ func (r *TopoServerReconciler) Reconcile( childSpan.End() } - // Reconcile PodDisruptionBudget + // Reconcile client Service { - ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcilePodDisruptionBudget") - if err := r.reconcilePodDisruptionBudget(ctx, toposerver); err != nil { + ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcileClientService") + if err := r.reconcileClientService(ctx, toposerver); err != nil { monitoring.RecordSpanError(childSpan, err) childSpan.End() - logger.Error(err, "Failed to reconcile PodDisruptionBudget") + logger.Error(err, "Failed to reconcile client Service") r.Recorder.Eventf( toposerver, "Warning", "FailedApply", - "Failed to reconcile PodDisruptionBudget: %v", + "Failed to reconcile client Service: %v", err, ) return ctrl.Result{}, err @@ -137,18 +120,18 @@ func (r *TopoServerReconciler) Reconcile( childSpan.End() } - // Reconcile headless Service + // Reconcile serving Certificate { - ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcileHeadlessService") - if err := r.reconcileHeadlessService(ctx, toposerver); err != nil { + ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcileCertificate") + if err := r.reconcileCertificate(ctx, toposerver); err != nil { monitoring.RecordSpanError(childSpan, err) childSpan.End() - logger.Error(err, "Failed to reconcile headless Service") + logger.Error(err, "Failed to reconcile serving Certificate") r.Recorder.Eventf( toposerver, "Warning", "FailedApply", - "Failed to reconcile headless Service: %v", + "Failed to reconcile serving Certificate: %v", err, ) return ctrl.Result{}, err @@ -156,18 +139,36 @@ func (r *TopoServerReconciler) Reconcile( childSpan.End() } - // Reconcile client Service + // An interrupted defrag may still be running on the server. Recover its + // reservation before allowing a pod rollout to take another member down. + if waiting, err := r.resumeMaintenance(ctx, toposerver); err != nil || waiting { + return ctrl.Result{RequeueAfter: statusRecheckDelay}, err + } + + // Validate StorageClass dependency before StatefulSet apply + if err := r.validateEtcdStorageClassDependency(ctx, toposerver); err != nil { + if isStorageClassDependencyError(err) { + logger.Info("StorageClass dependency is not ready; requeueing", + "after", storageClassDependencyRequeue) + return ctrl.Result{RequeueAfter: storageClassDependencyRequeue}, nil + } + monitoring.RecordSpanError(span, err) + logger.Error(err, "Failed to validate etcd StorageClass") + return ctrl.Result{}, err + } + + // Reconcile StatefulSet { - ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcileClientService") - if err := r.reconcileClientService(ctx, toposerver); err != nil { + ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcileStatefulSet") + if err := r.reconcileStatefulSet(ctx, toposerver); err != nil { monitoring.RecordSpanError(childSpan, err) childSpan.End() - logger.Error(err, "Failed to reconcile client Service") + logger.Error(err, "Failed to reconcile StatefulSet") r.Recorder.Eventf( toposerver, "Warning", "FailedApply", - "Failed to reconcile client Service: %v", + "Failed to reconcile StatefulSet: %v", err, ) return ctrl.Result{}, err @@ -175,18 +176,18 @@ func (r *TopoServerReconciler) Reconcile( childSpan.End() } - // Reconcile serving Certificate + // Reconcile PodDisruptionBudget { - ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcileCertificate") - if err := r.reconcileCertificate(ctx, toposerver); err != nil { + ctx, childSpan := monitoring.StartChildSpan(ctx, "TopoServer.ReconcilePodDisruptionBudget") + if err := r.reconcilePodDisruptionBudget(ctx, toposerver); err != nil { monitoring.RecordSpanError(childSpan, err) childSpan.End() - logger.Error(err, "Failed to reconcile serving Certificate") + logger.Error(err, "Failed to reconcile PodDisruptionBudget") r.Recorder.Eventf( toposerver, "Warning", "FailedApply", - "Failed to reconcile serving Certificate: %v", + "Failed to reconcile PodDisruptionBudget: %v", err, ) return ctrl.Result{}, err From fa072d9b35c318250293e5b7d2e0385126ff435c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ver=C3=B3nica=20L=C3=B3pez?= Date: Fri, 25 Sep 2026 06:24:50 -0600 Subject: [PATCH 3/3] test(toposerver): cover etcd maintenance failure cases MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Cover unexpected member counts, learners, and zero member IDs. - Verify status RPC failures abort maintenance without reserving a target. - Verify incomplete leadership transfers block defragmentation and retain the reservation. Signed-off-by: Verónica López --- .../controller/toposerver/maintenance_test.go | 125 ++++++++++++++++-- 1 file changed, 117 insertions(+), 8 deletions(-) diff --git a/pkg/resource-handler/controller/toposerver/maintenance_test.go b/pkg/resource-handler/controller/toposerver/maintenance_test.go index 714da5f5e..83123b0eb 100644 --- a/pkg/resource-handler/controller/toposerver/maintenance_test.go +++ b/pkg/resource-handler/controller/toposerver/maintenance_test.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "strings" "testing" "time" @@ -24,23 +25,29 @@ import ( ) type fakeEtcdMaintenance struct { - statuses map[string]*clientv3.StatusResponse - defragged []string - moved []uint64 - healthErr error - defragErr error - afterDefrag func() + statuses map[string]*clientv3.StatusResponse + members []*pb.Member + defragged []string + moved []uint64 + statusErr error + healthErr error + defragErr error + skipLeaderTransfer bool + afterDefrag func() } func (f *fakeEtcdMaintenance) Status( _ context.Context, ep string, ) (*clientv3.StatusResponse, error) { + if f.statusErr != nil { + return nil, f.statusErr + } return f.statuses[ep], nil } func (f *fakeEtcdMaintenance) Members(context.Context) (*clientv3.MemberListResponse, error) { - return &clientv3.MemberListResponse{Members: []*pb.Member{{ID: 1}, {ID: 2}, {ID: 3}}}, nil + return &clientv3.MemberListResponse{Members: f.members}, nil } func (f *fakeEtcdMaintenance) Health(context.Context, string) error { return f.healthErr } func (f *fakeEtcdMaintenance) Defragment(_ context.Context, ep string) error { @@ -53,6 +60,9 @@ func (f *fakeEtcdMaintenance) Defragment(_ context.Context, ep string) error { func (f *fakeEtcdMaintenance) MoveLeader(_ context.Context, _ string, id uint64) error { f.moved = append(f.moved, id) + if f.skipLeaderTransfer { + return nil + } for _, s := range f.statuses { s.Leader = id } @@ -87,7 +97,10 @@ func maintenanceFixture( UpdateRevision: "rev1", } objects := []client.Object{ts, sts} - f := &fakeEtcdMaintenance{statuses: map[string]*clientv3.StatusResponse{}} + f := &fakeEtcdMaintenance{ + statuses: map[string]*clientv3.StatusResponse{}, + members: []*pb.Member{{ID: 1}, {ID: 2}, {ID: 3}}, + } for i, ep := range maintenanceEndpoints(ts) { f.statuses[ep] = &clientv3.StatusResponse{ Header: &pb.ResponseHeader{ClusterId: 123, MemberId: uint64(i + 1)}, @@ -230,6 +243,102 @@ func TestEtcdMaintenanceHealthGates(t *testing.T) { } } +func TestEtcdMaintenanceRejectsInvalidMemberResponses(t *testing.T) { + statusErr := errors.New("status RPC unavailable") + for _, tc := range []struct { + name string + mutate func(*fakeEtcdMaintenance) + wantError string + wantCause error + }{ + { + name: "missing member", + mutate: func(f *fakeEtcdMaintenance) { f.members = f.members[:2] }, + wantError: "etcd membership differs from expected replicas", + }, + { + name: "extra member", + mutate: func(f *fakeEtcdMaintenance) { f.members = append(f.members, &pb.Member{ID: 4}) }, + wantError: "etcd membership differs from expected replicas", + }, + { + name: "learner member", + mutate: func(f *fakeEtcdMaintenance) { f.members[1].IsLearner = true }, + wantError: "etcd membership includes an unready voting member", + }, + { + name: "zero member ID", + mutate: func(f *fakeEtcdMaintenance) { f.members[1].ID = 0 }, + wantError: "etcd membership includes an unready voting member", + }, + { + name: "status RPC failure", + mutate: func(f *fakeEtcdMaintenance) { f.statusErr = statusErr }, + wantError: "status: status RPC unavailable", + wantCause: statusErr, + }, + } { + t.Run(tc.name, func(t *testing.T) { + r, ts, f := maintenanceFixture(t) + tc.mutate(f) + err := r.reconcileMaintenance(t.Context(), ts) + if err == nil || !strings.Contains(err.Error(), tc.wantError) { + t.Fatalf("expected %q, got %v", tc.wantError, err) + } + if tc.wantCause != nil && !errors.Is(err, tc.wantCause) { + t.Errorf("error %v does not wrap %v", err, tc.wantCause) + } + if len(f.moved) != 0 || len(f.defragged) != 0 { + t.Errorf( + "maintenance changed unhealthy members: transfers=%v defrags=%v", + f.moved, + f.defragged, + ) + } + fresh := &multigresv1alpha1.TopoServer{} + if err := r.Get(t.Context(), client.ObjectKeyFromObject(ts), fresh); err != nil { + t.Fatal(err) + } + if fresh.Status.EtcdMaintenance != nil { + t.Error("failed health checks created a maintenance reservation") + } + }) + } +} + +func TestEtcdMaintenanceIncompleteLeadershipTransfer(t *testing.T) { + r, ts, f := maintenanceFixture(t) + f.skipLeaderTransfer = true + err := r.reconcileMaintenance(t.Context(), ts) + if err == nil || err.Error() != "etcd leadership transfer has not completed" { + t.Fatalf("expected incomplete leadership transfer, got %v", err) + } + if len(f.moved) != 1 || f.moved[0] != 2 { + t.Errorf("leadership transfer attempts = %v, want [2]", f.moved) + } + if len(f.defragged) != 0 { + t.Errorf( + "defragmented a member whose leadership transfer did not complete: %v", + f.defragged, + ) + } + fresh := &multigresv1alpha1.TopoServer{} + if err := r.Get(t.Context(), client.ObjectKeyFromObject(ts), fresh); err != nil { + t.Fatal(err) + } + state := fresh.Status.EtcdMaintenance + if state == nil || !state.InProgress || state.Endpoint != maintenanceEndpoints(ts)[0] { + t.Fatalf("incomplete transfer did not retain the target reservation: %+v", state) + } + r2 := *r + if err := r2.reconcileMaintenance(t.Context(), fresh); err != nil { + t.Fatal(err) + } + if len(f.moved) != 1 || len(f.defragged) != 0 { + t.Error("maintenance restarted while the transfer reservation remained active") + } +} + func TestEtcdMaintenanceInterruptedOperation(t *testing.T) { r, ts, f := maintenanceFixture(t) f.defragErr = context.DeadlineExceeded