use crate::{prelude::*, queues::Queue};
#[derive(Debug)]
pub struct Buffered<C, Q> {
inner: C,
buffer: Q,
}
impl<C, Q> Buffered<C, Q> {
pub(crate) fn new(inner: C, buffer: Q) -> Self {
Self { inner, buffer }
}
pub fn into_inner(self) -> C {
self.inner
}
}
impl<C, Q> AsRef<C> for Buffered<C, Q> {
fn as_ref(&self) -> &C {
&self.inner
}
}
impl<C, Q> Buffered<C, Q>
where
C: Consumer<Item: Clone>,
Q: Queue<Item = C::Item>,
{
async fn write_buffer_to_inner(&mut self) -> Result<(), C::Error> {
loop {
let amount = self
.buffer
.expose_items(
async |items| match self.inner.consume_full_slice(items).await {
Ok(()) => (items.len(), Ok(items.len())),
Err(err) => (err.count, Err(err.reason)),
},
)
.await?;
if amount == 0 {
break;
}
}
debug_assert!(!self.buffer.is_full());
Ok(())
}
}
impl<C, Q> Consumer for Buffered<C, Q>
where
C: Consumer<Item: Clone>,
Q: Queue<Item = C::Item>,
{
type Item = C::Item;
type Final = C::Final;
type Error = C::Error;
async fn consume(&mut self, val: Either<Self::Item, Self::Final>) -> Result<(), Self::Error> {
match val {
Left(item) => match self.buffer.enqueue(item) {
None => Ok(()),
Some(item) => {
self.write_buffer_to_inner().await?;
let res = self.buffer.enqueue(item);
debug_assert!(
res.is_none(),
"Enqueueing into an empty queue must always succeed."
);
Ok(())
}
},
Right(fin) => {
self.write_buffer_to_inner().await?;
self.inner.consume_final(fin).await
}
}
}
async fn flush(&mut self) -> Result<(), Self::Error> {
self.write_buffer_to_inner().await?;
self.inner.flush().await
}
}
impl<C, Q> BulkConsumer for Buffered<C, Q>
where
C: Consumer<Item: Clone>,
Q: Queue<Item = C::Item>,
{
async fn expose_slots_gracefully<F, R>(&mut self, f: F) -> Result<R, (F, Self::Error)>
where
F: AsyncFnOnce(&mut [Self::Item]) -> (usize, R),
{
if self.buffer.is_full() {
if let Err(err) = self.write_buffer_to_inner().await {
return Err((f, err));
}
}
Ok(self.buffer.expose_slots(async |buffer_slots| {
debug_assert!(!buffer_slots.is_empty(), "A non-full queue must expose at least one item slot when expose_slots is invoked.");
f(buffer_slots).await
}).await)
}
}