Skip to content
Open
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
275 changes: 244 additions & 31 deletions src/handlers/http/resource_check.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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<Arc<AtomicBool>> = 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<f64>,
usage_micros: u64,
}

#[cfg(target_os = "linux")]
enum CgroupV2CpuMetrics {
Complete(CgroupCpuMetrics),
RootUsage(u64),
}

#[cfg(target_os = "linux")]
static CGROUP_CPU_USAGE_SAMPLE: LazyLock<Mutex<Option<CgroupCpuUsageSample>>> =
LazyLock::new(|| Mutex::new(None));

#[cfg(target_os = "linux")]
fn cpu_quota_cores(quota: &str, period: &str) -> Result<Option<f64>, ()> {
let quota = quota.trim();
Expand All @@ -74,49 +107,198 @@ fn cgroup_directory(pathname: &str, root: &str, mount_point: &Path) -> Option<Pa
}

#[cfg(target_os = "linux")]
fn cgroup_cpu_limit_cores() -> Result<Option<f64>, ()> {
fn cpu_usage_micros_from_stat(cpu_stat: &str) -> Result<u64, ()> {
for line in cpu_stat.lines() {
let mut fields = line.split_whitespace();
if fields.next() == Some("usage_usec") {
return fields.next().ok_or(())?.parse::<u64>().map_err(|_| ());
}
}
Err(())
}

#[cfg(target_os = "linux")]
fn read_cgroup_v2_cpu_metrics(directory: &Path) -> Result<CgroupV2CpuMetrics, ()> {
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<CgroupCpuMetrics, ()> {
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::<u64>()
.map_err(|_| ())?;

Ok(CgroupCpuMetrics {
limit_cores: cpu_quota_cores(&quota, &period)?,
usage_micros: usage_nanos / 1_000,
})
}

#[cfg(target_os = "linux")]
fn cgroup_cpu_metrics() -> Result<CgroupCpuMetrics, ()> {
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(&quota, &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<Option<f64>, ()> {
Err(())
#[cfg(target_os = "linux")]
Comment thread
coderabbitai[bot] marked this conversation as resolved.
fn cpu_usage_cores_between(
previous_usage_micros: u64,
current_usage_micros: u64,
elapsed_micros: u128,
) -> Option<f64> {
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<Option<f64>, ()> {
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)
}
}

Expand All @@ -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();
Expand Down Expand Up @@ -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() {
Expand All @@ -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
);
}
}
Loading
Loading