diff --git a/src/analytics.rs b/src/analytics.rs index 6ceb8ba83..ee6e0710c 100644 --- a/src/analytics.rs +++ b/src/analytics.rs @@ -41,6 +41,7 @@ use crate::{ modal::{NodeMetadata, NodeType}, }, }, + metrics::{ACTIVE_NODES, INACTIVE_NODES}, option::Mode, parseable::PARSEABLE, stats::{self, Stats}, @@ -114,8 +115,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_node_counts(NodeType::Querier).await?; // check liveness of indexers // get the count of active and inactive indexers @@ -129,17 +129,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 +172,42 @@ impl Report { } } +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_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(()) +} + +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 node in node_infos { + if check_liveness(&node.domain_name).await { + active_nodes += 1; + } else { + inactive_nodes += 1; + } + } + + Ok((active_nodes, inactive_nodes)) +} + /// build the node metrics for the node ingestor endpoint pub async fn get_analytics(_: HttpRequest) -> impl Responder { let json = NodeMetrics::build(); diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index fce8c70d3..135ac5e10 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -216,6 +216,22 @@ pub fn record_disk_metrics(disk_type: &str, path: &Path) { } } +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_NODES: Lazy = Lazy::new(|| { + IntGaugeVec::new( + Opts::new("inactive_nodes", "Number of inactive nodes").namespace(METRICS_NAMESPACE), + &["node_type"], + ) + .expect("metric can be created") +}); + pub static PROCESS_CPU_USAGE_PERCENT_AVG: Lazy = Lazy::new(|| { Gauge::with_opts( Opts::new( @@ -862,6 +878,12 @@ fn custom_metrics(registry: &Registry) { registry .register(Box::new(DISK_TOTAL_BYTES.clone())) .expect("metric can be registered"); + registry + .register(Box::new(ACTIVE_NODES.clone())) + .expect("metric can be registered"); + registry + .register(Box::new(INACTIVE_NODES.clone())) + .expect("metric can be registered"); registry .register(Box::new(PROCESS_CPU_USAGE_PERCENT_AVG.clone())) .expect("metric can be registered");