Skip to main content

subms_spsc_ring_buffer/features/
metrics.rs

1//! Per-instance metrics wrapper.
2//!
3//! Wraps a base `Producer` / `Consumer` pair behind `try_push` / `try_pop`
4//! and records:
5//!
6//! - `enqueue_success`, `enqueue_fail` (ring full)
7//! - `dequeue_success`, `dequeue_fail` (ring empty)
8//! - `max_depth_observed` (gauge; high-water mark of in-flight items)
9//! - `cas_retries` (only set by the mpmc-disruptor path; SPSC never CAS-retries)
10//!
11//! Counter overhead is one `fetch_add` per op - measurable on a hot loop, so
12//! reach for this when you're explicitly capturing operational stats, not
13//! when you need the absolute lowest latency.
14
15use std::sync::Arc;
16use std::sync::atomic::{AtomicU64, Ordering};
17
18use crate::{Consumer, Producer};
19
20/// Shared counter set. Hold an `Arc<RingMetrics>` across producer + consumer
21/// so the snapshot reflects both sides.
22pub struct RingMetrics {
23    enqueue_success: AtomicU64,
24    enqueue_fail: AtomicU64,
25    dequeue_success: AtomicU64,
26    dequeue_fail: AtomicU64,
27    max_depth_observed: AtomicU64,
28    cas_retries: AtomicU64,
29}
30
31impl Default for RingMetrics {
32    fn default() -> Self {
33        Self::new()
34    }
35}
36
37impl RingMetrics {
38    pub fn new() -> Self {
39        Self {
40            enqueue_success: AtomicU64::new(0),
41            enqueue_fail: AtomicU64::new(0),
42            dequeue_success: AtomicU64::new(0),
43            dequeue_fail: AtomicU64::new(0),
44            max_depth_observed: AtomicU64::new(0),
45            cas_retries: AtomicU64::new(0),
46        }
47    }
48
49    /// Bump the CAS-retry counter. Called by the mpmc-disruptor path; the
50    /// SPSC paths never retry under the wait-free invariant.
51    pub fn record_cas_retry(&self) {
52        self.cas_retries.fetch_add(1, Ordering::Relaxed);
53    }
54
55    /// Capture a point-in-time snapshot of every counter. Consistent enough
56    /// for monitoring; not a transactional read across all six values.
57    pub fn snapshot(&self) -> RingMetricsSnapshot {
58        RingMetricsSnapshot {
59            enqueue_success: self.enqueue_success.load(Ordering::Relaxed),
60            enqueue_fail: self.enqueue_fail.load(Ordering::Relaxed),
61            dequeue_success: self.dequeue_success.load(Ordering::Relaxed),
62            dequeue_fail: self.dequeue_fail.load(Ordering::Relaxed),
63            max_depth_observed: self.max_depth_observed.load(Ordering::Relaxed),
64            cas_retries: self.cas_retries.load(Ordering::Relaxed),
65        }
66    }
67
68    fn observe_depth(&self, d: u64) {
69        let mut cur = self.max_depth_observed.load(Ordering::Relaxed);
70        while d > cur {
71            match self.max_depth_observed.compare_exchange_weak(
72                cur,
73                d,
74                Ordering::Relaxed,
75                Ordering::Relaxed,
76            ) {
77                Ok(_) => break,
78                Err(latest) => cur = latest,
79            }
80        }
81    }
82}
83
84/// Snapshot value-type returned by `RingMetrics::snapshot`. `Clone + Copy` so
85/// callers can stash + compare across calls.
86#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
87pub struct RingMetricsSnapshot {
88    pub enqueue_success: u64,
89    pub enqueue_fail: u64,
90    pub dequeue_success: u64,
91    pub dequeue_fail: u64,
92    pub max_depth_observed: u64,
93    pub cas_retries: u64,
94}
95
96/// Helper that constructs a metrics-wrapped (producer, consumer, metrics)
97/// triple from a base SPSC pair.
98pub struct InstrumentedSpsc;
99
100impl InstrumentedSpsc {
101    /// Wrap an existing `(Producer, Consumer)` with counters. Returns the
102    /// instrumented sides and an `Arc<RingMetrics>` shared across both.
103    pub fn wrap<T>(
104        producer: Producer<T>,
105        consumer: Consumer<T>,
106    ) -> (
107        InstrumentedProducer<T>,
108        InstrumentedConsumer<T>,
109        Arc<RingMetrics>,
110    ) {
111        let metrics = Arc::new(RingMetrics::new());
112        let p = InstrumentedProducer {
113            inner: producer,
114            metrics: metrics.clone(),
115            local_depth: 0,
116        };
117        let c = InstrumentedConsumer {
118            inner: consumer,
119            metrics: metrics.clone(),
120        };
121        (p, c, metrics)
122    }
123}
124
125pub struct InstrumentedProducer<T> {
126    inner: Producer<T>,
127    metrics: Arc<RingMetrics>,
128    /// Producer-local view of in-flight count; +1 on success, no syscall.
129    local_depth: u64,
130}
131
132impl<T> InstrumentedProducer<T> {
133    pub fn try_push(&mut self, value: T) -> Result<(), T> {
134        match self.inner.try_push(value) {
135            Ok(()) => {
136                self.metrics.enqueue_success.fetch_add(1, Ordering::Relaxed);
137                self.local_depth += 1;
138                self.metrics.observe_depth(self.local_depth);
139                Ok(())
140            }
141            Err(v) => {
142                self.metrics.enqueue_fail.fetch_add(1, Ordering::Relaxed);
143                Err(v)
144            }
145        }
146    }
147
148    pub fn capacity(&self) -> usize {
149        self.inner.capacity()
150    }
151}
152
153pub struct InstrumentedConsumer<T> {
154    inner: Consumer<T>,
155    metrics: Arc<RingMetrics>,
156}
157
158impl<T> InstrumentedConsumer<T> {
159    pub fn try_pop(&mut self) -> Option<T> {
160        match self.inner.try_pop() {
161            Some(v) => {
162                self.metrics.dequeue_success.fetch_add(1, Ordering::Relaxed);
163                Some(v)
164            }
165            None => {
166                self.metrics.dequeue_fail.fetch_add(1, Ordering::Relaxed);
167                None
168            }
169        }
170    }
171
172    pub fn capacity(&self) -> usize {
173        self.inner.capacity()
174    }
175}
176
177#[cfg(test)]
178#[path = "metrics_tests.rs"]
179mod tests;