use std::sync::atomic::Ordering;
use crate::{Consumer, Producer};
impl<T: Copy> Producer<T> {
pub fn try_enqueue_bulk(&mut self, values: &[T]) -> usize {
if values.is_empty() {
return 0;
}
let tail = self.inner.tail.0.load(Ordering::Relaxed);
let cap = self.inner.capacity;
let mut free = cap - tail.wrapping_sub(self.cached_head);
if free < values.len() {
self.cached_head = self.inner.head.0.load(Ordering::Acquire);
free = cap - tail.wrapping_sub(self.cached_head);
}
let n = free.min(values.len());
if n == 0 {
return 0;
}
for (i, v) in values.iter().take(n).enumerate() {
unsafe {
(*self.inner.buf[(tail.wrapping_add(i)) & self.inner.mask].get()).write(*v);
}
}
self.inner
.tail
.0
.store(tail.wrapping_add(n), Ordering::Release);
n
}
}
impl<T: Copy> Consumer<T> {
pub fn try_dequeue_bulk(&mut self, out: &mut [T]) -> usize {
if out.is_empty() {
return 0;
}
let head = self.inner.head.0.load(Ordering::Relaxed);
let mut avail = self.cached_tail.wrapping_sub(head);
if avail < out.len() {
self.cached_tail = self.inner.tail.0.load(Ordering::Acquire);
avail = self.cached_tail.wrapping_sub(head);
}
let n = avail.min(out.len());
if n == 0 {
return 0;
}
for (i, slot) in out.iter_mut().take(n).enumerate() {
*slot = unsafe {
(*self.inner.buf[(head.wrapping_add(i)) & self.inner.mask].get()).assume_init_read()
};
}
self.inner
.head
.0
.store(head.wrapping_add(n), Ordering::Release);
n
}
}
#[cfg(test)]
#[path = "bulk_tests.rs"]
mod tests;