muxtop-core 0.3.1

Core data collection engine for muxtop
Documentation
use std::sync::OnceLock;
use std::time::{SystemTime, UNIX_EPOCH};

use bincode::{Decode, Encode};
use serde::{Deserialize, Serialize};

use crate::containers::ContainersSnapshot;
use crate::network::NetworkSnapshot;
use crate::process::ProcessInfo;

/// Interned per-core name table (PERF-L2).
///
/// `format!("cpu{i}")` previously ran on every collector tick for every core,
/// burning a few hundred small allocations per second on multi-core hosts.
/// The names never change for the lifetime of the process, so we mint them
/// once and clone the cached `String` on each tick (a `String::clone` for a
/// 4-6-byte payload reuses the small-string fast path on most allocators).
fn core_name(i: usize) -> String {
    static CORE_NAMES: OnceLock<std::sync::RwLock<Vec<String>>> = OnceLock::new();
    let lock = CORE_NAMES.get_or_init(|| std::sync::RwLock::new(Vec::new()));

    // Fast path: read lock + index lookup.
    {
        let table = lock.read().unwrap_or_else(|e| e.into_inner());
        if let Some(name) = table.get(i) {
            return name.clone();
        }
    }

    // Slow path: extend the table (rare — only on the first tick for that core
    // count, and on the first run a single contiguous extend).
    let mut table = lock.write().unwrap_or_else(|e| e.into_inner());
    while table.len() <= i {
        let idx = table.len();
        table.push(format!("cpu{idx}"));
    }
    table[i].clone()
}

/// Per-core CPU snapshot.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Encode, Decode)]
pub struct CoreSnapshot {
    pub name: String,
    pub usage: f32,
    pub frequency: u64,
}

/// Aggregated CPU snapshot with global usage and per-core data.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Encode, Decode)]
pub struct CpuSnapshot {
    pub global_usage: f32,
    pub cores: Vec<CoreSnapshot>,
}

/// Memory and swap snapshot (all values in bytes).
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Encode, Decode)]
pub struct MemorySnapshot {
    pub total: u64,
    pub used: u64,
    pub available: u64,
    pub swap_total: u64,
    pub swap_used: u64,
}

/// System load averages and uptime.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Encode, Decode)]
pub struct LoadSnapshot {
    pub one: f64,
    pub five: f64,
    pub fifteen: f64,
    pub uptime_secs: u64,
}

/// Full system snapshot aggregating all subsystems.
///
/// `containers` is `None` whenever the collector runs without a container
/// engine attached, or before the first container tick has completed.
/// Once set, it is `Some(ContainersSnapshot::unavailable())` to report
/// engine failure or `Some(..)` with the current fleet.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, Encode, Decode)]
pub struct SystemSnapshot {
    pub cpu: CpuSnapshot,
    pub memory: MemorySnapshot,
    pub load: LoadSnapshot,
    pub processes: Vec<ProcessInfo>,
    pub networks: NetworkSnapshot,
    pub containers: Option<ContainersSnapshot>,
    /// Milliseconds since Unix epoch.
    pub timestamp_ms: u64,
}

