#[cfg(not(alloc_frugal))]
use crate::compat::Condvar;
#[cfg(not(alloc_frugal))]
use crate::compat::Instant;
#[cfg(not(alloc_frugal))]
use crate::compat::Mutex;
#[cfg(not(alloc_frugal))]
use crate::compat::lock;
#[cfg(not(alloc_frugal))]
use alloc::collections::VecDeque;
#[cfg(not(alloc_frugal))]
use core::time::Duration;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QueueError {
Full,
Empty,
Closed,
}
#[cfg(not(alloc_frugal))]
#[derive(Debug)]
struct QueueState<T> {
items: VecDeque<T>,
closed: bool,
}
#[cfg(not(alloc_frugal))]
#[derive(Debug)]
pub struct BlockingQueue<T> {
state: Mutex<QueueState<T>>,
condvar: Condvar,
}
#[cfg(not(alloc_frugal))]
impl<T> BlockingQueue<T> {
pub fn new() -> Self {
Self {
state: Mutex::new(QueueState { items: VecDeque::new(), closed: false }),
condvar: Condvar::new(),
}
}
pub fn with_capacity(capacity: usize) -> Self {
Self {
state: Mutex::new(QueueState {
items: VecDeque::with_capacity(capacity),
closed: false,
}),
condvar: Condvar::new(),
}
}
pub fn push(&self, item: T) -> Result<(), QueueError> {
let mut state = lock(&self.state);
if state.closed {
return Err(QueueError::Closed);
}
state.items.push_back(item);
drop(state);
self.condvar.notify_one();
Ok(())
}
pub fn pop(&self) -> Result<T, QueueError> {
let mut state = lock(&self.state);
loop {
if let Some(item) = state.items.pop_front() {
return Ok(item);
}
if state.closed {
return Err(QueueError::Closed);
}
state = self.condvar.wait(state).unwrap_or_else(|e| e.into_inner());
}
}
pub fn pop_timeout(&self, timeout: Duration) -> Result<T, QueueError> {
let start = Instant::now();
let mut state = lock(&self.state);
loop {
if let Some(item) = state.items.pop_front() {
return Ok(item);
}
if state.closed {
return Err(QueueError::Closed);
}
let elapsed = start.elapsed();
if elapsed >= timeout {
return Err(QueueError::Empty);
}
let remaining = timeout - elapsed;
let result =
self.condvar.wait_timeout(state, remaining).unwrap_or_else(|e| e.into_inner());
state = result.0;
}
}
pub fn try_pop(&self) -> Option<T> {
lock(&self.state).items.pop_front()
}
pub fn close(&self) {
lock(&self.state).closed = true;
self.condvar.notify_all();
}
pub fn is_closed(&self) -> bool {
lock(&self.state).closed
}
pub fn len(&self) -> usize {
lock(&self.state).items.len()
}
pub fn is_empty(&self) -> bool {
lock(&self.state).items.is_empty()
}
pub fn clear(&self) {
lock(&self.state).items.clear();
}
}
#[cfg(not(alloc_frugal))]
impl<T> Default for BlockingQueue<T> {
fn default() -> Self {
Self::new()
}
}
#[cfg(all(test, not(alloc_frugal)))]
mod tests {
use super::*;
use crate::compat::Arc;
use core::sync::atomic::{AtomicBool, Ordering};
use std::sync::Barrier;
use std::thread;
use std::time::Duration;
#[test]
fn a_push_before_a_parked_waiter_is_not_lost() {
let q: Arc<BlockingQueue<i32>> = Arc::new(BlockingQueue::new());
let started = Arc::new(Barrier::new(2));
let waiter = {
let q = Arc::clone(&q);
let started = Arc::clone(&started);
thread::spawn(move || {
started.wait();
q.pop()
})
};
started.wait();
thread::sleep(Duration::from_millis(30));
q.push(7).expect("push succeeds before close");
let popped = waiter.join().expect("consumer thread joins");
assert_eq!(
popped,
Ok(7),
"the item pushed just before the waiter parked must be delivered"
);
}
#[test]
fn close_while_waiting_returns_closed() {
let q: Arc<BlockingQueue<i32>> = Arc::new(BlockingQueue::new());
let started = Arc::new(Barrier::new(2));
let waiter = {
let q = Arc::clone(&q);
let started = Arc::clone(&started);
thread::spawn(move || {
started.wait();
q.pop()
})
};
started.wait();
thread::sleep(Duration::from_millis(30));
q.close();
let result = waiter.join().expect("consumer thread joins");
assert_eq!(result, Err(QueueError::Closed), "close must unblock a waiting consumer");
}
#[test]
fn close_and_push_linearize() {
for _ in 0..64 {
let q: Arc<BlockingQueue<i32>> = Arc::new(BlockingQueue::new());
let go = Arc::new(Barrier::new(3));
let pushed_ok = Arc::new(AtomicBool::new(false));
let pusher = {
let q = Arc::clone(&q);
let go = Arc::clone(&go);
let pushed_ok = Arc::clone(&pushed_ok);
thread::spawn(move || {
go.wait();
if q.push(1).is_ok() {
pushed_ok.store(true, Ordering::SeqCst);
}
})
};
let closer = {
let q = Arc::clone(&q);
let go = Arc::clone(&go);
thread::spawn(move || {
go.wait();
q.close();
})
};
go.wait();
pusher.join().expect("producer thread joins");
closer.join().expect("closer thread joins");
assert_eq!(
q.push(2),
Err(QueueError::Closed),
"no push may succeed after close returns"
);
if pushed_ok.load(Ordering::SeqCst) {
assert_eq!(q.try_pop(), Some(1), "an accepted push must be queued, not dropped");
}
assert!(q.is_closed());
}
}
#[test]
fn pop_timeout_respects_the_deadline() {
let q: BlockingQueue<i32> = BlockingQueue::new();
let start = Instant::now();
let result = q.pop_timeout(Duration::from_millis(40));
let waited = start.elapsed();
assert_eq!(result, Err(QueueError::Empty), "an empty queue times out as `Empty`");
assert!(
waited >= Duration::from_millis(30),
"the call must have actually waited: {waited:?}"
);
}
#[test]
fn pop_timeout_returns_a_ready_item() {
let q = BlockingQueue::new();
q.push(11).unwrap();
let result = q.pop_timeout(Duration::from_millis(200));
assert_eq!(result, Ok(11));
}
#[test]
fn close_wakes_a_timed_waiter() {
let q: Arc<BlockingQueue<i32>> = Arc::new(BlockingQueue::new());
let started = Arc::new(Barrier::new(2));
let waiter = {
let q = Arc::clone(&q);
let started = Arc::clone(&started);
thread::spawn(move || {
started.wait();
q.pop_timeout(Duration::from_secs(10))
})
};
started.wait();
thread::sleep(Duration::from_millis(30));
q.close();
let result = waiter.join().expect("consumer thread joins");
assert_eq!(result, Err(QueueError::Closed), "close must short-circuit a timed wait");
}
#[test]
fn queued_items_survive_close() {
let q: BlockingQueue<i32> = BlockingQueue::new();
q.push(1).unwrap();
q.push(2).unwrap();
q.close();
assert_eq!(q.pop(), Ok(1));
assert_eq!(q.pop(), Ok(2));
assert_eq!(q.pop(), Err(QueueError::Closed));
}
}