use core::{cmp::min, fmt};
use alloc::collections::VecDeque;
use crate::queues::Queue;
pub struct Unbounded<T> {
buffer: VecDeque<T>,
amount: usize,
initialise_memory: fn() -> T,
}
impl<T> Unbounded<T> {
pub(crate) fn new(initialise_memory: fn() -> T) -> Self {
Self {
buffer: VecDeque::new(),
amount: 0,
initialise_memory,
}
}
}
impl<T> Unbounded<T> {
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)]
}
}
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
}
fn max_capacity(&self) -> Option<usize> {
None
}
fn enqueue(&mut self, item: T) -> Option<T> {
self.grow_if_needed();
self.buffer[self.amount] = item;
self.amount += 1;
None
}
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;
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);
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);
let _ = queue.enqueue(7);
let _ = queue.dequeue();
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());
});
}
}