From cb12959e2bfd9958a43bcb0aa7dc12b6cf3f6f49 Mon Sep 17 00:00:00 2001 From: Praveen K B Date: Fri, 7 Aug 2026 15:55:40 +0530 Subject: [PATCH 1/4] feat: expose averaged process CPU and memory per cluster node --- src/handlers/http/resource_check.rs | 18 +++++++ src/metrics/mod.rs | 82 ++++++++++++++++++++++++++++- src/metrics/prom_utils.rs | 10 ++++ 3 files changed, 109 insertions(+), 1 deletion(-) diff --git a/src/handlers/http/resource_check.rs b/src/handlers/http/resource_check.rs index aaf3595df..4cbcdb225 100644 --- a/src/handlers/http/resource_check.rs +++ b/src/handlers/http/resource_check.rs @@ -32,8 +32,11 @@ use tokio::{ use tracing::{info, trace, warn}; use crate::analytics::{SYS_INFO, refresh_sys_info}; +use crate::metrics::record_process_metrics_sample; use crate::parseable::PARSEABLE; +const PROCESS_METRICS_SAMPLE_INTERVAL: Duration = Duration::from_secs(10); + static RESOURCE_CHECK_ENABLED: LazyLock> = LazyLock::new(|| Arc::new(AtomicBool::new(false))); @@ -42,6 +45,7 @@ pub fn spawn_resource_monitor(shutdown_rx: tokio::sync::oneshot::Receiver<()>) { tokio::spawn(async move { let resource_check_interval = PARSEABLE.options.resource_check_interval; let mut check_interval = interval(Duration::from_secs(resource_check_interval)); + let mut process_metrics_interval = interval(PROCESS_METRICS_SAMPLE_INTERVAL); let mut shutdown_rx = shutdown_rx; let cpu_threshold = PARSEABLE.options.cpu_utilization_threshold; @@ -106,6 +110,20 @@ pub fn spawn_resource_monitor(shutdown_rx: tokio::sync::oneshot::Receiver<()>) { } } }, + _ = process_metrics_interval.tick() => { + refresh_sys_info(); + let process_metrics = tokio::task::spawn_blocking(|| { + let sys = SYS_INFO.lock().unwrap(); + sysinfo::get_current_pid() + .ok() + .and_then(|pid| sys.process(pid)) + .map(|process| (process.cpu_usage() as f64, process.memory())) + }).await.unwrap(); + + if let Some((cpu_usage, memory_bytes)) = process_metrics { + record_process_metrics_sample(cpu_usage, memory_bytes); + } + }, _ = &mut shutdown_rx => { trace!("Resource monitor shutting down"); break; diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index 57ca55a03..a705b89e2 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -25,7 +25,8 @@ use actix_web::Responder; use actix_web_prometheus::{PrometheusMetrics, PrometheusMetricsBuilder}; use error::MetricsError; use once_cell::sync::Lazy; -use prometheus::{HistogramOpts, HistogramVec, IntCounterVec, IntGaugeVec, Opts, Registry}; +use prometheus::{Gauge, HistogramOpts, HistogramVec, IntCounterVec, IntGaugeVec, Opts, Registry}; +use std::sync::atomic::{AtomicU64, Ordering}; pub const METRICS_NAMESPACE: &str = env!("CARGO_PKG_NAME"); @@ -175,6 +176,79 @@ pub static STAGING_FILES: Lazy = Lazy::new(|| { .expect("metric can be created") }); +pub static PROCESS_CPU_USAGE_PERCENT: Lazy = Lazy::new(|| { + Gauge::with_opts( + Opts::new( + "process_cpu_usage_percent", + "Current CPU usage percent for this Parseable process", + ) + .namespace(METRICS_NAMESPACE), + ) + .expect("metric can be created") +}); + +pub static PROCESS_MEMORY_BYTES: Lazy = Lazy::new(|| { + Gauge::with_opts( + Opts::new( + "process_memory_bytes", + "Current resident memory used by this Parseable process in bytes", + ) + .namespace(METRICS_NAMESPACE), + ) + .expect("metric can be created") +}); + +const CPU_USAGE_PRECISION: f64 = 1_000.0; + +#[derive(Default)] +struct ProcessMetricsAccumulator { + cpu_usage_sum: AtomicU64, + memory_bytes_sum: AtomicU64, + sample_count: AtomicU64, +} + +impl ProcessMetricsAccumulator { + fn record(&self, cpu_usage_percent: f64, memory_bytes: u64) -> (f64, f64) { + self.cpu_usage_sum.fetch_add( + (cpu_usage_percent * CPU_USAGE_PRECISION).round() as u64, + Ordering::Relaxed, + ); + self.memory_bytes_sum + .fetch_add(memory_bytes, Ordering::Relaxed); + let sample_count = self.sample_count.fetch_add(1, Ordering::Relaxed) + 1; + + ( + self.cpu_usage_sum.load(Ordering::Relaxed) as f64 + / sample_count as f64 + / CPU_USAGE_PRECISION, + self.memory_bytes_sum.load(Ordering::Relaxed) as f64 / sample_count as f64, + ) + } +} + +static PROCESS_METRICS_ACCUMULATOR: Lazy = + Lazy::new(ProcessMetricsAccumulator::default); + +pub fn record_process_metrics_sample(cpu_usage_percent: f64, memory_bytes: u64) { + let (average_cpu_usage, average_memory_bytes) = + PROCESS_METRICS_ACCUMULATOR.record(cpu_usage_percent, memory_bytes); + PROCESS_CPU_USAGE_PERCENT.set(average_cpu_usage); + PROCESS_MEMORY_BYTES.set(average_memory_bytes); +} + +#[cfg(test)] +mod process_metrics_tests { + use super::ProcessMetricsAccumulator; + + #[test] + fn averages_process_metric_samples() { + let accumulator = ProcessMetricsAccumulator::default(); + + assert_eq!(accumulator.record(10.0, 100), (10.0, 100.0)); + assert_eq!(accumulator.record(20.0, 300), (15.0, 200.0)); + } +} + pub static QUERY_EXECUTE_TIME: Lazy = Lazy::new(|| { HistogramVec::new( HistogramOpts::new("query_execute_time", "Query execute time").namespace(METRICS_NAMESPACE), @@ -663,6 +737,12 @@ fn custom_metrics(registry: &Registry) { registry .register(Box::new(STAGING_FILES.clone())) .expect("metric can be registered"); + registry + .register(Box::new(PROCESS_CPU_USAGE_PERCENT.clone())) + .expect("metric can be registered"); + registry + .register(Box::new(PROCESS_MEMORY_BYTES.clone())) + .expect("metric can be registered"); registry .register(Box::new(QUERY_EXECUTE_TIME.clone())) .expect("metric can be registered"); diff --git a/src/metrics/prom_utils.rs b/src/metrics/prom_utils.rs index 3f04d89f6..9d63bd78a 100644 --- a/src/metrics/prom_utils.rs +++ b/src/metrics/prom_utils.rs @@ -53,6 +53,8 @@ pub struct Metrics { event_time: NaiveDateTime, commit: String, staging: String, + process_cpu_usage_percent: f64, + process_memory_bytes: f64, } #[derive(Debug, Serialize, Default, Clone)] @@ -89,6 +91,8 @@ impl Default for Metrics { event_time: Utc::now().naive_utc(), commit: "".to_string(), staging: "".to_string(), + process_cpu_usage_percent: 0.0, + process_memory_bytes: 0.0, } } } @@ -113,6 +117,8 @@ impl Metrics { event_time: Utc::now().naive_utc(), commit: "".to_string(), staging: "".to_string(), + process_cpu_usage_percent: 0.0, + process_memory_bytes: 0.0, } } } @@ -187,6 +193,10 @@ impl Metrics { "process_resident_memory_bytes" => { prom_dress.process_resident_memory_bytes += val } + "parseable_process_cpu_usage_percent" => { + prom_dress.process_cpu_usage_percent += val + } + "parseable_process_memory_bytes" => prom_dress.process_memory_bytes += val, "parseable_storage_size" => { if sample.labels.get("type").expect("type is present") == "staging" { prom_dress.parseable_storage_size.staging += val; From b72ebfd3cd5850cc7a197b55259de25431251878 Mon Sep 17 00:00:00 2001 From: Praveen K B Date: Sun, 16 Aug 2026 16:26:56 +0530 Subject: [PATCH 2/4] rename: clarify process CPU/memory metrics as lifetime averages --- src/metrics/mod.rs | 20 ++++++++++---------- src/metrics/prom_utils.rs | 18 +++++++++--------- 2 files changed, 19 insertions(+), 19 deletions(-) diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index a705b89e2..2afb0d28f 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -176,22 +176,22 @@ pub static STAGING_FILES: Lazy = Lazy::new(|| { .expect("metric can be created") }); -pub static PROCESS_CPU_USAGE_PERCENT: Lazy = Lazy::new(|| { +pub static PROCESS_CPU_USAGE_PERCENT_AVG: Lazy = Lazy::new(|| { Gauge::with_opts( Opts::new( - "process_cpu_usage_percent", - "Current CPU usage percent for this Parseable process", + "process_cpu_usage_percent_avg", + "Lifetime average CPU usage percent for this Parseable process", ) .namespace(METRICS_NAMESPACE), ) .expect("metric can be created") }); -pub static PROCESS_MEMORY_BYTES: Lazy = Lazy::new(|| { +pub static PROCESS_MEMORY_BYTES_AVG: Lazy = Lazy::new(|| { Gauge::with_opts( Opts::new( - "process_memory_bytes", - "Current resident memory used by this Parseable process in bytes", + "process_memory_bytes_avg", + "Lifetime average resident memory used by this Parseable process in bytes", ) .namespace(METRICS_NAMESPACE), ) @@ -232,8 +232,8 @@ static PROCESS_METRICS_ACCUMULATOR: Lazy = pub fn record_process_metrics_sample(cpu_usage_percent: f64, memory_bytes: u64) { let (average_cpu_usage, average_memory_bytes) = PROCESS_METRICS_ACCUMULATOR.record(cpu_usage_percent, memory_bytes); - PROCESS_CPU_USAGE_PERCENT.set(average_cpu_usage); - PROCESS_MEMORY_BYTES.set(average_memory_bytes); + PROCESS_CPU_USAGE_PERCENT_AVG.set(average_cpu_usage); + PROCESS_MEMORY_BYTES_AVG.set(average_memory_bytes); } #[cfg(test)] @@ -738,10 +738,10 @@ fn custom_metrics(registry: &Registry) { .register(Box::new(STAGING_FILES.clone())) .expect("metric can be registered"); registry - .register(Box::new(PROCESS_CPU_USAGE_PERCENT.clone())) + .register(Box::new(PROCESS_CPU_USAGE_PERCENT_AVG.clone())) .expect("metric can be registered"); registry - .register(Box::new(PROCESS_MEMORY_BYTES.clone())) + .register(Box::new(PROCESS_MEMORY_BYTES_AVG.clone())) .expect("metric can be registered"); registry .register(Box::new(QUERY_EXECUTE_TIME.clone())) diff --git a/src/metrics/prom_utils.rs b/src/metrics/prom_utils.rs index 9d63bd78a..88a82787f 100644 --- a/src/metrics/prom_utils.rs +++ b/src/metrics/prom_utils.rs @@ -53,8 +53,8 @@ pub struct Metrics { event_time: NaiveDateTime, commit: String, staging: String, - process_cpu_usage_percent: f64, - process_memory_bytes: f64, + process_cpu_usage_percent_avg: f64, + process_memory_bytes_avg: f64, } #[derive(Debug, Serialize, Default, Clone)] @@ -91,8 +91,8 @@ impl Default for Metrics { event_time: Utc::now().naive_utc(), commit: "".to_string(), staging: "".to_string(), - process_cpu_usage_percent: 0.0, - process_memory_bytes: 0.0, + process_cpu_usage_percent_avg: 0.0, + process_memory_bytes_avg: 0.0, } } } @@ -117,8 +117,8 @@ impl Metrics { event_time: Utc::now().naive_utc(), commit: "".to_string(), staging: "".to_string(), - process_cpu_usage_percent: 0.0, - process_memory_bytes: 0.0, + process_cpu_usage_percent_avg: 0.0, + process_memory_bytes_avg: 0.0, } } } @@ -193,10 +193,10 @@ impl Metrics { "process_resident_memory_bytes" => { prom_dress.process_resident_memory_bytes += val } - "parseable_process_cpu_usage_percent" => { - prom_dress.process_cpu_usage_percent += val + "parseable_process_cpu_usage_percent_avg" => { + prom_dress.process_cpu_usage_percent_avg += val } - "parseable_process_memory_bytes" => prom_dress.process_memory_bytes += val, + "parseable_process_memory_bytes_avg" => prom_dress.process_memory_bytes_avg += val, "parseable_storage_size" => { if sample.labels.get("type").expect("type is present") == "staging" { prom_dress.parseable_storage_size.staging += val; From e7f7181757796cf96bd7f11f9dbf6ba616b62d95 Mon Sep 17 00:00:00 2001 From: Anant Vindal Date: Mon, 17 Aug 2026 10:50:01 +0530 Subject: [PATCH 3/4] EWMA instead of lifetime avg --- src/handlers/http/resource_check.rs | 2 +- src/metrics/mod.rs | 21 +++++++++++++-------- src/metrics/prom_utils.rs | 4 +++- 3 files changed, 17 insertions(+), 10 deletions(-) diff --git a/src/handlers/http/resource_check.rs b/src/handlers/http/resource_check.rs index 4cbcdb225..28e05c1db 100644 --- a/src/handlers/http/resource_check.rs +++ b/src/handlers/http/resource_check.rs @@ -35,7 +35,7 @@ use crate::analytics::{SYS_INFO, refresh_sys_info}; use crate::metrics::record_process_metrics_sample; use crate::parseable::PARSEABLE; -const PROCESS_METRICS_SAMPLE_INTERVAL: Duration = Duration::from_secs(10); +const PROCESS_METRICS_SAMPLE_INTERVAL: Duration = Duration::from_secs(5); static RESOURCE_CHECK_ENABLED: LazyLock> = LazyLock::new(|| Arc::new(AtomicBool::new(false))); diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index 2afb0d28f..50c336e41 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -204,7 +204,6 @@ const CPU_USAGE_PRECISION: f64 = 1_000.0; struct ProcessMetricsAccumulator { cpu_usage_sum: AtomicU64, memory_bytes_sum: AtomicU64, - sample_count: AtomicU64, } impl ProcessMetricsAccumulator { @@ -215,14 +214,20 @@ impl ProcessMetricsAccumulator { ); self.memory_bytes_sum .fetch_add(memory_bytes, Ordering::Relaxed); - let sample_count = self.sample_count.fetch_add(1, Ordering::Relaxed) + 1; - ( - self.cpu_usage_sum.load(Ordering::Relaxed) as f64 - / sample_count as f64 - / CPU_USAGE_PRECISION, - self.memory_bytes_sum.load(Ordering::Relaxed) as f64 / sample_count as f64, - ) + // Exponentially Weighted Moving Average is better than + // a lifetime average + // A spike which occurred 5 days ago should not affect the average utilization + // for the last minute + // α = 1 - exp(-Δt / τ) = 1 - exp(-5/60) ≈ 0.0800 + // S_new = S_old + α * (x_new - S_old) + let s_cpu_old = self.cpu_usage_sum.load(Ordering::Relaxed) as f64; + let s_cpu_new = s_cpu_old + 0.08 * (cpu_usage_percent - s_cpu_old); + + let s_mem_old = self.memory_bytes_sum.load(Ordering::Relaxed) as f64; + let s_mem_new = s_mem_old + 0.08 * (memory_bytes as f64 - s_mem_old); + + (s_cpu_new, s_mem_new) } } diff --git a/src/metrics/prom_utils.rs b/src/metrics/prom_utils.rs index 88a82787f..f0eff5e04 100644 --- a/src/metrics/prom_utils.rs +++ b/src/metrics/prom_utils.rs @@ -196,7 +196,9 @@ impl Metrics { "parseable_process_cpu_usage_percent_avg" => { prom_dress.process_cpu_usage_percent_avg += val } - "parseable_process_memory_bytes_avg" => prom_dress.process_memory_bytes_avg += val, + "parseable_process_memory_bytes_avg" => { + prom_dress.process_memory_bytes_avg += val + } "parseable_storage_size" => { if sample.labels.get("type").expect("type is present") == "staging" { prom_dress.parseable_storage_size.staging += val; From cf0e6fdd75d9da6e270ce1fe6e662a9b2b87835a Mon Sep 17 00:00:00 2001 From: Anant Vindal Date: Mon, 17 Aug 2026 12:25:00 +0530 Subject: [PATCH 4/4] coderabbit suggestion --- src/metrics/mod.rs | 53 ++++++++++++++++++++++++++++++---------------- 1 file changed, 35 insertions(+), 18 deletions(-) diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index 50c336e41..8253c5dca 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -17,6 +17,8 @@ */ pub mod prom_utils; +use std::sync::OnceLock; + use crate::{ handlers::{TelemetryType, http::metrics_path}, stats::FullStats, @@ -25,8 +27,10 @@ use actix_web::Responder; use actix_web_prometheus::{PrometheusMetrics, PrometheusMetricsBuilder}; use error::MetricsError; use once_cell::sync::Lazy; -use prometheus::{Gauge, HistogramOpts, HistogramVec, IntCounterVec, IntGaugeVec, Opts, Registry}; -use std::sync::atomic::{AtomicU64, Ordering}; +use prometheus::{ + Gauge, HistogramOpts, HistogramVec, IntCounterVec, IntGaugeVec, Opts, Registry, + core::{Atomic, AtomicF64}, +}; pub const METRICS_NAMESPACE: &str = env!("CARGO_PKG_NAME"); @@ -197,36 +201,41 @@ pub static PROCESS_MEMORY_BYTES_AVG: Lazy = Lazy::new(|| { ) .expect("metric can be created") }); - -const CPU_USAGE_PRECISION: f64 = 1_000.0; - -#[derive(Default)] +pub static PROCESS_METRICS_INIT: OnceLock<(f64, u64)> = OnceLock::new(); struct ProcessMetricsAccumulator { - cpu_usage_sum: AtomicU64, - memory_bytes_sum: AtomicU64, + cpu_usage_avg: AtomicF64, + memory_bytes_avg: AtomicF64, +} + +impl Default for ProcessMetricsAccumulator { + fn default() -> Self { + // PROCESS_METRICS_INIT must be initialized by now + let (cpu, mem) = *PROCESS_METRICS_INIT.get().unwrap(); + Self { + cpu_usage_avg: AtomicF64::new(cpu), + memory_bytes_avg: AtomicF64::new(mem as f64), + } + } } impl ProcessMetricsAccumulator { fn record(&self, cpu_usage_percent: f64, memory_bytes: u64) -> (f64, f64) { - self.cpu_usage_sum.fetch_add( - (cpu_usage_percent * CPU_USAGE_PRECISION).round() as u64, - Ordering::Relaxed, - ); - self.memory_bytes_sum - .fetch_add(memory_bytes, Ordering::Relaxed); - // Exponentially Weighted Moving Average is better than // a lifetime average // A spike which occurred 5 days ago should not affect the average utilization // for the last minute // α = 1 - exp(-Δt / τ) = 1 - exp(-5/60) ≈ 0.0800 // S_new = S_old + α * (x_new - S_old) - let s_cpu_old = self.cpu_usage_sum.load(Ordering::Relaxed) as f64; + let s_cpu_old = self.cpu_usage_avg.get(); let s_cpu_new = s_cpu_old + 0.08 * (cpu_usage_percent - s_cpu_old); - let s_mem_old = self.memory_bytes_sum.load(Ordering::Relaxed) as f64; + let s_mem_old = self.memory_bytes_avg.get(); let s_mem_new = s_mem_old + 0.08 * (memory_bytes as f64 - s_mem_old); + // update accumulator + self.cpu_usage_avg.set(s_cpu_new); + self.memory_bytes_avg.set(s_mem_new); + (s_cpu_new, s_mem_new) } } @@ -235,6 +244,10 @@ static PROCESS_METRICS_ACCUMULATOR: Lazy = Lazy::new(ProcessMetricsAccumulator::default); pub fn record_process_metrics_sample(cpu_usage_percent: f64, memory_bytes: u64) { + if PROCESS_METRICS_INIT.get().is_none() { + // first measurement + let _ = PROCESS_METRICS_INIT.set((cpu_usage_percent, memory_bytes)); + } let (average_cpu_usage, average_memory_bytes) = PROCESS_METRICS_ACCUMULATOR.record(cpu_usage_percent, memory_bytes); PROCESS_CPU_USAGE_PERCENT_AVG.set(average_cpu_usage); @@ -243,14 +256,18 @@ pub fn record_process_metrics_sample(cpu_usage_percent: f64, memory_bytes: u64) #[cfg(test)] mod process_metrics_tests { + use crate::metrics::PROCESS_METRICS_INIT; + use super::ProcessMetricsAccumulator; #[test] fn averages_process_metric_samples() { + // init PROCESS_METRICS_INIT + PROCESS_METRICS_INIT.get_or_init(|| (10.0, 100)); let accumulator = ProcessMetricsAccumulator::default(); assert_eq!(accumulator.record(10.0, 100), (10.0, 100.0)); - assert_eq!(accumulator.record(20.0, 300), (15.0, 200.0)); + assert_eq!(accumulator.record(20.0, 300), (10.8, 116.0)); } }