use futures::{stream::FuturesUnordered, Future, Stream, StreamExt};
use std::{
pin::Pin,
task::{Context, Poll, Waker},
};
pub struct FuturesStream<F> {
futures: FuturesUnordered<F>,
waker: Option<Waker>,
}
impl<F> Default for FuturesStream<F> {
fn default() -> FuturesStream<F> {
FuturesStream { futures: Default::default(), waker: None }
}
}
impl<F> FuturesStream<F> {
pub fn push(&mut self, future: F) {
self.futures.push(future);
if let Some(waker) = self.waker.take() {
waker.wake();
}
}
pub fn len(&self) -> usize {
self.futures.len()
}
}
impl<F: Future> Stream for FuturesStream<F> {
type Item = <F as Future>::Output;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
let Poll::Ready(Some(result)) = self.futures.poll_next_unpin(cx) else {
self.waker = Some(cx.waker().clone());
return Poll::Pending
};
Poll::Ready(Some(result))
}
}
#[cfg(test)]
mod tests {
use super::*;
use futures::future::{BoxFuture, FutureExt};
#[tokio::test]
async fn empty_futures_unordered_can_be_polled() {
let mut unordered = FuturesUnordered::<BoxFuture<()>>::default();
futures::future::poll_fn(|cx| {
assert_eq!(unordered.poll_next_unpin(cx), Poll::Ready(None));
assert_eq!(unordered.poll_next_unpin(cx), Poll::Ready(None));
Poll::Ready(())
})
.await;
}
#[tokio::test]
async fn deplenished_futures_unordered_can_be_polled() {
let mut unordered = FuturesUnordered::<BoxFuture<()>>::default();
unordered.push(futures::future::ready(()).boxed());
assert_eq!(unordered.next().await, Some(()));
futures::future::poll_fn(|cx| {
assert_eq!(unordered.poll_next_unpin(cx), Poll::Ready(None));
assert_eq!(unordered.poll_next_unpin(cx), Poll::Ready(None));
Poll::Ready(())
})
.await;
}
#[tokio::test]
async fn empty_futures_stream_yields_pending() {
let mut stream = FuturesStream::<BoxFuture<()>>::default();
futures::future::poll_fn(|cx| {
assert_eq!(stream.poll_next_unpin(cx), Poll::Pending);
Poll::Ready(())
})
.await;
}
#[tokio::test]
async fn futures_stream_resolves_futures_and_yields_pending() {
let mut stream = FuturesStream::default();
stream.push(futures::future::ready(17));
futures::future::poll_fn(|cx| {
assert_eq!(stream.poll_next_unpin(cx), Poll::Ready(Some(17)));
assert_eq!(stream.poll_next_unpin(cx), Poll::Pending);
Poll::Ready(())
})
.await;
}
}