lowlet 0.1.2

Low-latency IPC library using shared memory and lock-free structures
Documentation
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Instant;

#[derive(Debug, Clone, Copy)]
pub struct HealthSnapshot {
    pub utilization: f32,
    pub send_rate: u64,
    pub recv_rate: u64,
    pub pending: usize,
    pub capacity: usize,
    pub is_saturated: bool,
}

pub struct ChannelHealth {
    sends: AtomicU64,
    recvs: AtomicU64,
    last_sends: AtomicU64,
    last_recvs: AtomicU64,
    last_check: std::sync::Mutex<Instant>,
    capacity: usize,
    pending_fn: Box<dyn Fn() -> usize + Send + Sync>,
}

impl ChannelHealth {
    pub fn new<F>(capacity: usize, pending_fn: F) -> Self
    where
        F: Fn() -> usize + Send + Sync + 'static,
    {
        Self {
            sends: AtomicU64::new(0),
            recvs: AtomicU64::new(0),
            last_sends: AtomicU64::new(0),
            last_recvs: AtomicU64::new(0),
            last_check: std::sync::Mutex::new(Instant::now()),
            capacity,
            pending_fn: Box::new(pending_fn),
        }
    }

    #[inline]
    pub fn record_send(&self) {
        self.sends.fetch_add(1, Ordering::Relaxed);
    }

    #[inline]
    pub fn record_recv(&self) {
        self.recvs.fetch_add(1, Ordering::Relaxed);
    }

    pub fn snapshot(&self) -> HealthSnapshot {
        let mut last_check = self.last_check.lock().unwrap();
        let now = Instant::now();
        let elapsed = now.duration_since(*last_check).as_secs_f64();

        let sends = self.sends.load(Ordering::Relaxed);
        let recvs = self.recvs.load(Ordering::Relaxed);
        let last_sends = self.last_sends.swap(sends, Ordering::Relaxed);
        let last_recvs = self.last_recvs.swap(recvs, Ordering::Relaxed);

        let send_rate = if elapsed > 0.0 {
            ((sends - last_sends) as f64 / elapsed) as u64
        } else {
            0
        };

        let recv_rate = if elapsed > 0.0 {
            ((recvs - last_recvs) as f64 / elapsed) as u64
        } else {
            0
        };

        let pending = (self.pending_fn)();
        let utilization = pending as f32 / self.capacity as f32;
        let is_saturated = utilization > 0.9;

        *last_check = now;

        HealthSnapshot {
            utilization,
            send_rate,
            recv_rate,
            pending,
            capacity: self.capacity,
            is_saturated,
        }
    }

    pub fn reset(&self) {
        self.sends.store(0, Ordering::Relaxed);
        self.recvs.store(0, Ordering::Relaxed);
        self.last_sends.store(0, Ordering::Relaxed);
        self.last_recvs.store(0, Ordering::Relaxed);
        *self.last_check.lock().unwrap() = Instant::now();
    }

    #[inline]
    pub fn total_sends(&self) -> u64 {
        self.sends.load(Ordering::Relaxed)
    }

    #[inline]
    pub fn total_recvs(&self) -> u64 {
        self.recvs.load(Ordering::Relaxed)
    }
}