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,
messages_parsed: AtomicU64,
messages_relayed: AtomicU64,
messages_dropped_dup: AtomicU64,
serialization_calls: AtomicU64,
subscriber_fanout_total: AtomicU64,
ws_messages_received: AtomicU64,
ws_messages_sent: AtomicU64,
messages_dropped_ham: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
pub struct MetricsSnapshot {
pub dropped_sends: u64,
pub reaped_quorums: u64,
pub put_acks_seen: u64,
pub put_acks_quorum: u64,
pub messages_parsed: u64,
pub messages_relayed: u64,
pub messages_dropped_dup: u64,
pub serialization_calls: u64,
pub subscriber_fanout_total: u64,
pub ws_messages_received: u64,
pub ws_messages_sent: u64,
pub messages_dropped_ham: 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);
}
#[inline]
pub fn record_parsed(&self) {
self.messages_parsed.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_relayed(&self) {
self.messages_relayed.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_dropped_dup(&self) {
self.messages_dropped_dup.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_serialization(&self) {
self.serialization_calls.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_subscriber_fanout(&self, count: u64) {
self.subscriber_fanout_total
.fetch_add(count, Ordering::Relaxed);
}
#[inline]
pub fn record_ws_received(&self) {
self.ws_messages_received.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_ws_sent(&self) {
self.ws_messages_sent.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub fn record_dropped_ham(&self) {
self.messages_dropped_ham.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),
messages_parsed: self.messages_parsed.load(Ordering::Relaxed),
messages_relayed: self.messages_relayed.load(Ordering::Relaxed),
messages_dropped_dup: self.messages_dropped_dup.load(Ordering::Relaxed),
serialization_calls: self.serialization_calls.load(Ordering::Relaxed),
subscriber_fanout_total: self.subscriber_fanout_total.load(Ordering::Relaxed),
ws_messages_received: self.ws_messages_received.load(Ordering::Relaxed),
ws_messages_sent: self.ws_messages_sent.load(Ordering::Relaxed),
messages_dropped_ham: self.messages_dropped_ham.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);
assert_eq!(snap.messages_parsed, 0);
assert_eq!(snap.messages_relayed, 0);
assert_eq!(snap.messages_dropped_dup, 0);
assert_eq!(snap.serialization_calls, 0);
assert_eq!(snap.subscriber_fanout_total, 0);
assert_eq!(snap.ws_messages_received, 0);
assert_eq!(snap.ws_messages_sent, 0);
assert_eq!(snap.messages_dropped_ham, 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 record_parsed_increments_counter() {
let m = Metrics::new();
m.record_parsed();
m.record_parsed();
assert_eq!(m.snapshot().messages_parsed, 2);
}
#[test]
fn record_relayed_increments_counter() {
let m = Metrics::new();
m.record_relayed();
assert_eq!(m.snapshot().messages_relayed, 1);
}
#[test]
fn record_dropped_dup_increments_counter() {
let m = Metrics::new();
m.record_dropped_dup();
m.record_dropped_dup();
m.record_dropped_dup();
assert_eq!(m.snapshot().messages_dropped_dup, 3);
}
#[test]
fn record_dropped_ham_increments_counter() {
let m = Metrics::new();
m.record_dropped_ham();
m.record_dropped_ham();
assert_eq!(m.snapshot().messages_dropped_ham, 2);
}
#[test]
fn record_serialization_increments_counter() {
let m = Metrics::new();
m.record_serialization();
assert_eq!(m.snapshot().serialization_calls, 1);
}
#[test]
fn record_subscriber_fanout_accumulates() {
let m = Metrics::new();
m.record_subscriber_fanout(5);
m.record_subscriber_fanout(3);
m.record_subscriber_fanout(0); assert_eq!(m.snapshot().subscriber_fanout_total, 8);
}
#[test]
fn record_ws_received_increments_counter() {
let m = Metrics::new();
for _ in 0..500 {
m.record_ws_received();
}
assert_eq!(m.snapshot().ws_messages_received, 500);
}
#[test]
fn record_ws_sent_increments_counter() {
let m = Metrics::new();
m.record_ws_sent();
m.record_ws_sent();
assert_eq!(m.snapshot().ws_messages_sent, 2);
}
#[test]
fn snapshot_reflects_independent_increments() {
let m = Metrics::new();
m.record_dropped_send();
m.record_put_ack();
m.record_quorum_ack();
m.record_parsed();
m.record_relayed();
m.record_serialization();
m.record_dropped_ham();
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);
assert_eq!(snap.messages_parsed, 1);
assert_eq!(snap.messages_relayed, 1);
assert_eq!(snap.messages_dropped_dup, 0);
assert_eq!(snap.serialization_calls, 1);
assert_eq!(snap.messages_dropped_ham, 1);
}
#[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();
m.record_parsed();
assert_eq!(m2.snapshot().dropped_sends, 1);
assert_eq!(m2.snapshot().messages_parsed, 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 concurrent_hot_path_increments_are_not_lost() {
use std::thread;
let m = Arc::new(Metrics::new());
let mut handles = Vec::new();
for _ in 0..50 {
let m = Arc::clone(&m);
handles.push(thread::spawn(move || {
for _ in 0..2000 {
m.record_parsed();
}
}));
}
for h in handles {
h.join().unwrap();
}
assert_eq!(m.snapshot().messages_parsed, 100_000);
}
#[test]
fn snapshot_is_copy() {
let s = Metrics::new().snapshot();
let s2 = s; assert_eq!(s, s2);
}
#[test]
fn snapshot_serializes_to_json() {
let snap = Metrics::new().snapshot();
let json = serde_json::to_string(&snap).unwrap();
assert!(json.contains("dropped_sends"));
assert!(json.contains("messages_parsed"));
assert!(json.contains("ws_messages_sent"));
}
}