use std::time::Duration;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use crate::system::SystemSnapshot;
pub struct Collector {
sys: sysinfo::System,
networks: sysinfo::Networks,
interval: Duration,
}
impl Collector {
pub fn new(interval: Duration) -> Self {
Self {
sys: sysinfo::System::new_all(),
networks: sysinfo::Networks::new_with_refreshed_list(),
interval,
}
}
pub fn spawn(
self,
tx: mpsc::Sender<SystemSnapshot>,
token: CancellationToken,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(Self::run(self, tx, token))
}
async fn run(mut self, tx: mpsc::Sender<SystemSnapshot>, token: CancellationToken) {
let mut interval = tokio::time::interval(self.interval);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
interval.tick().await;
self.sys.refresh_all();
self.networks.refresh(true);
loop {
tokio::select! {
_ = interval.tick() => {
self.sys.refresh_all();
self.networks.refresh(false);
let snapshot = SystemSnapshot::collect(&self.sys, &self.networks);
match tx.try_send(snapshot) {
Ok(()) => {}
Err(mpsc::error::TrySendError::Full(_)) => {
tracing::trace!("channel full, dropping snapshot");
}
Err(mpsc::error::TrySendError::Closed(_)) => {
tracing::debug!("channel closed, stopping collector");
break;
}
}
}
_ = token.cancelled() => {
tracing::debug!("collector shutting down");
break;
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
fn make_collector(
cap: usize,
) -> (
mpsc::Receiver<SystemSnapshot>,
tokio::task::JoinHandle<()>,
CancellationToken,
) {
let (tx, rx) = mpsc::channel(cap);
let token = CancellationToken::new();
let collector = Collector::new(Duration::from_secs(1));
let handle = collector.spawn(tx, token.clone());
(rx, handle, token)
}
#[tokio::test]
async fn test_collector_produces_snapshots() {
let (mut rx, handle, token) = make_collector(4);
let mut count = 0usize;
let deadline = tokio::time::Instant::now() + Duration::from_secs(4);
loop {
match tokio::time::timeout_at(deadline, rx.recv()).await {
Ok(Some(_)) => {
count += 1;
if count >= 2 {
break;
}
}
Ok(None) => panic!("channel closed before receiving 2 snapshots"),
Err(_) => panic!("timeout: only received {count} snapshots within 4s"),
}
}
token.cancel();
handle.await.expect("collector task panicked");
assert!(count >= 2, "expected at least 2 snapshots, got {count}");
}
#[tokio::test]
async fn test_collector_snapshot_has_data() {
let (mut rx, handle, token) = make_collector(4);
let snapshot = tokio::time::timeout(Duration::from_secs(4), rx.recv())
.await
.expect("timeout waiting for snapshot")
.expect("channel closed before first snapshot");
token.cancel();
handle.await.expect("collector task panicked");
assert!(
!snapshot.processes.is_empty(),
"snapshot should contain processes"
);
assert!(
!snapshot.cpu.cores.is_empty(),
"snapshot should contain CPU cores"
);
}
#[tokio::test]
async fn test_collector_graceful_shutdown() {
let (mut rx, handle, token) = make_collector(4);
tokio::spawn(async move { while rx.recv().await.is_some() {} });
tokio::time::sleep(Duration::from_millis(500)).await;
token.cancel();
tokio::time::timeout(Duration::from_secs(2), handle)
.await
.expect("collector did not shut down within 2s")
.expect("collector task panicked");
}
#[tokio::test]
async fn test_collector_channel_backpressure() {
let (tx, _rx) = mpsc::channel::<SystemSnapshot>(1);
let token = CancellationToken::new();
let collector = Collector::new(Duration::from_secs(1));
let handle = collector.spawn(tx, token.clone());
tokio::time::sleep(Duration::from_secs(2)).await;
token.cancel();
tokio::time::timeout(Duration::from_secs(2), handle)
.await
.expect("collector did not shut down within 2s after backpressure test")
.expect("collector task panicked");
}
#[tokio::test]
async fn test_collector_respects_interval() {
let (mut rx, handle, token) = make_collector(4);
let first = tokio::time::timeout(Duration::from_secs(4), rx.recv())
.await
.expect("timeout waiting for first snapshot")
.expect("channel closed before first snapshot");
let second = tokio::time::timeout(Duration::from_secs(4), rx.recv())
.await
.expect("timeout waiting for second snapshot")
.expect("channel closed before second snapshot");
token.cancel();
handle.await.expect("collector task panicked");
let gap_ms = second.timestamp_ms.saturating_sub(first.timestamp_ms);
assert!(
(500..=1500).contains(&gap_ms),
"expected gap ~1000ms, got {gap_ms}ms"
);
}
}