File size: 6,195 Bytes
9ebf6d4 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 | use std::sync::Mutex;
use std::sync::Once;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::Ordering;
static GAUGES: Mutex<Vec<&'static Gauge>> = Mutex::new(Vec::new());
/// A process-wide gauge that registers itself the first time it is used.
pub struct Gauge {
name: &'static str,
value: AtomicU64,
registered: Once,
}
impl Gauge {
/// Creates a gauge suitable for use in a `static` declaration.
pub const fn new(name: &'static str) -> Self {
Self {
name,
value: AtomicU64::new(0),
registered: Once::new(),
}
}
/// Increments the gauge and registers it if needed.
pub fn increment(&'static self) {
self.register();
self.value.fetch_add(1, Ordering::Relaxed);
}
/// Decrements the gauge without allowing an underflow.
pub fn decrement(&self) {
let _ = self
.value
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |value| {
Some(value.saturating_sub(1))
});
}
/// Increments the gauge for the lifetime of the returned guard.
pub fn track(&'static self) -> GaugeGuard {
self.increment();
GaugeGuard { gauge: self }
}
fn register(&'static self) {
self.registered.call_once(|| {
GAUGES
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(self);
});
}
}
/// Decrements its gauge when the measured object is dropped.
pub struct GaugeGuard {
gauge: &'static Gauge,
}
impl Drop for GaugeGuard {
fn drop(&mut self) {
self.gauge.decrement();
}
}
/// The current value of one registered diagnostic gauge.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct GaugeSnapshot {
pub name: &'static str,
pub value: u64,
}
/// Best-effort operating-system measurements for the current process.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct ProcessSnapshot {
pub id: u32,
pub resident_memory_bytes: Option<u64>,
pub physical_footprint_bytes: Option<u64>,
}
/// Content-free diagnostic values contributed by this process.
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct DiagnosticsSnapshot {
pub process: ProcessSnapshot,
pub gauges: Vec<GaugeSnapshot>,
}
/// Collects built-in process measurements and every registered gauge.
pub fn snapshot() -> DiagnosticsSnapshot {
let mut gauges = GAUGES
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.iter()
.map(|gauge| GaugeSnapshot {
name: gauge.name,
value: gauge.value.load(Ordering::Relaxed),
})
.collect::<Vec<_>>();
gauges.sort_unstable_by_key(|gauge| gauge.name);
DiagnosticsSnapshot {
process: process_snapshot(),
gauges,
}
}
#[cfg(target_os = "macos")]
fn process_snapshot() -> ProcessSnapshot {
let usage = unsafe {
let mut usage = std::mem::MaybeUninit::<libc::rusage_info_v0>::zeroed();
// SAFETY: the kernel initializes this correctly sized buffer on success.
if libc::proc_pid_rusage(
libc::getpid(),
libc::RUSAGE_INFO_V0,
usage.as_mut_ptr().cast(),
) != 0
{
return empty_process_snapshot();
}
usage.assume_init()
};
ProcessSnapshot {
id: std::process::id(),
resident_memory_bytes: Some(usage.ri_resident_size),
physical_footprint_bytes: Some(usage.ri_phys_footprint),
}
}
#[cfg(target_os = "linux")]
fn process_snapshot() -> ProcessSnapshot {
// SAFETY: querying the system page size does not access caller-owned memory.
let page_size = u64::try_from(unsafe { libc::sysconf(libc::_SC_PAGESIZE) })
.ok()
.filter(|page_size| *page_size > 0);
let resident_pages = std::fs::read_to_string("/proc/self/statm")
.ok()
.and_then(|statm| statm.split_whitespace().nth(1)?.parse::<u64>().ok());
ProcessSnapshot {
id: std::process::id(),
resident_memory_bytes: resident_pages
.zip(page_size)
.map(|(pages, page_size)| pages.saturating_mul(page_size)),
physical_footprint_bytes: None,
}
}
#[cfg(target_os = "windows")]
fn process_snapshot() -> ProcessSnapshot {
#[repr(C)]
struct ProcessMemoryCounters {
size: u32,
page_fault_count: u32,
peak_working_set_size: usize,
working_set_size: usize,
quota_peak_paged_pool_usage: usize,
quota_paged_pool_usage: usize,
quota_peak_non_paged_pool_usage: usize,
quota_non_paged_pool_usage: usize,
pagefile_usage: usize,
peak_pagefile_usage: usize,
}
#[link(name = "kernel32")]
unsafe extern "system" {
fn GetCurrentProcess() -> *mut std::ffi::c_void;
fn K32GetProcessMemoryInfo(
process: *mut std::ffi::c_void,
counters: *mut ProcessMemoryCounters,
size: u32,
) -> i32;
}
let counters = unsafe {
let mut counters = std::mem::MaybeUninit::<ProcessMemoryCounters>::zeroed();
let size = u32::try_from(std::mem::size_of::<ProcessMemoryCounters>()).unwrap_or(u32::MAX);
// SAFETY: the pseudo-handle is valid and the kernel initializes this
// correctly sized writable buffer on success.
if K32GetProcessMemoryInfo(GetCurrentProcess(), counters.as_mut_ptr(), size) == 0 {
return empty_process_snapshot();
}
counters.assume_init()
};
ProcessSnapshot {
id: std::process::id(),
resident_memory_bytes: u64::try_from(counters.working_set_size).ok(),
physical_footprint_bytes: None,
}
}
#[cfg(not(any(target_os = "macos", target_os = "linux", target_os = "windows")))]
fn process_snapshot() -> ProcessSnapshot {
empty_process_snapshot()
}
#[cfg(not(target_os = "linux"))]
fn empty_process_snapshot() -> ProcessSnapshot {
ProcessSnapshot {
id: std::process::id(),
resident_memory_bytes: None,
physical_footprint_bytes: None,
}
}
#[cfg(test)]
#[path = "tests.rs"]
mod tests;
|