diff --git a/crates/common/src/telemetry/metrics/store.rs b/crates/common/src/telemetry/metrics/store.rs index 849a253..a509168 100644 --- a/crates/common/src/telemetry/metrics/store.rs +++ b/crates/common/src/telemetry/metrics/store.rs @@ -116,8 +116,19 @@ impl StoredMetric { /// changes (MON-4). Per process: a restart counts from the start. static LAST: Mutex>> = Mutex::new(None); -/// One tick's samples (MON-4 to MON-6). -pub fn sample() -> Vec { +/// Gauges that count the whole cluster's data, not this node's. Only the node +/// that computes them (the metrics-calculation role) has a true reading; on +/// the others the queue gauge only moves with local queue events and drifts +/// below zero, and the account and domain counts stay at 0. +const CLUSTER_GAUGES: [MetricType; 3] = [ + MetricType::QueueCount, + MetricType::UserCount, + MetricType::DomainCount, +]; + +/// One tick's samples (MON-4 to MON-6). `calculates` is whether this node +/// computes the cluster-wide gauges; a node that doesn't leaves them out. +pub fn sample(calculates: bool) -> Vec { let mut last_guard = LAST.lock().unwrap(); let last = last_guard.get_or_insert_with(AHashMap::new); let mut samples = Vec::new(); @@ -140,6 +151,9 @@ pub fn sample() -> Vec { // Gauges: the reading, always (MON-5) for gauge in Collector::collect_gauges() { + if !calculates && CLUSTER_GAUGES.contains(&gauge.id()) { + continue; + } samples.push(Metric::Gauge(MetricCount { count: gauge.get(), metric: gauge.id(), @@ -181,7 +195,7 @@ impl Server { if store.is_none() { return; } - let samples = sample(); + let samples = sample(self.core.network.roles.metrics_calculate); let count = samples.len(); let started = std::time::Instant::now(); match store.write_metrics(samples, now()).await { @@ -271,3 +285,41 @@ impl Server { } } } + +#[cfg(test)] +mod tests { + use super::*; + + fn gauges(samples: &[Metric]) -> Vec { + samples + .iter() + .filter_map(|m| match m { + Metric::Gauge(g) => Some(g.metric), + _ => None, + }) + .collect() + } + + #[test] + fn only_the_calculating_node_stores_cluster_gauges() { + let all = gauges(&sample(true)); + let local = gauges(&sample(false)); + for metric in CLUSTER_GAUGES { + assert!( + all.contains(&metric), + "{metric:?} missing on the calculating node" + ); + assert!( + !local.contains(&metric), + "{metric:?} stored by a node that doesn't compute it" + ); + } + // Per-node gauges are stored either way + for metric in [MetricType::ServerMemory, MetricType::HttpActiveConnections] { + assert!( + all.contains(&metric) && local.contains(&metric), + "{metric:?}" + ); + } + } +}