use crate::{MpscQueue, PopResult};
use std::sync::atomic::{AtomicU64, Ordering};
pub struct MetricsMpscQueue<T> {
inner: MpscQueue<T>,
enqueue_ok: AtomicU64,
enqueue_fail: AtomicU64,
dequeue_ok: AtomicU64,
dequeue_fail: AtomicU64,
batch_items: AtomicU64,
cas_retries: AtomicU64,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct QueueMetricsSnapshot {
pub enqueue_ok: u64,
pub enqueue_fail: u64,
pub dequeue_ok: u64,
pub dequeue_fail: u64,
pub batch_items: u64,
pub cas_retries: u64,
}
impl<T> MetricsMpscQueue<T> {
pub fn new() -> Self {
Self {
inner: MpscQueue::new(),
enqueue_ok: AtomicU64::new(0),
enqueue_fail: AtomicU64::new(0),
dequeue_ok: AtomicU64::new(0),
dequeue_fail: AtomicU64::new(0),
batch_items: AtomicU64::new(0),
cas_retries: AtomicU64::new(0),
}
}
pub fn push(&self, value: T) {
self.inner.push(value);
self.enqueue_ok.fetch_add(1, Ordering::Relaxed);
}
pub fn try_pop(&mut self) -> PopResult<T> {
let r = self.inner.try_pop();
match &r {
PopResult::Some(_) => {
self.dequeue_ok.fetch_add(1, Ordering::Relaxed);
}
PopResult::Empty | PopResult::Inconsistent => {
self.dequeue_fail.fetch_add(1, Ordering::Relaxed);
}
}
r
}
pub fn try_pop_batch(&mut self, out: &mut [Option<T>]) -> usize {
let mut n = 0;
while n < out.len() {
match self.inner.try_pop() {
PopResult::Some(v) => {
out[n] = Some(v);
n += 1;
self.dequeue_ok.fetch_add(1, Ordering::Relaxed);
}
PopResult::Empty | PopResult::Inconsistent => {
self.dequeue_fail.fetch_add(1, Ordering::Relaxed);
break;
}
}
}
self.batch_items.fetch_add(n as u64, Ordering::Relaxed);
n
}
pub fn peek(&mut self) -> Option<&T> {
self.inner.peek()
}
pub fn is_empty(&mut self) -> bool {
self.inner.is_empty()
}
pub fn len(&mut self) -> usize {
self.inner.len()
}
pub fn clear(&mut self) -> usize {
let n = self.inner.clear();
self.dequeue_ok.fetch_add(n as u64, Ordering::Relaxed);
n
}
pub fn record_enqueue_fail(&self) {
self.enqueue_fail.fetch_add(1, Ordering::Relaxed);
}
pub fn record_cas_retries(&self, n: u64) {
if n > 0 {
self.cas_retries.fetch_add(n, Ordering::Relaxed);
}
}
pub fn snapshot(&self) -> QueueMetricsSnapshot {
QueueMetricsSnapshot {
enqueue_ok: self.enqueue_ok.load(Ordering::Relaxed),
enqueue_fail: self.enqueue_fail.load(Ordering::Relaxed),
dequeue_ok: self.dequeue_ok.load(Ordering::Relaxed),
dequeue_fail: self.dequeue_fail.load(Ordering::Relaxed),
batch_items: self.batch_items.load(Ordering::Relaxed),
cas_retries: self.cas_retries.load(Ordering::Relaxed),
}
}
pub fn reset(&self) {
self.enqueue_ok.store(0, Ordering::Relaxed);
self.enqueue_fail.store(0, Ordering::Relaxed);
self.dequeue_ok.store(0, Ordering::Relaxed);
self.dequeue_fail.store(0, Ordering::Relaxed);
self.batch_items.store(0, Ordering::Relaxed);
self.cas_retries.store(0, Ordering::Relaxed);
}
}
impl<T> Default for MetricsMpscQueue<T> {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
#[path = "metrics_tests.rs"]
mod tests;