use alloc::vec::Vec;
use core::sync::atomic::{AtomicUsize, Ordering};
use crossbeam_queue::{ArrayQueue, SegQueue};
#[derive(Debug)]
pub enum QueueType<T> {
Segmented(SegQueue<T>),
Array(ArrayQueue<T>),
}
#[derive(Debug)]
pub struct Queue<T> {
queue: QueueType<T>,
count: AtomicUsize,
}
pub struct QueueIter<T> {
items: Vec<T>,
index: usize,
}
impl<T: Send + Clone> Queue<T> {
pub fn new_segmented() -> Self {
Self {
queue: QueueType::Segmented(SegQueue::new()),
count: AtomicUsize::new(0),
}
}
pub fn new_array(capacity: usize) -> Self {
Self {
queue: QueueType::Array(ArrayQueue::new(capacity)),
count: AtomicUsize::new(0),
}
}
pub fn length(&self) -> usize {
self.count.load(Ordering::Acquire)
}
pub fn is_empty(&self) -> bool {
self.length() == 0
}
pub fn enqueue(&self, item: T) -> Result<(), &'static str> {
match &self.queue {
QueueType::Segmented(q) => {
q.push(item);
self.count.fetch_add(1, Ordering::Release);
Ok(())
}
QueueType::Array(q) => {
if q.push(item).is_ok() {
self.count.fetch_add(1, Ordering::Release);
Ok(())
} else {
Err("Queue is full")
}
}
}
}
pub fn dequeue(&self) -> Option<T> {
let item = match &self.queue {
QueueType::Segmented(q) => q.pop(),
QueueType::Array(q) => q.pop(),
};
if item.is_some() {
self.count.fetch_sub(1, Ordering::Release);
}
item
}
pub fn try_iter(&self) -> Option<Vec<T>> {
let mut items = Vec::new();
let mut temp = Vec::new();
while let Some(item) = self.dequeue() {
items.push(item.clone());
temp.push(item);
}
for item in temp {
let _ = self.enqueue(item);
}
if items.is_empty() {
None
} else {
Some(items)
}
}
pub fn iter(&self) -> QueueIter<T> {
let items = self.try_iter().unwrap_or_default();
QueueIter { items, index: 0 }
}
}
impl<T: Clone> Iterator for QueueIter<T> {
type Item = T;
fn next(&mut self) -> Option<Self::Item> {
if self.index < self.items.len() {
let item = self.items[self.index].clone();
self.index += 1;
Some(item)
} else {
None
}
}
}
impl<T: Clone + Send> Clone for Queue<T> {
fn clone(&self) -> Self {
let items = self.try_iter().unwrap_or_default();
let new_queue = match &self.queue {
QueueType::Segmented(_) => Queue::new_segmented(),
QueueType::Array(q) => Queue::new_array(q.capacity()),
};
for item in items {
new_queue.enqueue(item).expect("Queue should not be full");
}
new_queue
}
}
impl<T: Send + Clone> Default for Queue<T> {
fn default() -> Self {
Self::new_segmented()
}
}