impl SystemSnapshot {
    /// Collect a full system snapshot from sysinfo.
    ///
    /// `containers` is passed through verbatim: callers either supply the
    /// latest snapshot from their container engine (via [`crate::Collector`])
    /// or `None` when running without one.
    pub fn collect(
        sys: &sysinfo::System,
        networks: &sysinfo::Networks,
        containers: Option<ContainersSnapshot>,
    ) -> Self {
        use sysinfo::System as SysSystem;

        let global_usage = sys.global_cpu_usage();
        let cores = sys
            .cpus()
            .iter()
            .enumerate()
            .map(|(i, cpu)| CoreSnapshot {
                name: core_name(i),
                usage: cpu.cpu_usage(),
                frequency: cpu.frequency(),
            })
            .collect();

        let cpu = CpuSnapshot {
            global_usage,
            cores,
        };

        let memory = MemorySnapshot {
            total: sys.total_memory(),
            used: sys.used_memory(),
            available: sys.available_memory(),
            swap_total: sys.total_swap(),
            swap_used: sys.used_swap(),
        };

        let load = {
            let avg = SysSystem::load_average();
            LoadSnapshot {
                one: avg.one,
                five: avg.five,
                fifteen: avg.fifteen,
                uptime_secs: SysSystem::uptime(),
            }
        };

        let total_mem = sys.total_memory();

        let processes = sys
            .processes()
            .iter()
            .map(|(pid, proc_info)| {
                let mem_pct = if total_mem > 0 {
                    ((proc_info.memory() as f64 / total_mem as f64) * 100.0).clamp(0.0, 100.0)
                        as f32
                } else {
                    0.0
                };

                let status = match proc_info.status() {
                    sysinfo::ProcessStatus::Run => "Running",
                    sysinfo::ProcessStatus::Sleep => "Sleeping",
                    sysinfo::ProcessStatus::Idle => "Idle",
                    sysinfo::ProcessStatus::Zombie => "Zombie",
                    sysinfo::ProcessStatus::Stop => "Stopped",
                    _ => "Unknown",
                };

                ProcessInfo {
                    pid: pid.as_u32(),
                    parent_pid: proc_info.parent().map(|p| p.as_u32()),
                    name: proc_info.name().to_string_lossy().into_owned(),
                    command: proc_info
                        .cmd()
                        .iter()
                        .map(|s| s.to_string_lossy().into_owned())
                        .collect::<Vec<_>>()
                        .join(" "),
                    user: proc_info
                        .user_id()
                        .map(|u| u.to_string())
                        .unwrap_or_default(),
                    cpu_percent: proc_info.cpu_usage(),
                    memory_bytes: proc_info.memory(),
                    memory_percent: mem_pct,
                    status: status.to_string(),
                }
            })
            .collect();

        let networks = NetworkSnapshot::collect(networks);

        Self {
            cpu,
            memory,
            load,
            processes,
            networks,
            containers,
            timestamp_ms: SystemTime::now()
                .duration_since(UNIX_EPOCH)
                .expect("system clock before Unix epoch")
                .as_millis() as u64,
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn test_snapshot_is_send_clone() {
        fn assert_send_clone<T: Send + Clone>() {}
        assert_send_clone::<CoreSnapshot>();
        assert_send_clone::<CpuSnapshot>();
        assert_send_clone::<MemorySnapshot>();
        assert_send_clone::<LoadSnapshot>();
        assert_send_clone::<SystemSnapshot>();
        assert_send_clone::<ProcessInfo>();
    }

    #[test]
    fn test_system_snapshot_from_sysinfo() {
        use sysinfo::System;
        let mut sys = System::new_all();
        // sysinfo needs a refresh to populate CPU data
        std::thread::sleep(std::time::Duration::from_millis(200));
        sys.refresh_all();

        let networks = sysinfo::Networks::new_with_refreshed_list();
        let snap = SystemSnapshot::collect(&sys, &networks, None);
        assert!(!snap.cpu.cores.is_empty(), "should have CPU cores");
        assert!(snap.memory.total > 0, "should have total memory");
        assert!(!snap.processes.is_empty(), "should have processes");
    }

    #[test]
    fn test_cpu_snapshot_has_global_and_cores() {
        use sysinfo::System;
        let mut sys = System::new_all();
        std::thread::sleep(std::time::Duration::from_millis(200));
        sys.refresh_all();

        let networks = sysinfo::Networks::new_with_refreshed_list();
        let snap = SystemSnapshot::collect(&sys, &networks, None);
        assert!(
            snap.cpu.global_usage >= 0.0 && snap.cpu.global_usage <= 100.0,
            "global CPU usage should be 0..=100, got {}",
            snap.cpu.global_usage
        );
        for core in &snap.cpu.cores {
            assert!(
                core.usage >= 0.0 && core.usage <= 100.0,
                "core usage should be 0..=100, got {}",
                core.usage
            );
        }
    }

    #[test]
    fn test_memory_snapshot_invariant() {
        use sysinfo::System;
        let mut sys = System::new_all();
        sys.refresh_all();

        let networks = sysinfo::Networks::new_with_refreshed_list();
        let snap = SystemSnapshot::collect(&sys, &networks, None);
        assert!(snap.memory.total > 0, "total memory should be positive");
        // used + available can slightly exceed total due to kernel accounting
        // but total should be >= used
        assert!(
            snap.memory.total >= snap.memory.used,
            "total ({}) should be >= used ({})",
            snap.memory.total,
            snap.memory.used
        );
    }

    #[test]
    fn test_system_snapshot_has_networks() {
        use sysinfo::System;
        let mut sys = System::new_all();
        std::thread::sleep(std::time::Duration::from_millis(200));
        sys.refresh_all();
        let networks = sysinfo::Networks::new_with_refreshed_list();

        let snap = SystemSnapshot::collect(&sys, &networks, None);
        assert!(
            !snap.networks.interfaces.is_empty(),
            "should have network interfaces"
        );
        assert_eq!(
            snap.networks.total_rx,
            snap.networks
                .interfaces
                .iter()
                .map(|i| i.bytes_rx)
                .sum::<u64>(),
            "total_rx should be consistent"
        );
    }

    #[test]
    fn test_core_name_returns_stable_label() {
        // PERF-L2: the interned table returns the canonical label and is
        // stable across calls (we don't observe `Arc` identity since the
        // interface returns owned `String`s, but the values must match the
        // legacy `format!("cpu{i}")` exactly so wire-format consumers don't
        // see a regression).
        assert_eq!(core_name(0), "cpu0");
        assert_eq!(core_name(7), "cpu7");
        assert_eq!(core_name(0), "cpu0"); // hits the fast path
    }

    #[test]
    fn test_all_structs_are_debug() {
        let core = CoreSnapshot {
            name: "cpu0".into(),
            usage: 50.0,
            frequency: 3600,
        };
        assert!(!format!("{core:?}").is_empty());

        let cpu = CpuSnapshot {
            global_usage: 25.0,
            cores: vec![core],
        };
        assert!(!format!("{cpu:?}").is_empty());

        let mem = MemorySnapshot {
            total: 16_000_000_000,
            used: 8_000_000_000,
            available: 8_000_000_000,
            swap_total: 4_000_000_000,
            swap_used: 1_000_000_000,
        };
        assert!(!format!("{mem:?}").is_empty());

        let load = LoadSnapshot {
            one: 1.5,
            five: 1.2,
            fifteen: 0.8,
            uptime_secs: 3600,
        };
        assert!(!format!("{load:?}").is_empty());
    }
}