pub mod mem_breakdown;
pub use mem_breakdown::{MemoryBreakdown, self_memory_breakdown};
use sysinfo::{Pid, ProcessRefreshKind, ProcessesToUpdate, RefreshKind, System};
#[cfg(target_os = "macos")]
pub fn physical_footprint_mb(pid: u32) -> Option<u64> {
Some(physical_footprint_bytes(pid)? / (1024 * 1024))
}
#[cfg(target_os = "macos")]
#[must_use]
pub fn physical_footprint_bytes(pid: u32) -> Option<u64> {
let mut info: libc::rusage_info_v0 = unsafe { std::mem::zeroed() };
let ret = unsafe {
libc::proc_pid_rusage(
pid as libc::c_int,
libc::RUSAGE_INFO_V0,
std::ptr::addr_of_mut!(info).cast(),
)
};
if ret != 0 {
return None;
}
Some(info.ri_phys_footprint)
}
#[must_use]
pub fn process_rss_mb(pid: u32) -> Option<u64> {
#[cfg(target_os = "macos")]
{
physical_footprint_mb(pid)
}
#[cfg(target_os = "linux")]
{
let status = std::fs::read_to_string(format!("/proc/{pid}/status")).ok()?;
for line in status.lines() {
let Some(rest) = line.strip_prefix("VmRSS:") else {
continue;
};
let kb: u64 = rest.split_whitespace().next()?.parse().ok()?;
return Some(kb / 1024);
}
None
}
#[cfg(not(any(target_os = "macos", target_os = "linux")))]
{
let _ = pid;
None
}
}
pub struct SysMetrics {
sys: System,
pid: Pid,
}
impl SysMetrics {
#[must_use]
pub fn new() -> Self {
let pid = Pid::from_u32(std::process::id());
let mut sys = System::new_with_specifics(
RefreshKind::nothing()
.with_processes(ProcessRefreshKind::nothing().with_memory().with_cpu()),
);
sys.refresh_processes_specifics(
ProcessesToUpdate::Some(&[pid]),
true,
ProcessRefreshKind::nothing().with_memory().with_cpu(),
);
Self { sys, pid }
}
pub fn sample(&mut self) -> (u64, f32) {
self.sys.refresh_processes_specifics(
ProcessesToUpdate::Some(&[self.pid]),
true,
ProcessRefreshKind::nothing().with_memory().with_cpu(),
);
let Some(proc) = self.sys.process(self.pid) else {
return (0, 0.0);
};
let sysinfo_rss_mb = proc.memory() / (1024 * 1024);
let cpu_pct = proc.cpu_usage();
#[cfg(target_os = "macos")]
let rss_mb = physical_footprint_mb(self.pid.as_u32()).unwrap_or(sysinfo_rss_mb);
#[cfg(not(target_os = "macos"))]
let rss_mb = sysinfo_rss_mb;
(rss_mb, cpu_pct)
}
}
impl Default for SysMetrics {
fn default() -> Self {
Self::new()
}
}
pub struct ProcessCpuSampler {
sys: System,
tracked: Vec<Pid>,
}
impl ProcessCpuSampler {
#[must_use]
pub fn new() -> Self {
Self {
sys: System::new_with_specifics(
RefreshKind::nothing().with_processes(process_refresh()),
),
tracked: Vec::new(),
}
}
pub fn track(&mut self, pid: u32) {
let pid = Pid::from_u32(pid);
if self.tracked.contains(&pid) {
return;
}
self.tracked.push(pid);
self.sys.refresh_processes_specifics(
ProcessesToUpdate::Some(&[pid]),
true,
process_refresh(),
);
}
pub fn untrack(&mut self, pid: u32) {
let pid = Pid::from_u32(pid);
self.tracked.retain(|p| *p != pid);
}
#[must_use]
pub fn is_tracked(&self, pid: u32) -> bool {
self.tracked.contains(&Pid::from_u32(pid))
}
#[must_use]
pub fn tracked_count(&self) -> usize {
self.tracked.len()
}
pub fn refresh(&mut self) {
if self.tracked.is_empty() {
return;
}
self.sys.refresh_processes_specifics(
ProcessesToUpdate::Some(&self.tracked),
true,
process_refresh(),
);
let sys = &self.sys;
self.tracked.retain(|pid| sys.process(*pid).is_some());
}
#[must_use]
pub fn cpu_pct(&self, pid: u32) -> Option<f32> {
let pid = Pid::from_u32(pid);
if !self.tracked.contains(&pid) {
return None;
}
self.sys.process(pid).map(sysinfo::Process::cpu_usage)
}
#[must_use]
pub fn rss_bytes(&self, pid: u32) -> Option<u64> {
let sysinfo_bytes = {
let key = Pid::from_u32(pid);
if !self.tracked.contains(&key) {
return None;
}
self.sys.process(key).map(sysinfo::Process::memory)?
};
#[cfg(target_os = "macos")]
{
Some(physical_footprint_bytes(pid).unwrap_or(sysinfo_bytes))
}
#[cfg(not(target_os = "macos"))]
{
Some(sysinfo_bytes)
}
}
}
fn process_refresh() -> ProcessRefreshKind {
ProcessRefreshKind::nothing().with_cpu().with_memory()
}
impl Default for ProcessCpuSampler {
fn default() -> Self {
Self::new()
}
}
const MAX_WALK_DEPTH: usize = 64;
const WALK_BUDGET: std::time::Duration = std::time::Duration::from_secs(30);
#[must_use]
pub fn dir_size_bytes(dir: &std::path::Path) -> u64 {
walk_summing(dir, "dir_size_bytes", apparent_len)
}
#[must_use]
pub fn dir_allocated_bytes(dir: &std::path::Path) -> u64 {
walk_summing(dir, "dir_allocated_bytes", allocated_len)
}
fn apparent_len(meta: &std::fs::Metadata) -> u64 {
meta.len()
}
fn allocated_len(meta: &std::fs::Metadata) -> u64 {
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt as _;
meta.blocks().saturating_mul(512)
}
#[cfg(not(unix))]
{
meta.len()
}
}
fn walk_summing(dir: &std::path::Path, label: &str, measure: fn(&std::fs::Metadata) -> u64) -> u64 {
let total = std::cell::Cell::new(0u64);
let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
walk_bounded(dir, &total, measure);
}));
if outcome.is_err() {
tracing::error!(
dir = %dir.display(),
walk = label,
"directory walk panicked (see preceding PANIC log for the \
payload); reporting the partial total"
);
}
total.get()
}
fn walk_bounded(
root: &std::path::Path,
total: &std::cell::Cell<u64>,
measure: fn(&std::fs::Metadata) -> u64,
) {
let started = std::time::Instant::now();
let mut stack: Vec<(std::path::PathBuf, usize)> = vec![(root.to_path_buf(), 0)];
while let Some((dir, depth)) = stack.pop() {
if started.elapsed() >= WALK_BUDGET {
tracing::warn!(
root = %root.display(),
pending = stack.len() + 1,
"directory walk exceeded its {WALK_BUDGET:?} budget; \
reporting the partial total"
);
return;
}
let Ok(entries) = std::fs::read_dir(&dir) else {
continue;
};
for entry in entries.flatten() {
let Ok(file_type) = entry.file_type() else {
continue;
};
if file_type.is_symlink() {
continue;
}
if file_type.is_dir() {
if depth < MAX_WALK_DEPTH {
stack.push((entry.path(), depth + 1));
}
continue;
}
if !file_type.is_file() {
continue;
}
if let Ok(meta) = entry.metadata() {
total.set(total.get().saturating_add(measure(&meta)));
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn sample_does_not_panic() {
let mut m = SysMetrics::new();
let (_rss, _cpu) = m.sample();
let (_rss2, cpu2) = m.sample();
assert!(cpu2 >= 0.0, "cpu usage must be non-negative, got {cpu2}");
}
#[test]
fn rss_is_plausible() {
let mut m = SysMetrics::new();
let (rss, _cpu) = m.sample();
assert!(
rss < 1024 * 1024,
"RSS implausibly large ({rss} MB) — unit must be MB"
);
}
#[test]
fn process_rss_mb_reports_own_process() {
let me = std::process::id();
match process_rss_mb(me) {
Some(mb) => assert!(
mb < 1024 * 1024,
"own RSS implausibly large ({mb} MB) — unit must be MB"
),
#[cfg(any(target_os = "macos", target_os = "linux"))]
None => panic!("process_rss_mb must resolve the current process on this platform"),
#[cfg(not(any(target_os = "macos", target_os = "linux")))]
None => {}
}
}
#[test]
fn process_rss_mb_is_none_for_absent_pid() {
assert_eq!(process_rss_mb(0), None);
}
#[cfg(unix)]
fn spawn_sleeper() -> std::process::Child {
std::process::Command::new("sleep")
.arg("30")
.spawn()
.expect("spawn a sleeping child")
}
#[cfg(unix)]
#[test]
fn process_cpu_sampler_measures_a_tracked_child() {
let mut child = spawn_sleeper();
let pid = child.id();
let mut sampler = ProcessCpuSampler::new();
sampler.track(pid);
sampler.refresh();
let cpu = sampler.cpu_pct(pid);
let _ = child.kill();
let _ = child.wait();
let cpu = cpu.expect("a live tracked pid must yield a measurement");
assert!(cpu >= 0.0, "cpu usage must be non-negative, got {cpu}");
}
#[cfg(unix)]
#[test]
fn process_cpu_sampler_measures_memory_of_a_tracked_child() {
let mut child = spawn_sleeper();
let pid = child.id();
let mut sampler = ProcessCpuSampler::new();
sampler.track(pid);
sampler.refresh();
let cpu = sampler.cpu_pct(pid);
let rss = sampler.rss_bytes(pid);
let _ = child.kill();
let _ = child.wait();
assert!(cpu.is_some(), "one refresh must serve both figures");
let rss = rss.expect("a live tracked pid must yield a memory measurement");
assert!(rss > 0, "a live process occupies memory, got {rss} bytes");
assert!(
rss < 1024 * 1024 * 1024 * 1024,
"implausibly large ({rss}) — the unit must be bytes"
);
}
#[test]
fn process_cpu_sampler_reports_no_memory_for_an_untracked_pid() {
let sampler = ProcessCpuSampler::new();
assert_eq!(sampler.rss_bytes(std::process::id()), None);
}
#[cfg(unix)]
#[test]
fn process_cpu_sampler_reports_none_for_a_vanished_pid() {
let mut child = spawn_sleeper();
let pid = child.id();
let mut sampler = ProcessCpuSampler::new();
sampler.track(pid);
sampler.refresh();
child.kill().expect("kill the sleeper");
child.wait().expect("reap the sleeper");
sampler.refresh();
assert_eq!(
sampler.cpu_pct(pid),
None,
"a vanished process must read as no-measurement, never as 0.0"
);
assert!(
!sampler.is_tracked(pid),
"a vanished pid must be dropped so a reused pid cannot be misread"
);
sampler.refresh();
assert_eq!(sampler.tracked_count(), 0);
}
#[test]
fn process_cpu_sampler_reports_none_for_an_untracked_pid() {
let mut sampler = ProcessCpuSampler::new();
assert_eq!(sampler.cpu_pct(std::process::id()), None);
sampler.refresh();
assert_eq!(sampler.cpu_pct(std::process::id()), None);
}
#[test]
fn process_cpu_sampler_untrack_removes_a_pid() {
let me = std::process::id();
let mut sampler = ProcessCpuSampler::new();
sampler.track(me);
sampler.track(me);
assert_eq!(sampler.tracked_count(), 1, "track must be idempotent");
assert!(sampler.is_tracked(me));
sampler.untrack(me);
assert_eq!(sampler.tracked_count(), 0);
assert!(!sampler.is_tracked(me));
assert_eq!(sampler.cpu_pct(me), None);
}
#[test]
fn dir_size_sums_files() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::write(tmp.path().join("a.txt"), vec![0u8; 100]).unwrap();
std::fs::write(tmp.path().join("b.txt"), vec![0u8; 250]).unwrap();
let sub = tmp.path().join("sub");
std::fs::create_dir(&sub).unwrap();
std::fs::write(sub.join("c.txt"), vec![0u8; 50]).unwrap();
assert_eq!(dir_size_bytes(tmp.path()), 400);
}
#[test]
fn dir_size_missing_dir_is_zero() {
let missing = std::path::Path::new("/nonexistent/trusty/path/xyz");
assert_eq!(dir_size_bytes(missing), 0);
}
#[test]
fn dir_allocated_bytes_matches_a_walked_fixture() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::write(tmp.path().join("a.bin"), vec![7u8; 100]).unwrap();
std::fs::write(tmp.path().join("b.bin"), vec![7u8; 40_000]).unwrap();
let sub = tmp.path().join("sub").join("deeper");
std::fs::create_dir_all(&sub).unwrap();
std::fs::write(sub.join("c.bin"), vec![7u8; 9_000]).unwrap();
let walked = walk_fixture_allocated(tmp.path());
assert!(walked > 0, "the fixture must allocate something");
assert_eq!(dir_allocated_bytes(tmp.path()), walked);
}
#[cfg(unix)]
fn walk_fixture_allocated(dir: &std::path::Path) -> u64 {
use std::os::unix::fs::MetadataExt as _;
let mut total = 0u64;
for entry in std::fs::read_dir(dir).expect("read_dir").flatten() {
let meta = entry.metadata().expect("metadata");
if meta.is_dir() {
total += walk_fixture_allocated(&entry.path());
} else if meta.is_file() {
total += meta.blocks() * 512;
}
}
total
}
#[cfg(not(unix))]
fn walk_fixture_allocated(dir: &std::path::Path) -> u64 {
let mut total = 0u64;
for entry in std::fs::read_dir(dir).expect("read_dir").flatten() {
let meta = entry.metadata().expect("metadata");
if meta.is_dir() {
total += walk_fixture_allocated(&entry.path());
} else if meta.is_file() {
total += meta.len();
}
}
total
}
#[cfg(unix)]
#[test]
fn dir_allocated_bytes_reads_block_allocation_not_file_length() {
let tmp = tempfile::tempdir().expect("tempdir");
for i in 0..16 {
std::fs::write(tmp.path().join(format!("{i}.bin")), b"x").unwrap();
}
let apparent = dir_size_bytes(tmp.path());
let allocated = dir_allocated_bytes(tmp.path());
assert_eq!(
apparent, 16,
"sixteen one-byte files are sixteen bytes long"
);
assert!(
allocated > apparent,
"block-rounded allocation must exceed the logical length: \
{allocated} vs {apparent}"
);
}
#[test]
fn dir_allocated_missing_dir_is_zero() {
let missing = std::path::Path::new("/nonexistent/trusty/path/xyz");
assert_eq!(dir_allocated_bytes(missing), 0);
}
#[test]
fn dir_size_survives_concurrent_mutation() {
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
const BRANCHES: u64 = 8;
const LEAF_BYTES: u64 = 64;
const TOP_BYTES: u64 = 32;
let tmp = tempfile::tempdir().expect("tempdir");
let root = tmp.path().to_path_buf();
for i in 0..BRANCHES {
let branch = root.join(format!("branch-{i}"));
std::fs::create_dir_all(branch.join("a/b/c")).expect("seed dirs");
std::fs::write(
branch.join("a/b/c/leaf.bin"),
vec![0u8; LEAF_BYTES as usize],
)
.expect("seed leaf");
std::fs::write(branch.join("a/top.bin"), vec![0u8; TOP_BYTES as usize])
.expect("seed top");
}
let stop = Arc::new(AtomicBool::new(false));
let mutators: Vec<_> = (0..3)
.map(|t| {
let root = root.clone();
let stop = Arc::clone(&stop);
std::thread::spawn(move || {
let mut n: u64 = 0;
while !stop.load(Ordering::Relaxed) {
let staged = root.join(format!("staged-{t}-{n}"));
let live = root.join(format!("live-{t}"));
if std::fs::create_dir_all(staged.join("nested")).is_ok() {
let _ = std::fs::write(staged.join("nested/data.bin"), vec![0u8; 128]);
let _ = std::fs::remove_dir_all(&live);
let _ = std::fs::rename(&staged, &live);
}
let _ = std::fs::remove_dir_all(&live);
let _ = std::fs::remove_dir_all(&staged);
n = n.wrapping_add(1);
}
})
})
.collect();
let mut min_observed = u64::MAX;
for _ in 0..30 {
min_observed = min_observed.min(dir_size_bytes(&root));
}
stop.store(true, Ordering::Relaxed);
for handle in mutators {
handle.join().expect("mutator thread must not panic");
}
let floor = BRANCHES * (LEAF_BYTES + TOP_BYTES);
assert!(
min_observed >= floor,
"a walk under concurrent mutation lost stable bytes: worst sample \
{min_observed}, floor {floor}"
);
}
#[test]
fn dir_size_depth_cap_boundary_is_exact() {
const TOP_BYTES: u64 = 7;
const AT_CAP_BYTES: u64 = 11;
const PAST_CAP_BYTES: u64 = 4096;
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::write(tmp.path().join("top.bin"), vec![0u8; TOP_BYTES as usize])
.expect("write top");
let mut at_cap = tmp.path().to_path_buf();
for _ in 0..MAX_WALK_DEPTH {
at_cap.push("d");
}
std::fs::create_dir_all(&at_cap).expect("create at-cap tree");
std::fs::write(at_cap.join("at-cap.bin"), vec![0u8; AT_CAP_BYTES as usize])
.expect("write at-cap file");
let past_cap = at_cap.join("d");
std::fs::create_dir(&past_cap).expect("create past-cap dir");
std::fs::write(
past_cap.join("past-cap.bin"),
vec![0u8; PAST_CAP_BYTES as usize],
)
.expect("write past-cap file");
let total = dir_size_bytes(tmp.path());
assert_eq!(
total,
TOP_BYTES + AT_CAP_BYTES,
"depth cap is off by one: {} means the cap fired a level early \
(the at-cap file was dropped); {} means it fired a level late \
(the past-cap file was counted)",
TOP_BYTES,
TOP_BYTES + AT_CAP_BYTES + PAST_CAP_BYTES
);
}
#[cfg(target_os = "macos")]
#[test]
fn self_physical_footprint_is_plausible() {
let pid = std::process::id();
let mb = physical_footprint_mb(pid).expect("proc_pid_rusage must resolve our own pid");
assert!(mb > 0, "physical footprint should be > 0 MB, got {mb}");
assert!(
mb < 1024 * 1024,
"physical footprint implausibly large ({mb} MB) — unit must be MB"
);
}
#[cfg(target_os = "macos")]
#[test]
fn physical_footprint_bytes_agrees_with_the_megabyte_reading() {
let pid = std::process::id();
let bytes = physical_footprint_bytes(pid).expect("bytes reading");
let mb = physical_footprint_mb(pid).expect("megabyte reading");
let bytes_as_mb = bytes / (1024 * 1024);
assert!(
bytes_as_mb.abs_diff(mb) <= 1,
"bytes ({bytes_as_mb} MB) and mb ({mb} MB) must be the same counter"
);
}
#[cfg(target_os = "macos")]
#[test]
fn physical_footprint_bogus_pid_returns_none() {
assert_eq!(physical_footprint_mb(u32::MAX), None);
}
#[cfg(target_os = "macos")]
#[test]
fn physical_footprint_tracks_real_allocation_growth() {
let pid = std::process::id();
let before = physical_footprint_mb(pid).expect("must resolve our own pid");
let mut touched: Vec<u8> = vec![0u8; 200 * 1024 * 1024];
for byte in touched.iter_mut().step_by(4096) {
*byte = 1;
}
let after = physical_footprint_mb(pid).expect("must resolve our own pid");
assert!(
after >= before + 100,
"expected footprint to grow by >= 100 MB after touching a 200 MB \
allocation; before={before} after={after}"
);
drop(touched);
}
}