lowlet 0.1.2

Low-latency IPC library using shared memory and lock-free structures
Documentation
mod health;
mod histogram;

pub use health::{ChannelHealth, HealthSnapshot};
pub use histogram::{HistogramSnapshot, LatencyHistogram};

use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};

pub struct ChannelStats {
    pub sends: AtomicU64,
    pub recvs: AtomicU64,
    pub failed_sends: AtomicU64,
    pub failed_recvs: AtomicU64,
}

impl ChannelStats {
    pub const fn new() -> Self {
        Self {
            sends: AtomicU64::new(0),
            recvs: AtomicU64::new(0),
            failed_sends: AtomicU64::new(0),
            failed_recvs: AtomicU64::new(0),
        }
    }

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

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

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

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

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

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

    pub fn reset(&self) {
        self.sends.store(0, Ordering::Relaxed);
        self.recvs.store(0, Ordering::Relaxed);
        self.failed_sends.store(0, Ordering::Relaxed);
        self.failed_recvs.store(0, Ordering::Relaxed);
    }
}

impl Default for ChannelStats {
    fn default() -> Self {
        Self::new()
    }
}

pub struct Metrics {
    total_bytes: AtomicUsize,
    peak_bytes: AtomicUsize,
}

impl Metrics {
    pub const fn new() -> Self {
        Self {
            total_bytes: AtomicUsize::new(0),
            peak_bytes: AtomicUsize::new(0),
        }
    }

    #[inline]
    pub fn record_allocation(&self, size: usize) {
        let total = self.total_bytes.fetch_add(size, Ordering::Relaxed) + size;
        let mut peak = self.peak_bytes.load(Ordering::Relaxed);
        while total > peak {
            match self.peak_bytes.compare_exchange_weak(
                peak,
                total,
                Ordering::Relaxed,
                Ordering::Relaxed,
            ) {
                Ok(_) => break,
                Err(p) => peak = p,
            }
        }
    }

    #[inline]
    pub fn total_bytes(&self) -> usize {
        self.total_bytes.load(Ordering::Relaxed)
    }

    #[inline]
    pub fn peak_bytes(&self) -> usize {
        self.peak_bytes.load(Ordering::Relaxed)
    }

    pub fn reset(&self) {
        self.total_bytes.store(0, Ordering::Relaxed);
        self.peak_bytes.store(0, Ordering::Relaxed);
    }
}

impl Default for Metrics {
    fn default() -> Self {
        Self::new()
    }
}

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

    #[test]
    fn test_channel_stats() {
        let stats = ChannelStats::new();
        stats.record_send(true);
        stats.record_send(false);
        stats.record_recv(true);

        assert_eq!(stats.get_sends(), 1);
        assert_eq!(stats.get_failed_sends(), 1);
        assert_eq!(stats.get_recvs(), 1);

        stats.reset();
        assert_eq!(stats.get_sends(), 0);
    }

    #[test]
    fn test_metrics() {
        let metrics = Metrics::new();
        metrics.record_allocation(100);
        metrics.record_allocation(200);

        assert_eq!(metrics.total_bytes(), 300);
        assert_eq!(metrics.peak_bytes(), 300);
    }

    #[test]
    fn test_histogram() {
        let hist = LatencyHistogram::<64>::new();
        for i in 0..100 {
            hist.record(i);
        }

        assert_eq!(hist.count(), 100);
        assert_eq!(hist.min(), 0);
        assert_eq!(hist.max(), 99);
    }
}