From 940cc9262c2c50545c697e649c572b3e4956f74c Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Sat, 3 Oct 2026 14:47:32 +0530 Subject: [PATCH 1/6] feat: add cgroup CPU usage metric --- src/handlers/http/resource_check.rs | 147 +++++++++++++++++++++++++++- src/metrics/mod.rs | 18 ++++ 2 files changed, 161 insertions(+), 4 deletions(-) diff --git a/src/handlers/http/resource_check.rs b/src/handlers/http/resource_check.rs index ddba90ce3..6cb131f72 100644 --- a/src/handlers/http/resource_check.rs +++ b/src/handlers/http/resource_check.rs @@ -18,7 +18,11 @@ #[cfg(target_os = "linux")] use std::path::{Path, PathBuf}; +#[cfg(target_os = "linux")] +use std::sync::Mutex; use std::sync::{Arc, LazyLock, atomic::AtomicBool}; +#[cfg(target_os = "linux")] +use std::time::Instant; use actix_web::{ body::MessageBody, @@ -35,19 +39,36 @@ use tokio::{ use tracing::{info, trace, warn}; use crate::analytics::{SYS_INFO, refresh_sys_info}; -use crate::metrics::{record_disk_metrics, record_process_metrics_sample}; +use crate::metrics::{ + record_disk_metrics, record_process_cpu_usage_cores, record_process_metrics_sample, +}; use crate::parseable::PARSEABLE; const PROCESS_METRICS_SAMPLE_INTERVAL: Duration = Duration::from_secs(5); #[cfg(target_os = "linux")] const CGROUP_V2_CPU_MAX_FILE: &str = "cpu.max"; #[cfg(target_os = "linux")] +const CGROUP_V2_CPU_STAT_FILE: &str = "cpu.stat"; +#[cfg(target_os = "linux")] const CGROUP_V1_CPU_QUOTA_FILE: &str = "cpu.cfs_quota_us"; #[cfg(target_os = "linux")] const CGROUP_V1_CPU_PERIOD_FILE: &str = "cpu.cfs_period_us"; +#[cfg(target_os = "linux")] +const CGROUP_V1_CPU_USAGE_FILE: &str = "cpuacct.usage"; static SERVER_OK: LazyLock> = LazyLock::new(|| Arc::new(AtomicBool::new(true))); +#[cfg(target_os = "linux")] +#[derive(Clone, Copy)] +struct CgroupCpuUsageSample { + usage_micros: u64, + sampled_at: Instant, +} + +#[cfg(target_os = "linux")] +static CGROUP_CPU_USAGE_SAMPLE: LazyLock>> = + LazyLock::new(|| Mutex::new(None)); + #[cfg(target_os = "linux")] fn cpu_quota_cores(quota: &str, period: &str) -> Result, ()> { let quota = quota.trim(); @@ -107,6 +128,82 @@ fn cgroup_cpu_limit_cores() -> Result, ()> { cpu_quota_cores("a, &period) } +#[cfg(target_os = "linux")] +fn cpu_usage_micros_from_stat(cpu_stat: &str) -> Result { + for line in cpu_stat.lines() { + let mut fields = line.split_whitespace(); + if fields.next() == Some("usage_usec") { + return fields.next().ok_or(())?.parse::().map_err(|_| ()); + } + } + Err(()) +} + +#[cfg(target_os = "linux")] +fn cgroup_cpu_usage_micros() -> Result { + let process = procfs::process::Process::myself().map_err(|_| ())?; + let cgroups = process.cgroups().map_err(|_| ())?.0; + let mounts = process.mountinfo().map_err(|_| ())?.0; + + if let (Some(cgroup), Some(mount)) = ( + cgroups.iter().find(|group| group.hierarchy == 0), + mounts.iter().find(|mount| mount.fs_type == "cgroup2"), + ) { + let directory = + cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; + let cpu_stat = + std::fs::read_to_string(directory.join(CGROUP_V2_CPU_STAT_FILE)).map_err(|_| ())?; + return cpu_usage_micros_from_stat(&cpu_stat); + } + + let cgroup = cgroups + .iter() + .find(|group| group.controllers.iter().any(|item| item == "cpuacct")) + .ok_or(())?; + let mount = mounts + .iter() + .find(|mount| mount.fs_type == "cgroup" && mount.super_options.contains_key("cpuacct")) + .ok_or(())?; + let directory = cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; + let usage_nanos = std::fs::read_to_string(directory.join(CGROUP_V1_CPU_USAGE_FILE)) + .map_err(|_| ())? + .trim() + .parse::() + .map_err(|_| ())?; + Ok(usage_nanos / 1_000) +} + +#[cfg(target_os = "linux")] +fn cpu_usage_cores_between( + previous_usage_micros: u64, + current_usage_micros: u64, + elapsed_micros: u128, +) -> Option { + let usage_delta = current_usage_micros.checked_sub(previous_usage_micros)?; + (elapsed_micros > 0).then_some(usage_delta as f64 / elapsed_micros as f64) +} + +#[cfg(target_os = "linux")] +fn cgroup_cpu_usage_cores() -> Result, ()> { + let current = CgroupCpuUsageSample { + usage_micros: cgroup_cpu_usage_micros()?, + sampled_at: Instant::now(), + }; + let mut previous = CGROUP_CPU_USAGE_SAMPLE.lock().map_err(|_| ())?; + let usage_cores = previous.and_then(|previous| { + cpu_usage_cores_between( + previous.usage_micros, + current.usage_micros, + current + .sampled_at + .duration_since(previous.sampled_at) + .as_micros(), + ) + }); + *previous = Some(current); + Ok(usage_cores) +} + #[cfg(not(target_os = "linux"))] fn cgroup_cpu_limit_cores() -> Result, ()> { Err(()) @@ -120,6 +217,18 @@ pub fn cpu_limit_cores() -> f64 { } } +fn cpu_usage_cores() -> f64 { + #[cfg(target_os = "linux")] + { + return cgroup_cpu_usage_cores().ok().flatten().unwrap_or_default(); + } + + #[cfg(not(target_os = "linux"))] + { + 0.0 + } +} + async fn sample_process_metrics() { refresh_sys_info(); let process_metrics = tokio::task::spawn_blocking(|| { @@ -137,13 +246,22 @@ async fn sample_process_metrics() { }; process_metrics.map(|(cpu_usage, memory_bytes, total_mem)| { - (cpu_usage, memory_bytes, total_mem, cpu_limit_cores()) + ( + cpu_usage, + memory_bytes, + total_mem, + cpu_limit_cores(), + cpu_usage_cores(), + ) }) }) .await .unwrap(); - if let Some((cpu_usage, memory_bytes, total_mem, cpu_limit_cores)) = process_metrics { + if let Some((cpu_usage, memory_bytes, total_mem, cpu_limit_cores, cpu_usage_cores)) = + process_metrics + { record_process_metrics_sample(cpu_usage, memory_bytes, total_mem, cpu_limit_cores); + record_process_cpu_usage_cores(cpu_usage_cores); } let staging_path = PARSEABLE.options.staging_dir().clone(); @@ -254,7 +372,7 @@ pub async fn check_resource_utilization_middleware( #[cfg(all(test, target_os = "linux"))] mod tests { - use super::cpu_quota_cores; + use super::{cpu_quota_cores, cpu_usage_cores_between, cpu_usage_micros_from_stat}; #[test] fn parses_limited_and_unlimited_cpu_quotas() { @@ -268,4 +386,25 @@ mod tests { assert_eq!(cpu_quota_cores("invalid", "100000"), Err(())); assert_eq!(cpu_quota_cores("50000", "0"), Err(())); } + + #[test] + fn parses_cgroup_v2_cpu_usage() { + assert_eq!( + cpu_usage_micros_from_stat("usage_usec 10200000\nuser_usec 8000000\n"), + Ok(10_200_000) + ); + assert_eq!(cpu_usage_micros_from_stat("user_usec 8000000\n"), Err(())); + } + + #[test] + fn calculates_cpu_usage_in_cores() { + assert_eq!( + cpu_usage_cores_between(10_000_000, 10_200_000, 5_000_000), + Some(0.04) + ); + assert_eq!( + cpu_usage_cores_between(10_200_000, 10_000_000, 5_000_000), + None + ); + } } diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index e2848c33f..1228b4c81 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -243,6 +243,17 @@ pub static PROCESS_CPU_USAGE_PERCENT_AVG: Lazy = Lazy::new(|| { .expect("metric can be created") }); +pub static PROCESS_CPU_USAGE_CORES: Lazy = Lazy::new(|| { + Gauge::with_opts( + Opts::new( + "process_cpu_usage_cores", + "CPU cores used by the Parseable cgroup, or zero when cgroup usage is unavailable", + ) + .namespace(METRICS_NAMESPACE), + ) + .expect("metric can be created") +}); + pub static PROCESS_CPU_LIMIT_CORES: Lazy = Lazy::new(|| { Gauge::with_opts( Opts::new( @@ -326,6 +337,10 @@ impl ProcessMetricsAccumulator { pub static PROCESS_METRICS_ACCUMULATOR: Lazy = Lazy::new(ProcessMetricsAccumulator::default); +pub fn record_process_cpu_usage_cores(cpu_usage_cores: f64) { + PROCESS_CPU_USAGE_CORES.set(cpu_usage_cores); +} + pub fn record_process_metrics_sample( cpu_usage_percent: f64, memory_bytes: u64, @@ -904,6 +919,9 @@ fn custom_metrics(registry: &Registry) { registry .register(Box::new(PROCESS_CPU_USAGE_PERCENT_AVG.clone())) .expect("metric can be registered"); + registry + .register(Box::new(PROCESS_CPU_USAGE_CORES.clone())) + .expect("metric can be registered"); registry .register(Box::new(PROCESS_CPU_LIMIT_CORES.clone())) .expect("metric can be registered"); From 8481c35453c867db7e412fcfd7d22ed56e83e00c Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Sat, 3 Oct 2026 15:28:26 +0530 Subject: [PATCH 2/6] fix: address Linux clippy warning --- src/handlers/http/resource_check.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/handlers/http/resource_check.rs b/src/handlers/http/resource_check.rs index 6cb131f72..f10765851 100644 --- a/src/handlers/http/resource_check.rs +++ b/src/handlers/http/resource_check.rs @@ -220,7 +220,7 @@ pub fn cpu_limit_cores() -> f64 { fn cpu_usage_cores() -> f64 { #[cfg(target_os = "linux")] { - return cgroup_cpu_usage_cores().ok().flatten().unwrap_or_default(); + cgroup_cpu_usage_cores().ok().flatten().unwrap_or_default() } #[cfg(not(target_os = "linux"))] From a17470be1b0b614249023e27c4b9814b51cd2056 Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Sat, 3 Oct 2026 15:36:42 +0530 Subject: [PATCH 3/6] fix: fall back to cgroup v1 metrics --- src/handlers/http/resource_check.rs | 32 +++++++++++++++++++---------- 1 file changed, 21 insertions(+), 11 deletions(-) diff --git a/src/handlers/http/resource_check.rs b/src/handlers/http/resource_check.rs index f10765851..e5168d2a0 100644 --- a/src/handlers/http/resource_check.rs +++ b/src/handlers/http/resource_check.rs @@ -104,12 +104,17 @@ fn cgroup_cpu_limit_cores() -> Result, ()> { cgroups.iter().find(|group| group.hierarchy == 0), mounts.iter().find(|mount| mount.fs_type == "cgroup2"), ) { - let directory = - cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; - let cpu_max = - std::fs::read_to_string(directory.join(CGROUP_V2_CPU_MAX_FILE)).map_err(|_| ())?; - let mut values = cpu_max.split_whitespace(); - return cpu_quota_cores(values.next().ok_or(())?, values.next().ok_or(())?); + let v2_limit = (|| { + let directory = + cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; + let cpu_max = + std::fs::read_to_string(directory.join(CGROUP_V2_CPU_MAX_FILE)).map_err(|_| ())?; + let mut values = cpu_max.split_whitespace(); + cpu_quota_cores(values.next().ok_or(())?, values.next().ok_or(())?) + })(); + if let Ok(limit) = v2_limit { + return Ok(limit); + } } let cgroup = cgroups @@ -149,11 +154,16 @@ fn cgroup_cpu_usage_micros() -> Result { cgroups.iter().find(|group| group.hierarchy == 0), mounts.iter().find(|mount| mount.fs_type == "cgroup2"), ) { - let directory = - cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; - let cpu_stat = - std::fs::read_to_string(directory.join(CGROUP_V2_CPU_STAT_FILE)).map_err(|_| ())?; - return cpu_usage_micros_from_stat(&cpu_stat); + let v2_usage = (|| { + let directory = + cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; + let cpu_stat = + std::fs::read_to_string(directory.join(CGROUP_V2_CPU_STAT_FILE)).map_err(|_| ())?; + cpu_usage_micros_from_stat(&cpu_stat) + })(); + if let Ok(usage_micros) = v2_usage { + return Ok(usage_micros); + } } let cgroup = cgroups From a70fba7fe6b8e0548283028dd0ab0243c986da47 Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Sat, 3 Oct 2026 16:21:14 +0530 Subject: [PATCH 4/6] fix: use consistent cgroup CPU scope --- src/handlers/http/resource_check.rs | 181 ++++++++++++++++------------ 1 file changed, 107 insertions(+), 74 deletions(-) diff --git a/src/handlers/http/resource_check.rs b/src/handlers/http/resource_check.rs index e5168d2a0..e523d1c1e 100644 --- a/src/handlers/http/resource_check.rs +++ b/src/handlers/http/resource_check.rs @@ -65,6 +65,12 @@ struct CgroupCpuUsageSample { sampled_at: Instant, } +#[cfg(target_os = "linux")] +struct CgroupCpuMetrics { + limit_cores: Option, + usage_micros: u64, +} + #[cfg(target_os = "linux")] static CGROUP_CPU_USAGE_SAMPLE: LazyLock>> = LazyLock::new(|| Mutex::new(None)); @@ -94,45 +100,6 @@ fn cgroup_directory(pathname: &str, root: &str, mount_point: &Path) -> Option Result, ()> { - let process = procfs::process::Process::myself().map_err(|_| ())?; - let cgroups = process.cgroups().map_err(|_| ())?.0; - let mounts = process.mountinfo().map_err(|_| ())?.0; - - if let (Some(cgroup), Some(mount)) = ( - cgroups.iter().find(|group| group.hierarchy == 0), - mounts.iter().find(|mount| mount.fs_type == "cgroup2"), - ) { - let v2_limit = (|| { - let directory = - cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; - let cpu_max = - std::fs::read_to_string(directory.join(CGROUP_V2_CPU_MAX_FILE)).map_err(|_| ())?; - let mut values = cpu_max.split_whitespace(); - cpu_quota_cores(values.next().ok_or(())?, values.next().ok_or(())?) - })(); - if let Ok(limit) = v2_limit { - return Ok(limit); - } - } - - let cgroup = cgroups - .iter() - .find(|group| group.controllers.iter().any(|item| item == "cpu")) - .ok_or(())?; - let mount = mounts - .iter() - .find(|mount| mount.fs_type == "cgroup" && mount.super_options.contains_key("cpu")) - .ok_or(())?; - let directory = cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; - let quota = - std::fs::read_to_string(directory.join(CGROUP_V1_CPU_QUOTA_FILE)).map_err(|_| ())?; - let period = - std::fs::read_to_string(directory.join(CGROUP_V1_CPU_PERIOD_FILE)).map_err(|_| ())?; - cpu_quota_cores("a, &period) -} - #[cfg(target_os = "linux")] fn cpu_usage_micros_from_stat(cpu_stat: &str) -> Result { for line in cpu_stat.lines() { @@ -145,7 +112,43 @@ fn cpu_usage_micros_from_stat(cpu_stat: &str) -> Result { } #[cfg(target_os = "linux")] -fn cgroup_cpu_usage_micros() -> Result { +fn read_cgroup_v2_cpu_metrics(directory: &Path) -> Result { + let cpu_max = + std::fs::read_to_string(directory.join(CGROUP_V2_CPU_MAX_FILE)).map_err(|_| ())?; + let mut values = cpu_max.split_whitespace(); + let limit_cores = cpu_quota_cores(values.next().ok_or(())?, values.next().ok_or(())?)?; + let cpu_stat = + std::fs::read_to_string(directory.join(CGROUP_V2_CPU_STAT_FILE)).map_err(|_| ())?; + + Ok(CgroupCpuMetrics { + limit_cores, + usage_micros: cpu_usage_micros_from_stat(&cpu_stat)?, + }) +} + +#[cfg(target_os = "linux")] +fn read_cgroup_v1_cpu_metrics( + limit_directory: &Path, + usage_directory: &Path, +) -> Result { + let quota = + std::fs::read_to_string(limit_directory.join(CGROUP_V1_CPU_QUOTA_FILE)).map_err(|_| ())?; + let period = + std::fs::read_to_string(limit_directory.join(CGROUP_V1_CPU_PERIOD_FILE)).map_err(|_| ())?; + let usage_nanos = std::fs::read_to_string(usage_directory.join(CGROUP_V1_CPU_USAGE_FILE)) + .map_err(|_| ())? + .trim() + .parse::() + .map_err(|_| ())?; + + Ok(CgroupCpuMetrics { + limit_cores: cpu_quota_cores("a, &period)?, + usage_micros: usage_nanos / 1_000, + }) +} + +#[cfg(target_os = "linux")] +fn cgroup_cpu_metrics() -> Result { let process = procfs::process::Process::myself().map_err(|_| ())?; let cgroups = process.cgroups().map_err(|_| ())?.0; let mounts = process.mountinfo().map_err(|_| ())?.0; @@ -154,33 +157,47 @@ fn cgroup_cpu_usage_micros() -> Result { cgroups.iter().find(|group| group.hierarchy == 0), mounts.iter().find(|mount| mount.fs_type == "cgroup2"), ) { - let v2_usage = (|| { - let directory = - cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; - let cpu_stat = - std::fs::read_to_string(directory.join(CGROUP_V2_CPU_STAT_FILE)).map_err(|_| ())?; - cpu_usage_micros_from_stat(&cpu_stat) - })(); - if let Ok(usage_micros) = v2_usage { - return Ok(usage_micros); + if let Some(directory) = cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point) + && let Ok(metrics) = read_cgroup_v2_cpu_metrics(&directory) + { + return Ok(metrics); } } - let cgroup = cgroups + let limit_cgroup = cgroups + .iter() + .find(|group| group.controllers.iter().any(|item| item == "cpu")) + .ok_or(())?; + let usage_cgroup = cgroups .iter() .find(|group| group.controllers.iter().any(|item| item == "cpuacct")) .ok_or(())?; - let mount = mounts + if limit_cgroup.pathname != usage_cgroup.pathname { + return Err(()); + } + + let limit_mount = mounts + .iter() + .find(|mount| mount.fs_type == "cgroup" && mount.super_options.contains_key("cpu")) + .ok_or(())?; + let usage_mount = mounts .iter() .find(|mount| mount.fs_type == "cgroup" && mount.super_options.contains_key("cpuacct")) .ok_or(())?; - let directory = cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point).ok_or(())?; - let usage_nanos = std::fs::read_to_string(directory.join(CGROUP_V1_CPU_USAGE_FILE)) - .map_err(|_| ())? - .trim() - .parse::() - .map_err(|_| ())?; - Ok(usage_nanos / 1_000) + let limit_directory = cgroup_directory( + &limit_cgroup.pathname, + &limit_mount.root, + &limit_mount.mount_point, + ) + .ok_or(())?; + let usage_directory = cgroup_directory( + &usage_cgroup.pathname, + &usage_mount.root, + &usage_mount.mount_point, + ) + .ok_or(())?; + + read_cgroup_v1_cpu_metrics(&limit_directory, &usage_directory) } #[cfg(target_os = "linux")] @@ -194,9 +211,9 @@ fn cpu_usage_cores_between( } #[cfg(target_os = "linux")] -fn cgroup_cpu_usage_cores() -> Result, ()> { +fn cgroup_cpu_usage_cores(current_usage_micros: u64) -> Result, ()> { let current = CgroupCpuUsageSample { - usage_micros: cgroup_cpu_usage_micros()?, + usage_micros: current_usage_micros, sampled_at: Instant::now(), }; let mut previous = CGROUP_CPU_USAGE_SAMPLE.lock().map_err(|_| ())?; @@ -214,28 +231,43 @@ fn cgroup_cpu_usage_cores() -> Result, ()> { Ok(usage_cores) } -#[cfg(not(target_os = "linux"))] -fn cgroup_cpu_limit_cores() -> Result, ()> { - Err(()) -} - pub fn cpu_limit_cores() -> f64 { - match cgroup_cpu_limit_cores() { - Ok(Some(limit)) => limit, - Ok(None) => num_cpus::get() as f64, - Err(()) => 0.0, + #[cfg(target_os = "linux")] + { + match cgroup_cpu_metrics() { + Ok(metrics) => metrics + .limit_cores + .unwrap_or_else(|| num_cpus::get() as f64), + Err(()) => 0.0, + } + } + + #[cfg(not(target_os = "linux"))] + { + 0.0 } } -fn cpu_usage_cores() -> f64 { +fn cpu_metrics() -> (f64, f64) { #[cfg(target_os = "linux")] { - cgroup_cpu_usage_cores().ok().flatten().unwrap_or_default() + match cgroup_cpu_metrics() { + Ok(metrics) => ( + metrics + .limit_cores + .unwrap_or_else(|| num_cpus::get() as f64), + cgroup_cpu_usage_cores(metrics.usage_micros) + .ok() + .flatten() + .unwrap_or_default(), + ), + Err(()) => (0.0, 0.0), + } } #[cfg(not(target_os = "linux"))] { - 0.0 + (0.0, 0.0) } } @@ -256,12 +288,13 @@ async fn sample_process_metrics() { }; process_metrics.map(|(cpu_usage, memory_bytes, total_mem)| { + let (cpu_limit_cores, cpu_usage_cores) = cpu_metrics(); ( cpu_usage, memory_bytes, total_mem, - cpu_limit_cores(), - cpu_usage_cores(), + cpu_limit_cores, + cpu_usage_cores, ) }) }) From 4a28ad87435382c0aa7909c0bbadafcaa333377f Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Sat, 3 Oct 2026 17:35:01 +0530 Subject: [PATCH 5/6] fix: handle cgroup v2 root metrics --- src/handlers/http/resource_check.rs | 117 ++++++++++++++++++---------- 1 file changed, 74 insertions(+), 43 deletions(-) diff --git a/src/handlers/http/resource_check.rs b/src/handlers/http/resource_check.rs index e523d1c1e..4626c7874 100644 --- a/src/handlers/http/resource_check.rs +++ b/src/handlers/http/resource_check.rs @@ -71,6 +71,12 @@ struct CgroupCpuMetrics { usage_micros: u64, } +#[cfg(target_os = "linux")] +enum CgroupV2CpuMetrics { + Complete(CgroupCpuMetrics), + RootUsage(u64), +} + #[cfg(target_os = "linux")] static CGROUP_CPU_USAGE_SAMPLE: LazyLock>> = LazyLock::new(|| Mutex::new(None)); @@ -112,18 +118,25 @@ fn cpu_usage_micros_from_stat(cpu_stat: &str) -> Result { } #[cfg(target_os = "linux")] -fn read_cgroup_v2_cpu_metrics(directory: &Path) -> Result { - let cpu_max = - std::fs::read_to_string(directory.join(CGROUP_V2_CPU_MAX_FILE)).map_err(|_| ())?; - let mut values = cpu_max.split_whitespace(); - let limit_cores = cpu_quota_cores(values.next().ok_or(())?, values.next().ok_or(())?)?; +fn read_cgroup_v2_cpu_metrics(directory: &Path) -> Result { let cpu_stat = std::fs::read_to_string(directory.join(CGROUP_V2_CPU_STAT_FILE)).map_err(|_| ())?; + let usage_micros = cpu_usage_micros_from_stat(&cpu_stat)?; - Ok(CgroupCpuMetrics { + let cpu_max = match std::fs::read_to_string(directory.join(CGROUP_V2_CPU_MAX_FILE)) { + Ok(cpu_max) => cpu_max, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => { + return Ok(CgroupV2CpuMetrics::RootUsage(usage_micros)); + } + Err(_) => return Err(()), + }; + let mut values = cpu_max.split_whitespace(); + let limit_cores = cpu_quota_cores(values.next().ok_or(())?, values.next().ok_or(())?)?; + + Ok(CgroupV2CpuMetrics::Complete(CgroupCpuMetrics { limit_cores, - usage_micros: cpu_usage_micros_from_stat(&cpu_stat)?, - }) + usage_micros, + })) } #[cfg(target_os = "linux")] @@ -152,52 +165,70 @@ fn cgroup_cpu_metrics() -> Result { let process = procfs::process::Process::myself().map_err(|_| ())?; let cgroups = process.cgroups().map_err(|_| ())?.0; let mounts = process.mountinfo().map_err(|_| ())?.0; + let mut v2_root_usage = None; if let (Some(cgroup), Some(mount)) = ( cgroups.iter().find(|group| group.hierarchy == 0), mounts.iter().find(|mount| mount.fs_type == "cgroup2"), - ) { - if let Some(directory) = cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point) - && let Ok(metrics) = read_cgroup_v2_cpu_metrics(&directory) - { - return Ok(metrics); + ) && let Some(directory) = + cgroup_directory(&cgroup.pathname, &mount.root, &mount.mount_point) + { + match read_cgroup_v2_cpu_metrics(&directory) { + Ok(CgroupV2CpuMetrics::Complete(metrics)) => return Ok(metrics), + Ok(CgroupV2CpuMetrics::RootUsage(usage_micros)) => { + v2_root_usage = Some(usage_micros); + } + Err(()) => {} } } - let limit_cgroup = cgroups - .iter() - .find(|group| group.controllers.iter().any(|item| item == "cpu")) + let v1_metrics = (|| { + let limit_cgroup = cgroups + .iter() + .find(|group| group.controllers.iter().any(|item| item == "cpu")) + .ok_or(())?; + let usage_cgroup = cgroups + .iter() + .find(|group| group.controllers.iter().any(|item| item == "cpuacct")) + .ok_or(())?; + if limit_cgroup.pathname != usage_cgroup.pathname { + return Err(()); + } + + let limit_mount = mounts + .iter() + .find(|mount| mount.fs_type == "cgroup" && mount.super_options.contains_key("cpu")) + .ok_or(())?; + let usage_mount = mounts + .iter() + .find(|mount| mount.fs_type == "cgroup" && mount.super_options.contains_key("cpuacct")) + .ok_or(())?; + let limit_directory = cgroup_directory( + &limit_cgroup.pathname, + &limit_mount.root, + &limit_mount.mount_point, + ) .ok_or(())?; - let usage_cgroup = cgroups - .iter() - .find(|group| group.controllers.iter().any(|item| item == "cpuacct")) + let usage_directory = cgroup_directory( + &usage_cgroup.pathname, + &usage_mount.root, + &usage_mount.mount_point, + ) .ok_or(())?; - if limit_cgroup.pathname != usage_cgroup.pathname { - return Err(()); + + read_cgroup_v1_cpu_metrics(&limit_directory, &usage_directory) + })(); + + if let Ok(metrics) = v1_metrics { + return Ok(metrics); } - let limit_mount = mounts - .iter() - .find(|mount| mount.fs_type == "cgroup" && mount.super_options.contains_key("cpu")) - .ok_or(())?; - let usage_mount = mounts - .iter() - .find(|mount| mount.fs_type == "cgroup" && mount.super_options.contains_key("cpuacct")) - .ok_or(())?; - let limit_directory = cgroup_directory( - &limit_cgroup.pathname, - &limit_mount.root, - &limit_mount.mount_point, - ) - .ok_or(())?; - let usage_directory = cgroup_directory( - &usage_cgroup.pathname, - &usage_mount.root, - &usage_mount.mount_point, - ) - .ok_or(())?; - - read_cgroup_v1_cpu_metrics(&limit_directory, &usage_directory) + v2_root_usage + .map(|usage_micros| CgroupCpuMetrics { + limit_cores: None, + usage_micros, + }) + .ok_or(()) } #[cfg(target_os = "linux")] From cc10985c98ff3d22d5b51f63be8a284da51d8876 Mon Sep 17 00:00:00 2001 From: Pratik Jadhav Date: Mon, 5 Oct 2026 11:21:15 +0530 Subject: [PATCH 6/6] add avg for last min for usage --- src/metrics/mod.rs | 43 +++++++++++++++++++++++++++++++++++++++---- 1 file changed, 39 insertions(+), 4 deletions(-) diff --git a/src/metrics/mod.rs b/src/metrics/mod.rs index 1228b4c81..1cfae8e1d 100644 --- a/src/metrics/mod.rs +++ b/src/metrics/mod.rs @@ -17,7 +17,11 @@ */ pub mod prom_utils; -use std::{path::Path, sync::OnceLock}; +use std::{ + collections::VecDeque, + path::Path, + sync::{Mutex, OnceLock}, +}; use crate::{ handlers::{TelemetryType, http::metrics_path}, @@ -247,7 +251,7 @@ pub static PROCESS_CPU_USAGE_CORES: Lazy = Lazy::new(|| { Gauge::with_opts( Opts::new( "process_cpu_usage_cores", - "CPU cores used by the Parseable cgroup, or zero when cgroup usage is unavailable", + "Average CPU cores used by the Parseable cgroup over the last minute", ) .namespace(METRICS_NAMESPACE), ) @@ -337,8 +341,23 @@ impl ProcessMetricsAccumulator { pub static PROCESS_METRICS_ACCUMULATOR: Lazy = Lazy::new(ProcessMetricsAccumulator::default); +const CPU_USAGE_CORES_WINDOW_SAMPLES: usize = 12; +static CPU_USAGE_CORES_SAMPLES: Lazy>> = + Lazy::new(|| Mutex::new(VecDeque::with_capacity(CPU_USAGE_CORES_WINDOW_SAMPLES))); + +fn update_cpu_usage_cores_average(samples: &mut VecDeque, sample: f64) -> f64 { + if samples.len() == CPU_USAGE_CORES_WINDOW_SAMPLES { + samples.pop_front(); + } + samples.push_back(sample); + + samples.iter().sum::() / samples.len() as f64 +} + pub fn record_process_cpu_usage_cores(cpu_usage_cores: f64) { - PROCESS_CPU_USAGE_CORES.set(cpu_usage_cores); + let mut samples = CPU_USAGE_CORES_SAMPLES.lock().unwrap(); + let average = update_cpu_usage_cores_average(&mut samples, cpu_usage_cores); + PROCESS_CPU_USAGE_CORES.set(average); } pub fn record_process_metrics_sample( @@ -363,7 +382,9 @@ pub fn record_process_metrics_sample( mod process_metrics_tests { use crate::metrics::PROCESS_METRICS_INIT; - use super::ProcessMetricsAccumulator; + use super::{ + CPU_USAGE_CORES_WINDOW_SAMPLES, ProcessMetricsAccumulator, update_cpu_usage_cores_average, + }; #[test] fn averages_process_metric_samples() { @@ -374,6 +395,20 @@ mod process_metrics_tests { assert_eq!(accumulator.record(10.0, 100), (10.0, 100.0)); assert_eq!(accumulator.record(20.0, 300), (10.8, 116.0)); } + + #[test] + fn cpu_usage_cores_window_is_one_minute_of_samples() { + let mut samples = std::collections::VecDeque::new(); + let mut average = 0.0; + for sample in 1..=CPU_USAGE_CORES_WINDOW_SAMPLES + 1 { + average = update_cpu_usage_cores_average(&mut samples, sample as f64); + } + + assert_eq!(samples.len(), CPU_USAGE_CORES_WINDOW_SAMPLES); + assert_eq!(samples.front(), Some(&2.0)); + assert_eq!(samples.back(), Some(&13.0)); + assert_eq!(average, 7.5); + } } pub static QUERY_EXECUTE_TIME: Lazy = Lazy::new(|| {