ufotofu 0.10.1

Abstractions for lazily consuming and producing sequences
Documentation
use core::{cmp::min, fmt};

use alloc::collections::VecDeque;

use crate::queues::Queue;

/// A queue of unbounded capacity.
///
/// Use the methods of the [Queue] trait to interact with the contents of the queue.
///
/// Created by the [`new_unbounded`](crate::queues::new_unbounded) and [`new_unbounded_with`](crate::queues::new_unbounded_with) functions.
pub struct Unbounded<T> {
    /// Buffered items.
    buffer: VecDeque<T>,
    ///  Number of items in the queue.
    amount: usize,
    /// The function we use to initialise buffer slots, or to reset them after having produced an item.
    initialise_memory: fn() -> T,
}

impl<T> Unbounded<T> {
    /// Creates a new unbounded queue.
    ///
    /// The `initialise_memory` function is used internally to ensure that all queue slots contain valid memory at all times. The specific choice of `T` returned by that funciton does not affect the observable semantics of the queue at all.
    pub(crate) fn new(initialise_memory: fn() -> T) -> Self {
        Self {
            buffer: VecDeque::new(),
            amount: 0,
            initialise_memory,
        }
    }
}

impl<T> Unbounded<T> {
    /// Returns a slice containing the next items that should be read.
    fn readable_slice(&mut self) -> &[T] {
        let (fst, snd) = self.buffer.as_slices();

        if fst.is_empty() {
            &snd[..min(snd.len(), self.amount)]
        } else {
            &fst[..min(fst.len(), self.amount)]
        }
    }

    /// Returns a slice containing the next slots that should be written to.
    fn writeable_slice(&mut self) -> &mut [T] {
        let (fst, snd) = self.buffer.as_mut_slices();

        let fst_len = fst.len();
        if fst_len > self.amount {
            &mut fst[self.amount..]
        } else {
            let snd_len = snd.len();
            &mut snd[min(snd_len, self.amount - fst_len)..]
        }
    }

    fn grow_if_needed(&mut self) {
        if self.amount == self.buffer.len() {
            let new_len = self.buffer.len() * 2 + 1;
            self.buffer.resize_with(new_len, self.initialise_memory);
        }
    }
}

impl<T> Queue for Unbounded<T> {
    type Item = T;

    fn len(&self) -> usize {
        self.amount
    }

    fn is_full(&self) -> bool {
        false
    }

    /// Always returns `None`, this queue is unbounded.
    fn max_capacity(&self) -> Option<usize> {
        None
    }

    /// Always returns `None`, this queue is unbounded.
    fn enqueue(&mut self, item: T) -> Option<T> {
        self.grow_if_needed();

        self.buffer[self.amount] = item;
        self.amount += 1;

        None
    }

    /// Never calls `f` with an empty slice, this queue is unbounded.
    async fn expose_slots<F, R>(&mut self, f: F) -> R
    where
        F: AsyncFnOnce(&mut [Self::Item]) -> (usize, R),
    {
        self.grow_if_needed();

        let (amount, ret) = f(self.writeable_slice()).await;
        self.amount += amount;
        ret
    }

    fn dequeue(&mut self) -> Option<T> {
        if self.amount == 0 {
            None
        } else {
            let tmp = self.buffer.pop_front();
            self.amount -= 1;

            Some(tmp.unwrap())
        }
    }

    async fn expose_items<F, R>(&mut self, f: F) -> R
    where
        F: AsyncFnOnce(&[Self::Item]) -> (usize, R),
    {
        let (amount, ret) = f(self.readable_slice()).await;
        self.amount -= amount;

        // TODO use truncate_front once stabilised: https://github.com/rust-lang/rust/issues/140667
        let rotate_by = amount;
        self.buffer.rotate_left(rotate_by);
        self.buffer.truncate(self.amount);

        ret
    }
}

impl<T: fmt::Debug> fmt::Debug for Unbounded<T> {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        f.debug_struct("Unbounded")
            .field("len", &self.amount)
            .field("data", &DataDebugger(self))
            .finish()
    }
}

