use gregg_protocol::v2::{
CommitMetrics, CpuMetricsV2, DiskIoPayload, DriveMetrics, MetricCapabilitiesV2, NetworkPayload,
StatusPayloadV2, StatusSnapshotV2, SwapMetrics as SwapMetricsV2, SCHEMA_VERSION_V2,
};
use gregg_protocol::{
CpuMetrics, LoadAverage, MemoryMetrics, MetricCapabilities, StatusSnapshot, SwapMetrics,
SystemIdentity,
};
mod drives;
pub(crate) mod rate;
#[cfg(target_os = "linux")]
pub mod linux;
#[cfg(target_os = "macos")]
pub mod macos;
#[cfg(target_os = "windows")]
pub mod windows;
pub mod error;
use error::{CollectError, CollectErrorKind};
const DRIVE_REFRESH_INTERVAL: std::time::Duration = std::time::Duration::from_secs(30);
const DRIVE_REFRESH_RETRY_START: std::time::Duration = std::time::Duration::from_millis(10);
#[derive(Debug)]
pub(crate) struct DriveRefreshCache {
request_tx: Option<std::sync::mpsc::SyncSender<()>>,
result_rx: std::sync::mpsc::Receiver<Result<Vec<DriveMetrics>, CollectError>>,
latest: Option<Vec<DriveMetrics>>,
}
impl DriveRefreshCache {
pub(crate) fn new<S, F>(source: S, collect: F) -> Self
where
S: Send + 'static,
F: Fn(&S) -> Result<Vec<DriveMetrics>, CollectError> + Send + 'static,
{
let (request_tx, request_rx) = std::sync::mpsc::sync_channel(1);
let (result_tx, result_rx) = std::sync::mpsc::sync_channel(1);
let worker = std::thread::Builder::new()
.name("greggd-drive-refresh".into())
.spawn(move || {
let mut retry_delay = std::time::Duration::ZERO;
loop {
let wait = if retry_delay.is_zero() {
DRIVE_REFRESH_INTERVAL
} else {
retry_delay
};
let request = request_rx.recv_timeout(wait);
if matches!(
request,
Err(std::sync::mpsc::RecvTimeoutError::Disconnected)
) {
break;
}
let caught =
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| collect(&source)));
let (result, panicked) = if let Ok(result) = caught {
(result, false)
} else {
tracing::warn!("drive refresh worker collector panicked; retrying");
(
Err(CollectError::new(
CollectErrorKind::SourceUnavailable,
"drive refresh collector panicked",
)),
true,
)
};
retry_delay = if panicked {
if retry_delay.is_zero() {
DRIVE_REFRESH_RETRY_START
} else {
retry_delay
.checked_mul(2)
.map_or(DRIVE_REFRESH_INTERVAL, |delay| {
delay.min(DRIVE_REFRESH_INTERVAL)
})
}
} else {
std::time::Duration::ZERO
};
if result_tx.send(result).is_err() {
break;
}
}
})
.expect("drive refresh worker spawn");
drop(worker);
let _ = request_tx.try_send(());
Self {
request_tx: Some(request_tx),
result_rx,
latest: None,
}
}
#[cfg(test)]
pub(crate) fn request(&self) {
if let Some(sender) = &self.request_tx {
let _ = sender.try_send(());
}
}
pub(crate) fn poll(&mut self) -> Option<Vec<DriveMetrics>> {
while let Ok(result) = self.result_rx.try_recv() {
match result {
Ok(drives) => self.latest = Some(drives),
Err(error) => tracing::debug!(kind = ?error.kind),
}
}
self.latest.clone()
}
}
impl Drop for DriveRefreshCache {
fn drop(&mut self) {
let _ = self.request_tx.take();
}
}
#[allow(clippy::cast_precision_loss, clippy::cast_possible_truncation)]
pub(crate) fn clamped_usage_pct(used_bytes: u64, total_bytes: u64) -> f32 {
if total_bytes == 0 {
0.0
} else if used_bytes >= total_bytes {
100.0
} else {
let pct = (used_bytes as f64 / total_bytes as f64) * 100.0;
finalize_percentage(pct).unwrap_or(0.0)
}
}
pub(crate) fn finalize_percentage(value: f64) -> Result<f32, CollectError> {
if !value.is_finite() {
return Err(CollectError::new(
CollectErrorKind::Numeric,
"percentage is not finite",
));
}
#[allow(clippy::cast_possible_truncation)]
let as_f32 = value.clamp(0.0, 100.0) as f32;
if !as_f32.is_finite() || !(0.0..=100.0).contains(&as_f32) {
return Err(CollectError::new(
CollectErrorKind::Numeric,
"percentage outside closed 0..=100 interval after conversion",
));
}
Ok(as_f32)
}
#[derive(Debug, Clone, PartialEq)]
pub struct CollectedMetrics {
pub logical_cores: u32,
pub cpu_usage_pct: Option<f32>,
pub cpu_iowait_pct: Option<f32>,
pub load: LoadAverage,
pub memory: MemoryMetrics,
pub swap: SwapMetrics,
pub commit: Option<CommitMetrics>,
pub drives: Option<Vec<DriveMetrics>>,
pub cpu_frequency_hz: Option<u64>,
pub disk_io: Option<DiskIoPayload>,
pub network: Option<NetworkPayload>,
}
impl CollectedMetrics {
pub fn into_snapshot(
self,
schema_version: u16,
observed_at_unix_ms: u64,
sample_interval_ms: u64,
capabilities: MetricCapabilities,
system: SystemIdentity,
) -> Result<StatusSnapshot, CollectError> {
let Some(cpu_usage_pct) = self.cpu_usage_pct.filter(|v| v.is_finite()) else {
return Err(CollectError::new(
CollectErrorKind::Numeric,
"cpu usage percentage is missing or non-finite",
));
};
let cpu_iowait_pct = if capabilities.cpu_iowait {
let Some(iowait_pct) = self.cpu_iowait_pct.filter(|v| v.is_finite()) else {
return Err(CollectError::new(
CollectErrorKind::Numeric,
"cpu iowait percentage is missing or non-finite",
));
};
Some(iowait_pct)
} else {
None
};
Ok(StatusSnapshot {
schema_version,
observed_at_unix_ms,
sample_interval_ms,
capabilities,
system,
cpu: CpuMetrics {
logical_cores: self.logical_cores,
usage_pct: cpu_usage_pct,
iowait_pct: cpu_iowait_pct,
},
load: self.load,
memory: self.memory,
swap: self.swap,
})
}
pub fn into_snapshot_v2(
self,
observed_at_unix_ms: u64,
sample_interval_ms: u64,
capabilities: MetricCapabilitiesV2,
system: SystemIdentity,
) -> Result<StatusSnapshotV2, CollectError> {
let Some(cpu_usage_pct) = self.cpu_usage_pct.filter(|v| v.is_finite()) else {
return Err(CollectError::new(
CollectErrorKind::Numeric,
"cpu usage percentage is missing or non-finite",
));
};
let cpu_iowait_pct = if capabilities.cpu_iowait {
let Some(iowait_pct) = self.cpu_iowait_pct.filter(|v| v.is_finite()) else {
return Err(CollectError::new(
CollectErrorKind::Numeric,
"cpu iowait percentage is missing or non-finite",
));
};
Some(iowait_pct)
} else {
None
};
let load = if capabilities.load_average {
Some(self.load)
} else {
None
};
let swap = if capabilities.swap {
Some(SwapMetricsV2 {
used_bytes: self.swap.used_bytes,
total_bytes: self.swap.total_bytes,
usage_pct: clamped_usage_pct(self.swap.used_bytes, self.swap.total_bytes),
})
} else {
None
};
Ok(StatusSnapshotV2 {
schema_version: SCHEMA_VERSION_V2,
observed_at_unix_ms,
sample_interval_ms,
capabilities,
system,
cpu: CpuMetricsV2 {
logical_cores: self.logical_cores,
usage_pct: cpu_usage_pct,
iowait_pct: cpu_iowait_pct,
},
load,
memory: self.memory,
swap,
commit: self.commit,
})
}
pub fn into_status_payload_v2(
self,
observed_at_unix_ms: u64,
sample_interval_ms: u64,
capabilities: MetricCapabilitiesV2,
system: SystemIdentity,
) -> Result<StatusPayloadV2, CollectError> {
let drives = self.drives.clone();
let cpu_frequency_hz = self.cpu_frequency_hz;
let disk_io = self.disk_io.clone();
let network = self.network.clone();
let snapshot = self.into_snapshot_v2(
observed_at_unix_ms,
sample_interval_ms,
capabilities,
system,
)?;
Ok(StatusPayloadV2 {
snapshot,
drives,
cpu_frequency_hz,
disk_io,
network,
})
}
}
pub trait SystemCollector: Send {
fn identity(&self) -> Result<SystemIdentity, error::CollectError>;
fn sample(&mut self) -> Result<CollectedMetrics, error::CollectError>;
fn capabilities(&self) -> MetricCapabilities;
fn capabilities_v2(&self) -> MetricCapabilitiesV2 {
let v1 = self.capabilities();
MetricCapabilitiesV2 {
cpu_iowait: v1.cpu_iowait,
load_average: true,
swap: true,
memory_commit: false,
}
}
fn supports_v1_snapshot(&self) -> bool {
true
}
}
#[cfg(test)]
mod tests {
use super::error::{CollectError, CollectErrorKind};
use super::{clamped_usage_pct, DriveRefreshCache};
use gregg_protocol::v2::DriveMetrics;
#[test]
fn large_byte_ratios_remain_finite_and_clamped() {
assert!((clamped_usage_pct(0, u64::MAX) - 0.0).abs() < f32::EPSILON);
assert!((clamped_usage_pct(u64::MAX, u64::MAX) - 100.0).abs() < f32::EPSILON);
assert!((clamped_usage_pct(u64::MAX, u64::MAX - 1) - 100.0).abs() < f32::EPSILON);
}
fn wait_until(mut condition: impl FnMut() -> bool) {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
while std::time::Instant::now() < deadline {
if condition() {
return;
}
std::thread::sleep(std::time::Duration::from_millis(1));
}
assert!(condition(), "worker did not reach expected state");
}
#[test]
fn blocked_drive_refresh_does_not_block_cache_drop() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
let started = Arc::new(AtomicBool::new(false));
let release = Arc::new(AtomicBool::new(false));
let started_for_worker = Arc::clone(&started);
let release_for_worker = Arc::clone(&release);
let mut cache = DriveRefreshCache::new((), move |()| {
started_for_worker.store(true, Ordering::Release);
while !release_for_worker.load(Ordering::Acquire) {
std::thread::yield_now();
}
Ok(Vec::new())
});
wait_until(|| started.load(Ordering::Acquire));
assert_eq!(cache.poll(), None);
let before = std::time::Instant::now();
drop(cache);
assert!(before.elapsed() < std::time::Duration::from_millis(100));
release.store(true, Ordering::Release);
}
#[test]
fn drive_refresh_retains_last_success_after_failure() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
let calls = Arc::new(AtomicUsize::new(0));
let calls_for_worker = Arc::clone(&calls);
let mut cache = DriveRefreshCache::new((), move |()| {
let call = calls_for_worker.fetch_add(1, Ordering::AcqRel);
if call == 0 {
Ok(vec![DriveMetrics {
name: "root".to_string(),
used_bytes: 1,
total_bytes: 2,
available_bytes: Some(1),
}])
} else {
Err(CollectError::new(
CollectErrorKind::SourceUnavailable,
"refresh failed",
))
}
});
wait_until(|| cache.poll().is_some());
let first = cache.poll().expect("first drive result");
assert_eq!(first[0].name, "root");
cache.request();
wait_until(|| calls.load(Ordering::Acquire) >= 2);
assert_eq!(
cache.poll().expect("last good drive result")[0].used_bytes,
1
);
}
#[test]
fn drive_refresh_does_not_drop_a_new_result_while_previous_is_queued() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
let calls = Arc::new(AtomicUsize::new(0));
let calls_for_worker = Arc::clone(&calls);
let mut cache = DriveRefreshCache::new((), move |()| {
let call = calls_for_worker.fetch_add(1, Ordering::AcqRel);
Ok(vec![DriveMetrics {
name: format!("drive-{call}"),
used_bytes: call as u64,
total_bytes: 10,
available_bytes: Some(10 - call as u64),
}])
});
wait_until(|| calls.load(Ordering::Acquire) >= 1);
cache.request();
wait_until(|| calls.load(Ordering::Acquire) >= 2);
wait_until(|| {
cache
.poll()
.is_some_and(|drives| drives[0].name == "drive-1")
});
}
#[test]
fn drive_refresh_recovers_after_collector_panic() {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
let calls = Arc::new(AtomicUsize::new(0));
let calls_for_worker = Arc::clone(&calls);
let mut cache = DriveRefreshCache::new((), move |()| {
assert_ne!(
calls_for_worker.fetch_add(1, Ordering::AcqRel),
0,
"injected drive refresh panic"
);
Ok(Vec::new())
});
wait_until(|| calls.load(Ordering::Acquire) >= 1);
wait_until(|| calls.load(Ordering::Acquire) >= 2);
wait_until(|| cache.poll().is_some());
assert_eq!(cache.poll(), Some(Vec::new()));
}
}