use std::{collections::VecDeque, sync::Arc};
use crate::async_runtime::lock::Mutex;
use crate::async_util::CondVar;
pub struct AsyncQueue<T> {
queue: Mutex<VecDeque<T>>,
capacity: usize,
not_empty: CondVar,
not_full: CondVar,
}
impl<T> AsyncQueue<T> {
pub fn new(capacity: usize) -> Arc<Self> {
assert!(capacity > 0, "AsyncQueue capacity must be > 0");
Arc::new(Self {
queue: Mutex::new(VecDeque::with_capacity(capacity)),
capacity,
not_empty: CondVar::new(),
not_full: CondVar::new(),
})
}
pub async fn push(&self, item: T) {
let mut queue = self.queue.lock().await;
while queue.len() >= self.capacity {
queue = self.not_full.wait(queue).await;
}
queue.push_back(item);
self.not_empty.signal();
}
pub async fn recv(&self) -> T {
let mut queue = self.queue.lock().await;
while queue.is_empty() {
queue = self.not_empty.wait(queue).await;
}
let item = queue.pop_front().expect("queue not empty");
self.not_full.signal();
item
}
pub async fn len(&self) -> usize {
self.queue.lock().await.len()
}
pub async fn is_empty(&self) -> bool {
self.queue.lock().await.is_empty()
}
pub fn capacity(&self) -> usize {
self.capacity
}
}