#![doc = include_str!("../README.md")]
use futures_util::{FutureExt, Stream, StreamExt, TryStream, future};
use tokio::{
sync::mpsc::{self, error::SendError},
task::JoinError,
};
use tokio_stream::wrappers::ReceiverStream;
pub trait TryStreamTranspose: TryStream {
fn transpose_with<Fut, F, T, E>(
self,
buf_size: usize,
mut f: F,
) -> impl Future<Output = Result<T, E>>
where
F: FnMut(ReceiverStream<Self::Ok>) -> Fut,
Fut: Future<Output = Result<T, E>>,
Self: Send
+ Sized
+ Stream<Item = Result<Self::Ok, Self::Error>>
+ 'static,
Self::Ok: Send,
Self::Error: Send,
E: Send
+ From<Self::Error>
+ From<JoinError>
+ From<SendError<Self::Ok>>
+ 'static,
{
let (sender, recver) = mpsc::channel(buf_size);
let send_handle = tokio::spawn(async move {
tokio::pin! {
let stream = self;
}
while let Some(line) = stream.next().await {
sender.send(line?).await?;
}
Ok(())
});
let recv_handle = f(ReceiverStream::new(recver));
future::join(send_handle, recv_handle).map(|result| match result {
(Err(err), _) => Err(err.into()),
(Ok(Err(err)), _) => Err(err),
(_, result) => result,
})
}
}
impl<S: TryStream> TryStreamTranspose for S {}