#[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)]
pub struct BlockingQueue<T> {
queue: Mutex<VecDeque<T>>,
condvar: Condvar,
closed: Mutex<bool>,
}
#[cfg(not(alloc_frugal))]
impl<T> BlockingQueue<T> {
pub fn new() -> Self {
Self {
queue: Mutex::new(VecDeque::new()),
condvar: Condvar::new(),
closed: Mutex::new(false),
}
}
pub fn with_capacity(capacity: usize) -> Self {
Self {
queue: Mutex::new(VecDeque::with_capacity(capacity)),
condvar: Condvar::new(),
closed: Mutex::new(false),
}
}
pub fn push(&self, item: T) -> Result<(), QueueError> {
if *lock(&self.closed) {
return Err(QueueError::Closed);
}
let mut queue = lock(&self.queue);
queue.push_back(item);
self.condvar.notify_one();
Ok(())
}
pub fn pop(&self) -> Result<T, QueueError> {
let mut queue = lock(&self.queue);
loop {
if let Some(item) = queue.pop_front() {
return Ok(item);
}
if *lock(&self.closed) {
return Err(QueueError::Closed);
}
queue = self.condvar.wait(queue).unwrap_or_else(|e| e.into_inner());
}
}
pub fn pop_timeout(&self, timeout: Duration) -> Result<T, QueueError> {
let start = Instant::now();
let mut queue = lock(&self.queue);
loop {
if let Some(item) = queue.pop_front() {
return Ok(item);
}
if *lock(&self.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(queue, remaining).unwrap_or_else(|e| e.into_inner());
queue = result.0;
}
}
pub fn try_pop(&self) -> Option<T> {
let mut queue = lock(&self.queue);
queue.pop_front()
}
pub fn close(&self) {
*lock(&self.closed) = true;
self.condvar.notify_all();
}
pub fn is_closed(&self) -> bool {
*lock(&self.closed)
}
pub fn len(&self) -> usize {
lock(&self.queue).len()
}
pub fn is_empty(&self) -> bool {
lock(&self.queue).is_empty()
}
pub fn clear(&self) {
lock(&self.queue).clear();
}
}
#[cfg(not(alloc_frugal))]
impl<T> Default for BlockingQueue<T> {
fn default() -> Self {
Self::new()
}
}