use futures::prelude::Stream;
use futures::stream::FusedStream;
use std::fmt::Debug;
use std::pin::Pin;
#[pin_project::pin_project]
pub struct DropFuse<S> {
#[pin]
stream: Option<S>,
}
impl<S: Stream> Stream for DropFuse<S> {
type Item = S::Item;
fn poll_next(
self: Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
let mut this = self.project();
let Some(stream) = this.stream.as_mut().as_pin_mut() else {
return std::task::Poll::Ready(None);
};
let item = futures::ready!(stream.poll_next(cx));
if let Some(v) = item {
return std::task::Poll::Ready(Some(v));
}
this.stream.set(None); std::task::Poll::Ready(None)
}
fn size_hint(&self) -> (usize, Option<usize>) {
if let Some(st) = self.stream.as_ref() {
st.size_hint()
} else {
(0, Some(0))
}
}
}
impl<S: Stream> FusedStream for DropFuse<S> {
fn is_terminated(&self) -> bool {
self.stream.is_none()
}
}
pub trait StreamUtils: Sized {
fn drop_fuse(self) -> DropFuse<Self>;
}
impl<S: Stream> StreamUtils for S {
fn drop_fuse(self) -> DropFuse<Self> {
DropFuse { stream: Some(self) }
}
}
#[pin_project::pin_project]
pub struct SelectOkStream<S>(#[pin] S);
impl<S: Stream<Item = Result<T, E>>, T, E: Debug> Stream for SelectOkStream<S> {
type Item = T;
fn poll_next(
self: Pin<&mut Self>,
cx: &mut std::task::Context<'_>,
) -> std::task::Poll<Option<Self::Item>> {
let mut this = self.project();
loop {
return match this.0.as_mut().poll_next(cx) {
std::task::Poll::Ready(Some(Ok(v))) => std::task::Poll::Ready(Some(v)),
std::task::Poll::Ready(Some(Err(_))) => continue,
std::task::Poll::Ready(None) => std::task::Poll::Ready(None),
std::task::Poll::Pending => std::task::Poll::Pending,
};
}
}
fn size_hint(&self) -> (usize, Option<usize>) {
(0, self.0.size_hint().1)
}
}
pub trait StreamErrors<T, E>: Sized {
fn select_ok(self) -> SelectOkStream<Self>;
}
impl<S: Stream<Item = Result<T, E>>, T, E: Debug> StreamErrors<T, E> for S {
fn select_ok(self) -> SelectOkStream<Self> {
SelectOkStream(self)
}
}