use std::io::{self, ErrorKind};
use either::{Either, Left, Right};
use crate::queues::Queue;
use futures_lite::io::{AsyncBufRead, AsyncBufReadExt, AsyncRead, AsyncReadExt};
use crate::prelude::*;
pub fn buf_reader_to_bulk_producer<R>(reader: R) -> BufReaderToBulkProducer<R> {
BufReaderToBulkProducer::new(reader)
}
#[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Debug)]
pub struct BufReaderToBulkProducer<R>(R);
impl<R> BufReaderToBulkProducer<R> {
fn new(reader: R) -> Self {
Self(reader)
}
pub fn into_inner(self) -> R {
self.0
}
}
impl<R> AsRef<R> for BufReaderToBulkProducer<R> {
fn as_ref(&self) -> &R {
&self.0
}
}
impl<R> AsMut<R> for BufReaderToBulkProducer<R> {
fn as_mut(&mut self) -> &mut R {
&mut self.0
}
}
impl<R> Producer for BufReaderToBulkProducer<R>
where
R: AsyncBufRead + Unpin,
{
type Item = u8;
type Final = ();
type Error = io::Error;
async fn produce(&mut self) -> Result<Either<Self::Item, Self::Final>, Self::Error> {
let mut buf = [0; 1];
match self.0.read_exact(&mut buf).await {
Err(err) => {
if err.kind() == ErrorKind::UnexpectedEof {
Ok(Right(()))
} else {
Err(err)
}
}
Ok(()) => Ok(Left(buf[0])),
}
}
async fn slurp(&mut self) -> Result<(), Self::Error> {
self.0.fill_buf().await?;
Ok(())
}
}
impl<R> BulkProducer for BufReaderToBulkProducer<R>
where
R: AsyncBufRead + Unpin,
{
async fn expose_items_gracefully<F, Ret>(
&mut self,
f: F,
) -> Result<Either<Ret, (F, Self::Final)>, (F, Self::Error)>
where
F: AsyncFnOnce(&[Self::Item]) -> (usize, Ret),
{
let buf = match self.0.fill_buf().await {
Ok(b) => b,
Err(err) => return Err((f, err)),
};
if buf.is_empty() {
Ok(Right((f, ())))
} else {
let (amount, ret) = f(buf).await;
self.0.consume(amount);
Ok(Left(ret))
}
}
}
pub fn reader_to_bulk_producer<R, Q>(reader: R, queue: Q) -> ReaderToBulkProducer<R, Q> {
ReaderToBulkProducer::new(reader, queue)
}
#[derive(Debug)]
pub struct ReaderToBulkProducer<R, Q> {
reader: R,
buffer: Q,
last: Option<Result<(), io::Error>>,
}
impl<R, Q> ReaderToBulkProducer<R, Q> {
fn new(reader: R, queue: Q) -> Self {
Self {
reader,
buffer: queue,
last: None,
}
}
pub fn into_inner(self) -> (R, Q) {
(self.reader, self.buffer)
}
}
impl<R, Q> AsRef<R> for ReaderToBulkProducer<R, Q> {
fn as_ref(&self) -> &R {
&self.reader
}
}
impl<R, Q> ReaderToBulkProducer<R, Q>
where
R: AsyncRead + Unpin,
Q: Queue<Item = u8>,
{
async fn fill_buffer_from_inner(&mut self) {
while self.last.is_none() && !self.buffer.is_full() {
match self
.buffer
.expose_slots(async |slots| match self.reader.read(slots).await {
Ok(amount) => (amount, Ok(amount)),
Err(err) => (0, Err(err)),
})
.await
{
Ok(0) => {
self.last = Some(Ok(()));
break;
}
Ok(_) => { }
Err(err) => {
self.last = Some(Err(err));
break;
}
}
}
debug_assert!(self.last.is_some() || !self.buffer.is_empty());
}
fn check_last(&mut self) -> Option<Result<(), io::Error>> {
if !self.buffer.is_empty() {
None
} else {
self.last.take()
}
}
}
impl<R, Q> Producer for ReaderToBulkProducer<R, Q>
where
R: AsyncRead + Unpin,
Q: Queue<Item = u8>,
{
type Item = u8;
type Final = ();
type Error = io::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 => Ok(()),
}
}
}
impl<R, Q> BulkProducer for ReaderToBulkProducer<R, Q>
where
R: AsyncRead + Unpin,
Q: Queue<Item = u8>,
{
async fn expose_items_gracefully<F, Ret>(
&mut self,
f: F,
) -> Result<Either<Ret, (F, Self::Final)>, (F, Self::Error)>
where
F: AsyncFnOnce(&[Self::Item]) -> (usize, Ret),
{
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))
}
}
}
}