From b5d9df2a9ddb0c3b2e7b2c64e976e5acffe5888d Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Thu, 24 Sep 2026 10:46:44 +0530 Subject: [PATCH 1/5] expose node counts for enterpise analytics --- src/analytics.rs | 49 +++++++++++++++++++++++++++++++++++------------- 1 file changed, 36 insertions(+), 13 deletions(-) diff --git a/src/analytics.rs b/src/analytics.rs index 6ceb8ba83..f0110f09e 100644 --- a/src/analytics.rs +++ b/src/analytics.rs @@ -93,6 +93,13 @@ pub struct Report { metrics: HashMap, } +pub struct ClusterNodeCounts { + pub active_ingestors: u64, + pub inactive_ingestors: u64, + pub active_queriers: u64, + pub inactive_queriers: u64, +} + impl Report { pub async fn new() -> anyhow::Result { let mut upt: f64 = 0.0; @@ -114,8 +121,7 @@ impl Report { let ingestor_metrics = fetch_ingestors_metrics().await?; let mut active_indexers = 0; let mut inactive_indexers = 0; - let mut active_queriers = 0; - let mut inactive_queriers = 0; + let (active_queriers, inactive_queriers) = fetch_querier_metrics().await?; // check liveness of indexers // get the count of active and inactive indexers @@ -129,17 +135,6 @@ impl Report { } } - // check liveness of queriers - // get the count of active and inactive queriers - let query_infos: Vec = - cluster::get_node_info(NodeType::Querier, &None).await?; - for query in query_infos { - if check_liveness(&query.domain_name).await { - active_queriers += 1; - } else { - inactive_queriers += 1; - } - } Ok(Self { deployment_id: storage::StorageMetadata::global().deployment_id, uptime: upt, @@ -183,6 +178,34 @@ impl Report { } } +pub async fn fetch_cluster_node_counts() -> anyhow::Result { + let ingestor_metrics = fetch_ingestors_metrics().await?; + let (active_queriers, inactive_queriers) = fetch_querier_metrics().await?; + + Ok(ClusterNodeCounts { + active_ingestors: ingestor_metrics.0, + inactive_ingestors: ingestor_metrics.1, + active_queriers, + inactive_queriers, + }) +} + +async fn fetch_querier_metrics() -> anyhow::Result<(u64, u64)> { + let mut active_queriers = 0; + let mut inactive_queriers = 0; + let query_infos: Vec = cluster::get_node_info(NodeType::Querier, &None).await?; + + for query in query_infos { + if check_liveness(&query.domain_name).await { + active_queriers += 1; + } else { + inactive_queriers += 1; + } + } + + Ok((active_queriers, inactive_queriers)) +} + /// build the node metrics for the node ingestor endpoint pub async fn get_analytics(_: HttpRequest) -> impl Responder { let json = NodeMetrics::build(); From 84f5403e105005ce2a07b119691bd3494ca2aed0 Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Thu, 24 Sep 2026 17:19:39 +0530 Subject: [PATCH 2/5] fix: count cluster nodes using liveness --- src/analytics.rs | 28 ++++++++++++++-------------- 1 file changed, 14 insertions(+), 14 deletions(-) diff --git a/src/analytics.rs b/src/analytics.rs index f0110f09e..2044c519d 100644 --- a/src/analytics.rs +++ b/src/analytics.rs @@ -121,7 +121,7 @@ impl Report { let ingestor_metrics = fetch_ingestors_metrics().await?; let mut active_indexers = 0; let mut inactive_indexers = 0; - let (active_queriers, inactive_queriers) = fetch_querier_metrics().await?; + let (active_queriers, inactive_queriers) = fetch_node_counts(NodeType::Querier).await?; // check liveness of indexers // get the count of active and inactive indexers @@ -179,31 +179,31 @@ impl Report { } pub async fn fetch_cluster_node_counts() -> anyhow::Result { - let ingestor_metrics = fetch_ingestors_metrics().await?; - let (active_queriers, inactive_queriers) = fetch_querier_metrics().await?; + let (active_ingestors, inactive_ingestors) = fetch_node_counts(NodeType::Ingestor).await?; + let (active_queriers, inactive_queriers) = fetch_node_counts(NodeType::Querier).await?; Ok(ClusterNodeCounts { - active_ingestors: ingestor_metrics.0, - inactive_ingestors: ingestor_metrics.1, + active_ingestors, + inactive_ingestors, active_queriers, inactive_queriers, }) } -async fn fetch_querier_metrics() -> anyhow::Result<(u64, u64)> { - let mut active_queriers = 0; - let mut inactive_queriers = 0; - let query_infos: Vec = cluster::get_node_info(NodeType::Querier, &None).await?; +async fn fetch_node_counts(node_type: NodeType) -> anyhow::Result<(u64, u64)> { + let mut active_nodes = 0; + let mut inactive_nodes = 0; + let node_infos: Vec = cluster::get_node_info(node_type, &None).await?; - for query in query_infos { - if check_liveness(&query.domain_name).await { - active_queriers += 1; + for node in node_infos { + if check_liveness(&node.domain_name).await { + active_nodes += 1; } else { - inactive_queriers += 1; + inactive_nodes += 1; } } - Ok((active_queriers, inactive_queriers)) + Ok((active_nodes, inactive_nodes)) } /// build the node metrics for the node ingestor endpoint From 98b05fd96d3d49dd8e980dfb4073a7b581224f9b Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Fri, 25 Sep 2026 15:44:51 +0530 Subject: [PATCH 3/5] add cluster node counts to metrics registry --- src/analytics.rs | 6 ++++++ src/metrics/mod.rs | 47 +++++++++++++++++++++++++++++++++++++++++++++- 2 files changed, 52 insertions(+), 1 deletion(-) diff --git a/src/analytics.rs b/src/analytics.rs index 2044c519d..5874bc1fb 100644 --- a/src/analytics.rs +++ b/src/analytics.rs @@ -41,6 +41,7 @@ use crate::{ modal::{NodeMetadata, NodeType}, }, }, + metrics::{ACTIVE_INGESTORS, ACTIVE_QUERIERS, INACTIVE_INGESTORS, INACTIVE_QUERIERS}, option::Mode, parseable::PARSEABLE, stats::{self, Stats}, @@ -182,6 +183,11 @@ pub async fn fetch_cluster_node_counts() -> anyhow::Result { let (active_ingestors, inactive_ingestors) = fetch_node_counts(NodeType::Ingestor).await?; let (active_queriers, inactive_queriers) = fetch_node_counts(NodeType::Querier).await?; + ACTIVE_INGESTORS.set(i64::try_from(active_ingestors).unwrap_or(i64::MAX)); + INACTIVE_INGESTORS.set(i64::try_from(inactive_ingestors).unwrap_or(i64::MAX)); + ACTIVE_QUERIERS.set(i64::try_from(active_queriers).unwrap_or(i64::MAX)); + INACTIVE_QUERIERS.set(i64::try_from(inactive_queriers).unwrap_or(i64::MAX)); + Ok(ClusterNodeCounts { active_ingestors, inactive_ingestors, diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index c06cc407f..5de91d643 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -29,7 +29,8 @@ use actix_web_prometheus::{PrometheusMetrics, PrometheusMetricsBuilder}; use error::MetricsError; use once_cell::sync::Lazy; use prometheus::{ - Gauge, GaugeVec, HistogramOpts, HistogramVec, IntCounterVec, IntGaugeVec, Opts, Registry, + Gauge, GaugeVec, HistogramOpts, HistogramVec, IntCounterVec, IntGauge, IntGaugeVec, Opts, + Registry, core::{Atomic, AtomicF64}, }; @@ -208,6 +209,38 @@ pub fn record_disk_metrics(disk_type: &str, path: &Path) { } } +pub static ACTIVE_INGESTORS: Lazy = Lazy::new(|| { + IntGauge::with_opts( + Opts::new("active_ingestors", "Number of active ingestor nodes") + .namespace(METRICS_NAMESPACE), + ) + .expect("metric can be created") +}); + +pub static INACTIVE_INGESTORS: Lazy = Lazy::new(|| { + IntGauge::with_opts( + Opts::new("inactive_ingestors", "Number of inactive ingestor nodes") + .namespace(METRICS_NAMESPACE), + ) + .expect("metric can be created") +}); + +pub static ACTIVE_QUERIERS: Lazy = Lazy::new(|| { + IntGauge::with_opts( + Opts::new("active_queriers", "Number of active querier nodes") + .namespace(METRICS_NAMESPACE), + ) + .expect("metric can be created") +}); + +pub static INACTIVE_QUERIERS: Lazy = Lazy::new(|| { + IntGauge::with_opts( + Opts::new("inactive_queriers", "Number of inactive querier nodes") + .namespace(METRICS_NAMESPACE), + ) + .expect("metric can be created") +}); + pub static PROCESS_CPU_USAGE_PERCENT_AVG: Lazy = Lazy::new(|| { Gauge::with_opts( Opts::new( @@ -826,6 +859,18 @@ fn custom_metrics(registry: &Registry) { registry .register(Box::new(DISK_TOTAL_BYTES.clone())) .expect("metric can be registered"); + registry + .register(Box::new(ACTIVE_INGESTORS.clone())) + .expect("metric can be registered"); + registry + .register(Box::new(INACTIVE_INGESTORS.clone())) + .expect("metric can be registered"); + registry + .register(Box::new(ACTIVE_QUERIERS.clone())) + .expect("metric can be registered"); + registry + .register(Box::new(INACTIVE_QUERIERS.clone())) + .expect("metric can be registered"); registry .register(Box::new(PROCESS_CPU_USAGE_PERCENT_AVG.clone())) .expect("metric can be registered"); From 84d4e6a6db9b84effd2cc439c4bc2af8b85d1f1b Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Fri, 25 Sep 2026 15:56:29 +0530 Subject: [PATCH 4/5] removed unwanted struct --- src/analytics.rs | 16 ++-------------- src/metrics/mod.rs | 3 +-- 2 files changed, 3 insertions(+), 16 deletions(-) diff --git a/src/analytics.rs b/src/analytics.rs index 5874bc1fb..0065f447b 100644 --- a/src/analytics.rs +++ b/src/analytics.rs @@ -94,13 +94,6 @@ pub struct Report { metrics: HashMap, } -pub struct ClusterNodeCounts { - pub active_ingestors: u64, - pub inactive_ingestors: u64, - pub active_queriers: u64, - pub inactive_queriers: u64, -} - impl Report { pub async fn new() -> anyhow::Result { let mut upt: f64 = 0.0; @@ -179,7 +172,7 @@ impl Report { } } -pub async fn fetch_cluster_node_counts() -> anyhow::Result { +pub async fn fetch_cluster_node_counts() -> anyhow::Result<()> { let (active_ingestors, inactive_ingestors) = fetch_node_counts(NodeType::Ingestor).await?; let (active_queriers, inactive_queriers) = fetch_node_counts(NodeType::Querier).await?; @@ -188,12 +181,7 @@ pub async fn fetch_cluster_node_counts() -> anyhow::Result { ACTIVE_QUERIERS.set(i64::try_from(active_queriers).unwrap_or(i64::MAX)); INACTIVE_QUERIERS.set(i64::try_from(inactive_queriers).unwrap_or(i64::MAX)); - Ok(ClusterNodeCounts { - active_ingestors, - inactive_ingestors, - active_queriers, - inactive_queriers, - }) + Ok(()) } async fn fetch_node_counts(node_type: NodeType) -> anyhow::Result<(u64, u64)> { diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index 5de91d643..2491b243f 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -227,8 +227,7 @@ pub static INACTIVE_INGESTORS: Lazy = Lazy::new(|| { pub static ACTIVE_QUERIERS: Lazy = Lazy::new(|| { IntGauge::with_opts( - Opts::new("active_queriers", "Number of active querier nodes") - .namespace(METRICS_NAMESPACE), + Opts::new("active_queriers", "Number of active querier nodes").namespace(METRICS_NAMESPACE), ) .expect("metric can be created") }); From 3ee09f356e0cab419e35fa604d91e918921e55d8 Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Mon, 28 Sep 2026 13:23:06 +0530 Subject: [PATCH 5/5] refactor: use generic cluster node metrics --- src/analytics.rs | 18 +++++++++++++----- src/metrics/mod.rs | 44 +++++++++++--------------------------------- 2 files changed, 24 insertions(+), 38 deletions(-) diff --git a/src/analytics.rs b/src/analytics.rs index 0065f447b..ee6e0710c 100644 --- a/src/analytics.rs +++ b/src/analytics.rs @@ -41,7 +41,7 @@ use crate::{ modal::{NodeMetadata, NodeType}, }, }, - metrics::{ACTIVE_INGESTORS, ACTIVE_QUERIERS, INACTIVE_INGESTORS, INACTIVE_QUERIERS}, + metrics::{ACTIVE_NODES, INACTIVE_NODES}, option::Mode, parseable::PARSEABLE, stats::{self, Stats}, @@ -176,10 +176,18 @@ pub async fn fetch_cluster_node_counts() -> anyhow::Result<()> { let (active_ingestors, inactive_ingestors) = fetch_node_counts(NodeType::Ingestor).await?; let (active_queriers, inactive_queriers) = fetch_node_counts(NodeType::Querier).await?; - ACTIVE_INGESTORS.set(i64::try_from(active_ingestors).unwrap_or(i64::MAX)); - INACTIVE_INGESTORS.set(i64::try_from(inactive_ingestors).unwrap_or(i64::MAX)); - ACTIVE_QUERIERS.set(i64::try_from(active_queriers).unwrap_or(i64::MAX)); - INACTIVE_QUERIERS.set(i64::try_from(inactive_queriers).unwrap_or(i64::MAX)); + ACTIVE_NODES + .with_label_values(&["ingestor"]) + .set(i64::try_from(active_ingestors).unwrap_or(i64::MAX)); + INACTIVE_NODES + .with_label_values(&["ingestor"]) + .set(i64::try_from(inactive_ingestors).unwrap_or(i64::MAX)); + ACTIVE_NODES + .with_label_values(&["querier"]) + .set(i64::try_from(active_queriers).unwrap_or(i64::MAX)); + INACTIVE_NODES + .with_label_values(&["querier"]) + .set(i64::try_from(inactive_queriers).unwrap_or(i64::MAX)); Ok(()) } diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index 2491b243f..3b774c5a0 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -29,8 +29,7 @@ use actix_web_prometheus::{PrometheusMetrics, PrometheusMetricsBuilder}; use error::MetricsError; use once_cell::sync::Lazy; use prometheus::{ - Gauge, GaugeVec, HistogramOpts, HistogramVec, IntCounterVec, IntGauge, IntGaugeVec, Opts, - Registry, + Gauge, GaugeVec, HistogramOpts, HistogramVec, IntCounterVec, IntGaugeVec, Opts, Registry, core::{Atomic, AtomicF64}, }; @@ -209,33 +208,18 @@ pub fn record_disk_metrics(disk_type: &str, path: &Path) { } } -pub static ACTIVE_INGESTORS: Lazy = Lazy::new(|| { - IntGauge::with_opts( - Opts::new("active_ingestors", "Number of active ingestor nodes") - .namespace(METRICS_NAMESPACE), - ) - .expect("metric can be created") -}); - -pub static INACTIVE_INGESTORS: Lazy = Lazy::new(|| { - IntGauge::with_opts( - Opts::new("inactive_ingestors", "Number of inactive ingestor nodes") - .namespace(METRICS_NAMESPACE), - ) - .expect("metric can be created") -}); - -pub static ACTIVE_QUERIERS: Lazy = Lazy::new(|| { - IntGauge::with_opts( - Opts::new("active_queriers", "Number of active querier nodes").namespace(METRICS_NAMESPACE), +pub static ACTIVE_NODES: Lazy = Lazy::new(|| { + IntGaugeVec::new( + Opts::new("active_nodes", "Number of active nodes").namespace(METRICS_NAMESPACE), + &["node_type"], ) .expect("metric can be created") }); -pub static INACTIVE_QUERIERS: Lazy = Lazy::new(|| { - IntGauge::with_opts( - Opts::new("inactive_queriers", "Number of inactive querier nodes") - .namespace(METRICS_NAMESPACE), +pub static INACTIVE_NODES: Lazy = Lazy::new(|| { + IntGaugeVec::new( + Opts::new("inactive_nodes", "Number of inactive nodes").namespace(METRICS_NAMESPACE), + &["node_type"], ) .expect("metric can be created") }); @@ -859,16 +843,10 @@ fn custom_metrics(registry: &Registry) { .register(Box::new(DISK_TOTAL_BYTES.clone())) .expect("metric can be registered"); registry - .register(Box::new(ACTIVE_INGESTORS.clone())) - .expect("metric can be registered"); - registry - .register(Box::new(INACTIVE_INGESTORS.clone())) - .expect("metric can be registered"); - registry - .register(Box::new(ACTIVE_QUERIERS.clone())) + .register(Box::new(ACTIVE_NODES.clone())) .expect("metric can be registered"); registry - .register(Box::new(INACTIVE_QUERIERS.clone())) + .register(Box::new(INACTIVE_NODES.clone())) .expect("metric can be registered"); registry .register(Box::new(PROCESS_CPU_USAGE_PERCENT_AVG.clone()))