subms_spsc_ring_buffer/features/
metrics.rs1use std::sync::Arc;
16use std::sync::atomic::{AtomicU64, Ordering};
17
18use crate::{Consumer, Producer};
19
20pub 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 pub fn record_cas_retry(&self) {
52 self.cas_retries.fetch_add(1, Ordering::Relaxed);
53 }
54
55 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#[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
96pub struct InstrumentedSpsc;
99
100impl InstrumentedSpsc {
101 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 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;