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);
}
}