use std::cell::{Cell, UnsafeCell};
use std::mem::MaybeUninit;
use std::ptr;
use std::sync::atomic::{AtomicBool, AtomicPtr, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
pub(crate) const CHUNK_SIZE: usize = 64;
const MAX_CHUNKS_PER_SWEEP: usize = 1024;
struct Chunk<M> {
slots: [UnsafeCell<MaybeUninit<M>>; CHUNK_SIZE],
len: AtomicUsize,
next: AtomicPtr<Chunk<M>>,
}
impl<M> Chunk<M> {
fn new_raw() -> *mut Chunk<M> {
Box::into_raw(Box::new(Chunk {
slots: [const { UnsafeCell::new(MaybeUninit::uninit()) }; CHUNK_SIZE],
len: AtomicUsize::new(0),
next: AtomicPtr::new(ptr::null_mut()),
}))
}
}
pub(crate) struct EventQueue<M> {
head: AtomicPtr<Chunk<M>>,
consumed: AtomicUsize,
closed: AtomicBool,
}
unsafe impl<M: Send> Send for EventQueue<M> {}
unsafe impl<M: Send> Sync for EventQueue<M> {}
impl<M> EventQueue<M> {
fn drain_into(&self, out: &mut Vec<M>, max_chunks: usize) -> bool {
let mut chunk_ptr = self.head.load(Ordering::Relaxed);
let mut consumed = self.consumed.load(Ordering::Relaxed);
let mut chunks_walked = 0;
let mut reached_tail = true;
loop {
let chunk = unsafe { &*chunk_ptr };
let len = chunk.len.load(Ordering::Acquire);
for i in consumed..len {
out.push(unsafe { (*chunk.slots[i].get()).assume_init_read() });
}
consumed = len;
if len == CHUNK_SIZE {
let next = chunk.next.load(Ordering::Acquire);
if !next.is_null() {
if chunks_walked >= max_chunks {
reached_tail = false;
break;
}
unsafe { drop(Box::from_raw(chunk_ptr)) };
chunk_ptr = next;
consumed = 0;
chunks_walked += 1;
continue;
}
}
break;
}
self.head.store(chunk_ptr, Ordering::Relaxed);
self.consumed.store(consumed, Ordering::Relaxed);
reached_tail
}
}
impl<M> Drop for EventQueue<M> {
fn drop(&mut self) {
let mut chunk_ptr = *self.head.get_mut();
let mut consumed = *self.consumed.get_mut();
while !chunk_ptr.is_null() {
let chunk = unsafe { &mut *chunk_ptr };
let len = *chunk.len.get_mut();
for i in consumed..len {
unsafe { (*chunk.slots[i].get()).assume_init_drop() };
}
consumed = 0;
let next = *chunk.next.get_mut();
unsafe { drop(Box::from_raw(chunk_ptr)) };
chunk_ptr = next;
}
}
}
pub(crate) struct EventProducer<M> {
queue: Arc<EventQueue<M>>,
tail: Cell<*mut Chunk<M>>,
len: Cell<usize>,
}
impl<M> EventProducer<M> {
#[inline]
pub(crate) fn push(&self, m: M) {
let tail = self.tail.get();
let i = self.len.get();
unsafe {
(*tail).slots[i].get().write(MaybeUninit::new(m));
(*tail).len.store(i + 1, Ordering::Release);
}
if i + 1 == CHUNK_SIZE {
let new = Chunk::new_raw();
unsafe { (*tail).next.store(new, Ordering::Release) };
self.tail.set(new);
self.len.set(0);
} else {
self.len.set(i + 1);
}
}
}
impl<M> Drop for EventProducer<M> {
fn drop(&mut self) {
self.queue.closed.store(true, Ordering::Release);
}
}
pub(crate) struct EventQueueRegistry<M> {
active: AtomicBool,
queues: Mutex<Vec<Arc<EventQueue<M>>>>,
}
impl<M: Send> EventQueueRegistry<M> {
pub(crate) const fn new() -> Self {
Self {
active: AtomicBool::new(false),
queues: Mutex::new(Vec::new()),
}
}
#[inline]
pub(crate) fn is_active(&self) -> bool {
self.active.load(Ordering::Relaxed)
}
pub(crate) fn set_active(&self, active: bool) {
self.active.store(active, Ordering::Release);
}
pub(crate) fn register(&self) -> EventProducer<M> {
let first = Chunk::new_raw();
let queue = Arc::new(EventQueue {
head: AtomicPtr::new(first),
consumed: AtomicUsize::new(0),
closed: AtomicBool::new(false),
});
if let Ok(mut queues) = self.queues.lock() {
queues.push(Arc::clone(&queue));
}
EventProducer {
queue,
tail: Cell::new(first),
len: Cell::new(0),
}
}
pub(crate) fn sweep(&self, out: &mut Vec<M>) {
self.sweep_inner(out, MAX_CHUNKS_PER_SWEEP);
}
pub(crate) fn drain_all(&self, out: &mut Vec<M>) {
self.sweep_inner(out, usize::MAX);
}
fn sweep_inner(&self, out: &mut Vec<M>, max_chunks: usize) {
if let Ok(mut queues) = self.queues.lock() {
queues.retain(|queue| {
let closed = queue.closed.load(Ordering::Acquire);
let reached_tail = queue.drain_into(out, max_chunks);
!(closed && reached_tail)
});
}
}
}