subms_spsc_ring_buffer/features/
bulk.rs1use std::sync::atomic::Ordering;
15
16use crate::{Consumer, Producer};
17
18impl<T: Copy> Producer<T> {
19 pub fn try_enqueue_bulk(&mut self, values: &[T]) -> usize {
22 if values.is_empty() {
23 return 0;
24 }
25 let tail = self.inner.tail.0.load(Ordering::Relaxed);
26 let cap = self.inner.capacity;
27
28 let mut free = cap - tail.wrapping_sub(self.cached_head);
29 if free < values.len() {
30 self.cached_head = self.inner.head.0.load(Ordering::Acquire);
33 free = cap - tail.wrapping_sub(self.cached_head);
34 }
35 let n = free.min(values.len());
36 if n == 0 {
37 return 0;
38 }
39 for (i, v) in values.iter().take(n).enumerate() {
40 unsafe {
41 (*self.inner.buf[(tail.wrapping_add(i)) & self.inner.mask].get()).write(*v);
42 }
43 }
44 self.inner
46 .tail
47 .0
48 .store(tail.wrapping_add(n), Ordering::Release);
49 n
50 }
51}
52
53impl<T: Copy> Consumer<T> {
54 pub fn try_dequeue_bulk(&mut self, out: &mut [T]) -> usize {
57 if out.is_empty() {
58 return 0;
59 }
60 let head = self.inner.head.0.load(Ordering::Relaxed);
61
62 let mut avail = self.cached_tail.wrapping_sub(head);
63 if avail < out.len() {
64 self.cached_tail = self.inner.tail.0.load(Ordering::Acquire);
65 avail = self.cached_tail.wrapping_sub(head);
66 }
67 let n = avail.min(out.len());
68 if n == 0 {
69 return 0;
70 }
71 for (i, slot) in out.iter_mut().take(n).enumerate() {
72 *slot = unsafe {
73 (*self.inner.buf[(head.wrapping_add(i)) & self.inner.mask].get()).assume_init_read()
74 };
75 }
76 self.inner
77 .head
78 .0
79 .store(head.wrapping_add(n), Ordering::Release);
80 n
81 }
82}
83
84#[cfg(test)]
85#[path = "bulk_tests.rs"]
86mod tests;