use core::sync::atomic::{AtomicU32, AtomicUsize, Ordering};
#[repr(C)]
struct AccessEvent {
data: AtomicU32,
}
impl AccessEvent {
const EMPTY: u32 = u32::MAX;
fn new_empty() -> Self {
Self {
data: AtomicU32::new(Self::EMPTY),
}
}
#[inline]
fn pack(bucket_idx: usize, slot_idx: u8) -> u32 {
((bucket_idx as u32) << 8) | (slot_idx as u32)
}
#[allow(dead_code)]
#[inline]
fn unpack(val: u32) -> (usize, u8) {
let bucket_idx = (val >> 8) as usize;
let slot_idx = (val & 0xFF) as u8;
(bucket_idx, slot_idx)
}
}
pub struct AccessBuffer {
buffer: Vec<AccessEvent>,
mask: usize,
head: AtomicUsize, tail: AtomicUsize, }
impl AccessBuffer {
pub fn new(capacity: usize) -> Self {
let cap = capacity.max(64).next_power_of_two();
let buffer = (0..cap).map(|_| AccessEvent::new_empty()).collect();
Self {
buffer,
mask: cap - 1,
head: AtomicUsize::new(0),
tail: AtomicUsize::new(0),
}
}
#[inline]
pub fn push(&self, bucket_idx: usize, slot_idx: u8) -> bool {
let head = self.head.load(Ordering::Relaxed);
let tail = self.tail.load(Ordering::Relaxed);
if head.wrapping_sub(tail) > self.mask {
return false;
}
let idx = head & self.mask;
let packed = AccessEvent::pack(bucket_idx, slot_idx);
self.buffer[idx].data.store(packed, Ordering::Relaxed);
self.head.store(head.wrapping_add(1), Ordering::Release);
true
}
#[allow(dead_code)]
#[inline]
pub fn drain(&self, max_events: usize, mut f: impl FnMut(usize, u8)) {
let tail = self.tail.load(Ordering::Relaxed);
let head = self.head.load(Ordering::Acquire);
let available = head.wrapping_sub(tail);
let to_drain = available.min(max_events);
for i in 0..to_drain {
let idx = (tail.wrapping_add(i)) & self.mask;
let val = self.buffer[idx].data.load(Ordering::Relaxed);
if val != AccessEvent::EMPTY {
let (bucket_idx, slot_idx) = AccessEvent::unpack(val);
f(bucket_idx, slot_idx);
self.buffer[idx]
.data
.store(AccessEvent::EMPTY, Ordering::Relaxed);
}
}
self.tail
.store(tail.wrapping_add(to_drain), Ordering::Release);
}
}
unsafe impl Send for AccessBuffer {}
unsafe impl Sync for AccessBuffer {}