struct DataDebugger<'q, T>(&'q Unbounded<T>);

impl<T: fmt::Debug> fmt::Debug for DataDebugger<'_, T> {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        let mut list = f.debug_list();

        for item in self.0.buffer.iter().take(self.0.amount) {
            list.entry(item);
        }

        list.finish()
    }
}

#[cfg(test)]
mod tests {
    use core::array;

    use alloc::format;

    use crate::queues::QueueExt;

    use super::*;

    #[test]
    fn enqueues_and_dequeues_with_correct_amount() {
        let mut queue: Unbounded<u8> = Unbounded::new(Default::default);

        assert_eq!(queue.enqueue(7), None);
        assert_eq!(queue.enqueue(21), None);
        assert_eq!(queue.enqueue(196), None);
        assert_eq!(queue.len(), 3);

        assert_eq!(queue.enqueue(233), None);
        assert_eq!(queue.len(), 4);

        // Queue should be first-in, first-out.
        assert_eq!(queue.dequeue(), Some(7));
        assert_eq!(queue.len(), 3);
    }

    #[test]
    fn bulk_enqueues_and_dequeues_with_correct_amount() {
        pollster::block_on(async {
            let mut queue: Unbounded<u8> = Unbounded::new(Default::default);
            let mut buf = [0; 4];

            let enqueue_amount = queue.bulk_enqueue(b"ufo").await;
            let dequeue_amount = queue.bulk_dequeue(&mut buf).await;

            assert_eq!(enqueue_amount, dequeue_amount);
        });
    }

    #[test]
    fn returns_none_on_dequeue_when_queue_is_empty() {
        let mut queue: Unbounded<u8> = Unbounded::new(Default::default);

        // Enqueue and then dequeue an item.
        let _ = queue.enqueue(7);
        let _ = queue.dequeue();

        // The queue is now empty.
        assert!(queue.dequeue().is_none());
    }

    #[test]
    fn test_debug_impl() {
        let mut queue: Unbounded<u8> = Unbounded::new(Default::default);

        assert_eq!(queue.enqueue(7), None);
        assert_eq!(queue.enqueue(21), None);
        assert_eq!(queue.enqueue(196), None);
        assert_eq!(
            format!("{queue:?}"),
            "Unbounded { len: 3, data: [7, 21, 196] }"
        );

        assert_eq!(queue.dequeue(), Some(7));
        assert_eq!(
            format!("{queue:?}"),
            "Unbounded { len: 2, data: [21, 196] }"
        );

        assert_eq!(queue.dequeue(), Some(21));
        assert_eq!(format!("{queue:?}"), "Unbounded { len: 1, data: [196] }");

        assert_eq!(queue.enqueue(33), None);
        assert_eq!(
            format!("{queue:?}"),
            "Unbounded { len: 2, data: [196, 33] }"
        );

        assert_eq!(queue.enqueue(17), None);
        assert_eq!(
            format!("{queue:?}"),
            "Unbounded { len: 3, data: [196, 33, 17] }"
        );

        assert_eq!(queue.enqueue(200), None);
        assert_eq!(
            format!("{queue:?}"),
            "Unbounded { len: 4, data: [196, 33, 17, 200] }"
        );
    }

    #[test]
    fn regression() {
        pollster::block_on(async {
            let mut queue: Unbounded<u8> = Unbounded::new(Default::default);

            let count = 34usize;
            for i in 0..count {
                assert_eq!(queue.enqueue(i as u8), None);
            }

            let mut buf = [99; 99];

            let deq_amount_1 = 2;
            let deq_amount_2 = count - deq_amount_1;

            assert_eq!(
                queue.bulk_dequeue(&mut buf[..deq_amount_1]).await,
                deq_amount_1
            );
            assert_eq!(
                queue.bulk_dequeue(&mut buf[deq_amount_1..]).await,
                deq_amount_2
            );

            let expected: [u8; 34] = array::from_fn(|i| i as u8);

            assert_eq!(&buf[..count], expected.as_slice());
        });
    }
}