From 0bcf65fcc851e5ca93b252d366fe0d3cd3379b8d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ver=C3=B3nica=20L=C3=B3pez?= Date: Fri, 25 Sep 2026 13:55:09 -0600 Subject: [PATCH 1/2] fix(topology): report quorum health and failover readiness MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Verónica López --- api/v1alpha1/toposerver_types.go | 4 + api/v1alpha1/zz_generated.deepcopy.go | 4 + .../crd/bases/multigres.com_toposervers.yaml | 5 + config/deploy-observability/prometheus.yaml | 26 ++ config/monitoring/prometheus-rules.yaml | 125 ++++++- docs/monitoring/runbooks/TopologyHealth.md | 137 ++++++++ docs/observability.md | 13 + .../controller/multigrescluster/health.go | 197 +++++++++++ .../multigrescluster/health_test.go | 305 ++++++++++++++++++ .../multigrescluster_controller.go | 31 +- .../multigrescluster/reconcile_topology.go | 33 +- .../controller/multigrescluster/status.go | 32 ++ .../multigrescluster/status_test.go | 27 ++ pkg/monitoring/alerts_test.go | 44 +++ pkg/monitoring/metrics.go | 23 +- pkg/monitoring/testdata/topology_alerts.yaml | 147 +++++++++ pkg/monitoring/topology.go | 158 +++++++++ pkg/monitoring/topology_test.go | 56 ++++ .../controller/toposerver/health.go | 204 ++++++++++++ .../controller/toposerver/health_etcd_test.go | 63 ++++ .../controller/toposerver/health_test.go | 162 ++++++++++ .../toposerver/maintenance_etcd_test.go | 9 +- .../maintenance_status_integration_test.go | 19 ++ .../toposerver/toposerver_controller.go | 8 + 24 files changed, 1785 insertions(+), 47 deletions(-) create mode 100644 docs/monitoring/runbooks/TopologyHealth.md create mode 100644 pkg/cluster-handler/controller/multigrescluster/health.go create mode 100644 pkg/cluster-handler/controller/multigrescluster/health_test.go create mode 100644 pkg/monitoring/alerts_test.go create mode 100644 pkg/monitoring/testdata/topology_alerts.yaml create mode 100644 pkg/monitoring/topology.go create mode 100644 pkg/monitoring/topology_test.go create mode 100644 pkg/resource-handler/controller/toposerver/health.go create mode 100644 pkg/resource-handler/controller/toposerver/health_etcd_test.go create mode 100644 pkg/resource-handler/controller/toposerver/health_test.go diff --git a/api/v1alpha1/toposerver_types.go b/api/v1alpha1/toposerver_types.go index 52c9aa8ed..ac7776b97 100644 --- a/api/v1alpha1/toposerver_types.go +++ b/api/v1alpha1/toposerver_types.go @@ -71,6 +71,10 @@ type TopoServerSpec struct { // TopoServerStatus defines the observed state of TopoServer. type TopoServerStatus struct { + // HealthCheckedAt is the completion time of the latest direct etcd health probe. + // +optional + HealthCheckedAt *metav1.Time `json:"healthCheckedAt,omitempty"` + // ObservedGeneration is the most recent generation observed. // +optional ObservedGeneration int64 `json:"observedGeneration,omitempty"` diff --git a/api/v1alpha1/zz_generated.deepcopy.go b/api/v1alpha1/zz_generated.deepcopy.go index dc88876c7..76b42e366 100644 --- a/api/v1alpha1/zz_generated.deepcopy.go +++ b/api/v1alpha1/zz_generated.deepcopy.go @@ -2419,6 +2419,10 @@ func (in *TopoServerSpec) DeepCopy() *TopoServerSpec { // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *TopoServerStatus) DeepCopyInto(out *TopoServerStatus) { *out = *in + if in.HealthCheckedAt != nil { + in, out := &in.HealthCheckedAt, &out.HealthCheckedAt + *out = (*in).DeepCopy() + } if in.Conditions != nil { in, out := &in.Conditions, &out.Conditions *out = make([]v1.Condition, len(*in)) diff --git a/config/crd/bases/multigres.com_toposervers.yaml b/config/crd/bases/multigres.com_toposervers.yaml index d71d98929..3257c1346 100644 --- a/config/crd/bases/multigres.com_toposervers.yaml +++ b/config/crd/bases/multigres.com_toposervers.yaml @@ -1520,6 +1520,11 @@ spec: - inProgress - lastAttemptTime type: object + healthCheckedAt: + description: HealthCheckedAt is the completion time of the latest + direct etcd health probe. + format: date-time + type: string message: description: Message provides details about the current phase. type: string diff --git a/config/deploy-observability/prometheus.yaml b/config/deploy-observability/prometheus.yaml index 53858f4b2..eeb7e2e0c 100644 --- a/config/deploy-observability/prometheus.yaml +++ b/config/deploy-observability/prometheus.yaml @@ -60,6 +60,9 @@ spec: - exemplar-storage - native-histograms enableOTLPReceiver: true + additionalScrapeConfigs: + name: kubelet-cadvisor-scrape + key: scrape-configs.yaml serviceMonitorSelector: matchLabels: app.kubernetes.io/name: multigres-operator @@ -125,3 +128,26 @@ spec: targetPort: 9090 name: http type: ClusterIP + +--- +# cAdvisor supplies per-container memory usage for the topology pressure alert. +apiVersion: v1 +kind: Secret +metadata: + name: kubelet-cadvisor-scrape +stringData: + scrape-configs.yaml: | + - job_name: kubelet-cadvisor + scheme: https + metrics_path: /metrics/cadvisor + scrape_interval: 30s + kubernetes_sd_configs: + - role: node + authorization: + credentials_file: /var/run/secrets/kubernetes.io/serviceaccount/token + tls_config: + insecure_skip_verify: true + metric_relabel_configs: + - source_labels: [__name__, container] + regex: 'container_memory_working_set_bytes;etcd' + action: keep diff --git a/config/monitoring/prometheus-rules.yaml b/config/monitoring/prometheus-rules.yaml index c93896db6..77b25cca3 100644 --- a/config/monitoring/prometheus-rules.yaml +++ b/config/monitoring/prometheus-rules.yaml @@ -22,6 +22,9 @@ spec: Controller {{ $labels.controller }} has had a non-zero reconcile error rate for {{ $labels.namespace }}/{{ $labels.name }} over the last 5 minutes (current: {{ $value | humanize }}/s). + Inspect TopologyReady for topology errors and FailoverReady for their + impact on automatic failover; MultigresFailoverUnavailable reports + these failures separately at critical severity. runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/MultigresClusterReconcileErrors.md" - alert: MultigresClusterDegraded @@ -78,6 +81,126 @@ spec: (current: {{ $value | humanize }}/s). runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/MultigresWebhookErrors.md" + # ── Topology and failover ────────────────────────────────── + + - alert: MultigresTopologyQuorumUnavailable + expr: multigres_operator_toposerver_quorum_available == 0 + for: 1m + labels: + severity: critical + annotations: + summary: "Topology quorum unavailable for {{ $labels.target_namespace }}/{{ $labels.name }}" + description: >- + {{ $labels.reason }}: the operator cannot confirm quorum for one consistent etcd cluster. + Cluster {{ $labels.cluster }} has lost verified failover protection; + existing SQL connections may still serve traffic. + runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" + + - alert: MultigresTopologyHealthUnknown + expr: >- + (multigres_operator_toposerver_quorum_available < 0) + or (time() - multigres_operator_toposerver_health_checked_timestamp_seconds > 120) + for: 1m + labels: + severity: warning + annotations: + summary: "Topology health observation missing for {{ $labels.target_namespace }}/{{ $labels.name }}" + description: "Etcd quorum checks are failing or stale. Failover protection cannot be verified for cluster {{ $labels.cluster }}." + runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" + + - alert: MultigresTopologyMemberUnavailable + expr: multigres_operator_toposerver_member_up == 0 + for: 2m + labels: + severity: warning + annotations: + summary: "Etcd member {{ $labels.member }} is unavailable" + description: >- + Member {{ $labels.target_namespace }}/{{ $labels.member }} cannot complete + a linearizable read. Check the quorum alert before restarting other members. + runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" + + - alert: MultigresTopologyBackendNearQuota + expr: >- + multigres_operator_toposerver_backend_bytes + / multigres_operator_toposerver_backend_quota_bytes > 0.8 + for: 5m + labels: + severity: warning + annotations: + summary: "Etcd backend nearing quota on {{ $labels.member }}" + description: >- + Backend usage on {{ $labels.target_namespace }}/{{ $labels.member }} + is {{ $value | humanizePercentage }} of its quota. Exhausting the quota + blocks topology writes and puts failover protection at risk. + runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" + + - alert: MultigresTopologyMemoryPressure + expr: >- + label_replace(label_replace( + max by (namespace, pod) (container_memory_working_set_bytes{container="etcd", pod!=""}), + "target_namespace", "$1", "namespace", "(.+)"), + "member", "$1", "pod", "(.+)") + / on (target_namespace, member) group_left (cluster, name) + (max by (cluster, name, target_namespace, member) + (multigres_operator_toposerver_memory_limit_bytes) > 0) + > 0.8 + for: 1m + labels: + severity: warning + annotations: + summary: "Etcd memory pressure on {{ $labels.member }}" + description: >- + Member {{ $labels.target_namespace }}/{{ $labels.member }} is using + {{ $value | humanizePercentage }} of its memory limit. An OOM kill + can remove a quorum member and interrupt failover protection for {{ $labels.cluster }}. + runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" + + - alert: MultigresTopologyMemberOOMKilled + expr: >- + (time() - multigres_operator_toposerver_member_oom_timestamp_seconds < 300) + and (multigres_operator_toposerver_member_oom_timestamp_seconds > 0) + labels: + severity: critical + annotations: + summary: "Etcd member {{ $labels.member }} was OOM-killed" + description: "An OOM kill removed member {{ $labels.target_namespace }}/{{ $labels.member }}. Check quorum and FailoverReady for cluster {{ $labels.cluster }} before restarting another member." + runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" + + - alert: MultigresTopologyMemberRestarting + expr: increase(multigres_operator_toposerver_member_restarts_total[10m]) >= 3 + labels: + severity: warning + annotations: + summary: "Etcd member {{ $labels.member }} is repeatedly restarting" + description: "Member {{ $labels.target_namespace }}/{{ $labels.member }} restarted at least three times in ten minutes. Check termination reasons, memory pressure, and quorum health." + runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" + + - alert: MultigresFailoverUnavailable + expr: multigres_operator_cluster_failover_ready == 0 + for: 1m + labels: + severity: critical + annotations: + summary: "Failover protection unavailable for {{ $labels.target_namespace }}/{{ $labels.cluster }}" + description: >- + FailoverReady is false: {{ $labels.reason }}. Topology access or shard + orchestrators are unavailable. SQL availability does not imply that + automatic failover is working. + runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" + + - alert: MultigresFailoverHealthUnknown + expr: >- + (multigres_operator_cluster_failover_ready < 0) + or (time() - multigres_operator_cluster_failover_checked_timestamp_seconds > 120) + for: 2m + labels: + severity: warning + annotations: + summary: "Failover protection unverified for {{ $labels.target_namespace }}/{{ $labels.cluster }}" + description: "Failover observations are unknown or stale. Inspect TopologyQuorumAvailable, TopologyReady, and shard OrchReady conditions." + runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" + # ── Backup & Drain ──────────────────────────────────────── - alert: MultigresBackupStale @@ -138,7 +261,7 @@ spec: # ── Saturation ──────────────────────────────────────────── - alert: MultigresControllerSaturated - expr: workqueue_depth{name=~"multigrescluster|tablegroup|shard|cell|toposerver"} > 50 + expr: workqueue_depth{name=~"multigrescluster|tablegroup|shard|cell|toposerver|toposerver-health"} > 50 for: 10m labels: severity: warning diff --git a/docs/monitoring/runbooks/TopologyHealth.md b/docs/monitoring/runbooks/TopologyHealth.md new file mode 100644 index 000000000..a02099421 --- /dev/null +++ b/docs/monitoring/runbooks/TopologyHealth.md @@ -0,0 +1,137 @@ +# Topology health and failover protection + +These alerts distinguish topology and failover failures from SQL availability. +An existing primary may continue serving queries while topology access or +multiorch is unavailable. `Available=True` alone does not confirm that a new +primary can be elected. + +## Conditions and probes + +The operator checks each managed etcd member every 30 seconds with a status RPC +and a read-only, linearizable Get. The probes run independently of resource +reconciliation and defragmentation, with a five-second deadline. A successful +linearizable read confirms quorum even when the operator cannot reach every +member. Status RPCs alone do not confirm quorum. + +`TopoServer.status.conditions[QuorumAvailable]` reports the result. +`TopologyUnreachable` means no member answered the operator; it does not prove +that the members have lost quorum among themselves. `QuorumUnavailable` means +members answered status requests but none completed the read. `ClusterMismatch` +means the endpoints reported different etcd cluster IDs. Client setup +failures, including invalid TLS credentials, produce `Unknown` with reason +`ProbeFailed`. `status.healthCheckedAt` records the latest probe completion. +The TopoServer `Ready` condition continues to report StatefulSet readiness. + +On a MultigresCluster: + +- `TopologyReady` reports the result of topology registration, including failures + during the startup grace period. +- `TopologyQuorumAvailable` aggregates managed global and cell-local topology + servers. Missing servers report false; observations older than two minutes or + from an earlier generation report unknown. +- `FailoverReady` requires successful topology access, current managed quorum + observations, and `OrchReady` for every desired shard at its current generation. + A missing or unready shard reports false. Unknown observations prevent a true + result. A previously initialized cluster with lost failover readiness is degraded. + +External topology is not probed directly. Its quorum condition is unknown with +reason `ExternalTopology`; failover readiness uses topology registration and shard +orchestrator readiness. Monitor external etcd through its owner. + +## Alert thresholds + +| Alert | Threshold | Severity | +| --- | --- | --- | +| `MultigresTopologyQuorumUnavailable` | No successful quorum read for one minute | critical | +| `MultigresTopologyMemberUnavailable` | A member cannot complete reads for two minutes | warning | +| `MultigresTopologyHealthUnknown` | Unknown probe or observation over two minutes old, for one minute | warning | +| `MultigresTopologyBackendNearQuota` | Backend over 80% of quota for five minutes | warning | +| `MultigresTopologyMemoryPressure` | Working set over 80% of memory limit for one minute | warning | +| `MultigresTopologyMemberOOMKilled` | An observed OOM termination within the last five minutes | critical | +| `MultigresTopologyMemberRestarting` | At least three restarts in ten minutes | warning | +| `MultigresFailoverUnavailable` | FailoverReady false for one minute | critical | +| `MultigresFailoverHealthUnknown` | Unknown failover readiness or an observation over two minutes old, for two minutes | warning | + +Alert delays are in addition to the probe, scrape, and evaluation intervals. +The pressure alerts can warn before resource exhaustion; a rapid memory spike +can reach the limit between scrapes. + +## Investigate + +Read the conditions and their messages, then inspect the affected members: + +```bash +kubectl -n get multigrescluster -o yaml +kubectl -n get toposervers -l multigres.com/cluster= -o yaml +kubectl -n get shards -l multigres.com/cluster= -o yaml +kubectl -n describe pod +kubectl -n logs -c etcd --previous +``` + +For quorum or connectivity failures, check member logs, headless Service and DNS, +network reachability from the operator, and TLS Secrets. Check all members before +restarting one: taking down another voter can turn a single-member failure into +quorum loss. When topology is healthy but `FailoverReady` names an orchestrator, +inspect the affected shard's multiorch Deployments and pod readiness. + +For backend pressure, compare total backend bytes with bytes in use. Review +compaction and defragmentation settings in [Topology maintenance](../../topology-maintenance.md). +For memory pressure or OOMs, inspect working-set history and the running pod's +memory limit. Raising the backend quota does not increase available memory. + +## Metrics and queries + +Topology metrics use `cluster`, `name`, and `target_namespace`. Member metrics also +include `member` (the pod name). These labels avoid collisions with the operator +scrape target's own `namespace` and `pod` labels. + +All names below begin with `multigres_operator_`: + +| Suffix | Observation | +| --- | --- | +| `toposerver_quorum_available` | 1 true, 0 false, -1 unknown; includes `reason` | +| `toposerver_health_checked_timestamp_seconds` | Last completed probe | +| `toposerver_member_up` | Member completed a linearizable read | +| `toposerver_backend_bytes` | Total backend size | +| `toposerver_backend_in_use_bytes` | Backend bytes in use | +| `toposerver_backend_quota_bytes` | Quota configured on the running pod | +| `toposerver_revision` | Current MVCC revision | +| `toposerver_memory_limit_bytes` | Running pod's memory limit; zero means unlimited | +| `toposerver_member_restarts_total` | Observed container restart count, reset on pod replacement | +| `toposerver_member_oom_timestamp_seconds` | Latest OOM termination still recorded in pod status, or zero | +| `cluster_failover_ready` | 1 true, 0 false, -1 unknown; labels `cluster`, `target_namespace`, `reason` | +| `cluster_failover_checked_timestamp_seconds` | Last failover readiness observation | + +Backend and revision series are removed when a member cannot be observed. They +are not carried forward as current measurements. Kubernetes retains only the +current and last container termination, so the OOM timestamp is not a complete +history of OOM events. Restart counts are exported as gauges because they come +from pod status; `increase` handles their resets. + +Revision growth per second: + +```promql +deriv(multigres_operator_toposerver_revision[5m]) +``` + +Backend quota utilization: + +```promql +multigres_operator_toposerver_backend_bytes +/ multigres_operator_toposerver_backend_quota_bytes +``` + +The memory alert requires kubelet/cAdvisor +`container_memory_working_set_bytes{container="etcd"}` with `namespace` and `pod` +labels. The local `deploy-observability` overlay configures this scrape using the +Prometheus service account. Other installations must enable their kubelet scrape. +A missing scrape or an unlimited container produces no memory-pressure alert; +check scrape health and configure limits when enabling this alert. The local +scrape accepts the kubelet's self-signed serving certificate; production scrape +configuration should use the cluster's trusted kubelet certificates. + +To run the alert evaluation fixtures: + +```bash +PROMTOOL=/path/to/promtool go test ./pkg/monitoring -run TestTopologyAlertRules +``` diff --git a/docs/observability.md b/docs/observability.md index 210fa62e7..cb7076716 100644 --- a/docs/observability.md +++ b/docs/observability.md @@ -54,6 +54,19 @@ kubectl apply -f config/monitoring/prometheus-rules.yaml Each alert links to a dedicated runbook with investigation steps, PromQL queries, and remediation actions. +## Topology health + +Managed etcd members are probed every 30 seconds. The operator exports quorum, +backend size and quota, MVCC revision, memory limits, restarts, and observed OOM +terminations. `TopologyQuorumAvailable` and `FailoverReady` on MultigresCluster +separate topology and orchestrator health from SQL availability. Quorum loss and +lost failover readiness raise critical alerts after one minute. + +The [topology health runbook](monitoring/runbooks/TopologyHealth.md) lists the +metrics, thresholds, condition meanings, and investigation steps. Memory-pressure +alerts require kubelet/cAdvisor metrics; the local observability overlay includes +that scrape. + ## Grafana Dashboards Three Grafana dashboards are included in `config/monitoring/`: diff --git a/pkg/cluster-handler/controller/multigrescluster/health.go b/pkg/cluster-handler/controller/multigrescluster/health.go new file mode 100644 index 000000000..7bb3288a6 --- /dev/null +++ b/pkg/cluster-handler/controller/multigrescluster/health.go @@ -0,0 +1,197 @@ +package multigrescluster + +import ( + "context" + "fmt" + "time" + + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + "github.com/multigres/multigres-operator/pkg/data-handler/topo" + "github.com/multigres/multigres-operator/pkg/monitoring" + "github.com/multigres/multigres-operator/pkg/resolver" + "github.com/multigres/multigres-operator/pkg/util/metadata" + "github.com/multigres/multigres-operator/pkg/util/name" +) + +const ( + conditionTopologyQuorumAvailable = "TopologyQuorumAvailable" + conditionFailoverReady = "FailoverReady" + clusterHealthInterval = 30 * time.Second + topologyHealthMaxAge = 2 * time.Minute +) + +type shardHealthIdentity struct { + database multigresv1alpha1.DatabaseName + tableGroup multigresv1alpha1.TableGroupName + shard multigresv1alpha1.ShardName +} + +func (r *MultigresClusterReconciler) updateHealthConditions( + ctx context.Context, + cluster *multigresv1alpha1.MultigresCluster, + globalTopoSpec *multigresv1alpha1.GlobalTopoServerSpec, + servers []multigresv1alpha1.TopoServer, +) { + quorum := r.topologyQuorumCondition(ctx, cluster, globalTopoSpec, servers) + meta.SetStatusCondition(&cluster.Status.Conditions, quorum) + failover := metav1.Condition{ + Type: conditionFailoverReady, + ObservedGeneration: cluster.Generation, + Status: metav1.ConditionTrue, + Reason: "FailoverReady", + Message: "Topology access and all shard orchestrators are ready", + } + orch := failover + shards := &multigresv1alpha1.ShardList{} + if err := r.childStatusReader(). + List(ctx, shards, client.InNamespace(cluster.Namespace), client.MatchingLabels{metadata.LabelMultigresCluster: cluster.Name}); err != nil { + orch.Status, orch.Reason, orch.Message = metav1.ConditionUnknown, "ObservationFailed", "Could not observe shard orchestrators: "+err.Error() + } else { + orch = orchestratorCondition(cluster, shards.Items, orch) + } + access := meta.FindStatusCondition(cluster.Status.Conditions, conditionTopologyReady) + switch { + case quorum.Status == metav1.ConditionFalse: + failover.Status, failover.Reason, failover.Message = quorum.Status, quorum.Reason, quorum.Message + case access != nil && access.ObservedGeneration == cluster.Generation && access.Status == metav1.ConditionFalse: + failover.Status, failover.Reason = metav1.ConditionFalse, access.Reason + failover.Message = "Failover protection is unavailable: " + access.Message + case orch.Status == metav1.ConditionFalse: + failover = orch + case quorum.Status == metav1.ConditionUnknown && quorum.Reason != "ExternalTopology": + failover.Status, failover.Reason, failover.Message = quorum.Status, quorum.Reason, quorum.Message + case access == nil || access.ObservedGeneration != cluster.Generation || access.Status != metav1.ConditionTrue: + failover.Status, failover.Reason, failover.Message = metav1.ConditionUnknown, "TopologyUnchecked", "Topology access has not been verified for the current configuration" + default: + failover = orch + } + meta.SetStatusCondition(&cluster.Status.Conditions, failover) + monitoring.SetFailoverHealth(cluster.Name, cluster.Namespace, failover) +} + +func (r *MultigresClusterReconciler) topologyQuorumCondition( + ctx context.Context, + cluster *multigresv1alpha1.MultigresCluster, + global *multigresv1alpha1.GlobalTopoServerSpec, + servers []multigresv1alpha1.TopoServer, +) metav1.Condition { + condition := metav1.Condition{ + Type: conditionTopologyQuorumAvailable, + ObservedGeneration: cluster.Generation, + Status: metav1.ConditionTrue, + Reason: "QuorumAvailable", + Message: "All managed topology servers have quorum", + } + var expected []string + external := global.External != nil + if global.Etcd != nil { + expected = append(expected, cluster.Name+"-global-topo") + } + res := resolver.NewResolver(r.Client, cluster.Namespace) + for _, cell := range cluster.Spec.Cells { + cell.CellTemplate = cluster.Spec.EffectiveCellTemplate(cell.CellTemplate) + _, _, local, err := res.ResolveCell(ctx, cluster, &cell) + if err != nil { + condition.Status, condition.Reason, condition.Message = metav1.ConditionUnknown, "ObservationFailed", "Could not resolve local topology: "+err.Error() + return condition + } + if local != nil && local.Etcd != nil { + cellName := name.JoinWithConstraints( + name.DefaultConstraints, + cluster.Name, + string(cell.Name), + ) + expected = append(expected, topo.ManagedLocalTopoServerName(cellName)) + } + if local != nil && local.External != nil { + external = true + } + } + byName := make(map[string]*multigresv1alpha1.TopoServer, len(servers)) + for i := range servers { + byName[servers[i].Name] = &servers[i] + } + for _, serverName := range expected { + server := byName[serverName] + if server == nil || !server.DeletionTimestamp.IsZero() { + condition.Status, condition.Reason = metav1.ConditionFalse, "TopologyMissing" + condition.Message = fmt.Sprintf( + "Topology server %s is missing or being deleted; failover protection is unavailable", + serverName, + ) + return condition + } + observed := meta.FindStatusCondition(server.Status.Conditions, "QuorumAvailable") + if observed == nil || observed.ObservedGeneration != server.Generation || + server.Status.HealthCheckedAt == nil || + time.Since(server.Status.HealthCheckedAt.Time) > topologyHealthMaxAge { + condition.Status, condition.Reason, condition.Message = metav1.ConditionUnknown, "TopologyHealthStale", fmt.Sprintf( + "No recent quorum observation for %s", + serverName, + ) + continue + } + if observed.Status == metav1.ConditionFalse { + condition.Status, condition.Reason, condition.Message = observed.Status, observed.Reason, fmt.Sprintf( + "%s: %s", + serverName, + observed.Message, + ) + return condition + } + if observed.Status == metav1.ConditionUnknown { + condition.Status, condition.Reason, condition.Message = observed.Status, observed.Reason, fmt.Sprintf( + "%s: %s", + serverName, + observed.Message, + ) + } + } + if external && condition.Status == metav1.ConditionTrue { + condition.Status, condition.Reason, condition.Message = metav1.ConditionUnknown, "ExternalTopology", "External topology quorum is not monitored; topology registration reports access separately" + } + return condition +} + +func orchestratorCondition( + cluster *multigresv1alpha1.MultigresCluster, + shards []multigresv1alpha1.Shard, + condition metav1.Condition, +) metav1.Condition { + observed := make(map[shardHealthIdentity]*multigresv1alpha1.Shard, len(shards)) + for i := range shards { + s := &shards[i] + observed[shardHealthIdentity{s.Spec.DatabaseName, s.Spec.TableGroupName, s.Spec.ShardName}] = s + } + for _, db := range cluster.Spec.Databases { + for _, tg := range db.TableGroups { + for _, desired := range tg.Shards { + s := observed[shardHealthIdentity{db.Name, tg.Name, desired.Name}] + if s != nil && s.DeletionTimestamp.IsZero() && + s.Status.ObservedGeneration != s.Generation { + condition.Status, condition.Reason, condition.Message = metav1.ConditionUnknown, "OrchestratorHealthStale", fmt.Sprintf( + "Shard %s readiness predates its current configuration", + s.Name, + ) + continue + } + if s == nil || !s.DeletionTimestamp.IsZero() || !s.Status.OrchReady { + condition.Status, condition.Reason = metav1.ConditionFalse, "OrchestratorUnavailable" + condition.Message = fmt.Sprintf( + "Shard %s/%s/%s has no ready multiorch observation; failover protection is unavailable", + db.Name, + tg.Name, + desired.Name, + ) + return condition + } + + } + } + } + return condition +} diff --git a/pkg/cluster-handler/controller/multigrescluster/health_test.go b/pkg/cluster-handler/controller/multigrescluster/health_test.go new file mode 100644 index 000000000..fbe988b08 --- /dev/null +++ b/pkg/cluster-handler/controller/multigrescluster/health_test.go @@ -0,0 +1,305 @@ +package multigrescluster + +import ( + "errors" + "testing" + "time" + + "github.com/multigres/multigres/go/common/topoclient" + "github.com/stretchr/testify/require" + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/client-go/tools/record" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + "github.com/multigres/multigres-operator/pkg/testutil" + "github.com/multigres/multigres-operator/pkg/util/metadata" +) + +func TestFailoverHealthConditions(t *testing.T) { + for _, tc := range []struct { + name string + quorum metav1.ConditionStatus + stale, missing, orchDown, shardStale, shardMissing, accessFailed, external, observationFailed bool + want metav1.ConditionStatus + reason string + }{ + {name: "failed shard observation clears readiness", quorum: metav1.ConditionTrue, observationFailed: true, want: metav1.ConditionUnknown, reason: "ObservationFailed"}, + {name: "healthy", quorum: metav1.ConditionTrue, want: metav1.ConditionTrue, reason: "FailoverReady"}, + {name: "quorum lost with ready pods", quorum: metav1.ConditionFalse, want: metav1.ConditionFalse, reason: "QuorumUnavailable"}, + {name: "stale topology observation", quorum: metav1.ConditionTrue, stale: true, want: metav1.ConditionUnknown, reason: "TopologyHealthStale"}, + {name: "missing topology", missing: true, want: metav1.ConditionFalse, reason: "TopologyMissing"}, + {name: "known orchestrator failure takes precedence over stale quorum", quorum: metav1.ConditionTrue, stale: true, orchDown: true, want: metav1.ConditionFalse, reason: "OrchestratorUnavailable"}, + {name: "multiorch not ready", quorum: metav1.ConditionTrue, orchDown: true, want: metav1.ConditionFalse, reason: "OrchestratorUnavailable"}, + {name: "missing desired shard", quorum: metav1.ConditionTrue, shardMissing: true, want: metav1.ConditionFalse, reason: "OrchestratorUnavailable"}, + {name: "old unready shard generation", quorum: metav1.ConditionTrue, shardStale: true, orchDown: true, want: metav1.ConditionUnknown, reason: "OrchestratorHealthStale"}, + {name: "old shard generation", quorum: metav1.ConditionTrue, shardStale: true, want: metav1.ConditionUnknown, reason: "OrchestratorHealthStale"}, + {name: "topology writes fail despite quorum", quorum: metav1.ConditionTrue, accessFailed: true, want: metav1.ConditionFalse, reason: "TopologyUnavailable"}, + {name: "external topology uses registration and orch", external: true, want: metav1.ConditionTrue, reason: "FailoverReady"}, + } { + t.Run(tc.name, func(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, multigresv1alpha1.AddToScheme(scheme)) + global := &multigresv1alpha1.GlobalTopoServerSpec{Etcd: &multigresv1alpha1.EtcdSpec{}} + if tc.external { + global = &multigresv1alpha1.GlobalTopoServerSpec{ + External: &multigresv1alpha1.ExternalTopoServerSpec{ + Endpoints: []multigresv1alpha1.EndpointUrl{"http://external:2379"}, + }, + } + } + cluster := &multigresv1alpha1.MultigresCluster{ + ObjectMeta: metav1.ObjectMeta{Name: "example", Namespace: "default", Generation: 1}, + Spec: multigresv1alpha1.MultigresClusterSpec{ + GlobalTopoServer: global, + Databases: []multigresv1alpha1.DatabaseConfig{ + { + Name: "postgres", + TableGroups: []multigresv1alpha1.TableGroupConfig{ + { + Name: "default", + Shards: []multigresv1alpha1.ShardConfig{{Name: "0-inf"}}, + }, + }, + }, + }, + }, + } + access := metav1.Condition{ + Type: conditionTopologyReady, + ObservedGeneration: 1, + Status: metav1.ConditionTrue, + Reason: "TopoConnected", + Message: "Topology reachable", + } + if tc.accessFailed { + access.Status, access.Reason, access.Message = metav1.ConditionFalse, "TopologyUnavailable", "Topology write failed" + } + cluster.Status.Conditions = []metav1.Condition{ + access, + {Type: "Available", Status: metav1.ConditionTrue, Reason: "CellsReady"}, + } + checked := metav1.Now() + if tc.stale { + checked = metav1.NewTime(time.Now().Add(-3 * time.Minute)) + } + servers := []multigresv1alpha1.TopoServer{ + { + ObjectMeta: metav1.ObjectMeta{Name: "example-global-topo", Generation: 1}, + Status: multigresv1alpha1.TopoServerStatus{ + Phase: multigresv1alpha1.PhaseHealthy, + HealthCheckedAt: &checked, + Conditions: []metav1.Condition{ + { + Type: "QuorumAvailable", + ObservedGeneration: 1, + Status: tc.quorum, + Reason: "QuorumUnavailable", + Message: "Quorum observation", + }, + }, + }, + }, + } + if tc.missing || tc.external { + servers = nil + } + shard := &multigresv1alpha1.Shard{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example-shard", + Namespace: "default", + Generation: 1, + Labels: map[string]string{metadata.LabelMultigresCluster: "example"}, + }, + Spec: multigresv1alpha1.ShardSpec{ + DatabaseName: "postgres", + TableGroupName: "default", + ShardName: "0-inf", + }, + Status: multigresv1alpha1.ShardStatus{ + ObservedGeneration: 1, + OrchReady: !tc.orchDown, + }, + } + if tc.shardStale { + shard.Status.ObservedGeneration = 0 + } + objects := []client.Object{cluster} + if !tc.shardMissing { + objects = append(objects, shard) + } + c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(objects...).Build() + r := &MultigresClusterReconciler{Client: c, APIReader: c} + if tc.observationFailed { + r.APIReader = testutil.NewFakeClientWithFailures( + c, + &testutil.FailureConfig{OnList: func(list client.ObjectList) error { + if _, ok := list.(*multigresv1alpha1.ShardList); ok { + return errors.New("shard list unavailable") + } + return nil + }}, + ) + } + + r.updateHealthConditions(t.Context(), cluster, global, servers) + condition := meta.FindStatusCondition(cluster.Status.Conditions, conditionFailoverReady) + require.Equal(t, tc.want, condition.Status) + require.Equal(t, tc.reason, condition.Reason) + require.True( + t, + meta.IsStatusConditionTrue(cluster.Status.Conditions, "Available"), + "SQL availability remains separately reported", + ) + if tc.external { + require.Equal( + t, + metav1.ConditionUnknown, + meta.FindStatusCondition( + cluster.Status.Conditions, + conditionTopologyQuorumAvailable, + ).Status, + ) + } + }) + } +} + +func TestTopologyUnavailableClearsPreviousReadinessDuringGracePeriod(t *testing.T) { + scheme := runtime.NewScheme() + require.NoError(t, multigresv1alpha1.AddToScheme(scheme)) + cluster := &multigresv1alpha1.MultigresCluster{ + ObjectMeta: metav1.ObjectMeta{ + Name: "example", + Namespace: "default", + CreationTimestamp: metav1.Now(), + }, + Status: multigresv1alpha1.MultigresClusterStatus{ + Conditions: []metav1.Condition{ + {Type: conditionTopologyReady, Status: metav1.ConditionTrue}, + }, + }, + } + c := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(cluster). + WithStatusSubresource(cluster). + Build() + r := &MultigresClusterReconciler{Client: c, Recorder: record.NewFakeRecorder(10)} + result, err := r.handleTopoUnavailable( + t.Context(), + cluster, + errors.New("linearizable read failed"), + discardHealthLogger{}, + ) + require.NoError(t, err) + require.Equal(t, topoUnavailableRequeueDelay, result.RequeueAfter) + fresh := &multigresv1alpha1.MultigresCluster{} + require.NoError(t, c.Get(t.Context(), client.ObjectKeyFromObject(cluster), fresh)) + condition := meta.FindStatusCondition(fresh.Status.Conditions, conditionTopologyReady) + require.Equal(t, metav1.ConditionFalse, condition.Status) + require.Equal(t, "linearizable read failed", condition.Message) +} + +type discardHealthLogger struct{} + +func (discardHealthLogger) Info(string, ...any) {} +func (discardHealthLogger) Error(error, string, ...any) {} + +func TestReconcileRefreshesFailoverStatusOnTopologyRequeue(t *testing.T) { + scheme := setupScheme() + core, cell, shard, cluster, _, _ := setupFixtures(t) + cluster.CreationTimestamp = metav1.Now() + cluster.Status.Conditions = []metav1.Condition{ + { + Type: conditionTopologyReady, + Status: metav1.ConditionTrue, + Reason: "TopoConnected", + Message: "Topology reachable", + }, + { + Type: conditionFailoverReady, + Status: metav1.ConditionTrue, + Reason: "FailoverReady", + Message: "Orchestrators ready", + }, + } + c := fake.NewClientBuilder(). + WithScheme(scheme). + WithObjects(core, cell, shard, cluster). + WithStatusSubresource(cluster). + Build() + r := &MultigresClusterReconciler{ + Client: c, Scheme: scheme, Recorder: record.NewFakeRecorder(100), + CreateTopoStore: func(multigresv1alpha1.GlobalTopoServerRef) (topoclient.Store, error) { + return nil, errors.New("UNAVAILABLE: quorum read failed") + }, + } + result, err := r.Reconcile( + t.Context(), + ctrl.Request{NamespacedName: client.ObjectKeyFromObject(cluster)}, + ) + require.NoError(t, err) + require.Equal(t, topoUnavailableRequeueDelay, result.RequeueAfter) + fresh := &multigresv1alpha1.MultigresCluster{} + require.NoError(t, c.Get(t.Context(), client.ObjectKeyFromObject(cluster), fresh)) + require.True(t, meta.IsStatusConditionFalse(fresh.Status.Conditions, conditionTopologyReady)) + require.True(t, meta.IsStatusConditionFalse(fresh.Status.Conditions, conditionFailoverReady)) + require.Equal(t, multigresv1alpha1.PhaseDegraded, fresh.Status.Phase) +} + +func TestManagedLocalTopologyQuorumIsRequired(t *testing.T) { + scheme := setupScheme() + global := &multigresv1alpha1.GlobalTopoServerSpec{ + External: &multigresv1alpha1.ExternalTopoServerSpec{ + Endpoints: []multigresv1alpha1.EndpointUrl{"http://external:2379"}, + }, + } + cluster := &multigresv1alpha1.MultigresCluster{ + ObjectMeta: metav1.ObjectMeta{Name: "example", Namespace: "default"}, + Spec: multigresv1alpha1.MultigresClusterSpec{ + GlobalTopoServer: global, + Cells: []multigresv1alpha1.CellConfig{ + { + Name: "local", + Spec: &multigresv1alpha1.CellInlineSpec{ + LocalTopoServer: &multigresv1alpha1.LocalTopoServerSpec{ + Etcd: &multigresv1alpha1.EtcdSpec{}, + }, + }, + }, + }, + }, + } + c := fake.NewClientBuilder().WithScheme(scheme).WithObjects(cluster).Build() + r := &MultigresClusterReconciler{Client: c} + condition := r.topologyQuorumCondition(t.Context(), cluster, global, nil) + require.Equal(t, metav1.ConditionFalse, condition.Status) + require.Equal(t, "TopologyMissing", condition.Reason) + local := healthyManagedLocalTopoServer( + expectedManagedLocalCell(cluster.Name, "local", cluster.Namespace), + ) + checked := metav1.Now() + local.Status.HealthCheckedAt = &checked + local.Status.Conditions = []metav1.Condition{ + { + Type: "QuorumAvailable", + ObservedGeneration: local.Generation, + Status: metav1.ConditionFalse, + Reason: "QuorumUnavailable", + Message: "No linearizable read succeeded", + }, + } + condition = r.topologyQuorumCondition( + t.Context(), + cluster, + global, + []multigresv1alpha1.TopoServer{*local}, + ) + require.Equal(t, metav1.ConditionFalse, condition.Status) + require.Equal(t, "QuorumUnavailable", condition.Reason) +} diff --git a/pkg/cluster-handler/controller/multigrescluster/multigrescluster_controller.go b/pkg/cluster-handler/controller/multigrescluster/multigrescluster_controller.go index bb4336806..7cd3a5797 100644 --- a/pkg/cluster-handler/controller/multigrescluster/multigrescluster_controller.go +++ b/pkg/cluster-handler/controller/multigrescluster/multigrescluster_controller.go @@ -100,6 +100,7 @@ func (r *MultigresClusterReconciler) Reconcile( err = r.Get(ctx, req.NamespacedName, cluster) if err != nil { if errors.IsNotFound(err) { + monitoring.DeleteFailoverHealth(req.Name, req.Namespace) if r.PoolerClientCache != nil { r.PoolerClientCache.ForgetCluster(req.NamespacedName) } @@ -137,6 +138,7 @@ func (r *MultigresClusterReconciler) Reconcile( } if !cluster.DeletionTimestamp.IsZero() { + monitoring.DeleteFailoverHealth(cluster.Name, cluster.Namespace) return r.handleDeletion(ctx, cluster) } if r.PoolerClientCache != nil { @@ -290,6 +292,15 @@ func (r *MultigresClusterReconciler) Reconcile( { ctx, childSpan := monitoring.StartChildSpan(ctx, "MultigresCluster.ReconcileTopology") result, err := r.reconcileTopology(ctx, cluster, res, pendingCells) + if err != nil || result.RequeueAfter > 0 { + if statusErr := r.updateStatus(ctx, cluster); statusErr != nil { + childSpan.End() + return ctrl.Result{}, fmt.Errorf( + "refreshing status after topology reconciliation: %w", + statusErr, + ) + } + } if err != nil { monitoring.RecordSpanError(childSpan, err) childSpan.End() @@ -343,24 +354,6 @@ func (r *MultigresClusterReconciler) Reconcile( childSpan.End() } - // Emit cluster-level metrics - monitoring.SetClusterInfo( - cluster.Name, - cluster.Namespace, - string(cluster.Status.Phase), - cluster.Status.InitializedAt != nil, - ) - var totalShards int - for _, db := range cluster.Status.Databases { - totalShards += int(db.TotalShards) - } - monitoring.SetClusterTopology( - cluster.Name, - cluster.Namespace, - len(cluster.Status.Cells), - totalShards, - ) - if pendingCells || pendingDBs { l.V(1).Info("Pending graceful deletions, requeueing", "duration", time.Since(start).String()) @@ -369,7 +362,7 @@ func (r *MultigresClusterReconciler) Reconcile( l.V(1).Info("reconcile complete", "duration", time.Since(start).String()) r.Recorder.Event(cluster, "Normal", "Synced", "Successfully reconciled MultigresCluster") - return ctrl.Result{}, nil + return ctrl.Result{RequeueAfter: clusterHealthInterval}, nil } func (r *MultigresClusterReconciler) ensureClusterFinalizer( diff --git a/pkg/cluster-handler/controller/multigrescluster/reconcile_topology.go b/pkg/cluster-handler/controller/multigrescluster/reconcile_topology.go index 2e98359a5..28e7782a3 100644 --- a/pkg/cluster-handler/controller/multigrescluster/reconcile_topology.go +++ b/pkg/cluster-handler/controller/multigrescluster/reconcile_topology.go @@ -50,8 +50,17 @@ func (r *MultigresClusterReconciler) reconcileTopology( cluster *multigresv1alpha1.MultigresCluster, res *resolver.Resolver, preservePendingDeletionCells ...bool, -) (ctrl.Result, error) { +) (result ctrl.Result, reconcileErr error) { logger := log.FromContext(ctx) + failureReason := "TopologyReconcileFailed" + defer func() { + if reconcileErr != nil { + if topo.IsTopoUnavailable(reconcileErr) { + failureReason = "TopologyUnavailable" + } + r.markTopologyFailed(ctx, cluster, failureReason, reconcileErr, logger) + } + }() preservePendingCells := len(preservePendingDeletionCells) > 0 && preservePendingDeletionCells[0] @@ -89,7 +98,7 @@ func (r *MultigresClusterReconciler) reconcileTopology( store, err := r.openTopoStore(ctx, cluster.Namespace, globalTopoRef) if err != nil { if topo.IsTopoUnavailable(err) { - return r.handleTopoUnavailable(cluster, logger) + return r.handleTopoUnavailable(ctx, cluster, err, logger) } // A missing or unusable client credential is a configuration error, not a // transient outage, so record it in the cluster's status as well as an @@ -97,7 +106,7 @@ func (r *MultigresClusterReconciler) reconcileTopology( // stays not ready is visible without reading the operator logs. r.Recorder.Eventf(cluster, "Warning", "TopoConnectFailed", "Failed to connect to topology server: %v", err) - r.markTopologyFailed(ctx, cluster, "TopoConnectFailed", err, logger) + failureReason = "TopoConnectFailed" return ctrl.Result{}, fmt.Errorf("failed to open topology store: %w", err) } defer func() { _ = store.Close() }() @@ -122,7 +131,7 @@ func (r *MultigresClusterReconciler) reconcileTopology( ), ); err != nil { if topo.IsTopoUnavailable(err) { - return r.handleTopoUnavailable(cluster, logger) + return r.handleTopoUnavailable(ctx, cluster, err, logger) } return ctrl.Result{}, fmt.Errorf("failed to register cell '%s' in topology: %w", cellCfg.Name, err) @@ -143,7 +152,7 @@ func (r *MultigresClusterReconciler) reconcileTopology( cluster.Spec.DurabilityPolicy, ); err != nil { if topo.IsTopoUnavailable(err) { - return r.handleTopoUnavailable(cluster, logger) + return r.handleTopoUnavailable(ctx, cluster, err, logger) } return ctrl.Result{}, fmt.Errorf("failed to register database '%s' in topology: %w", dbConfig.Name, err) @@ -160,7 +169,7 @@ func (r *MultigresClusterReconciler) reconcileTopology( topoCtx, store, r.Recorder, cluster, specDBNames, ); err != nil { if topo.IsTopoUnavailable(err) { - return r.handleTopoUnavailable(cluster, logger) + return r.handleTopoUnavailable(ctx, cluster, err, logger) } return ctrl.Result{}, fmt.Errorf("failed to prune databases: %w", err) } @@ -175,7 +184,7 @@ func (r *MultigresClusterReconciler) reconcileTopology( } if err := topo.PruneCells(topoCtx, store, r.Recorder, cluster, cellNames); err != nil { if topo.IsTopoUnavailable(err) { - return r.handleTopoUnavailable(cluster, logger) + return r.handleTopoUnavailable(ctx, cluster, err, logger) } return ctrl.Result{}, fmt.Errorf("failed to prune cells: %w", err) } @@ -336,11 +345,17 @@ func isPruningEnabled(cluster *multigresv1alpha1.MultigresCluster) bool { // During the grace period after cluster creation, it silently requeues. // After the grace period, it returns an error. func (r *MultigresClusterReconciler) handleTopoUnavailable( + ctx context.Context, cluster *multigresv1alpha1.MultigresCluster, - logger interface{ Info(string, ...any) }, + cause error, + logger interface { + Info(string, ...any) + Error(error, string, ...any) + }, ) (ctrl.Result, error) { resourceAge := time.Since(cluster.CreationTimestamp.Time) if resourceAge < topoUnavailableGracePeriod { + r.markTopologyFailed(ctx, cluster, "TopologyUnavailable", cause, logger) logger.Info("Topology server not available yet, requeueing", "resourceAge", resourceAge.Round(time.Second).String(), "gracePeriod", topoUnavailableGracePeriod.String(), @@ -350,5 +365,5 @@ func (r *MultigresClusterReconciler) handleTopoUnavailable( resourceAge.Round(time.Second)) return ctrl.Result{RequeueAfter: topoUnavailableRequeueDelay}, nil } - return ctrl.Result{}, fmt.Errorf("topology server unavailable") + return ctrl.Result{}, fmt.Errorf("topology server unavailable: %w", cause) } diff --git a/pkg/cluster-handler/controller/multigrescluster/status.go b/pkg/cluster-handler/controller/multigrescluster/status.go index a20c0b83a..eb3901f7c 100644 --- a/pkg/cluster-handler/controller/multigrescluster/status.go +++ b/pkg/cluster-handler/controller/multigrescluster/status.go @@ -13,6 +13,7 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" + "github.com/multigres/multigres-operator/pkg/monitoring" "github.com/multigres/multigres-operator/pkg/resolver" ) @@ -282,6 +283,16 @@ func (r *MultigresClusterReconciler) updateStatus( } } + r.updateHealthConditions(ctx, cluster, globalTopoSpec, topoServers.Items) + failover := meta.FindStatusCondition(cluster.Status.Conditions, conditionFailoverReady) + if failover.Status != metav1.ConditionTrue { + allHealthy = false + } + if failover.Status == metav1.ConditionFalse && + (cluster.Status.InitializedAt != nil || (failover.Reason != "TopologyMissing" && failover.Reason != "OrchestratorUnavailable")) { + anyDegraded = true + } + if len(expectedCells) > 0 || len(expectedTableGroups) > 0 || expectsGlobalTopoServer { allHealthy = false } @@ -290,6 +301,9 @@ func (r *MultigresClusterReconciler) updateStatus( case anyDegraded: cluster.Status.Phase = multigresv1alpha1.PhaseDegraded cluster.Status.Message = "Cluster is degraded" + if failover.Status == metav1.ConditionFalse { + cluster.Status.Message = failover.Message + } case allHealthy: cluster.Status.Phase = multigresv1alpha1.PhaseHealthy cluster.Status.Message = "Ready" @@ -470,5 +484,23 @@ func (r *MultigresClusterReconciler) updateStatus( return fmt.Errorf("failed to patch status: %w", err) } + // Emit cluster-level metrics + monitoring.SetClusterInfo( + cluster.Name, + cluster.Namespace, + string(cluster.Status.Phase), + cluster.Status.InitializedAt != nil, + ) + var totalShards int + for _, db := range cluster.Status.Databases { + totalShards += int(db.TotalShards) + } + monitoring.SetClusterTopology( + cluster.Name, + cluster.Namespace, + len(cluster.Status.Cells), + totalShards, + ) + return nil } diff --git a/pkg/cluster-handler/controller/multigrescluster/status_test.go b/pkg/cluster-handler/controller/multigrescluster/status_test.go index 997e15f39..cce82274b 100644 --- a/pkg/cluster-handler/controller/multigrescluster/status_test.go +++ b/pkg/cluster-handler/controller/multigrescluster/status_test.go @@ -470,6 +470,15 @@ func TestUpdateStatus_ExpectedChildren(t *testing.T) { Labels: map[string]string{"multigres.com/cluster": "test-cluster"}, }, Status: multigresv1alpha1.TopoServerStatus{ + HealthCheckedAt: func() *metav1.Time { now := metav1.Now(); return &now }(), + Conditions: []metav1.Condition{ + { + Type: "QuorumAvailable", + Status: metav1.ConditionTrue, + Reason: "QuorumAvailable", + Message: "Quorum confirmed", + }, + }, Phase: multigresv1alpha1.PhaseHealthy, }, } @@ -529,6 +538,15 @@ func TestUpdateStatus_ExpectedChildren(t *testing.T) { for name, tt := range tests { t.Run(name, func(t *testing.T) { + meta.SetStatusCondition( + &tt.cluster.Status.Conditions, + metav1.Condition{ + Type: conditionTopologyReady, + Status: metav1.ConditionTrue, + Reason: "TopoConnected", + Message: "Topology reachable", + }, + ) objects := append([]client.Object{tt.cluster}, tt.children...) fakeClient := fake.NewClientBuilder(). WithScheme(scheme). @@ -604,6 +622,15 @@ func TestUpdateStatus_InitializedAtSticky(t *testing.T) { }, } + meta.SetStatusCondition( + &cluster.Status.Conditions, + metav1.Condition{ + Type: conditionTopologyReady, + Status: metav1.ConditionTrue, + Reason: "TopoConnected", + Message: "Topology reachable", + }, + ) fakeClient := fake.NewClientBuilder(). WithScheme(scheme). WithObjects(cluster, cell, tableGroup). diff --git a/pkg/monitoring/alerts_test.go b/pkg/monitoring/alerts_test.go new file mode 100644 index 000000000..548ff7845 --- /dev/null +++ b/pkg/monitoring/alerts_test.go @@ -0,0 +1,44 @@ +package monitoring + +import ( + "os" + "os/exec" + "path/filepath" + "testing" + + "github.com/stretchr/testify/require" + "sigs.k8s.io/yaml" +) + +// PROMTOOL can point to a standalone binary when Prometheus is not installed. +func TestTopologyAlertRules(t *testing.T) { + binary := os.Getenv("PROMTOOL") + if binary == "" { + var err error + binary, err = exec.LookPath("promtool") + if err != nil { + t.Skip("set PROMTOOL or install promtool to test alert evaluation") + } + } + data, err := os.ReadFile("../../config/monitoring/prometheus-rules.yaml") + require.NoError(t, err) + var resource struct { + Spec struct { + Groups []any `json:"groups"` + } `json:"spec"` + } + require.NoError(t, yaml.Unmarshal(data, &resource)) + rules, err := yaml.Marshal(resource.Spec) + require.NoError(t, err) + dir := t.TempDir() + require.NoError(t, os.WriteFile(filepath.Join(dir, "rules.yaml"), rules, 0o600)) + fixture, err := os.ReadFile("testdata/topology_alerts.yaml") + require.NoError(t, err) + // #nosec G703 -- dir comes from t.TempDir and the filename is constant. + require.NoError(t, os.WriteFile(filepath.Join(dir, "tests.yaml"), fixture, 0o600)) + // #nosec G204 G702 -- PROMTOOL is an explicit test-runner binary, never application input. + cmd := exec.CommandContext(t.Context(), binary, "test", "rules", "tests.yaml") + cmd.Dir = dir + output, err := cmd.CombinedOutput() + require.NoError(t, err, "%s", output) +} diff --git a/pkg/monitoring/metrics.go b/pkg/monitoring/metrics.go index 798b21c97..01ad881fa 100644 --- a/pkg/monitoring/metrics.go +++ b/pkg/monitoring/metrics.go @@ -127,28 +127,13 @@ var ( ) func init() { - metrics.Registry.MustRegister( - clusterInfo, - clusterCellsTotal, - clusterShardsTotal, - cellGatewayReplicas, - shardPoolReplicas, - poolPodsDrifted, - toposerverReplicas, - webhookRequestTotal, - webhookRequestDuration, - lastBackupAgeSeconds, - drainOperationsTotal, - rollingUpdateInProgress, - reconcileErrorsTotal, - shardPostureInconsistent, - ) + metrics.Registry.MustRegister(Collectors()...) } // Collectors returns all registered metric collectors. This is useful for // testing that metrics are properly registered. func Collectors() []prometheus.Collector { - return []prometheus.Collector{ + collectors := []prometheus.Collector{ clusterInfo, clusterCellsTotal, clusterShardsTotal, @@ -164,6 +149,10 @@ func Collectors() []prometheus.Collector { reconcileErrorsTotal, shardPostureInconsistent, } + for _, collector := range topologyCollectors { + collectors = append(collectors, collector) + } + return collectors } // RecordReconcileError increments the per-object reconcile error counter. diff --git a/pkg/monitoring/testdata/topology_alerts.yaml b/pkg/monitoring/testdata/topology_alerts.yaml new file mode 100644 index 000000000..0b9f37b1c --- /dev/null +++ b/pkg/monitoring/testdata/topology_alerts.yaml @@ -0,0 +1,147 @@ +evaluation_interval: 30s +rule_files: +- rules.yaml +tests: +- input_series: + - series: multigres_operator_toposerver_quorum_available{cluster="example", name="example-global-topo", + target_namespace="database",reason="QuorumUnavailable"} + values: 1x1 0x4 1x3 + interval: 30s + name: quorum loss fires in one minute and clears on recovery + promql_expr_test: + - eval_time: 90s + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyQuorumUnavailable",alertstate="firing"}) + - eval_time: 120s + exp_samples: + - labels: '{}' + value: 1 + expr: count(ALERTS{alertname="MultigresTopologyQuorumUnavailable",alertstate="firing"}) + - eval_time: 210s + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyQuorumUnavailable",alertstate="firing"}) +- input_series: + - series: container_memory_working_set_bytes{namespace="database",pod="example-global-topo-0",container="etcd",instance="node"} + values: 85x8 + - series: multigres_operator_toposerver_memory_limit_bytes{cluster="example", name="example-global-topo", + target_namespace="database", member="example-global-topo-0",namespace="operator",pod="controller-manager",instance="operator"} + values: 100x8 + - series: multigres_operator_toposerver_member_oom_timestamp_seconds{cluster="example", + name="example-global-topo", target_namespace="database", member="example-global-topo-0"} + values: "0x8" + interval: 30s + name: memory pressure fires before any OOM despite exporter label collisions + promql_expr_test: + - eval_time: 30s + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyMemoryPressure",alertstate="firing"}) + - eval_time: 60s + exp_samples: + - labels: '{}' + value: 1 + expr: count(ALERTS{alertname="MultigresTopologyMemoryPressure",alertstate="firing"}) + - eval_time: 120s + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyMemberOOMKilled",alertstate="firing"}) +- input_series: + - series: container_memory_working_set_bytes{namespace="database",pod="example-global-topo-0",container="etcd"} + values: 1000x8 + - series: multigres_operator_toposerver_memory_limit_bytes{cluster="example", name="example-global-topo", + target_namespace="database", member="example-global-topo-0"} + values: "0x8" + interval: 30s + name: unlimited memory and unrelated containers do not alert + promql_expr_test: + - eval_time: 120s + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyMemoryPressure",alertstate="firing"}) +- input_series: + - series: multigres_operator_toposerver_backend_bytes{cluster="example", name="example-global-topo", + target_namespace="database", member="example-global-topo-0"} + values: 90x12 + - series: multigres_operator_toposerver_backend_quota_bytes{cluster="example", name="example-global-topo", + target_namespace="database", member="example-global-topo-0"} + values: 100x12 + interval: 30s + name: backend pressure uses configured quota + promql_expr_test: + - eval_time: 4m + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyBackendNearQuota",alertstate="firing"}) + - eval_time: 5m + exp_samples: + - labels: '{}' + value: 1 + expr: count(ALERTS{alertname="MultigresTopologyBackendNearQuota",alertstate="firing"}) +- input_series: + - series: multigres_operator_toposerver_backend_bytes{cluster="example", name="example-global-topo", + target_namespace="database", member="example-global-topo-0"} + values: 90x1 stale _x12 + - series: multigres_operator_toposerver_backend_quota_bytes{cluster="example", name="example-global-topo", + target_namespace="database", member="example-global-topo-0"} + values: 100x14 + interval: 30s + name: lost backend observations do not become zeros or stay healthy + promql_expr_test: + - eval_time: 6m + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyBackendNearQuota",alertstate="firing"}) +- input_series: + - series: multigres_operator_cluster_failover_ready{cluster="example",target_namespace="database",reason="OrchestratorUnavailable"} + values: "0x8" + interval: 30s + name: failover failure retains its cause and critical severity + promql_expr_test: + - eval_time: 1m + exp_samples: + - labels: '{}' + value: 1 + expr: count(ALERTS{alertname="MultigresFailoverUnavailable",alertstate="firing",reason="OrchestratorUnavailable",severity="critical"}) +- input_series: + - series: multigres_operator_toposerver_health_checked_timestamp_seconds{cluster="example", + name="example-global-topo", target_namespace="database"} + values: "0x12" + - series: multigres_operator_toposerver_quorum_available{cluster="example", name="example-global-topo", + target_namespace="database",reason="QuorumAvailable"} + values: 1x12 + interval: 30s + name: old successful probes become unknown + promql_expr_test: + - eval_time: 2m + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyHealthUnknown",alertstate="firing"}) + - eval_time: 4m + exp_samples: + - labels: '{}' + value: 1 + expr: count(ALERTS{alertname="MultigresTopologyHealthUnknown",alertstate="firing"}) +- input_series: + - series: multigres_operator_toposerver_member_oom_timestamp_seconds{cluster="example", + name="example-global-topo", target_namespace="database", member="example-global-topo-0"} + values: 0 30x12 + interval: 30s + name: recent OOMs alert and expire + promql_expr_test: + - eval_time: 0s + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyMemberOOMKilled",alertstate="firing"}) + - eval_time: 60s + exp_samples: + - labels: '{}' + value: 1 + expr: count(ALERTS{alertname="MultigresTopologyMemberOOMKilled",alertstate="firing"}) + - eval_time: 6m + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyMemberOOMKilled",alertstate="firing"}) +- input_series: + - series: multigres_operator_toposerver_member_restarts_total{cluster="example", + name="example-global-topo", target_namespace="database", member="example-global-topo-0"} + values: 0 1 2 3 4 4x10 + interval: 30s + name: repeated restarts alert + promql_expr_test: + - eval_time: 2m + exp_samples: + - labels: '{}' + value: 1 + expr: count(ALERTS{alertname="MultigresTopologyMemberRestarting",alertstate="firing"}) diff --git a/pkg/monitoring/topology.go b/pkg/monitoring/topology.go new file mode 100644 index 000000000..fc1539941 --- /dev/null +++ b/pkg/monitoring/topology.go @@ -0,0 +1,158 @@ +package monitoring + +import ( + "time" + + "github.com/prometheus/client_golang/prometheus" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +// target_namespace and member identify the monitored workload without colliding +// with the namespace and pod labels attached to the operator's scrape target. +var ( + topologyLabels = []string{"cluster", "name", "target_namespace"} + memberLabels = []string{"cluster", "name", "target_namespace", "member"} +) + +func topologyGauge(name, help string, labels []string) *prometheus.GaugeVec { + return prometheus.NewGaugeVec(prometheus.GaugeOpts{ + Name: "multigres_operator_" + name, Help: help, + }, labels) +} + +var ( + topoQuorum = topologyGauge( + "toposerver_quorum_available", + "Linearizable etcd reads succeed: 1 yes, 0 no, -1 unknown.", + append(append([]string{}, topologyLabels...), "reason"), + ) + topoChecked = topologyGauge( + "toposerver_health_checked_timestamp_seconds", + "Unix timestamp of the latest completed etcd probe.", + topologyLabels, + ) + topoMemberUp = topologyGauge( + "toposerver_member_up", + "Whether this member answered a linearizable read.", + memberLabels, + ) + topoBackend = topologyGauge( + "toposerver_backend_bytes", + "Etcd backend size reported by the member.", + memberLabels, + ) + topoBackendInUse = topologyGauge( + "toposerver_backend_in_use_bytes", + "Etcd backend bytes in use reported by the member.", + memberLabels, + ) + topoRevision = topologyGauge( + "toposerver_revision", + "Current etcd MVCC revision; use deriv to observe growth.", + memberLabels, + ) + topoQuota = topologyGauge( + "toposerver_backend_quota_bytes", + "Backend quota configured on the running etcd pod.", + memberLabels, + ) + topoMemoryLimit = topologyGauge( + "toposerver_memory_limit_bytes", + "Memory limit of the etcd container, or zero when unlimited.", + memberLabels, + ) + topoRestarts = topologyGauge( + "toposerver_member_restarts_total", + "Observed etcd container restart count; resets when the pod is replaced.", + memberLabels, + ) + topoOOM = topologyGauge( + "toposerver_member_oom_timestamp_seconds", + "Unix timestamp of the most recent observed OOM termination, or zero.", + memberLabels, + ) + failoverReady = topologyGauge( + "cluster_failover_ready", + "Topology access and multiorch readiness: 1 ready, 0 unavailable, -1 unknown.", + []string{"cluster", "target_namespace", "reason"}, + ) + failoverChecked = topologyGauge( + "cluster_failover_checked_timestamp_seconds", + "Unix timestamp of the latest failover readiness observation.", + []string{"cluster", "target_namespace"}, + ) +) + +var topologyCollectors = []*prometheus.GaugeVec{ + topoQuorum, topoChecked, topoMemberUp, topoBackend, topoBackendInUse, + topoRevision, topoQuota, topoMemoryLimit, topoRestarts, topoOOM, failoverReady, failoverChecked, +} + +// TopologyMemberMetrics contains the latest observations of a managed etcd pod. +// Nil values mean an observation failed, rather than a zero-sized backend. +type TopologyMemberMetrics struct { + Name string + Up bool + BackendBytes, BackendInUseBytes, Revision *int64 + QuotaBytes, MemoryLimitBytes, Restarts, OOMTimestamp *int64 +} + +func conditionValue(status metav1.ConditionStatus) float64 { + switch status { + case metav1.ConditionTrue: + return 1 + case metav1.ConditionFalse: + return 0 + default: + return -1 + } +} + +// DeleteTopologyHealth removes observations after deletion and before replacement. +func DeleteTopologyHealth(name, namespace string) { + for _, metric := range []*prometheus.GaugeVec{topoQuorum, topoChecked, topoMemberUp, topoBackend, topoBackendInUse, topoRevision, topoQuota, topoMemoryLimit, topoRestarts, topoOOM} { + metric.DeletePartialMatch(prometheus.Labels{"name": name, "target_namespace": namespace}) + } +} + +func SetTopologyHealth( + cluster, name, namespace string, + condition metav1.Condition, + checked time.Time, + members []TopologyMemberMetrics, +) { + DeleteTopologyHealth(name, namespace) + topoQuorum.WithLabelValues(cluster, name, namespace, condition.Reason). + Set(conditionValue(condition.Status)) + topoChecked.WithLabelValues(cluster, name, namespace).Set(float64(checked.Unix())) + for _, member := range members { + labels := []string{cluster, name, namespace, member.Name} + up := 0.0 + if member.Up { + up = 1 + } + topoMemberUp.WithLabelValues(labels...).Set(up) + for metric, value := range map[*prometheus.GaugeVec]*int64{ + topoBackend: member.BackendBytes, topoBackendInUse: member.BackendInUseBytes, + topoRevision: member.Revision, topoQuota: member.QuotaBytes, + topoMemoryLimit: member.MemoryLimitBytes, topoRestarts: member.Restarts, topoOOM: member.OOMTimestamp, + } { + if value != nil { + metric.WithLabelValues(labels...).Set(float64(*value)) + } + } + } +} + +func DeleteFailoverHealth(cluster, namespace string) { + labels := prometheus.Labels{"cluster": cluster, "target_namespace": namespace} + failoverReady.DeletePartialMatch(labels) + failoverChecked.DeletePartialMatch(labels) +} + +func SetFailoverHealth(cluster, namespace string, condition metav1.Condition) { + DeleteFailoverHealth(cluster, namespace) + failoverReady.WithLabelValues(cluster, namespace, condition.Reason). + Set(conditionValue(condition.Status)) + failoverChecked.WithLabelValues(cluster, namespace).Set(float64(time.Now().Unix())) +} diff --git a/pkg/monitoring/topology_test.go b/pkg/monitoring/topology_test.go new file mode 100644 index 000000000..56abdd7f2 --- /dev/null +++ b/pkg/monitoring/topology_test.go @@ -0,0 +1,56 @@ +package monitoring + +import ( + "testing" + "time" + + "github.com/prometheus/client_golang/prometheus" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/utils/ptr" +) + +func TestTopologyMetricsReplaceFailedObservationsAndDelete(t *testing.T) { + name, ns := "metrics-topo", "metrics-test" + t.Cleanup(func() { DeleteTopologyHealth(name, ns) }) + condition := metav1.Condition{Status: metav1.ConditionTrue, Reason: "QuorumAvailable"} + members := []TopologyMemberMetrics{ + { + Name: "metrics-topo-0", + Up: true, + BackendBytes: ptr.To(int64(100)), + Revision: ptr.To(int64(42)), + }, + } + SetTopologyHealth("example", name, ns, condition, time.Now(), members) + require.Equal( + t, + float64(100), + gaugeValue(t, topoBackend, "example", name, ns, "metrics-topo-0"), + ) + condition.Status, condition.Reason = metav1.ConditionFalse, "TopologyUnreachable" + members[0].Up, members[0].BackendBytes, members[0].Revision = false, nil, nil + SetTopologyHealth("example", name, ns, condition, time.Now(), members) + registry := prometheus.NewRegistry() + registry.MustRegister(topoBackend, topoRevision) + families, err := registry.Gather() + require.NoError(t, err) + for _, family := range families { + for _, metric := range family.Metric { + for _, label := range metric.Label { + require.False( + t, + label.GetName() == "name" && label.GetValue() == name, + "failed probes must not retain backend/revision observations", + ) + } + } + } + require.Equal( + t, + float64(0), + gaugeValue(t, topoMemberUp, "example", name, ns, "metrics-topo-0"), + ) + DeleteTopologyHealth(name, ns) + require.False(t, topoQuorum.DeleteLabelValues("example", name, ns, "TopologyUnreachable")) +} diff --git a/pkg/resource-handler/controller/toposerver/health.go b/pkg/resource-handler/controller/toposerver/health.go new file mode 100644 index 000000000..91a2163ee --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/health.go @@ -0,0 +1,204 @@ +package toposerver + +import ( + "context" + "fmt" + "strconv" + "sync" + "time" + + appsv1 "k8s.io/api/apps/v1" + corev1 "k8s.io/api/core/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/meta" + 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/monitoring" + "github.com/multigres/multigres-operator/pkg/util/metadata" +) + +const healthProbeTimeout = 5 * time.Second + +// reconcileHealth runs independently of resource reconciliation: a failed apply +// or an interrupted maintenance operation must not stop health observations. +func (r *TopoServerReconciler) reconcileHealth( + ctx context.Context, + req ctrl.Request, +) (ctrl.Result, error) { + ts := &multigresv1alpha1.TopoServer{} + if err := r.maintenanceReader().Get(ctx, req.NamespacedName, ts); err != nil { + if apierrors.IsNotFound(err) { + monitoring.DeleteTopologyHealth(req.Name, req.Namespace) + } + return ctrl.Result{}, client.IgnoreNotFound(err) + } + if !ts.DeletionTimestamp.IsZero() { + monitoring.DeleteTopologyHealth(req.Name, req.Namespace) + return ctrl.Result{}, nil + } + condition, members := r.probeHealth(ctx, ts) + r.observeMemberPods(ctx, ts, members) + checked := metav1.Now() + monitoring.SetTopologyHealth( + ts.Labels[metadata.LabelMultigresCluster], + ts.Name, + ts.Namespace, + condition, + checked.Time, + members, + ) + meta.SetStatusCondition(&ts.Status.Conditions, condition) + condition = *meta.FindStatusCondition(ts.Status.Conditions, "QuorumAvailable") + patch := &multigresv1alpha1.TopoServer{ + TypeMeta: metav1.TypeMeta{ + APIVersion: multigresv1alpha1.GroupVersion.String(), + Kind: "TopoServer", + }, + ObjectMeta: metav1.ObjectMeta{Name: ts.Name, Namespace: ts.Namespace}, + Status: multigresv1alpha1.TopoServerStatus{ + Conditions: []metav1.Condition{condition}, + HealthCheckedAt: &checked, + }, + } + err := r.Status(). + Patch(ctx, patch, client.Apply, client.FieldOwner("multigres-topology-health"), client.ForceOwnership) + return ctrl.Result{RequeueAfter: statusRecheckDelay}, err +} + +func (r *TopoServerReconciler) probeHealth( + ctx context.Context, + ts *multigresv1alpha1.TopoServer, +) (metav1.Condition, []monitoring.TopologyMemberMetrics) { + condition := metav1.Condition{ + Type: "QuorumAvailable", + ObservedGeneration: ts.Generation, + Status: metav1.ConditionUnknown, + Reason: "ProbeFailed", + Message: "Could not check etcd quorum", + } + endpoints := maintenanceEndpoints(ts) + members := make([]monitoring.TopologyMemberMetrics, len(endpoints)) + for i := range members { + members[i].Name = fmt.Sprintf("%s-%d", ts.Name, i) + } + probeCtx, cancel := context.WithTimeout(ctx, healthProbeTimeout) + defer cancel() + c, err := r.maintenanceClient(probeCtx, ts) + if err != nil { + condition.Message = fmt.Sprintf("Could not create etcd health client: %v", err) + return condition, members + } + defer c.Close() + clusterIDs := make([]uint64, len(endpoints)) + var wg sync.WaitGroup + for i, endpoint := range endpoints { + wg.Go(func() { + s, err := c.Status(probeCtx, endpoint) + if err != nil || s == nil || s.Header == nil || s.Header.ClusterId == 0 || + s.Header.MemberId == 0 { + return + } + clusterIDs[i] = s.Header.ClusterId + members[i].BackendBytes = ptr.To(s.DbSize) + members[i].BackendInUseBytes = ptr.To(s.DbSizeInUse) + members[i].Revision = ptr.To(s.Header.Revision) + members[i].Up = c.Health(probeCtx, endpoint) == nil + }) + } + wg.Wait() + var clusterID uint64 + for _, observedID := range clusterIDs { + if observedID == 0 { + continue + } + if clusterID != 0 && observedID != clusterID { + condition.Status, condition.Reason = metav1.ConditionFalse, "ClusterMismatch" + condition.Message = "Etcd endpoints report different cluster IDs; topology membership is inconsistent and failover protection is unavailable" + return condition, members + } + clusterID = observedID + } + responding, readable := 0, 0 + for _, member := range members { + if member.BackendBytes != nil { + responding++ + } + if member.Up { + readable++ + } + } + switch { + case readable > 0: + // A successful linearizable read requires a quorum round trip. Counting + // reachable endpoints alone would misclassify asymmetric network failures. + condition.Status, condition.Reason = metav1.ConditionTrue, "QuorumAvailable" + condition.Message = fmt.Sprintf( + "Linearizable reads succeeded through %d/%d etcd members", + readable, + len(endpoints), + ) + case responding > 0: + condition.Status, condition.Reason = metav1.ConditionFalse, "QuorumUnavailable" + condition.Message = "Etcd members respond to status requests, but no member completed a linearizable read; failover protection is unavailable" + default: + condition.Status, condition.Reason = metav1.ConditionFalse, "TopologyUnreachable" + condition.Message = "No etcd member responded to the operator; quorum cannot be verified and failover protection is unavailable" + } + return condition, members +} + +func (r *TopoServerReconciler) observeMemberPods( + ctx context.Context, + ts *multigresv1alpha1.TopoServer, + members []monitoring.TopologyMemberMetrics, +) { + sts := &appsv1.StatefulSet{} + if err := r.maintenanceReader().Get(ctx, client.ObjectKeyFromObject(ts), sts); err != nil { + return + } + for i := range members { + pod := &corev1.Pod{} + if err := r.maintenanceReader(). + Get(ctx, client.ObjectKey{Namespace: ts.Namespace, Name: members[i].Name}, pod); err != nil { + continue + } + if !metav1.IsControlledBy(pod, sts) { + continue + } + for _, container := range pod.Spec.Containers { + if container.Name != "etcd" { + continue + } + members[i].MemoryLimitBytes = ptr.To(container.Resources.Limits.Memory().Value()) + // Older managed pods use etcd's default quota until their template rolls. + members[i].QuotaBytes = ptr.To(int64(2 * 1024 * 1024 * 1024)) + for _, env := range container.Env { + if env.Name == "ETCD_QUOTA_BACKEND_BYTES" { + quota, err := strconv.ParseInt(env.Value, 10, 64) + if err == nil { + members[i].QuotaBytes = "a + } else { + members[i].QuotaBytes = nil + } + } + } + } + for _, state := range pod.Status.ContainerStatuses { + if state.Name != "etcd" { + continue + } + members[i].Restarts = ptr.To(int64(state.RestartCount)) + members[i].OOMTimestamp = ptr.To(int64(0)) + for _, terminated := range []*corev1.ContainerStateTerminated{state.LastTerminationState.Terminated, state.State.Terminated} { + if terminated != nil && terminated.Reason == "OOMKilled" && + !terminated.FinishedAt.IsZero() { + members[i].OOMTimestamp = ptr.To(terminated.FinishedAt.Unix()) + } + } + } + } +} diff --git a/pkg/resource-handler/controller/toposerver/health_etcd_test.go b/pkg/resource-handler/controller/toposerver/health_etcd_test.go new file mode 100644 index 000000000..bc13dd2d4 --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/health_etcd_test.go @@ -0,0 +1,63 @@ +//go:build integration + +package toposerver + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/require" + clientv3 "go.etcd.io/etcd/client/v3" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" +) + +type mappedHealthClient struct { + etcdMaintenanceClient + endpoints map[string]string +} + +func (c mappedHealthClient) Status( + ctx context.Context, + endpoint string, +) (*clientv3.StatusResponse, error) { + return c.etcdMaintenanceClient.Status(ctx, c.endpoints[endpoint]) +} + +func (c mappedHealthClient) Health(ctx context.Context, endpoint string) error { + return c.etcdMaintenanceClient.Health(ctx, c.endpoints[endpoint]) +} +func (c mappedHealthClient) Close() {} + +func TestLiveEtcdQuorumHealth(t *testing.T) { + c, endpoints, stop := startTestEtcd(t) + ts := certTestTopoServer(nil) + mapping := map[string]string{} + for i, endpoint := range maintenanceEndpoints(ts) { + mapping[endpoint] = endpoints[i] + } + r := &TopoServerReconciler{ + newMaintenanceClient: func(context.Context, *multigresv1alpha1.TopoServer) (etcdMaintenanceClient, error) { + return mappedHealthClient{c, mapping}, nil + }, + } + condition, _ := r.probeHealth(t.Context(), ts) + require.Equal(t, metav1.ConditionTrue, condition.Status) + stop[2]() + require.Eventually(t, func() bool { + condition, _ := r.probeHealth(t.Context(), ts) + return condition.Status == metav1.ConditionTrue + }, 20*time.Second, 100*time.Millisecond, "two surviving members retain quorum") + stop[1]() + require.Eventually(t, func() bool { + condition, _ := r.probeHealth(t.Context(), ts) + return condition.Status == metav1.ConditionFalse && condition.Reason == "QuorumUnavailable" + }, 20*time.Second, 100*time.Millisecond, "one live member can answer status but cannot serve linearizable reads") + stop[0]() + require.Eventually(t, func() bool { + condition, _ := r.probeHealth(t.Context(), ts) + return condition.Reason == "TopologyUnreachable" + }, 20*time.Second, 100*time.Millisecond) +} diff --git a/pkg/resource-handler/controller/toposerver/health_test.go b/pkg/resource-handler/controller/toposerver/health_test.go new file mode 100644 index 000000000..512e6a74e --- /dev/null +++ b/pkg/resource-handler/controller/toposerver/health_test.go @@ -0,0 +1,162 @@ +package toposerver + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/stretchr/testify/require" + clientv3 "go.etcd.io/etcd/client/v3" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/meta" + "k8s.io/apimachinery/pkg/api/resource" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + ctrl "sigs.k8s.io/controller-runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" +) + +type healthTestClient struct { + *fakeEtcdMaintenance + unavailable map[string]bool + blockedReads map[string]bool +} + +func (c *healthTestClient) Status( + ctx context.Context, + endpoint string, +) (*clientv3.StatusResponse, error) { + if c.unavailable[endpoint] { + return nil, errors.New("connection refused") + } + return c.fakeEtcdMaintenance.Status(ctx, endpoint) +} + +func (c *healthTestClient) Health(_ context.Context, endpoint string) error { + if c.blockedReads[endpoint] { + return errors.New("linearizable read timed out") + } + return nil +} + +func TestProbeTopologyQuorum(t *testing.T) { + for _, tc := range []struct { + name string + unavailable, blocked []int + clientError, clusterMismatch bool + want metav1.ConditionStatus + reason string + }{ + {name: "different etcd clusters", clusterMismatch: true, want: metav1.ConditionFalse, reason: "ClusterMismatch"}, + {name: "healthy", want: metav1.ConditionTrue, reason: "QuorumAvailable"}, + {name: "one unreachable member", unavailable: []int{0}, want: metav1.ConditionTrue, reason: "QuorumAvailable"}, + {name: "only one endpoint reachable but quorum read succeeds", unavailable: []int{0, 1}, want: metav1.ConditionTrue, reason: "QuorumAvailable"}, + {name: "status works without quorum", blocked: []int{0, 1, 2}, want: metav1.ConditionFalse, reason: "QuorumUnavailable"}, + {name: "all members unreachable", unavailable: []int{0, 1, 2}, want: metav1.ConditionFalse, reason: "TopologyUnreachable"}, + {name: "credential error", clientError: true, want: metav1.ConditionUnknown, reason: "ProbeFailed"}, + } { + t.Run(tc.name, func(t *testing.T) { + r, ts, f := maintenanceFixture(t) + c := &healthTestClient{ + fakeEtcdMaintenance: f, + unavailable: map[string]bool{}, + blockedReads: map[string]bool{}, + } + endpoints := maintenanceEndpoints(ts) + if tc.clusterMismatch { + f.statuses[endpoints[1]].Header.ClusterId = 456 + } + for _, i := range tc.unavailable { + c.unavailable[endpoints[i]] = true + } + for _, i := range tc.blocked { + c.blockedReads[endpoints[i]] = true + } + r.newMaintenanceClient = func(ctx context.Context, _ *multigresv1alpha1.TopoServer) (etcdMaintenanceClient, error) { + deadline, ok := ctx.Deadline() + require.True(t, ok) + require.LessOrEqual(t, time.Until(deadline), healthProbeTimeout) + if tc.clientError { + return nil, errors.New("invalid TLS certificate") + } + return c, nil + } + condition, members := r.probeHealth(t.Context(), ts) + require.Equal(t, tc.want, condition.Status) + require.Equal(t, tc.reason, condition.Reason) + for _, i := range tc.unavailable { + require.False(t, members[i].Up) + require.Nil(t, members[i].BackendBytes) + } + require.Empty(t, f.defragged) + require.Empty(t, f.moved) + }) + } +} + +func TestHealthObservationDuringMaintenance(t *testing.T) { + r, ts, f := maintenanceFixture(t) + state := &multigresv1alpha1.EtcdMaintenanceStatus{ + LastAttemptTime: metav1.Now(), + Endpoint: maintenanceEndpoints(ts)[0], + InProgress: true, + } + require.NoError(t, r.saveMaintenance(t.Context(), ts, state)) + f.healthErr = errors.New("quorum unavailable") + result, err := r.reconcileHealth( + t.Context(), + ctrl.Request{NamespacedName: client.ObjectKeyFromObject(ts)}, + ) + require.NoError(t, err) + require.Equal(t, statusRecheckDelay, result.RequeueAfter) + fresh := &multigresv1alpha1.TopoServer{} + require.NoError(t, r.Get(t.Context(), client.ObjectKeyFromObject(ts), fresh)) + require.Equal( + t, + metav1.ConditionFalse, + meta.FindStatusCondition(fresh.Status.Conditions, "QuorumAvailable").Status, + ) + require.True(t, fresh.Status.EtcdMaintenance.InProgress) + require.NotNil(t, fresh.Status.HealthCheckedAt) +} + +func TestObserveRunningEtcdPodLimitsAndOOM(t *testing.T) { + r, ts, _ := maintenanceFixture(t) + _, members := r.probeHealth(t.Context(), ts) + pod := &corev1.Pod{} + require.NoError( + t, + r.Get(t.Context(), client.ObjectKey{Namespace: ts.Namespace, Name: members[0].Name}, pod), + ) + pod.Spec.Containers = []corev1.Container{ + { + Name: "etcd", + Resources: corev1.ResourceRequirements{ + Limits: corev1.ResourceList{corev1.ResourceMemory: resource.MustParse("512Mi")}, + }, + Env: []corev1.EnvVar{{Name: "ETCD_QUOTA_BACKEND_BYTES", Value: "1073741824"}}, + }, + } + require.NoError(t, r.Update(t.Context(), pod)) + finished := metav1.NewTime(time.Now().Add(-time.Minute)) + pod.Status.ContainerStatuses = []corev1.ContainerStatus{ + { + Name: "etcd", + RestartCount: 3, + LastTerminationState: corev1.ContainerState{ + Terminated: &corev1.ContainerStateTerminated{ + Reason: "OOMKilled", + FinishedAt: finished, + }, + }, + }, + } + require.NoError(t, r.Status().Update(t.Context(), pod)) + r.observeMemberPods(t.Context(), ts, members) + require.EqualValues(t, 512<<20, *members[0].MemoryLimitBytes) + require.EqualValues(t, 1<<30, *members[0].QuotaBytes) + require.EqualValues(t, 3, *members[0].Restarts) + require.Equal(t, finished.Unix(), *members[0].OOMTimestamp) +} diff --git a/pkg/resource-handler/controller/toposerver/maintenance_etcd_test.go b/pkg/resource-handler/controller/toposerver/maintenance_etcd_test.go index 3316d8e77..23bf93007 100644 --- a/pkg/resource-handler/controller/toposerver/maintenance_etcd_test.go +++ b/pkg/resource-handler/controller/toposerver/maintenance_etcd_test.go @@ -26,6 +26,11 @@ import ( // 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) { + c, endpoints, _ := startTestEtcd(t) + return c, endpoints +} + +func startTestEtcd(t *testing.T) (*memberClients, []string, []func()) { t.Helper() binary := filepath.Join(os.Getenv("KUBEBUILDER_ASSETS"), "etcd") if _, err := os.Stat(binary); err != nil { @@ -53,6 +58,7 @@ func startMaintenanceEtcd(t *testing.T) (*memberClients, []string) { peers[i] = allocateURL() cluster[i] = fmt.Sprintf("member-%d=%s", i, peers[i]) } + stop := make([]func(), len(endpoints)) for i := range endpoints { dir := t.TempDir() logFile, err := os.Create(filepath.Join(dir, "etcd.log")) @@ -91,6 +97,7 @@ func startMaintenanceEtcd(t *testing.T) (*memberClients, []string) { _ = logFile.Close() t.Fatal(err) } + stop[i] = func() { _ = cmd.Process.Kill() } t.Cleanup(func() { _ = cmd.Process.Kill() _ = cmd.Wait() @@ -122,7 +129,7 @@ func startMaintenanceEtcd(t *testing.T) (*memberClients, []string) { } time.Sleep(100 * time.Millisecond) } - return c, endpoints + return c, endpoints, stop } func TestLiveEtcdCompactionAndMaintenance(t *testing.T) { diff --git a/pkg/resource-handler/controller/toposerver/maintenance_status_integration_test.go b/pkg/resource-handler/controller/toposerver/maintenance_status_integration_test.go index d959add02..f4acce41e 100644 --- a/pkg/resource-handler/controller/toposerver/maintenance_status_integration_test.go +++ b/pkg/resource-handler/controller/toposerver/maintenance_status_integration_test.go @@ -3,9 +3,14 @@ package toposerver import ( + "context" + "errors" "path/filepath" "testing" + "k8s.io/apimachinery/pkg/api/meta" + ctrl "sigs.k8s.io/controller-runtime" + multigresv1alpha1 "github.com/multigres/multigres-operator/api/v1alpha1" "github.com/multigres/multigres-operator/pkg/testutil" appsv1 "k8s.io/api/apps/v1" @@ -57,6 +62,16 @@ func TestMaintenanceReservationSurvivesStatusApply(t *testing.T) { if err := r.saveMaintenance(t.Context(), stale, state); !apierrors.IsConflict(err) { t.Fatalf("stale reservation should conflict, got %v", err) } + + r.newMaintenanceClient = func(context.Context, *multigresv1alpha1.TopoServer) (etcdMaintenanceClient, error) { + return nil, errors.New("credential unavailable") + } + if _, err := r.reconcileHealth( + t.Context(), + ctrl.Request{NamespacedName: client.ObjectKeyFromObject(ts)}, + ); err != nil { + t.Fatal(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 { @@ -66,6 +81,10 @@ func TestMaintenanceReservationSurvivesStatusApply(t *testing.T) { if err := c.Get(t.Context(), client.ObjectKeyFromObject(ts), fresh); err != nil { t.Fatal(err) } + if fresh.Status.HealthCheckedAt == nil || + meta.FindStatusCondition(fresh.Status.Conditions, "QuorumAvailable") == nil { + t.Fatal("resource status apply removed health observation") + } 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/toposerver_controller.go b/pkg/resource-handler/controller/toposerver/toposerver_controller.go index 7a8a5ef83..8d74bbfe4 100644 --- a/pkg/resource-handler/controller/toposerver/toposerver_controller.go +++ b/pkg/resource-handler/controller/toposerver/toposerver_controller.go @@ -551,6 +551,14 @@ func (r *TopoServerReconciler) SetupWithManagerReconciler( controllerOpts = opts[0] } + if err := ctrl.NewControllerManagedBy(mgr). + Named("toposerver-health"). + For(&multigresv1alpha1.TopoServer{}, builder.WithPredicates(predicate.GenerationChangedPredicate{})). + WithOptions(controllerOpts). + Complete(reconcile.Func(r.reconcileHealth)); err != nil { + return err + } + return ctrl.NewControllerManagedBy(mgr). For(&multigresv1alpha1.TopoServer{}, builder.WithPredicates( predicate.GenerationChangedPredicate{}, From 4388a45823e1e01814b44239ef241edaa4afe587 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ver=C3=B3nica=20L=C3=B3pez?= Date: Fri, 25 Sep 2026 14:35:02 -0600 Subject: [PATCH 2/2] fix(topology): block failover readiness on etcd alarms MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Verónica López --- config/monitoring/prometheus-rules.yaml | 4 +- docs/monitoring/runbooks/TopologyHealth.md | 11 +++- .../multigrescluster/health_test.go | 9 +++ pkg/monitoring/testdata/topology_alerts.yaml | 26 +++++++++ pkg/monitoring/topology.go | 2 +- .../controller/toposerver/health.go | 22 ++++++++ .../controller/toposerver/health_etcd_test.go | 55 +++++++++++++++++++ .../controller/toposerver/health_test.go | 49 +++++++++++++++++ 8 files changed, 172 insertions(+), 6 deletions(-) diff --git a/config/monitoring/prometheus-rules.yaml b/config/monitoring/prometheus-rules.yaml index 77b25cca3..eda83df00 100644 --- a/config/monitoring/prometheus-rules.yaml +++ b/config/monitoring/prometheus-rules.yaml @@ -89,9 +89,9 @@ spec: labels: severity: critical annotations: - summary: "Topology quorum unavailable for {{ $labels.target_namespace }}/{{ $labels.name }}" + summary: "Topology unavailable for failover in {{ $labels.target_namespace }}/{{ $labels.name }}" description: >- - {{ $labels.reason }}: the operator cannot confirm quorum for one consistent etcd cluster. + {{ $labels.reason }}: the operator detected an etcd status error or cannot confirm quorum for one consistent etcd cluster. Cluster {{ $labels.cluster }} has lost verified failover protection; existing SQL connections may still serve traffic. runbook_url: "https://github.com/multigres/multigres-operator/blob/main/docs/monitoring/runbooks/TopologyHealth.md" diff --git a/docs/monitoring/runbooks/TopologyHealth.md b/docs/monitoring/runbooks/TopologyHealth.md index a02099421..0e2e1db88 100644 --- a/docs/monitoring/runbooks/TopologyHealth.md +++ b/docs/monitoring/runbooks/TopologyHealth.md @@ -11,7 +11,10 @@ The operator checks each managed etcd member every 30 seconds with a status RPC and a read-only, linearizable Get. The probes run independently of resource reconciliation and defragmentation, with a five-second deadline. A successful linearizable read confirms quorum even when the operator cannot reach every -member. Status RPCs alone do not confirm quorum. +member. Status RPCs alone do not confirm quorum. A status error from any member +sets `QuorumAvailable=False` with reason `EtcdStatusError`, even when reads succeed. +For example, a NOSPACE alarm allows reads but blocks the topology writes required +during failover. The condition message identifies each member reporting an error. `TopoServer.status.conditions[QuorumAvailable]` reports the result. `TopologyUnreachable` means no member answered the operator; it does not prove @@ -42,7 +45,7 @@ orchestrator readiness. Monitor external etcd through its owner. | Alert | Threshold | Severity | | --- | --- | --- | -| `MultigresTopologyQuorumUnavailable` | No successful quorum read for one minute | critical | +| `MultigresTopologyQuorumUnavailable` | Failed quorum check, inconsistent cluster IDs, or etcd status errors for one minute | critical | | `MultigresTopologyMemberUnavailable` | A member cannot complete reads for two minutes | warning | | `MultigresTopologyHealthUnknown` | Unknown probe or observation over two minutes old, for one minute | warning | | `MultigresTopologyBackendNearQuota` | Backend over 80% of quota for five minutes | warning | @@ -76,6 +79,8 @@ inspect the affected shard's multiorch Deployments and pod readiness. For backend pressure, compare total backend bytes with bytes in use. Review compaction and defragmentation settings in [Topology maintenance](../../topology-maintenance.md). +If the condition reports NOSPACE, reclaim backend space before disarming the alarm +with `etcdctl alarm disarm`. Successful reads alone do not confirm recovery. For memory pressure or OOMs, inspect working-set history and the running pod's memory limit. Raising the backend quota does not increase available memory. @@ -91,7 +96,7 @@ All names below begin with `multigres_operator_`: | --- | --- | | `toposerver_quorum_available` | 1 true, 0 false, -1 unknown; includes `reason` | | `toposerver_health_checked_timestamp_seconds` | Last completed probe | -| `toposerver_member_up` | Member completed a linearizable read | +| `toposerver_member_up` | Member completed a linearizable read; may remain 1 during a NOSPACE alarm | | `toposerver_backend_bytes` | Total backend size | | `toposerver_backend_in_use_bytes` | Backend bytes in use | | `toposerver_backend_quota_bytes` | Quota configured on the running pod | diff --git a/pkg/cluster-handler/controller/multigrescluster/health_test.go b/pkg/cluster-handler/controller/multigrescluster/health_test.go index fbe988b08..e4905966a 100644 --- a/pkg/cluster-handler/controller/multigrescluster/health_test.go +++ b/pkg/cluster-handler/controller/multigrescluster/health_test.go @@ -24,6 +24,7 @@ func TestFailoverHealthConditions(t *testing.T) { for _, tc := range []struct { name string quorum metav1.ConditionStatus + quorumReason string stale, missing, orchDown, shardStale, shardMissing, accessFailed, external, observationFailed bool want metav1.ConditionStatus reason string @@ -31,6 +32,7 @@ func TestFailoverHealthConditions(t *testing.T) { {name: "failed shard observation clears readiness", quorum: metav1.ConditionTrue, observationFailed: true, want: metav1.ConditionUnknown, reason: "ObservationFailed"}, {name: "healthy", quorum: metav1.ConditionTrue, want: metav1.ConditionTrue, reason: "FailoverReady"}, {name: "quorum lost with ready pods", quorum: metav1.ConditionFalse, want: metav1.ConditionFalse, reason: "QuorumUnavailable"}, + {name: "etcd alarm with ready pods and successful registration", quorum: metav1.ConditionFalse, quorumReason: "EtcdStatusError", want: metav1.ConditionFalse, reason: "EtcdStatusError"}, {name: "stale topology observation", quorum: metav1.ConditionTrue, stale: true, want: metav1.ConditionUnknown, reason: "TopologyHealthStale"}, {name: "missing topology", missing: true, want: metav1.ConditionFalse, reason: "TopologyMissing"}, {name: "known orchestrator failure takes precedence over stale quorum", quorum: metav1.ConditionTrue, stale: true, orchDown: true, want: metav1.ConditionFalse, reason: "OrchestratorUnavailable"}, @@ -105,6 +107,10 @@ func TestFailoverHealthConditions(t *testing.T) { }, }, } + if tc.quorumReason != "" { + servers[0].Status.Conditions[0].Reason = tc.quorumReason + servers[0].Status.Conditions[0].Message = "Etcd member reports NOSPACE" + } if tc.missing || tc.external { servers = nil } @@ -150,6 +156,9 @@ func TestFailoverHealthConditions(t *testing.T) { condition := meta.FindStatusCondition(cluster.Status.Conditions, conditionFailoverReady) require.Equal(t, tc.want, condition.Status) require.Equal(t, tc.reason, condition.Reason) + if tc.quorumReason != "" { + require.Contains(t, condition.Message, "NOSPACE") + } require.True( t, meta.IsStatusConditionTrue(cluster.Status.Conditions, "Available"), diff --git a/pkg/monitoring/testdata/topology_alerts.yaml b/pkg/monitoring/testdata/topology_alerts.yaml index 0b9f37b1c..147087ddb 100644 --- a/pkg/monitoring/testdata/topology_alerts.yaml +++ b/pkg/monitoring/testdata/topology_alerts.yaml @@ -2,6 +2,32 @@ evaluation_interval: 30s rule_files: - rules.yaml tests: +- input_series: + - series: multigres_operator_toposerver_quorum_available{cluster="example", name="example-global-topo", + target_namespace="database",reason="EtcdStatusError"} + values: 0x4 stale + - series: multigres_operator_cluster_failover_ready{cluster="example", target_namespace="database",reason="EtcdStatusError"} + values: 0x4 stale + - series: multigres_operator_toposerver_member_up{cluster="example", name="example-global-topo", + target_namespace="database",member="example-global-topo-0"} + values: 1x6 + interval: 30s + name: etcd alarms trigger topology and failover alerts despite successful reads + promql_expr_test: + - eval_time: 30s + exp_samples: [] + expr: count(ALERTS{reason="EtcdStatusError",alertstate="firing"}) + - eval_time: 60s + exp_samples: + - labels: '{}' + value: 2 + expr: count(ALERTS{alertname=~"MultigresTopologyQuorumUnavailable|MultigresFailoverUnavailable",reason="EtcdStatusError",alertstate="firing"}) + - eval_time: 120s + exp_samples: [] + expr: count(ALERTS{alertname="MultigresTopologyMemberUnavailable",alertstate="firing"}) + - eval_time: 150s + exp_samples: [] + expr: count(ALERTS{reason="EtcdStatusError",alertstate="firing"}) - input_series: - series: multigres_operator_toposerver_quorum_available{cluster="example", name="example-global-topo", target_namespace="database",reason="QuorumUnavailable"} diff --git a/pkg/monitoring/topology.go b/pkg/monitoring/topology.go index fc1539941..46481a7cc 100644 --- a/pkg/monitoring/topology.go +++ b/pkg/monitoring/topology.go @@ -23,7 +23,7 @@ func topologyGauge(name, help string, labels []string) *prometheus.GaugeVec { var ( topoQuorum = topologyGauge( "toposerver_quorum_available", - "Linearizable etcd reads succeed: 1 yes, 0 no, -1 unknown.", + "Etcd quorum is verified with no reported status errors: 1 yes, 0 no, -1 unknown.", append(append([]string{}, topologyLabels...), "reason"), ) topoChecked = topologyGauge( diff --git a/pkg/resource-handler/controller/toposerver/health.go b/pkg/resource-handler/controller/toposerver/health.go index 91a2163ee..35f2f8557 100644 --- a/pkg/resource-handler/controller/toposerver/health.go +++ b/pkg/resource-handler/controller/toposerver/health.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "strconv" + "strings" "sync" "time" @@ -94,6 +95,7 @@ func (r *TopoServerReconciler) probeHealth( } defer c.Close() clusterIDs := make([]uint64, len(endpoints)) + statusErrors := make([][]string, len(endpoints)) var wg sync.WaitGroup for i, endpoint := range endpoints { wg.Go(func() { @@ -103,6 +105,7 @@ func (r *TopoServerReconciler) probeHealth( return } clusterIDs[i] = s.Header.ClusterId + statusErrors[i] = s.Errors members[i].BackendBytes = ptr.To(s.DbSize) members[i].BackendInUseBytes = ptr.To(s.DbSizeInUse) members[i].Revision = ptr.To(s.Header.Revision) @@ -122,6 +125,25 @@ func (r *TopoServerReconciler) probeHealth( } clusterID = observedID } + var failures []string + for i, errors := range statusErrors { + if len(errors) > 0 { + failures = append( + failures, + fmt.Sprintf("%s: %s", members[i].Name, strings.Join(errors, ", ")), + ) + } + } + if len(failures) > 0 { + // NOSPACE can reject writes cluster-wide while linearizable reads still + // succeed. A read through another member must not mask a reported error. + condition.Status, condition.Reason = metav1.ConditionFalse, "EtcdStatusError" + condition.Message = fmt.Sprintf( + "Etcd members report status errors (%s); failover protection is unavailable", + strings.Join(failures, "; "), + ) + return condition, members + } responding, readable := 0, 0 for _, member := range members { if member.BackendBytes != nil { diff --git a/pkg/resource-handler/controller/toposerver/health_etcd_test.go b/pkg/resource-handler/controller/toposerver/health_etcd_test.go index bc13dd2d4..8e0ce3507 100644 --- a/pkg/resource-handler/controller/toposerver/health_etcd_test.go +++ b/pkg/resource-handler/controller/toposerver/health_etcd_test.go @@ -7,7 +7,10 @@ import ( "testing" "time" + "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + pb "go.etcd.io/etcd/api/v3/etcdserverpb" + "go.etcd.io/etcd/api/v3/v3rpc/rpctypes" clientv3 "go.etcd.io/etcd/client/v3" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -31,6 +34,58 @@ func (c mappedHealthClient) Health(ctx context.Context, endpoint string) error { } func (c mappedHealthClient) Close() {} +func TestLiveEtcdNOSPACEHealth(t *testing.T) { + c, endpoints, _ := startTestEtcd(t) + ts := certTestTopoServer(nil) + mapping := map[string]string{} + for i, endpoint := range maintenanceEndpoints(ts) { + mapping[endpoint] = endpoints[i] + } + r := &TopoServerReconciler{ + newMaintenanceClient: func(context.Context, *multigresv1alpha1.TopoServer) (etcdMaintenanceClient, error) { + return mappedHealthClient{c, mapping}, nil + }, + } + ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) + defer cancel() + status, err := c.Status(ctx, endpoints[0]) + require.NoError(t, err) + _, err = pb.NewMaintenanceClient(c.first.ActiveConnection()).Alarm(ctx, &pb.AlarmRequest{ + Action: pb.AlarmRequest_ACTIVATE, + MemberID: status.Header.MemberId, + Alarm: pb.AlarmType_NOSPACE, + }) + require.NoError(t, err) + require.EventuallyWithT(t, func(t *assert.CollectT) { + condition, _ := r.probeHealth(ctx, ts) + assert.Equal(t, metav1.ConditionFalse, condition.Status) + assert.Equal(t, "EtcdStatusError", condition.Reason) + assert.Contains(t, condition.Message, "NOSPACE") + }, 10*time.Second, 100*time.Millisecond) + for _, endpoint := range endpoints { + require.NoError(t, c.Health(ctx, endpoint), "NOSPACE still permits linearizable reads") + _, err = c.clients[endpoint].Put(ctx, "/health-test", "blocked") + require.ErrorIs( + t, + err, + rpctypes.ErrNoSpace, + "a single member alarm blocks writes through every endpoint", + ) + } + _, err = c.first.AlarmDisarm(ctx, &clientv3.AlarmMember{ + MemberID: status.Header.MemberId, Alarm: pb.AlarmType_NOSPACE, + }) + require.NoError(t, err) + require.EventuallyWithT(t, func(t *assert.CollectT) { + condition, _ := r.probeHealth(ctx, ts) + assert.Equal(t, metav1.ConditionTrue, condition.Status) + assert.Equal(t, "QuorumAvailable", condition.Reason) + assert.NotContains(t, condition.Message, "NOSPACE") + }, 10*time.Second, 100*time.Millisecond) + _, err = c.first.Put(ctx, "/health-test", "recovered") + require.NoError(t, err) +} + func TestLiveEtcdQuorumHealth(t *testing.T) { c, endpoints, stop := startTestEtcd(t) ts := certTestTopoServer(nil) diff --git a/pkg/resource-handler/controller/toposerver/health_test.go b/pkg/resource-handler/controller/toposerver/health_test.go index 512e6a74e..fa454023f 100644 --- a/pkg/resource-handler/controller/toposerver/health_test.go +++ b/pkg/resource-handler/controller/toposerver/health_test.go @@ -45,12 +45,17 @@ func TestProbeTopologyQuorum(t *testing.T) { for _, tc := range []struct { name string unavailable, blocked []int + statusErrors map[int][]string clientError, clusterMismatch bool want metav1.ConditionStatus reason string }{ {name: "different etcd clusters", clusterMismatch: true, want: metav1.ConditionFalse, reason: "ClusterMismatch"}, {name: "healthy", want: metav1.ConditionTrue, reason: "QuorumAvailable"}, + {name: "NOSPACE despite successful reads", statusErrors: map[int][]string{0: {"NOSPACE"}}, want: metav1.ConditionFalse, reason: "EtcdStatusError"}, + {name: "healthy members cannot mask another member alarm", statusErrors: map[int][]string{2: {"NOSPACE"}}, want: metav1.ConditionFalse, reason: "EtcdStatusError"}, + {name: "multiple status errors", statusErrors: map[int][]string{0: {"NOSPACE", "CORRUPT"}, 1: {"unexpected health error"}}, want: metav1.ConditionFalse, reason: "EtcdStatusError"}, + {name: "alarm on member with failed reads", blocked: []int{0}, statusErrors: map[int][]string{0: {"NOSPACE"}}, want: metav1.ConditionFalse, reason: "EtcdStatusError"}, {name: "one unreachable member", unavailable: []int{0}, want: metav1.ConditionTrue, reason: "QuorumAvailable"}, {name: "only one endpoint reachable but quorum read succeeds", unavailable: []int{0, 1}, want: metav1.ConditionTrue, reason: "QuorumAvailable"}, {name: "status works without quorum", blocked: []int{0, 1, 2}, want: metav1.ConditionFalse, reason: "QuorumUnavailable"}, @@ -74,6 +79,9 @@ func TestProbeTopologyQuorum(t *testing.T) { for _, i := range tc.blocked { c.blockedReads[endpoints[i]] = true } + for i, errors := range tc.statusErrors { + f.statuses[endpoints[i]].Errors = errors + } r.newMaintenanceClient = func(ctx context.Context, _ *multigresv1alpha1.TopoServer) (etcdMaintenanceClient, error) { deadline, ok := ctx.Deadline() require.True(t, ok) @@ -90,12 +98,53 @@ func TestProbeTopologyQuorum(t *testing.T) { require.False(t, members[i].Up) require.Nil(t, members[i].BackendBytes) } + for i, errors := range tc.statusErrors { + require.Contains(t, condition.Message, members[i].Name) + for _, err := range errors { + require.Contains(t, condition.Message, err) + } + require.Equal(t, !c.blockedReads[endpoints[i]], members[i].Up) + require.Equal(t, f.statuses[endpoints[i]].DbSize, *members[i].BackendBytes) + require.Equal( + t, + f.statuses[endpoints[i]].DbSizeInUse, + *members[i].BackendInUseBytes, + ) + require.Equal(t, f.statuses[endpoints[i]].Header.Revision, *members[i].Revision) + } require.Empty(t, f.defragged) require.Empty(t, f.moved) }) } } +func TestHealthObservationRecoversAfterEtcdAlarm(t *testing.T) { + r, ts, f := maintenanceFixture(t) + for _, alarmed := range []bool{false, true, false} { + status := f.statuses[maintenanceEndpoints(ts)[0]] + status.Errors = nil + want, reason := metav1.ConditionTrue, "QuorumAvailable" + if alarmed { + status.Errors = []string{"NOSPACE"} + want, reason = metav1.ConditionFalse, "EtcdStatusError" + } + _, err := r.reconcileHealth( + t.Context(), + ctrl.Request{NamespacedName: client.ObjectKeyFromObject(ts)}, + ) + require.NoError(t, err) + fresh := &multigresv1alpha1.TopoServer{} + require.NoError(t, r.Get(t.Context(), client.ObjectKeyFromObject(ts), fresh)) + condition := meta.FindStatusCondition(fresh.Status.Conditions, "QuorumAvailable") + require.NotNil(t, condition) + require.Equal(t, want, condition.Status) + require.Equal(t, reason, condition.Reason) + if !alarmed { + require.NotContains(t, condition.Message, "NOSPACE") + } + } +} + func TestHealthObservationDuringMaintenance(t *testing.T) { r, ts, f := maintenanceFixture(t) state := &multigresv1alpha1.EtcdMaintenanceStatus{