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::cluster_engine::ClusterEngine;
use crate::container_engine::ContainerEngine;
use crate::containers::ContainersSnapshot;
use crate::gpu::GpusSnapshot;
use crate::gpu_engine::GpuEngine;
use crate::kube::KubeSnapshot;
use crate::system::SystemSnapshot;
const CONTAINER_INTERVAL: Duration = Duration::from_secs(2);
const CLUSTER_INTERVAL: Duration = Duration::from_secs(5);
const GPU_INTERVAL: Duration = Duration::from_secs(1);
pub struct Collector {
sys: sysinfo::System,
networks: sysinfo::Networks,
interval: Duration,
container_engine: Option<Arc<dyn ContainerEngine + Send + Sync>>,
cluster_engine: Option<Arc<dyn ClusterEngine + Send + Sync>>,
gpu_engine: Option<Arc<dyn GpuEngine + Send + Sync>>,
}
impl Collector {
pub fn new(interval: Duration) -> Self {
Self::with_all_engines(interval, None, None, None)
}
pub fn with_container_engine(
interval: Duration,
container_engine: Option<Arc<dyn ContainerEngine + Send + Sync>>,
) -> Self {
Self::with_all_engines(interval, container_engine, None, None)
}
pub fn with_engines(
interval: Duration,
container_engine: Option<Arc<dyn ContainerEngine + Send + Sync>>,
cluster_engine: Option<Arc<dyn ClusterEngine + Send + Sync>>,
) -> Self {
Self::with_all_engines(interval, container_engine, cluster_engine, None)
}
pub fn with_all_engines(
interval: Duration,
container_engine: Option<Arc<dyn ContainerEngine + Send + Sync>>,
cluster_engine: Option<Arc<dyn ClusterEngine + Send + Sync>>,
gpu_engine: Option<Arc<dyn GpuEngine + Send + Sync>>,
) -> Self {
Self {
sys: sysinfo::System::new_all(),
networks: sysinfo::Networks::new_with_refreshed_list(),
interval,
container_engine,
cluster_engine,
gpu_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 initial_kube = if self.cluster_engine.is_some() {
None
} else {
Some(KubeSnapshot::unavailable())
};
let last_kube: Arc<Mutex<Option<KubeSnapshot>>> = Arc::new(Mutex::new(initial_kube));
let initial_gpu = if self.gpu_engine.is_some() {
None
} else {
Some(GpusSnapshot::unavailable())
};
let last_gpu: Arc<Mutex<Option<GpusSnapshot>>> = Arc::new(Mutex::new(initial_gpu));
let container_task = self.container_engine.take().map(|engine| {
spawn_container_loop(engine, Arc::clone(&last_containers), token.clone())
});
let cluster_task = self
.cluster_engine
.take()
.map(|engine| spawn_cluster_loop(engine, Arc::clone(&last_kube), token.clone()));
let gpu_task = self
.gpu_engine
.take()
.map(|engine| spawn_gpu_loop(engine, Arc::clone(&last_gpu), 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 kube = last_kube.lock().await.clone();
let gpu = last_gpu.lock().await.clone();
let snapshot =
SystemSnapshot::collect(&self.sys, &self.networks, containers, kube, gpu);
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;
}
if let Some(handle) = cluster_task {
let _ = handle.await;
}
if let Some(handle) = gpu_task {
let _ = handle.await;
}
}
}
fn spawn_gpu_loop(
engine: Arc<dyn GpuEngine + Send + Sync>,
slot: Arc<Mutex<Option<GpusSnapshot>>>,
token: CancellationToken,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut interval = tokio::time::interval(GPU_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = interval.tick() => {
match engine.snapshot().await {
Ok(snapshot) => {
*slot.lock().await = Some(snapshot);
}
Err(err) => {
tracing::warn!(target: "muxtop::gpu", error = %err, "GPU engine failed");
*slot.lock().await =
Some(GpusSnapshot::unavailable_with(err.to_string()));
}
}
}
_ = token.cancelled() => {
tracing::debug!("gpu loop shutting down");
break;
}
}
}
})
}
fn spawn_cluster_loop(
engine: Arc<dyn ClusterEngine + Send + Sync>,
slot: Arc<Mutex<Option<KubeSnapshot>>>,
token: CancellationToken,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut interval = tokio::time::interval(CLUSTER_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = interval.tick() => {
match engine.snapshot().await {
Ok(snapshot) => {
*slot.lock().await = Some(snapshot);
}
Err(err) => {
tracing::warn!(error = %err, "cluster engine failed");
*slot.lock().await = Some(KubeSnapshot::unavailable());
}
}
}
_ = token.cancelled() => {
tracing::debug!("cluster loop shutting down");
break;
}
}
}
})
}
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);
}
struct MockGpuEngine {
call_count: AtomicUsize,
fail_mode: AtomicBool,
}
impl MockGpuEngine {
fn new() -> Self {
Self {
call_count: AtomicUsize::new(0),
fail_mode: AtomicBool::new(false),
}
}
fn sample_device() -> crate::gpu::GpuDeviceSnapshot {
crate::gpu::GpuDeviceSnapshot {
index: 0,
vendor: crate::gpu::GpuVendor::Nvidia,
backend: crate::gpu::GpuBackend::Nvml,
name: "Mock GPU".into(),
bus_id: "0000:01:00.0".into(),
driver_version: Some("999.99".into()),
utilization_pct: Some(50.0),
mem_utilization_pct: Some(20.0),
mem_used_bytes: Some(1024),
mem_total_bytes: Some(4096),
temperature_c: Some(55.0),
power_watts: Some(100.0),
power_limit_watts: Some(200.0),
graphics_clock_mhz: Some(1800),
memory_clock_mhz: Some(7000),
fan_pct: Some(30.0),
encoder_pct: None,
decoder_pct: None,
supports_process_stats: true,
}
}
}
#[async_trait]
impl crate::gpu_engine::GpuEngine for MockGpuEngine {
async fn snapshot(&self) -> Result<crate::gpu::GpusSnapshot, crate::gpu_engine::GpuError> {
self.call_count.fetch_add(1, Ordering::Relaxed);
if self.fail_mode.load(Ordering::Relaxed) {
return Err(crate::gpu_engine::GpuError::Query("mock failure".into()));
}
Ok(crate::gpu::GpusSnapshot {
backends: vec![crate::gpu::GpuBackend::Nvml],
available: true,
devices: vec![Self::sample_device()],
processes: vec![crate::gpu::GpuProcessSnapshot {
pid: std::process::id(),
device_index: 0,
name: String::new(),
kind: crate::gpu::GpuProcessKind::Compute,
mem_bytes: Some(512),
}],
detail: String::new(),
})
}
fn backend(&self) -> crate::gpu::GpuBackend {
crate::gpu::GpuBackend::Nvml
}
}
async fn drain_until(
rx: &mut mpsc::Receiver<SystemSnapshot>,
secs: u64,
pred: impl Fn(&SystemSnapshot) -> bool,
) -> Option<SystemSnapshot> {
let deadline = tokio::time::Instant::now() + Duration::from_secs(secs);
loop {
match tokio::time::timeout_at(deadline, rx.recv()).await {
Ok(Some(snap)) if pred(&snap) => return Some(snap),
Ok(Some(_)) => continue,
Ok(None) | Err(_) => return None,
}
}
}
#[tokio::test]
async fn test_collector_populates_gpu_with_engine() {
let engine: Arc<dyn crate::gpu_engine::GpuEngine + Send + Sync> =
Arc::new(MockGpuEngine::new());
let (tx, mut rx) = mpsc::channel(8);
let token = CancellationToken::new();
let collector =
Collector::with_all_engines(Duration::from_millis(300), None, None, Some(engine));
let handle = collector.spawn(tx, token.clone());
let populated = drain_until(&mut rx, 5, |s| s.gpu.is_some()).await;
token.cancel();
let _ = tokio::time::timeout(Duration::from_secs(3), handle).await;
let snap = populated.expect("no snapshot with gpu populated within 5 s");
let gpu = snap.gpu.unwrap();
assert!(gpu.available);
assert_eq!(gpu.devices.len(), 1);
assert_eq!(gpu.devices[0].name, "Mock GPU");
}
#[tokio::test]
async fn test_collector_resolves_gpu_process_names() {
let engine: Arc<dyn crate::gpu_engine::GpuEngine + Send + Sync> =
Arc::new(MockGpuEngine::new());
let (tx, mut rx) = mpsc::channel(8);
let token = CancellationToken::new();
let collector =
Collector::with_all_engines(Duration::from_millis(300), None, None, Some(engine));
let handle = collector.spawn(tx, token.clone());
let populated = drain_until(&mut rx, 5, |s| {
s.gpu.as_ref().is_some_and(|g| !g.processes.is_empty())
})
.await;
token.cancel();
let _ = tokio::time::timeout(Duration::from_secs(3), handle).await;
let snap = populated.expect("no snapshot with gpu processes within 5 s");
let gpu = snap.gpu.unwrap();
assert!(
!gpu.processes[0].name.is_empty(),
"collector should have resolved the PID to a name"
);
}
#[tokio::test]
async fn test_collector_publishes_gpu_unavailable_on_engine_error() {
let engine = Arc::new(MockGpuEngine::new());
engine.fail_mode.store(true, Ordering::Relaxed);
let engine_dyn: Arc<dyn crate::gpu_engine::GpuEngine + Send + Sync> = engine.clone();
let (tx, mut rx) = mpsc::channel(8);
let token = CancellationToken::new();
let collector =
Collector::with_all_engines(Duration::from_millis(300), None, None, Some(engine_dyn));
let handle = collector.spawn(tx, token.clone());
let observed =
drain_until(&mut rx, 5, |s| s.gpu.as_ref().is_some_and(|g| !g.available)).await;
token.cancel();
let _ = tokio::time::timeout(Duration::from_secs(3), handle).await;
let snap = observed.expect("no gpu-populated snapshot within 5 s");
let gpu = snap.gpu.unwrap();
assert!(!gpu.available);
assert!(gpu.devices.is_empty());
assert!(
gpu.detail.contains("mock failure"),
"the failure reason should reach the UI, got {:?}",
gpu.detail
);
assert!(engine.call_count.load(Ordering::Relaxed) >= 1);
}
#[tokio::test]
async fn test_collector_without_gpu_engine_seeds_unavailable() {
let (mut rx, handle, token) = make_collector(4);
let snap = 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");
let gpu = snap.gpu.expect("gpu slot should be seeded, not None");
assert!(!gpu.available);
assert!(gpu.devices.is_empty());
}
#[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");
}
}