use std::mem::MaybeUninit;
use std::ptr::{self, NonNull};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Mutex;
use super::allocator::CacheAlignedAllocator;
pub struct UnifiedRingBuffer<T> {
buffer: NonNull<MaybeUninit<T>>,
capacity: usize,
mask: usize,
head: AtomicUsize,
tail: AtomicUsize,
write_lock: Mutex<()>,
read_lock: Mutex<()>,
}
impl<T> UnifiedRingBuffer<T> {
pub fn new(capacity: usize) -> Option<Self> {
let capacity = capacity.next_power_of_two().max(2);
let buffer = CacheAlignedAllocator::allocate::<MaybeUninit<T>>(capacity)?;
Some(Self {
buffer,
capacity,
mask: capacity - 1,
head: AtomicUsize::new(0),
tail: AtomicUsize::new(0),
write_lock: Mutex::new(()),
read_lock: Mutex::new(()),
})
}
pub fn try_push(&self, item: T) -> Result<(), T> {
let _guard = self.write_lock.lock().unwrap();
let head = self.head.load(Ordering::Relaxed);
let next_head = (head + 1) & self.mask;
let tail = self.tail.load(Ordering::Acquire);
if next_head == tail {
return Err(item);
}
unsafe {
let slot = self.buffer.as_ptr().add(head & self.mask);
ptr::write((*slot).as_mut_ptr(), item);
}
self.head.store(next_head, Ordering::Release);
Ok(())
}
pub fn try_pop(&self) -> Option<T> {
let _guard = self.read_lock.lock().unwrap();
let tail = self.tail.load(Ordering::Relaxed);
let head = self.head.load(Ordering::Acquire);
if tail == head {
return None;
}
let item = unsafe {
let slot = self.buffer.as_ptr().add(tail & self.mask);
ptr::read((*slot).as_ptr())
};
self.tail.store((tail + 1) & self.mask, Ordering::Release);
Some(item)
}
pub fn capacity(&self) -> usize {
self.capacity
}
pub fn is_empty(&self) -> bool {
let head = self.head.load(Ordering::Acquire);
let tail = self.tail.load(Ordering::Acquire);
head == tail
}
pub fn len(&self) -> usize {
let head = self.head.load(Ordering::Acquire);
let tail = self.tail.load(Ordering::Acquire);
(head.wrapping_sub(tail)) & self.mask
}
}
impl<T> Drop for UnifiedRingBuffer<T> {
fn drop(&mut self) {
#[allow(clippy::redundant_pattern_matching)]
while let Some(_) = self.try_pop() {
}
unsafe {
CacheAlignedAllocator::deallocate(self.buffer, self.capacity);
}
}
}
unsafe impl<T: Send> Send for UnifiedRingBuffer<T> {}
unsafe impl<T: Send> Sync for UnifiedRingBuffer<T> {}