use parking_lot::Mutex;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Clone, Debug)]
pub struct RingBufferSlot<T: Clone> {
value: Option<T>,
}
pub struct RingBuffer<T: Clone> {
buffer: Vec<RingBufferSlot<T>>,
write_idx: usize, read_idx: usize, count: usize, metrics: Arc<RingBufferMetrics>,
}
#[derive(Debug, Default)]
pub struct RingBufferMetrics {
pub total_pushed: AtomicU64,
pub total_dropped: AtomicU64,
pub total_popped: AtomicU64,
pub max_depth: AtomicU64,
}
impl<T: Clone> RingBuffer<T> {
pub fn new(capacity: usize) -> Self {
let buffer = vec![RingBufferSlot { value: None }; capacity.max(1)];
Self {
buffer,
write_idx: 0,
read_idx: 0,
count: 0,
metrics: Arc::new(RingBufferMetrics::default()),
}
}
pub fn push(&mut self, item: T) -> bool {
let was_full = self.is_full();
if was_full {
self.buffer[self.read_idx].value = None;
self.read_idx = (self.read_idx + 1) % self.buffer.len();
self.count = self.count.saturating_sub(1);
self.metrics.total_dropped.fetch_add(1, Ordering::Relaxed);
}
self.buffer[self.write_idx].value = Some(item);
self.write_idx = (self.write_idx + 1) % self.buffer.len();
self.count += 1;
self.metrics.total_pushed.fetch_add(1, Ordering::Relaxed);
let current_depth = self.count as u64;
let _ = self
.metrics
.max_depth
.fetch_max(current_depth, Ordering::Relaxed);
!was_full
}
pub fn pop(&mut self) -> Option<T> {
if self.count == 0 {
return None;
}
let slot = &mut self.buffer[self.read_idx];
let item = slot.value.take();
self.read_idx = (self.read_idx + 1) % self.buffer.len();
self.count = self.count.saturating_sub(1);
if item.is_some() {
self.metrics.total_popped.fetch_add(1, Ordering::Relaxed);
}
item
}
pub fn is_empty(&self) -> bool {
self.count == 0
}
pub fn is_full(&self) -> bool {
self.count >= self.buffer.len()
}
pub fn len(&self) -> usize {
self.count
}
pub fn capacity(&self) -> usize {
self.buffer.len()
}
pub fn metrics(&self) -> Arc<RingBufferMetrics> {
Arc::clone(&self.metrics)
}
pub fn clear(&mut self) {
for slot in &mut self.buffer {
slot.value = None;
}
self.write_idx = 0;
self.read_idx = 0;
self.count = 0;
}
pub fn drain_all(&mut self) -> Vec<T> {
let mut result = Vec::with_capacity(self.count);
while let Some(item) = self.pop() {
result.push(item);
}
result
}
}
pub struct ConcurrentRingBuffer<T: Clone + Send + 'static> {
inner: Arc<Mutex<RingBuffer<T>>>,
}
impl<T: Clone + Send + 'static> ConcurrentRingBuffer<T> {
pub fn new(capacity: usize) -> Self {
Self {
inner: Arc::new(Mutex::new(RingBuffer::new(capacity))),
}
}
pub fn push(&self, item: T) -> bool {
let mut buf = self.inner.lock();
buf.push(item)
}
pub fn pop(&self) -> Option<T> {
let mut buf = self.inner.lock();
buf.pop()
}
pub fn len(&self) -> usize {
self.inner.lock().len()
}
pub fn is_empty(&self) -> bool {
self.inner.lock().is_empty()
}
pub fn capacity(&self) -> usize {
self.inner.lock().capacity()
}
pub fn metrics(&self) -> Arc<RingBufferMetrics> {
self.inner.lock().metrics()
}
pub fn drain_all(&self) -> Vec<T> {
let mut buf = self.inner.lock();
buf.drain_all()
}
pub fn clone_ref(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_ring_buffer_basic() {
let mut buf = RingBuffer::new(3);
assert!(buf.is_empty());
assert!(!buf.is_full());
assert!(buf.push(1));
assert!(!buf.is_empty());
assert!(!buf.is_full());
assert!(buf.push(2));
assert!(buf.push(3));
assert!(buf.is_full());
assert!(!buf.push(4));
assert_eq!(buf.len(), 3);
assert_eq!(buf.metrics().total_dropped.load(Ordering::Relaxed), 1);
assert_eq!(buf.pop(), Some(2));
assert_eq!(buf.pop(), Some(3));
assert_eq!(buf.pop(), Some(4));
assert_eq!(buf.pop(), None);
}
#[test]
fn test_ring_buffer_metrics() {
let mut buf = RingBuffer::new(2);
buf.push(1);
buf.push(2);
assert_eq!(buf.metrics().total_pushed.load(Ordering::Relaxed), 2);
assert_eq!(buf.metrics().max_depth.load(Ordering::Relaxed), 2);
buf.push(3); assert_eq!(buf.metrics().total_dropped.load(Ordering::Relaxed), 1);
buf.pop(); buf.pop(); buf.pop(); assert_eq!(buf.metrics().total_popped.load(Ordering::Relaxed), 2);
}
#[test]
fn test_concurrent_ring_buffer() {
let buf = ConcurrentRingBuffer::new(2);
assert!(buf.push(1));
assert!(buf.push(2));
assert!(!buf.push(3));
assert_eq!(buf.pop(), Some(2));
assert_eq!(buf.len(), 1);
}
#[test]
fn test_drain_all_basic() {
let mut buf = RingBuffer::new(10);
buf.push(1);
buf.push(2);
buf.push(3);
let drained = buf.drain_all();
assert_eq!(drained, vec![1, 2, 3]);
assert!(buf.is_empty());
assert_eq!(buf.len(), 0);
}
#[test]
fn test_drain_all_empty() {
let mut buf: RingBuffer<i32> = RingBuffer::new(10);
let drained = buf.drain_all();
assert!(drained.is_empty());
}
#[test]
fn test_drain_all_after_overflow() {
let mut buf = RingBuffer::new(3);
buf.push(1);
buf.push(2);
buf.push(3);
buf.push(4);
let drained = buf.drain_all();
assert_eq!(drained, vec![2, 3, 4]);
assert!(buf.is_empty());
}
#[test]
fn test_concurrent_drain_all() {
let buf = ConcurrentRingBuffer::new(100);
for i in 0..50 {
buf.push(i);
}
assert_eq!(buf.len(), 50);
let drained = buf.drain_all();
assert_eq!(drained.len(), 50);
assert_eq!(drained[0], 0);
assert_eq!(drained[49], 49);
assert!(buf.is_empty());
}
#[test]
fn test_concurrent_drain_via_clone_ref() {
let buf = ConcurrentRingBuffer::new(1000);
let buf_clone = buf.clone_ref();
for i in 0..364 {
buf_clone.push(i);
}
let drained = buf.drain_all();
assert_eq!(drained.len(), 364);
assert_eq!(drained[0], 0);
assert_eq!(drained[363], 363);
assert!(buf.is_empty());
}
}