nimbusqueue 0.2.7

fifo collection
Documentation
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>),
}

/// A thread-safe queue implementation for no_std environments.
#[derive(Debug)]
pub struct Queue<T> {
    queue: QueueType<T>,
    count: AtomicUsize,
}

/// Iterator for Queue elements
pub struct QueueIter<T> {
    items: Vec<T>,
    index: usize,
}

impl<T: Send + Clone> Queue<T> {
    /// Creates a new empty segmented queue with unlimited capacity.
    pub fn new_segmented() -> Self {
        Self {
            queue: QueueType::Segmented(SegQueue::new()),
            count: AtomicUsize::new(0),
        }
    }

    /// Creates a new empty array queue with specified capacity.
    pub fn new_array(capacity: usize) -> Self {
        Self {
            queue: QueueType::Array(ArrayQueue::new(capacity)),
            count: AtomicUsize::new(0),
        }
    }

    /// Returns the current number of elements in the queue.
    pub fn length(&self) -> usize {
        self.count.load(Ordering::Acquire)
    }

    /// Returns true if the queue is empty.
    pub fn is_empty(&self) -> bool {
        self.length() == 0
    }

    /// Adds an item to the queue.
    /// Returns Err if the queue is at capacity.
    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")
                }
            }
        }
    }

    /// Removes and returns the first item in the queue.
    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
    }

    /// Creates a vector containing all items in the queue.
    /// Returns None if the queue is empty.
    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);
        }

        // Restore items to queue
        for item in temp {
            let _ = self.enqueue(item);
        }

        if items.is_empty() {
            None
        } else {
            Some(items)
        }
    }

    /// Returns an iterator over the elements of the queue.
    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()
    }
}