Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 38 additions & 13 deletions src/analytics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ use crate::{
modal::{NodeMetadata, NodeType},
},
},
metrics::{ACTIVE_NODES, INACTIVE_NODES},
option::Mode,
parseable::PARSEABLE,
stats::{self, Stats},
Expand Down Expand Up @@ -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
Expand All @@ -129,17 +129,6 @@ impl Report {
}
}

// check liveness of queriers
// get the count of active and inactive queriers
let query_infos: Vec<NodeMetadata> =
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,
Expand Down Expand Up @@ -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<NodeMetadata> = 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();
Expand Down
22 changes: 22 additions & 0 deletions src/metrics/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,22 @@ pub fn record_disk_metrics(disk_type: &str, path: &Path) {
}
}

pub static ACTIVE_NODES: Lazy<IntGaugeVec> = 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<IntGaugeVec> = 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<Gauge> = Lazy::new(|| {
Gauge::with_opts(
Opts::new(
Expand Down Expand Up @@ -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");
Expand Down
Loading