ufotofu 0.12.5

Abstractions for lazily consuming and producing sequences
Documentation
use crate::{prelude::*, queues::Queue};

/// A bulk consumer wrapper which collects consumed items in an internal buffer before flushing them all at once into the wrapped bulk consumer.
///
/// More efficient than [`Buffered`](super::Buffered) (which wraps regular consumers, not bulk consumers).
///
/// The internal buffer can be any value implementing the [`queues::Queue`](crate::queues::Queue) trait.
///
/// Use the `AsRef<C>` impl to access the wrapped bulk consumer.
///
/// Created via [`BulkConsumerExt::to_bulk_buffered`].
///
/// <br/>Counterpart: the [producer::BulkBuffered] type.
#[derive(Debug)]

pub struct BulkBuffered<C, Q> {
    inner: C,
    buffer: Q,
}

impl<C, Q> BulkBuffered<C, Q> {
    pub(crate) fn new(inner: C, buffer: Q) -> Self {
        Self { inner, buffer }
    }

    /// Consumes `self` and returns the wrapped bulk consumer.
    pub fn into_inner(self) -> C {
        self.inner
    }
}

impl<C, Q> AsRef<C> for BulkBuffered<C, Q> {
    fn as_ref(&self) -> &C {
        &self.inner
    }
}

impl<C, Q> BulkBuffered<C, Q>
where
    C: BulkConsumer<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.bulk_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 BulkBuffered<C, Q>
where
    C: BulkConsumer<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 BulkBuffered<C, Q>
where
    C: BulkConsumer<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)
    }
}