use crate::{prelude::*, queues::Queue};
#[derive(Debug)]
pub struct BulkBuffered<P, Final, Error, Q> {
inner: P,
buffer: Q,
last: Option<Result<Final, Error>>,
}
impl<P, Final, Error, Q> BulkBuffered<P, Final, Error, Q> {
pub(crate) fn new(inner: P, buffer: Q) -> Self {
Self {
inner,
buffer,
last: None,
}
}
pub fn into_inner(self) -> P {
self.inner
}
}
impl<P, Final, Error, Q> AsRef<P> for BulkBuffered<P, Final, Error, Q> {
fn as_ref(&self) -> &P {
&self.inner
}
}
impl<P, Final, Error, Q> BulkBuffered<P, Final, Error, Q>
where
P: BulkProducer<Item: Clone, Final = Final, Error = Error>,
Q: Queue<Item = P::Item>,
{
async fn fill_buffer_from_inner(&mut self) {
while self.last.is_none() && !self.buffer.is_full() {
match self
.buffer
.expose_slots(async |items| {
match self.inner.bulk_overwrite_full_slice(items).await {
Ok(()) => (items.len(), Ok(items.len())),
Err(err) => (err.count, Err(err.reason)),
}
})
.await
{
Ok(0) => break,
Ok(_) => { }
Err(last) => {
self.last = Some(last);
break;
}
}
}
debug_assert!(self.last.is_some() || !self.buffer.is_empty());
}
fn check_last(&mut self) -> Option<Result<Final, Error>> {
if !self.buffer.is_empty() {
None
} else {
self.last.take()
}
}
}
impl<P, Final, Error, Q> Producer for BulkBuffered<P, Final, Error, Q>
where
P: BulkProducer<Item: Clone, Final = Final, Error = Error>,
Q: Queue<Item = P::Item>,
{
type Item = P::Item;
type Final = P::Final;
type Error = P::Error;
async fn produce(&mut self) -> Result<Either<Self::Item, Self::Final>, Self::Error> {
match self.check_last() {
Some(Ok(fin)) => Ok(Right(fin)),
Some(Err(err)) => Err(err),
None => match self.buffer.dequeue() {
Some(item) => Ok(Left(item)),
None => {
self.fill_buffer_from_inner().await;
match self.check_last() {
Some(Ok(fin)) => Ok(Right(fin)),
Some(Err(err)) => Err(err),
None => {
Ok(Left(self.buffer.dequeue().expect(
"Dequeueing from a non-empty queue must always suceed.",
)))
}
}
}
},
}
}
async fn slurp(&mut self) -> Result<(), Self::Error> {
match self.check_last() {
Some(Ok(fin)) => {
self.last = Some(Ok(fin));
Ok(())
}
Some(Err(err)) => {
debug_assert!(self.buffer.is_empty());
Err(err)
}
None => {
self.fill_buffer_from_inner().await;
if self.last.is_none() {
match self.inner.slurp().await {
Ok(()) => Ok(()),
Err(err) => {
if self.buffer.is_empty() {
Err(err)
} else {
self.last = Some(Err(err));
Ok(())
}
}
}
} else {
Ok(())
}
}
}
}
}
impl<P, Final, Error, Q> BulkProducer for BulkBuffered<P, Final, Error, Q>
where
P: BulkProducer<Item: Clone, Final = Final, Error = Error>,
Q: Queue<Item = P::Item>,
{
async fn expose_items_gracefully<F, R>(
&mut self,
f: F,
) -> Result<Either<R, (F, Self::Final)>, (F, Self::Error)>
where
F: AsyncFnOnce(&[Self::Item]) -> (usize, R),
{
match self.check_last() {
Some(Ok(fin)) => Ok(Right((f, fin))),
Some(Err(err)) => Err((f, err)),
None => {
if self.buffer.is_empty() {
self.fill_buffer_from_inner().await;
match self.check_last() {
Some(Ok(fin)) => return Ok(Right((f, fin))),
Some(Err(err)) => return Err((f, err)),
None => { }
}
}
Ok(Left(self.buffer.expose_items(async |buffer_items| {
debug_assert!(!buffer_items.is_empty(), "A non-empty queue must expose at least one item when expose_items is invoked");
f(buffer_items).await
}).await))
}
}
}
}