use std::sync::Arc;
use std::time::Duration;
use sysinfo::{MemoryRefreshKind, ProcessRefreshKind, ProcessesToUpdate};
use tokio::sync::Mutex;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use crate::container_engine::ContainerEngine;
use crate::containers::ContainersSnapshot;
use crate::system::SystemSnapshot;
const CONTAINER_INTERVAL: Duration = Duration::from_secs(2);
pub struct Collector {
sys: sysinfo::System,
networks: sysinfo::Networks,
interval: Duration,
container_engine: Option<Arc<dyn ContainerEngine + Send + Sync>>,
}
impl Collector {
pub fn new(interval: Duration) -> Self {
Self::with_container_engine(interval, None)
}
pub fn with_container_engine(
interval: Duration,
container_engine: Option<Arc<dyn ContainerEngine + Send + Sync>>,
) -> Self {
Self {
sys: sysinfo::System::new_all(),
networks: sysinfo::Networks::new_with_refreshed_list(),
interval,
container_engine,
}
}
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);
let last_containers: Arc<Mutex<Option<ContainersSnapshot>>> = Arc::new(Mutex::new(None));
let container_task = self.container_engine.take().map(|engine| {
spawn_container_loop(engine, Arc::clone(&last_containers), token.clone())
});
loop {
tokio::select! {
_ = interval.tick() => {
self.sys.refresh_memory_specifics(MemoryRefreshKind::everything());
self.sys.refresh_cpu_usage();
self.sys.refresh_processes_specifics(
ProcessesToUpdate::All,
true,
ProcessRefreshKind::everything(),
);
self.networks.refresh(false);
let containers = last_containers.lock().await.clone();
let snapshot =
SystemSnapshot::collect(&self.sys, &self.networks, containers);
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;
}
}
}
if let Some(handle) = container_task {
let _ = handle.await;
}
}
}
fn spawn_container_loop(
engine: Arc<dyn ContainerEngine + Send + Sync>,
slot: Arc<Mutex<Option<ContainersSnapshot>>>,
token: CancellationToken,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut interval = tokio::time::interval(CONTAINER_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = interval.tick() => {
match engine.list_and_stats().await {
Ok(containers) => {
let snapshot = ContainersSnapshot {
engine: engine.kind(),
daemon_up: true,
containers,
};
*slot.lock().await = Some(snapshot);
}
Err(err) => {
tracing::warn!(error = %err, "container engine failed");
*slot.lock().await = Some(ContainersSnapshot::unavailable());
}
}
}
_ = token.cancelled() => {
tracing::debug!("container loop shutting down");
break;
}
}
}
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::container_engine::EngineError;
use crate::containers::{ContainerSnapshot, ContainerState, EngineKind};
use async_trait::async_trait;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::time::Duration;
struct MockEngine {
kind: EngineKind,
call_count: AtomicUsize,
fail_mode: AtomicBool,
}
impl MockEngine {
fn new(kind: EngineKind) -> Self {
Self {
kind,
call_count: AtomicUsize::new(0),
fail_mode: AtomicBool::new(false),
}
}
fn sample_container() -> ContainerSnapshot {
ContainerSnapshot {
id: "abc123".into(),
id_full: "abc123".to_string() + &"0".repeat(58),
name: "mock-svc".into(),
image: "mock:latest".into(),
state: ContainerState::Running,
status_text: "Up 1 minute".into(),
cpu_pct: 2.5,
mem_used_bytes: 128 * 1024 * 1024,
mem_limit_bytes: 512 * 1024 * 1024,
net_rx_bytes: 1024,
net_tx_bytes: 512,
block_read_bytes: 0,
block_write_bytes: 0,
started_at_ms: 1_700_000_000_000,
}
}
}
#[async_trait]
impl ContainerEngine for MockEngine {
async fn list_and_stats(&self) -> Result<Vec<ContainerSnapshot>, EngineError> {
self.call_count.fetch_add(1, Ordering::Relaxed);
if self.fail_mode.load(Ordering::Relaxed) {
return Err(EngineError::ConnectFailed("mock failure".into()));
}
Ok(vec![Self::sample_container()])
}
async fn stop(&self, _id: &str, _t: Option<u64>) -> Result<(), EngineError> {
Ok(())
}
async fn kill(&self, _id: &str) -> Result<(), EngineError> {
Ok(())
}
async fn restart(&self, _id: &str) -> Result<(), EngineError> {
Ok(())
}
fn kind(&self) -> EngineKind {
self.kind
}
}
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"
);
assert!(snapshot.containers.is_none());
}
#[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"
);
}
#[tokio::test]
async fn test_collector_populates_containers_with_engine() {
let engine: Arc<dyn ContainerEngine + Send + Sync> =
Arc::new(MockEngine::new(EngineKind::Docker));
let (tx, mut rx) = mpsc::channel(8);
let token = CancellationToken::new();
let collector = Collector::with_container_engine(Duration::from_millis(300), Some(engine));
let handle = collector.spawn(tx, token.clone());
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
let mut populated: Option<SystemSnapshot> = None;
loop {
match tokio::time::timeout_at(deadline, rx.recv()).await {
Ok(Some(snap)) => {
if snap.containers.is_some() {
populated = Some(snap);
break;
}
}
Ok(None) => break,
Err(_) => break,
}
}
token.cancel();
let _ = tokio::time::timeout(Duration::from_secs(3), handle).await;
let snap = populated.expect("no snapshot with containers populated within 5 s");
let cs = snap.containers.unwrap();
assert!(cs.daemon_up);
assert_eq!(cs.engine, EngineKind::Docker);
assert_eq!(cs.containers.len(), 1);
assert_eq!(cs.containers[0].name, "mock-svc");
}
#[tokio::test]
async fn test_collector_publishes_unavailable_on_engine_error() {
let engine = Arc::new(MockEngine::new(EngineKind::Docker));
engine.fail_mode.store(true, Ordering::Relaxed);
let engine_dyn: Arc<dyn ContainerEngine + Send + Sync> = engine.clone();
let (tx, mut rx) = mpsc::channel(8);
let token = CancellationToken::new();
let collector =
Collector::with_container_engine(Duration::from_millis(300), Some(engine_dyn));
let handle = collector.spawn(tx, token.clone());
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
let mut observed: Option<SystemSnapshot> = None;
loop {
match tokio::time::timeout_at(deadline, rx.recv()).await {
Ok(Some(snap)) => {
if snap.containers.is_some() {
observed = Some(snap);
break;
}
}
Ok(None) => break,
Err(_) => break,
}
}
token.cancel();
let _ = tokio::time::timeout(Duration::from_secs(3), handle).await;
let snap = observed.expect("no containers-populated snapshot within 5 s");
let cs = snap.containers.unwrap();
assert!(!cs.daemon_up);
assert!(cs.containers.is_empty());
assert_eq!(cs.engine, EngineKind::Unknown);
assert!(engine.call_count.load(Ordering::Relaxed) >= 1);
}
#[tokio::test]
async fn test_collector_without_engine_keeps_containers_none() {
let (mut rx, handle, token) = make_collector(4);
for _ in 0..3 {
if let Some(snap) = tokio::time::timeout(Duration::from_secs(4), rx.recv())
.await
.expect("timeout")
{
assert!(
snap.containers.is_none(),
"unexpected containers snapshot without engine"
);
} else {
break;
}
}
token.cancel();
handle.await.expect("collector task panicked");
}
}