mod cgroup;
mod gpu;
use serde_json::{Value, json};
use std::sync::Arc;
use std::time::{Duration, Instant};
use sysinfo::{DiskRefreshKind, Disks, ProcessRefreshKind, ProcessesToUpdate, System};
use tokio::sync::Mutex;
const DEADLINE: Duration = Duration::from_secs(2);
#[derive(Default)]
pub(super) struct Resources {
sampler: Arc<Mutex<Sampler>>,
}
#[derive(Default)]
struct Sampler {
system: System,
disks: Disks,
previous: Option<Instant>,
}
impl Resources {
pub(super) async fn sample(&self) -> Value {
let Ok(mut sampler) = Arc::clone(&self.sampler).try_lock_owned() else {
return unavailable("busy");
};
let worker = tokio::task::spawn_blocking(move || {
let value = sampler.sample();
(sampler, value)
});
let (measured, gpu) = tokio::join!(tokio::time::timeout(DEADLINE, worker), gpu::sample());
let mut measured = match measured {
Ok(Ok((_sampler, value))) => value,
Ok(Err(_)) => unavailable("measurement_failed"),
Err(_) => unavailable("timeout"),
};
if gpu["status"] != "unavailable" || measured["gpu"].is_null() {
measured["gpu"] = gpu;
}
measured
}
}
impl Sampler {
fn sample(&mut self) -> Value {
let now = Instant::now();
let interval = self.previous.map(|previous| now.duration_since(previous));
let ready = interval.is_some_and(|elapsed| elapsed >= sysinfo::MINIMUM_CPU_UPDATE_INTERVAL);
if interval.is_none() || ready {
self.system.refresh_cpu_usage();
self.previous = Some(now);
}
self.system.refresh_memory();
let pid = sysinfo::get_current_pid().ok();
let updated = if let Some(pid) = pid {
if interval.is_none() || ready {
self.system.refresh_processes_specifics(
ProcessesToUpdate::Some(&[pid]),
true,
ProcessRefreshKind::nothing()
.without_tasks()
.with_memory()
.with_disk_usage()
.with_cpu(),
) > 0
} else {
self.system.process(pid).is_some()
}
} else {
false
};
self.disks
.refresh_specifics(true, DiskRefreshKind::nothing().with_io_usage());
let gateway = pid
.filter(|_| updated)
.and_then(|pid| self.system.process(pid))
.map(|process| {
let io = process.disk_usage();
json!({
"cpu_percent": ready.then(|| process.cpu_usage()),
"cpu_time_ms": process.accumulated_cpu_time(),
"resident_memory_bytes": process.memory(),
"read_bytes_total": io.total_read_bytes,
"written_bytes_total": io.total_written_bytes,
})
});
let mut seen = std::collections::BTreeSet::new();
let disks: Vec<_> = self.disks.iter().filter(|disk| seen.insert(disk.name())).take(64).map(|disk| {
let io = disk.usage();
json!({"device": disk.name().to_string_lossy(),
"read_bytes_total": io.total_read_bytes, "written_bytes_total": io.total_written_bytes})
}).collect();
json!({
"version": 1,
"status": "ok",
"measured_at_ms": chrono::Utc::now().timestamp_millis(),
"sample_interval_ms": interval.map(|elapsed| elapsed.as_millis()),
"host": {
"cpu_percent": ready.then(|| self.system.global_cpu_usage()),
"logical_cpus": self.system.cpus().len(),
"memory_total_bytes": self.system.total_memory(),
"memory_available_bytes": self.system.available_memory(),
"swap_total_bytes": self.system.total_swap(),
"swap_used_bytes": self.system.used_swap(),
"uptime_seconds": System::uptime(),
"disks": disks,
},
"gateway": gateway,
"container": cgroup::sample(),
"gpu": gpu::fallback(),
})
}
}
fn unavailable(reason: &str) -> Value {
json!({"status": "unavailable", "reason": reason})
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn resources_warm_up_and_bound_concurrent_measurements() {
let resources = Resources::default();
let held = resources.sampler.lock().await;
assert_eq!(resources.sample().await["reason"], "busy");
drop(held);
let first = resources.sample().await;
assert_eq!(first["status"], "ok");
assert!(first["host"]["cpu_percent"].is_null());
assert!(first["gateway"]["resident_memory_bytes"].as_u64().unwrap() > 0);
tokio::time::sleep(sysinfo::MINIMUM_CPU_UPDATE_INTERVAL).await;
let second = resources.sample().await;
assert!(second["host"]["cpu_percent"].as_f64().is_some());
assert!(second["host"]["memory_total_bytes"].as_u64().unwrap() > 0);
assert!(second["gpu"]["status"].is_string());
}
#[test]
fn too_close_samples_preserve_process_counters_and_cpu_baseline() {
let mut sampler = Sampler::default();
let first = sampler.sample();
let allocation = std::hint::black_box(vec![42_u8; 16 * 1024 * 1024]);
let baseline = Instant::now();
sampler.previous = Some(baseline);
let quick = sampler.sample();
assert_eq!(quick["gateway"], first["gateway"]);
assert_eq!(sampler.previous, Some(baseline));
assert!(quick["host"]["cpu_percent"].is_null());
drop(allocation);
}
#[tokio::test]
async fn timed_out_os_work_keeps_the_sampler_busy_until_it_finishes() {
let resources = Resources::default();
let guard = Arc::clone(&resources.sampler).try_lock_owned().unwrap();
let (release, wait) = std::sync::mpsc::channel::<()>();
let worker = tokio::task::spawn_blocking(move || {
wait.recv().unwrap();
guard
});
assert!(tokio::time::timeout(Duration::ZERO, worker).await.is_err());
assert_eq!(resources.sample().await["reason"], "busy");
release.send(()).unwrap();
let guard = tokio::time::timeout(DEADLINE, resources.sampler.lock()).await;
assert!(guard.is_ok());
}
}