use std::time::Duration;
use sysinfo::{ProcessesToUpdate, System};
pub(crate) struct ResourceSampler {
handle: tokio::task::JoinHandle<()>,
}
impl ResourceSampler {
pub(crate) fn spawn(interval: Duration) -> Self {
let handle = tokio::spawn(async move {
let Ok(pid) = sysinfo::get_current_pid() else {
return;
};
let mut system = System::new();
loop {
tokio::time::sleep(interval).await;
system.refresh_processes(ProcessesToUpdate::Some(&[pid]), true);
let Some(process) = system.process(pid) else {
continue;
};
let disk = process.disk_usage();
let instance = crate::observability::instance();
metrics::gauge!("pigeon_resource_cpu_percent", "instance" => instance)
.set(process.cpu_usage() as f64);
metrics::gauge!("pigeon_resource_mem_bytes", "instance" => instance)
.set(process.memory() as f64);
metrics::counter!("pigeon_resource_disk_read_bytes_total", "instance" => instance)
.absolute(disk.total_read_bytes);
metrics::counter!("pigeon_resource_disk_written_bytes_total", "instance" => instance)
.absolute(disk.total_written_bytes);
}
});
Self { handle }
}
}
impl Drop for ResourceSampler {
fn drop(&mut self) {
self.handle.abort();
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn dropping_the_sampler_stops_its_sampling_task() {
let sampler = ResourceSampler::spawn(Duration::from_millis(5));
tokio::time::sleep(Duration::from_millis(20)).await;
let abort_handle = sampler.handle.abort_handle();
drop(sampler);
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(abort_handle.is_finished());
}
}