Skip to main content

subms_mpsc_queue/features/
metrics.rs

1//! Per-instance metrics wrapper.
2//!
3//! Wraps the base [`MpscQueue`] with relaxed atomic counters for
4//! enqueue success/fail, dequeue success/fail, total batch items
5//! drained, and CAS retries. Counters are relaxed because they're
6//! advisory diagnostics, not ordering primitives.
7//!
8//! All counters are zero-cost when the wrapper isn't used; the
9//! feature flag keeps them out of the base build.
10
11use crate::{MpscQueue, PopResult};
12use std::sync::atomic::{AtomicU64, Ordering};
13
14/// Wrapping queue that tracks per-instance counters.
15pub struct MetricsMpscQueue<T> {
16    inner: MpscQueue<T>,
17    enqueue_ok: AtomicU64,
18    enqueue_fail: AtomicU64,
19    dequeue_ok: AtomicU64,
20    dequeue_fail: AtomicU64,
21    batch_items: AtomicU64,
22    cas_retries: AtomicU64,
23}
24
25/// Immutable snapshot of the counters at one instant.
26#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
27pub struct QueueMetricsSnapshot {
28    pub enqueue_ok: u64,
29    pub enqueue_fail: u64,
30    pub dequeue_ok: u64,
31    pub dequeue_fail: u64,
32    pub batch_items: u64,
33    pub cas_retries: u64,
34}
35
36impl<T> MetricsMpscQueue<T> {
37    pub fn new() -> Self {
38        Self {
39            inner: MpscQueue::new(),
40            enqueue_ok: AtomicU64::new(0),
41            enqueue_fail: AtomicU64::new(0),
42            dequeue_ok: AtomicU64::new(0),
43            dequeue_fail: AtomicU64::new(0),
44            batch_items: AtomicU64::new(0),
45            cas_retries: AtomicU64::new(0),
46        }
47    }
48
49    /// Push always succeeds for the unbounded base; the fail counter
50    /// is only bumped via [`Self::record_enqueue_fail`] from a bounded
51    /// composition wrapper.
52    pub fn push(&self, value: T) {
53        self.inner.push(value);
54        self.enqueue_ok.fetch_add(1, Ordering::Relaxed);
55    }
56
57    /// Single-consumer pop. Bumps `dequeue_ok` or `dequeue_fail`.
58    pub fn try_pop(&mut self) -> PopResult<T> {
59        let r = self.inner.try_pop();
60        match &r {
61            PopResult::Some(_) => {
62                self.dequeue_ok.fetch_add(1, Ordering::Relaxed);
63            }
64            PopResult::Empty | PopResult::Inconsistent => {
65                self.dequeue_fail.fetch_add(1, Ordering::Relaxed);
66            }
67        }
68        r
69    }
70
71    /// Bulk drain into `out`. Returns the number drained and bumps
72    /// `batch_items` by the count.
73    pub fn try_pop_batch(&mut self, out: &mut [Option<T>]) -> usize {
74        let mut n = 0;
75        while n < out.len() {
76            match self.inner.try_pop() {
77                PopResult::Some(v) => {
78                    out[n] = Some(v);
79                    n += 1;
80                    self.dequeue_ok.fetch_add(1, Ordering::Relaxed);
81                }
82                PopResult::Empty | PopResult::Inconsistent => {
83                    self.dequeue_fail.fetch_add(1, Ordering::Relaxed);
84                    break;
85                }
86            }
87        }
88        self.batch_items.fetch_add(n as u64, Ordering::Relaxed);
89        n
90    }
91
92    /// Borrow the next value without consuming it. Does not touch the
93    /// counters: a peek is not a dequeue.
94    pub fn peek(&mut self) -> Option<&T> {
95        self.inner.peek()
96    }
97
98    /// See [`MpscQueue::is_empty`].
99    pub fn is_empty(&mut self) -> bool {
100        self.inner.is_empty()
101    }
102
103    /// See [`MpscQueue::len`]. O(n) in the backlog.
104    pub fn len(&mut self) -> usize {
105        self.inner.len()
106    }
107
108    /// Drain everything reachable and return the count. The drained items
109    /// count as successful dequeues, so a cleared backlog still shows up in
110    /// the snapshot rather than vanishing from the totals.
111    pub fn clear(&mut self) -> usize {
112        let n = self.inner.clear();
113        self.dequeue_ok.fetch_add(n as u64, Ordering::Relaxed);
114        n
115    }
116
117    /// External hook for callers that combine this with a bounded
118    /// upstream (or any path where an enqueue can be rejected).
119    pub fn record_enqueue_fail(&self) {
120        self.enqueue_fail.fetch_add(1, Ordering::Relaxed);
121    }
122
123    /// External hook used by MPMC compositions to log retry counts.
124    pub fn record_cas_retries(&self, n: u64) {
125        if n > 0 {
126            self.cas_retries.fetch_add(n, Ordering::Relaxed);
127        }
128    }
129
130    /// Atomic-load snapshot. Counters may move between loads (relaxed
131    /// across atomics), so this is a point-in-time approximation.
132    pub fn snapshot(&self) -> QueueMetricsSnapshot {
133        QueueMetricsSnapshot {
134            enqueue_ok: self.enqueue_ok.load(Ordering::Relaxed),
135            enqueue_fail: self.enqueue_fail.load(Ordering::Relaxed),
136            dequeue_ok: self.dequeue_ok.load(Ordering::Relaxed),
137            dequeue_fail: self.dequeue_fail.load(Ordering::Relaxed),
138            batch_items: self.batch_items.load(Ordering::Relaxed),
139            cas_retries: self.cas_retries.load(Ordering::Relaxed),
140        }
141    }
142
143    /// Reset all counters to zero. Useful for cycle-bounded
144    /// reporting.
145    pub fn reset(&self) {
146        self.enqueue_ok.store(0, Ordering::Relaxed);
147        self.enqueue_fail.store(0, Ordering::Relaxed);
148        self.dequeue_ok.store(0, Ordering::Relaxed);
149        self.dequeue_fail.store(0, Ordering::Relaxed);
150        self.batch_items.store(0, Ordering::Relaxed);
151        self.cas_retries.store(0, Ordering::Relaxed);
152    }
153}
154
155impl<T> Default for MetricsMpscQueue<T> {
156    fn default() -> Self {
157        Self::new()
158    }
159}
160
161#[cfg(test)]
162#[path = "metrics_tests.rs"]
163mod tests;