use crate::{
QueueError,
owned::MpmcQueue,
traits::{QueueConsumer, QueueFactory, QueueProducer},
};
use std::{
fmt,
sync::{
Arc,
atomic::{AtomicUsize, Ordering},
},
};
pub struct QueuePack<T, I = u32, const G: usize = 4, const K: usize = 16, const N: usize = 0>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
queues: Vec<Arc<MpmcQueue<T, I, N>>>,
writer_counter: AtomicUsize,
reader_counter: AtomicUsize,
}
impl<T, I, const G: usize, const K: usize, const N: usize> fmt::Debug for QueuePack<T, I, G, K, N>
where
T: Copy + Send + Sync + Default + fmt::Debug,
I: Copy + Into<u128> + fmt::Debug,
{
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("QueuePack")
.field("queue_count", &G)
.field("scan_threshold", &K)
.field("queue_capacity", &self.queue_capacity())
.field("total_len", &self.len())
.field("total_capacity", &self.capacity())
.finish()
}
}
#[derive(Debug, Clone)]
pub struct QueuePackBuilder<T, I = u32, const G: usize = 4, const K: usize = 16>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
queue_capacity: Option<usize>,
_phantom: std::marker::PhantomData<(T, I)>,
}
impl<T, I, const G: usize, const K: usize> Default for QueuePackBuilder<T, I, G, K>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
fn default() -> Self {
Self::new()
}
}
impl<T, I, const G: usize, const K: usize> QueuePackBuilder<T, I, G, K>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
pub const fn new() -> Self {
Self {
queue_capacity: None,
_phantom: std::marker::PhantomData,
}
}
#[must_use]
pub const fn queue_capacity(mut self, capacity: usize) -> Self {
self.queue_capacity = Some(capacity);
self
}
pub fn build(self) -> Result<Arc<QueuePack<T, I, G, K>>, QueueError> {
let capacity = self.queue_capacity.ok_or(QueueError::InvalidCapacity)?;
Ok(Arc::new(QueuePack::new(capacity)?))
}
pub fn build_static<const N: usize>(self) -> Result<Arc<QueuePack<T, I, G, K, N>>, QueueError> {
let capacity = self.queue_capacity.unwrap_or(N);
Ok(Arc::new(QueuePack::new(capacity)?))
}
pub fn channels(
self,
) -> Result<(PackProducer<T, I, G, K>, PackConsumer<T, I, G, K>), QueueError> {
let pack = self.build()?;
Ok((pack.producer(), pack.consumer()))
}
pub fn channels_static<const N: usize>(
self,
) -> Result<(PackProducer<T, I, G, K, N>, PackConsumer<T, I, G, K, N>), QueueError> {
let pack = self.build_static::<N>()?;
Ok((pack.producer(), pack.consumer()))
}
}
pub fn queue_pack<T, const G: usize, const K: usize>() -> QueuePackBuilder<T, u32, G, K>
where
T: Copy + Send + Sync + Default,
{
QueuePackBuilder::new()
}
pub fn queue_pack_with_index<T, I, const G: usize, const K: usize>() -> QueuePackBuilder<T, I, G, K>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
QueuePackBuilder::new()
}
impl<T, I, const G: usize, const K: usize, const N: usize> QueuePack<T, I, G, K, N>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
pub fn new(queue_capacity: usize) -> Result<Self, QueueError> {
if G == 0 {
return Err(QueueError::InvalidCapacity);
}
let mut queues = Vec::with_capacity(G);
for _ in 0..G {
queues.push(Arc::new(MpmcQueue::new(queue_capacity)?));
}
Ok(Self {
queues,
writer_counter: AtomicUsize::new(0),
reader_counter: AtomicUsize::new(0),
})
}
pub const fn queue_count() -> usize {
G
}
pub const fn scan_threshold() -> usize {
K
}
pub fn queue_capacity(&self) -> usize {
self.queues[0].capacity()
}
pub fn capacity(&self) -> usize {
self.queues.len() * self.queue_capacity()
}
pub fn len(&self) -> usize {
self.queues.iter().map(|q| q.len()).sum()
}
pub fn is_empty(&self) -> bool {
self.queues.iter().all(|q| q.is_empty())
}
pub fn is_full(&self) -> bool {
self.queues.iter().all(|q| q.is_full())
}
pub fn queue_stats(&self) -> Vec<QueueStats> {
self.queues
.iter()
.enumerate()
.map(|(index, queue)| QueueStats {
index,
len: queue.len(),
capacity: queue.capacity(),
is_empty: queue.is_empty(),
is_full: queue.is_full(),
})
.collect()
}
pub fn try_push_to(&self, queue_index: usize, value: T) -> Result<(), (T, QueueError)> {
if queue_index >= G {
return Err((value, QueueError::InvalidCapacity));
}
self.queues[queue_index].try_push(value)
}
pub fn try_pop_from(&self, queue_index: usize) -> Result<T, QueueError> {
if queue_index >= G {
return Err(QueueError::InvalidCapacity);
}
self.queues[queue_index].try_pop()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QueueStats {
pub index: usize,
pub len: usize,
pub capacity: usize,
pub is_empty: bool,
pub is_full: bool,
}
pub type PackProducer<T, I = u32, const G: usize = 4, const K: usize = 16, const N: usize = 0> =
PackProducerHandle<T, I, G, K, N>;
pub type PackConsumer<T, I = u32, const G: usize = 4, const K: usize = 16, const N: usize = 0> =
PackConsumerHandle<T, I, G, K, N>;
#[derive(Debug)]
pub struct PackProducerHandle<
T,
I = u32,
const G: usize = 4,
const K: usize = 16,
const N: usize = 0,
> where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
pack: Arc<QueuePack<T, I, G, K, N>>,
queue_index: usize,
}
impl<T, I, const G: usize, const K: usize, const N: usize> Clone
for PackProducerHandle<T, I, G, K, N>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
fn clone(&self) -> Self {
self.pack.producer()
}
}
impl<T, I, const G: usize, const K: usize, const N: usize> PackProducerHandle<T, I, G, K, N>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
const fn new(pack: Arc<QueuePack<T, I, G, K, N>>, queue_index: usize) -> Self {
Self { pack, queue_index }
}
pub const fn queue_index(&self) -> usize {
self.queue_index
}
pub fn try_push(&self, value: T) -> Result<(), (T, QueueError)> {
self.pack.queues[self.queue_index].try_push(value)
}
pub fn queue_stats(&self) -> QueueStats {
let queue = &self.pack.queues[self.queue_index];
QueueStats {
index: self.queue_index,
len: queue.len(),
capacity: queue.capacity(),
is_empty: queue.is_empty(),
is_full: queue.is_full(),
}
}
}
impl<T, I, const G: usize, const K: usize, const N: usize> QueueProducer<T>
for PackProducerHandle<T, I, G, K, N>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
fn try_push(&self, value: T) -> Result<(), (T, QueueError)> {
self.pack.queues[self.queue_index].try_push(value)
}
fn push(&self, value: T) -> Result<(), QueueError> {
self.pack.queues[self.queue_index].push(value)
}
fn push_with_seq(&self, value: T) -> Result<usize, QueueError> {
let local_index = self.pack.queues[self.queue_index].push_impl(value, true)?;
let global_seq = (self.queue_index << 24) | (local_index & 0x00FF_FFFF);
Ok(global_seq)
}
}
#[derive(Debug)]
pub struct PackConsumerHandle<
T,
I = u32,
const G: usize = 4,
const K: usize = 16,
const N: usize = 0,
> where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
pack: Arc<QueuePack<T, I, G, K, N>>,
preferred_queue_index: AtomicUsize,
pop_count: AtomicUsize,
}
impl<T, I, const G: usize, const K: usize, const N: usize> Clone
for PackConsumerHandle<T, I, G, K, N>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
fn clone(&self) -> Self {
self.pack.consumer()
}
}
impl<T, I, const G: usize, const K: usize, const N: usize> PackConsumerHandle<T, I, G, K, N>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
const fn new(pack: Arc<QueuePack<T, I, G, K, N>>, preferred_queue_index: usize) -> Self {
Self {
pack,
preferred_queue_index: AtomicUsize::new(preferred_queue_index),
pop_count: AtomicUsize::new(0),
}
}
pub fn preferred_queue_index(&self) -> usize {
self.preferred_queue_index.load(Ordering::Relaxed)
}
pub fn pop_count(&self) -> usize {
self.pop_count.load(Ordering::Relaxed)
}
pub fn try_pop(&self) -> Result<T, QueueError> {
let queue_idx = self.preferred_queue_index.load(Ordering::Relaxed);
match self.pack.queues[queue_idx].try_pop() {
Ok(value) => {
let count = self.pop_count.fetch_add(1, Ordering::Relaxed) + 1;
if count >= K {
self.pop_count.store(0, Ordering::Relaxed);
let next_idx = (queue_idx + 1) % G;
self.preferred_queue_index
.store(next_idx, Ordering::Relaxed);
}
return Ok(value);
},
Err(QueueError::Empty) => {
self.pop_count.store(0, Ordering::Relaxed);
},
Err(e) => return Err(e),
}
for i in 1..G {
let scan_idx = (queue_idx + i) % G;
match self.pack.queues[scan_idx].try_pop() {
Ok(value) => {
self.preferred_queue_index
.store(scan_idx, Ordering::Relaxed);
self.pop_count.store(1, Ordering::Relaxed);
return Ok(value);
},
Err(QueueError::Empty) => {},
Err(e) => return Err(e),
}
}
Err(QueueError::Empty)
}
pub fn scan_stats(&self) -> Vec<QueueStats> {
self.pack.queue_stats()
}
}
impl<T, I, const G: usize, const K: usize, const N: usize> QueueConsumer<T>
for PackConsumerHandle<T, I, G, K, N>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
fn try_pop(&self) -> Result<T, QueueError> {
let queue_idx = self.preferred_queue_index.load(Ordering::Relaxed);
self.pack.queues[queue_idx].try_pop()
}
fn pop(&self) -> Result<T, QueueError> {
let queue_idx = self.preferred_queue_index.load(Ordering::Relaxed);
match self.pack.queues[queue_idx].pop() {
Ok(value) => {
let count = self.pop_count.fetch_add(1, Ordering::Relaxed) + 1;
if count >= K {
self.pop_count.store(0, Ordering::Relaxed);
let next_idx = (queue_idx + 1) % G;
self.preferred_queue_index
.store(next_idx, Ordering::Relaxed);
}
return Ok(value);
},
Err(QueueError::Empty) => {
self.pop_count.store(0, Ordering::Relaxed);
},
Err(e) => return Err(e),
}
for i in 1..G {
let scan_idx = (queue_idx + i) % G;
match self.pack.queues[scan_idx].pop() {
Ok(value) => {
self.preferred_queue_index
.store(scan_idx, Ordering::Relaxed);
self.pop_count.store(1, Ordering::Relaxed);
return Ok(value);
},
Err(QueueError::Empty) => {},
Err(e) => return Err(e),
}
}
Err(QueueError::Empty)
}
fn pop_with_seq(&self) -> Result<(T, usize), QueueError> {
let queue_idx = self.preferred_queue_index.load(Ordering::Relaxed);
match self.pack.queues[queue_idx].pop_impl(true) {
Ok((value, local_index)) => {
let count = self.pop_count.fetch_add(1, Ordering::Relaxed) + 1;
if count >= K {
self.pop_count.store(0, Ordering::Relaxed);
let next_idx = (queue_idx + 1) % G;
self.preferred_queue_index
.store(next_idx, Ordering::Relaxed);
}
let global_seq = (queue_idx << 24) | (local_index & 0x00FF_FFFF);
return Ok((value, global_seq));
},
Err(QueueError::Empty) => {
self.pop_count.store(0, Ordering::Relaxed);
},
Err(e) => return Err(e),
}
for i in 1..G {
let scan_idx = (queue_idx + i) % G;
match self.pack.queues[scan_idx].pop_impl(true) {
Ok((value, local_index)) => {
self.preferred_queue_index
.store(scan_idx, Ordering::Relaxed);
self.pop_count.store(1, Ordering::Relaxed);
let global_seq = (scan_idx << 24) | (local_index & 0x00FF_FFFF);
return Ok((value, global_seq));
},
Err(QueueError::Empty) => {},
Err(e) => return Err(e),
}
}
Err(QueueError::Empty)
}
fn peek(&self) -> Result<T, QueueError> {
let queue_idx = self.preferred_queue_index.load(Ordering::Relaxed);
match self.pack.queues[queue_idx].peek() {
Ok(value) => return Ok(value),
Err(QueueError::Empty) => {},
Err(e) => return Err(e),
}
for i in 1..G {
let scan_idx = (queue_idx + i) % G;
match self.pack.queues[scan_idx].peek() {
Ok(value) => return Ok(value),
Err(QueueError::Empty) => {},
Err(e) => return Err(e),
}
}
Err(QueueError::Empty)
}
fn peek_with_seq(&self) -> Result<(T, usize), QueueError> {
let queue_idx = self.preferred_queue_index.load(Ordering::Relaxed);
match self.pack.queues[queue_idx].peek() {
Ok(value) => {
let seq = queue_idx << 24;
return Ok((value, seq));
},
Err(QueueError::Empty) => {},
Err(e) => return Err(e),
}
for i in 1..G {
let scan_idx = (queue_idx + i) % G;
match self.pack.queues[scan_idx].peek() {
Ok(value) => {
let seq = scan_idx << 24;
return Ok((value, seq));
},
Err(QueueError::Empty) => {},
Err(e) => return Err(e),
}
}
Err(QueueError::Empty)
}
fn pop_if<F>(&self, mut predicate: F) -> Result<T, QueueError>
where
F: FnMut(&T, usize) -> bool,
{
let queue_idx = self.preferred_queue_index.load(Ordering::Relaxed);
match self.pack.queues[queue_idx].peek() {
Ok(value) => {
let seq = queue_idx << 24;
if predicate(&value, seq) {
match self.pack.queues[queue_idx].pop() {
Ok(popped) => {
let count = self.pop_count.fetch_add(1, Ordering::Relaxed) + 1;
if count >= K {
self.pop_count.store(0, Ordering::Relaxed);
let next_idx = (queue_idx + 1) % G;
self.preferred_queue_index
.store(next_idx, Ordering::Relaxed);
}
return Ok(popped);
},
Err(QueueError::Empty) => {
self.pop_count.store(0, Ordering::Relaxed);
},
Err(e) => return Err(e),
}
}
},
Err(QueueError::Empty) => {
self.pop_count.store(0, Ordering::Relaxed);
},
Err(e) => return Err(e),
}
for i in 1..G {
let scan_idx = (queue_idx + i) % G;
match self.pack.queues[scan_idx].peek() {
Ok(value) => {
let seq = scan_idx << 24;
if predicate(&value, seq) {
match self.pack.queues[scan_idx].pop() {
Ok(popped) => {
self.preferred_queue_index
.store(scan_idx, Ordering::Relaxed);
self.pop_count.store(1, Ordering::Relaxed);
return Ok(popped);
},
Err(QueueError::Empty) => {},
Err(e) => return Err(e),
}
}
},
Err(QueueError::Empty) => {},
Err(e) => return Err(e),
}
}
Err(QueueError::Empty)
}
fn consume<F>(&self, mut consumer_fn: F) -> usize
where
F: FnMut(T, usize) -> bool,
{
let queue_idx = self.preferred_queue_index.load(Ordering::Relaxed);
let mut total_consumed = 0;
for i in 0..G {
let scan_idx = (queue_idx + i) % G;
while let Ok((value, local_index)) = self.pack.queues[scan_idx].pop_impl(true) {
let global_seq = (scan_idx << 24) | (local_index & 0x00FF_FFFF);
total_consumed += 1;
if scan_idx == queue_idx {
let count = self.pop_count.fetch_add(1, Ordering::Relaxed) + 1;
if count >= K {
self.pop_count.store(0, Ordering::Relaxed);
let next_idx = (queue_idx + 1) % G;
self.preferred_queue_index
.store(next_idx, Ordering::Relaxed);
}
} else {
self.preferred_queue_index
.store(scan_idx, Ordering::Relaxed);
self.pop_count.store(1, Ordering::Relaxed);
}
if consumer_fn(value, global_seq) {
return total_consumed;
}
}
}
total_consumed
}
fn is_empty(&self) -> bool {
self.pack.is_empty()
}
fn size(&self) -> usize {
self.pack.len()
}
}
impl<T, I, const G: usize, const K: usize, const N: usize> QueueFactory<T>
for Arc<QueuePack<T, I, G, K, N>>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
type Producer = PackProducerHandle<T, I, G, K, N>;
type Consumer = PackConsumerHandle<T, I, G, K, N>;
fn producer(&self) -> Self::Producer {
let queue_index = self.writer_counter.fetch_add(1, Ordering::Relaxed) % G;
PackProducerHandle::new(self.clone(), queue_index)
}
fn consumer(&self) -> Self::Consumer {
let preferred_index = self.reader_counter.fetch_add(1, Ordering::Relaxed) % G;
PackConsumerHandle::new(self.clone(), preferred_index)
}
}
unsafe impl<T, I, const G: usize, const K: usize, const N: usize> Send for QueuePack<T, I, G, K, N>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
}
unsafe impl<T, I, const G: usize, const K: usize, const N: usize> Sync for QueuePack<T, I, G, K, N>
where
T: Copy + Send + Sync + Default,
I: Copy + Into<u128>,
{
}
#[cfg(test)]
mod tests {
use super::*;
use std::{
collections::HashSet,
sync::atomic::{AtomicUsize, Ordering},
time::Instant,
};
#[test]
fn test_builder_pattern() {
let pack = queue_pack::<u32, 4, 16>()
.queue_capacity(32)
.build()
.unwrap();
assert_eq!(pack.queue_capacity(), 32);
assert_eq!(pack.capacity(), 128); }
#[test]
fn test_channels() {
let (producer, consumer) = queue_pack::<u64, 2, 8>()
.queue_capacity(16)
.channels()
.unwrap();
producer.push(42).unwrap();
producer.push(43).unwrap();
assert_eq!(consumer.pop().unwrap(), 42);
assert_eq!(consumer.pop().unwrap(), 43);
assert!(consumer.try_pop().is_err());
}
#[test]
fn test_producer_assignment() {
let pack = queue_pack::<u32, 3, 4>().queue_capacity(8).build().unwrap();
let p1 = pack.producer();
let p2 = pack.producer();
let p3 = pack.producer();
let p4 = pack.producer();
assert_eq!(p1.queue_index(), 0);
assert_eq!(p2.queue_index(), 1);
assert_eq!(p3.queue_index(), 2);
assert_eq!(p4.queue_index(), 0); }
#[test]
fn test_consumer_scanning() {
let pack = queue_pack::<u64, 3, 2>().queue_capacity(8).build().unwrap();
let p1 = pack.producer(); let p2 = pack.producer(); let p3 = pack.producer();
let consumer = pack.consumer();
p1.push(100).unwrap();
p2.push(200).unwrap();
p3.push(300).unwrap();
let mut values = Vec::new();
for _ in 0..3 {
values.push(consumer.pop().unwrap());
}
values.sort_unstable();
assert_eq!(values, vec![100, 200, 300]);
}
#[test]
fn test_queue_stats() {
let pack = queue_pack::<u32, 2, 4>().queue_capacity(4).build().unwrap();
let p1 = pack.producer(); let p2 = pack.producer();
p1.push(1).unwrap();
p1.push(2).unwrap();
p2.push(10).unwrap();
let stats = pack.queue_stats();
assert_eq!(stats.len(), 2);
assert_eq!(stats[0].index, 0);
assert_eq!(stats[0].len, 2);
assert_eq!(stats[0].capacity, 4);
assert!(!stats[0].is_empty);
assert!(!stats[0].is_full);
assert_eq!(stats[1].index, 1);
assert_eq!(stats[1].len, 1);
assert_eq!(stats[1].capacity, 4);
assert!(!stats[1].is_empty);
assert!(!stats[1].is_full);
}
#[test]
fn test_sequence_numbers() {
let (producer, consumer) = queue_pack::<u32, 2, 4>()
.queue_capacity(8)
.channels()
.unwrap();
let seq1 = producer.push_with_seq(777).unwrap();
let seq2 = producer.push_with_seq(888).unwrap();
let (val1, _) = consumer.pop_with_seq().unwrap();
let (val2, _) = consumer.pop_with_seq().unwrap();
assert_eq!(val1, 777);
assert_eq!(val2, 888);
assert_eq!(seq1 >> 24, producer.queue_index());
assert_eq!(seq2 >> 24, producer.queue_index());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 8)]
async fn stress_test_pack() {
const QUEUE_COUNT: usize = 4;
const SCAN_THRESHOLD: usize = 16;
const PRODUCERS: usize = 4;
const CONSUMERS: usize = 4;
const ITEMS_PER_PRODUCER: usize = 10_000;
let (producer, consumer) = queue_pack::<u64, QUEUE_COUNT, SCAN_THRESHOLD>()
.queue_capacity(256)
.channels()
.unwrap();
let total_items = PRODUCERS * ITEMS_PER_PRODUCER;
let consumed_count = Arc::new(AtomicUsize::new(0));
let seen = Arc::new(tokio::sync::Mutex::new(HashSet::<u64>::new()));
let mut consumer_handles = Vec::new();
for _ in 0..CONSUMERS {
let consumer = consumer.clone();
let consumed_clone = consumed_count.clone();
let seen_clone = seen.clone();
let handle = tokio::task::spawn(async move {
loop {
if consumed_clone.load(Ordering::SeqCst) >= total_items {
break;
}
match consumer.try_pop() {
Ok(value) => {
assert!(
seen_clone.lock().await.insert(value),
"Duplicate value found: {value}"
);
consumed_clone.fetch_add(1, Ordering::SeqCst);
},
Err(QueueError::Empty) => {
tokio::task::yield_now().await;
},
Err(e) => panic!("Unexpected error: {e:?}"),
}
}
});
consumer_handles.push(handle);
}
let mut producer_handles = Vec::new();
let start = Instant::now();
for producer_id in 0..PRODUCERS {
let producer = producer.clone();
let handle = tokio::task::spawn(async move {
for item_id in 0..ITEMS_PER_PRODUCER {
let value = ((producer_id as u64) << 32) | (item_id as u64);
loop {
match producer.try_push(value) {
Ok(()) => break,
Err((_, QueueError::Full)) => {
tokio::task::yield_now().await;
},
Err((_, e)) => panic!("Unexpected error: {e:?}"),
}
}
}
});
producer_handles.push(handle);
}
for handle in producer_handles {
handle.await.unwrap();
}
while consumed_count.load(Ordering::SeqCst) < total_items {
tokio::time::sleep(tokio::time::Duration::from_millis(1)).await;
}
for handle in consumer_handles {
handle.await.unwrap();
}
let elapsed = start.elapsed();
let throughput = (total_items as f64) / elapsed.as_secs_f64();
println!(
"Pack stress test: {QUEUE_COUNT} queues, {PRODUCERS} producers, {CONSUMERS} consumers, {ITEMS_PER_PRODUCER} items = {total_items} total in {elapsed:?} ({throughput:.0} ops/sec)"
);
assert_eq!(consumed_count.load(Ordering::SeqCst), total_items);
let final_seen_count = seen.lock().await.len();
assert_eq!(final_seen_count, total_items);
}
#[test]
fn test_pop_if() {
let (producer, consumer) = queue_pack::<u32, 2, 4>()
.queue_capacity(8)
.channels()
.unwrap();
producer.push(2).unwrap();
producer.push(1).unwrap();
producer.push(3).unwrap();
let result = consumer.pop_if(|&value, _seq| value % 2 == 0);
assert_eq!(result.unwrap(), 2);
assert_eq!(consumer.pop().unwrap(), 1);
assert_eq!(consumer.pop().unwrap(), 3);
}
#[test]
fn test_consume_function() {
let (producer, consumer) = queue_pack::<u32, 2, 4>()
.queue_capacity(8)
.channels()
.unwrap();
for i in 0..5 {
producer.push(i).unwrap();
}
let mut collected = Vec::new();
let consumed = consumer.consume(|value, _seq| {
collected.push(value);
value >= 2 });
assert!(consumed > 0);
assert!(!collected.is_empty());
assert!(collected.contains(&2) || collected.iter().any(|&x| x > 2));
}
}