diff --git a/src/handlers/http/resource_check.rs b/src/handlers/http/resource_check.rs index ddba90ce3..4626c7874 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,48 @@ 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")] +struct CgroupCpuMetrics { + limit_cores: Option, + 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)); + #[cfg(target_os = "linux")] fn cpu_quota_cores(quota: &str, period: &str) -> Result, ()> { let quota = quota.trim(); @@ -74,49 +107,198 @@ fn cgroup_directory(pathname: &str, root: &str, mount_point: &Path) -> Option Result, ()> { +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 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)?; + + 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, + })) +} + +#[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; + 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"), - ) { - 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 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 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 mount = mounts - .iter() - .find(|mount| mount.fs_type == "cgroup" && mount.super_options.contains_key("cpu")) + let usage_directory = cgroup_directory( + &usage_cgroup.pathname, + &usage_mount.root, + &usage_mount.mount_point, + ) .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) + + read_cgroup_v1_cpu_metrics(&limit_directory, &usage_directory) + })(); + + if let Ok(metrics) = v1_metrics { + return Ok(metrics); + } + + v2_root_usage + .map(|usage_micros| CgroupCpuMetrics { + limit_cores: None, + usage_micros, + }) + .ok_or(()) } -#[cfg(not(target_os = "linux"))] -fn cgroup_cpu_limit_cores() -> Result, ()> { - Err(()) +#[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(current_usage_micros: u64) -> Result, ()> { + let current = CgroupCpuUsageSample { + usage_micros: current_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) } 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_metrics() -> (f64, f64) { + #[cfg(target_os = "linux")] + { + 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) } } @@ -137,13 +319,23 @@ async fn sample_process_metrics() { }; process_metrics.map(|(cpu_usage, memory_bytes, total_mem)| { - (cpu_usage, memory_bytes, total_mem, cpu_limit_cores()) + let (cpu_limit_cores, cpu_usage_cores) = cpu_metrics(); + ( + 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 +446,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 +460,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..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}, @@ -243,6 +247,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", + "Average CPU cores used by the Parseable cgroup over the last minute", + ) + .namespace(METRICS_NAMESPACE), + ) + .expect("metric can be created") +}); + pub static PROCESS_CPU_LIMIT_CORES: Lazy = Lazy::new(|| { Gauge::with_opts( Opts::new( @@ -326,6 +341,25 @@ 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) { + 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( cpu_usage_percent: f64, memory_bytes: u64, @@ -348,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() { @@ -359,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(|| { @@ -904,6 +954,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");