use std::cell::UnsafeCell;
use std::os::raw::c_void;
use std::sync::Arc;
use crate::bench_harness::{
BenchQueueHandleFactory, BenchQueueThreadOps, LogQueueHandleFactory, LogQueueThreadOps,
LogRecord,
};
unsafe extern "C" {
fn mc_queue_new() -> *mut c_void;
fn mc_queue_free(handle: *mut c_void);
fn mc_queue_create_producer_token(handle: *mut c_void) -> *mut c_void;
fn mc_queue_free_producer_token(token: *mut c_void);
fn mc_queue_enqueue_with_token(handle: *mut c_void, token: *mut c_void, value: u64) -> bool;
fn mc_queue_enqueue_bulk_with_token(
handle: *mut c_void,
token: *mut c_void,
items: *const u64,
count: usize,
) -> bool;
fn mc_queue_create_consumer_token(handle: *mut c_void) -> *mut c_void;
fn mc_queue_free_consumer_token(token: *mut c_void);
fn mc_queue_try_dequeue_with_token(
handle: *mut c_void,
token: *mut c_void,
out: *mut u64,
) -> bool;
fn mc_queue_try_dequeue_bulk_with_token(
handle: *mut c_void,
token: *mut c_void,
out: *mut u64,
max: usize,
) -> usize;
}
pub struct MoodycamelQueue {
handle: *mut c_void,
}
unsafe impl Send for MoodycamelQueue {}
unsafe impl Sync for MoodycamelQueue {}
impl MoodycamelQueue {
pub fn new_handle() -> Arc<Self> {
let handle = unsafe { mc_queue_new() };
assert!(!handle.is_null(), "moodycamel queue allocation failed");
Arc::new(Self { handle })
}
}
impl Drop for MoodycamelQueue {
fn drop(&mut self) {
unsafe { mc_queue_free(self.handle) };
}
}
struct ProducerTokenHandle(*mut c_void);
unsafe impl Send for ProducerTokenHandle {}
impl Drop for ProducerTokenHandle {
fn drop(&mut self) {
unsafe { mc_queue_free_producer_token(self.0) };
}
}
struct ConsumerTokenHandle(*mut c_void);
unsafe impl Send for ConsumerTokenHandle {}
impl Drop for ConsumerTokenHandle {
fn drop(&mut self) {
unsafe { mc_queue_free_consumer_token(self.0) };
}
}
enum MoodycamelThreadRole {
Producer {
token: ProducerTokenHandle,
batch: UnsafeCell<Vec<u64>>,
},
Consumer {
token: ConsumerTokenHandle,
batch: UnsafeCell<Vec<u64>>,
},
}
pub struct MoodycamelThreadHandle {
queue: Arc<MoodycamelQueue>,
role: MoodycamelThreadRole,
}
fn new_producer_handle(queue: &Arc<MoodycamelQueue>) -> MoodycamelThreadHandle {
let producer_token = unsafe { mc_queue_create_producer_token(queue.handle) };
assert!(
!producer_token.is_null(),
"moodycamel producer token allocation failed"
);
MoodycamelThreadHandle {
queue: queue.clone(),
role: MoodycamelThreadRole::Producer {
token: ProducerTokenHandle(producer_token),
batch: UnsafeCell::new(Vec::new()),
},
}
}
fn new_consumer_handle(queue: &Arc<MoodycamelQueue>) -> MoodycamelThreadHandle {
let consumer_token = unsafe { mc_queue_create_consumer_token(queue.handle) };
assert!(
!consumer_token.is_null(),
"moodycamel consumer token allocation failed"
);
MoodycamelThreadHandle {
queue: queue.clone(),
role: MoodycamelThreadRole::Consumer {
token: ConsumerTokenHandle(consumer_token),
batch: UnsafeCell::new(Vec::new()),
},
}
}
impl BenchQueueHandleFactory for MoodycamelQueue {
type ThreadHandle = MoodycamelThreadHandle;
fn thread_handle(self: &Arc<Self>) -> Self::ThreadHandle {
panic!("moodycamel handles must be requested for a producer or consumer role")
}
fn producer_thread_handle(self: &Arc<Self>) -> Self::ThreadHandle {
new_producer_handle(self)
}
fn consumer_thread_handle(self: &Arc<Self>) -> Self::ThreadHandle {
new_consumer_handle(self)
}
}
impl LogQueueHandleFactory for MoodycamelQueue {
type ThreadHandle = MoodycamelThreadHandle;
fn log_thread_handle(self: &Arc<Self>) -> Self::ThreadHandle {
panic!("moodycamel log handles must be requested for a producer or consumer role")
}
fn log_producer_thread_handle(self: &Arc<Self>) -> Self::ThreadHandle {
new_producer_handle(self)
}
fn log_consumer_thread_handle(self: &Arc<Self>) -> Self::ThreadHandle {
new_consumer_handle(self)
}
}
impl BenchQueueThreadOps for MoodycamelThreadHandle {
fn try_send_value(&self, value: u64) -> bool {
let MoodycamelThreadRole::Producer { token, .. } = &self.role else {
panic!("attempted to send through a moodycamel consumer handle");
};
unsafe { mc_queue_enqueue_with_token(self.queue.handle, token.0, value) }
}
fn send_batch(&self, base: u64, offsets: std::ops::Range<usize>) {
let MoodycamelThreadRole::Producer { token, batch } = &self.role else {
panic!("attempted to send a batch through a moodycamel consumer handle");
};
let items = unsafe { &mut *batch.get() };
items.clear();
items.extend(offsets.map(|offset| base + offset as u64));
let backoff = crossbeam_utils::Backoff::new();
while !unsafe {
mc_queue_enqueue_bulk_with_token(
self.queue.handle,
token.0,
items.as_ptr(),
items.len(),
)
} {
backoff.snooze();
}
}
fn try_recv_value(&self) -> Option<u64> {
let MoodycamelThreadRole::Consumer { token, .. } = &self.role else {
panic!("attempted to receive through a moodycamel producer handle");
};
let mut value = 0_u64;
let ok = unsafe { mc_queue_try_dequeue_with_token(self.queue.handle, token.0, &mut value) };
ok.then_some(value)
}
fn try_recv_batch(&self, request_size: usize) -> usize {
let MoodycamelThreadRole::Consumer { token, batch } = &self.role else {
panic!("attempted to receive a batch through a moodycamel producer handle");
};
let buf = unsafe { &mut *batch.get() };
buf.resize(request_size, 0);
unsafe {
mc_queue_try_dequeue_bulk_with_token(
self.queue.handle,
token.0,
buf.as_mut_ptr(),
request_size,
)
}
}
}
impl LogQueueThreadOps for MoodycamelThreadHandle {
fn send_log(&self, record: LogRecord) {
let MoodycamelThreadRole::Producer { token, .. } = &self.role else {
panic!("attempted to send a log through a moodycamel consumer handle");
};
let ptr = Box::into_raw(Box::new(record)) as usize as u64;
let backoff = crossbeam_utils::Backoff::new();
while !unsafe { mc_queue_enqueue_with_token(self.queue.handle, token.0, ptr) } {
backoff.snooze();
}
}
fn try_recv_log(&self) -> Option<LogRecord> {
let MoodycamelThreadRole::Consumer { token, .. } = &self.role else {
panic!("attempted to receive a log through a moodycamel producer handle");
};
let mut ptr = 0_u64;
let ok = unsafe { mc_queue_try_dequeue_with_token(self.queue.handle, token.0, &mut ptr) };
ok.then(|| *unsafe { Box::from_raw(ptr as usize as *mut LogRecord) })
}
}