subms_mpsc_queue/features/
metrics.rs1use crate::{MpscQueue, PopResult};
12use std::sync::atomic::{AtomicU64, Ordering};
13
14pub 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#[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 pub fn push(&self, value: T) {
53 self.inner.push(value);
54 self.enqueue_ok.fetch_add(1, Ordering::Relaxed);
55 }
56
57 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 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 pub fn peek(&mut self) -> Option<&T> {
95 self.inner.peek()
96 }
97
98 pub fn is_empty(&mut self) -> bool {
100 self.inner.is_empty()
101 }
102
103 pub fn len(&mut self) -> usize {
105 self.inner.len()
106 }
107
108 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 pub fn record_enqueue_fail(&self) {
120 self.enqueue_fail.fetch_add(1, Ordering::Relaxed);
121 }
122
123 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 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 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;