use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Debug, Default)]
pub struct Metrics {
dropped_sends: AtomicU64,
reaped_quorums: AtomicU64,
put_acks_seen: AtomicU64,
put_acks_quorum: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MetricsSnapshot {
pub dropped_sends: u64,
pub reaped_quorums: u64,
pub put_acks_seen: u64,
pub put_acks_quorum: u64,
}
impl Metrics {
pub fn new() -> Self {
Self::default()
}
#[inline]
pub fn record_dropped_send(&self) {
self.dropped_sends.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_reaped_quorum(&self) {
self.reaped_quorums.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_put_ack(&self) {
self.put_acks_seen.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_quorum_ack(&self) {
self.put_acks_quorum.fetch_add(1, Ordering::Relaxed);
}
pub fn snapshot(&self) -> MetricsSnapshot {
MetricsSnapshot {
dropped_sends: self.dropped_sends.load(Ordering::Relaxed),
reaped_quorums: self.reaped_quorums.load(Ordering::Relaxed),
put_acks_seen: self.put_acks_seen.load(Ordering::Relaxed),
put_acks_quorum: self.put_acks_quorum.load(Ordering::Relaxed),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
#[test]
fn default_is_all_zero() {
let snap = Metrics::new().snapshot();
assert_eq!(snap.dropped_sends, 0);
assert_eq!(snap.reaped_quorums, 0);
assert_eq!(snap.put_acks_seen, 0);
assert_eq!(snap.put_acks_quorum, 0);
}
#[test]
fn record_dropped_send_increments_counter() {
let m = Metrics::new();
m.record_dropped_send();
m.record_dropped_send();
m.record_dropped_send();
assert_eq!(m.snapshot().dropped_sends, 3);
}
#[test]
fn record_reaped_quorum_increments_counter() {
let m = Metrics::new();
m.record_reaped_quorum();
assert_eq!(m.snapshot().reaped_quorums, 1);
}
#[test]
fn record_put_ack_increments_counter() {
let m = Metrics::new();
m.record_put_ack();
m.record_put_ack();
assert_eq!(m.snapshot().put_acks_seen, 2);
}
#[test]
fn record_quorum_ack_increments_counter() {
let m = Metrics::new();
m.record_quorum_ack();
assert_eq!(m.snapshot().put_acks_quorum, 1);
}
#[test]
fn snapshot_reflects_independent_increments() {
let m = Metrics::new();
m.record_dropped_send();
m.record_put_ack();
m.record_quorum_ack();
let snap = m.snapshot();
assert_eq!(snap.dropped_sends, 1);
assert_eq!(snap.put_acks_seen, 1);
assert_eq!(snap.put_acks_quorum, 1);
assert_eq!(snap.reaped_quorums, 0);
}
#[test]
fn counters_are_monotonic() {
let m = Metrics::new();
for i in 1..=1000 {
m.record_dropped_send();
assert_eq!(m.snapshot().dropped_sends, i);
}
}
#[test]
fn shared_metrics_via_arc() {
let m: Arc<Metrics> = Arc::new(Metrics::new());
let m2 = Arc::clone(&m);
m.record_dropped_send();
assert_eq!(m2.snapshot().dropped_sends, 1);
}
#[test]
fn concurrent_increments_are_not_lost() {
use std::thread;
let m = Arc::new(Metrics::new());
let mut handles = Vec::new();
for _ in 0..100 {
let m = Arc::clone(&m);
handles.push(thread::spawn(move || {
for _ in 0..1000 {
m.record_dropped_send();
}
}));
}
for h in handles {
h.join().unwrap();
}
assert_eq!(m.snapshot().dropped_sends, 100_000);
}
#[test]
fn snapshot_is_copy() {
let s = Metrics::new().snapshot();
let s2 = s; assert_eq!(s, s2);
}
}