From 8471e0a0ee7f6c973fb517ccb7efcf6c7e2ddc6f Mon Sep 17 00:00:00 2001 From: Carl Lerche Date: Sat, 11 Jan 2020 12:32:19 -0800 Subject: [PATCH] stream: add `empty()` and `pending()` (#2092) `stream::empty()` is the asynchronous equivalent to `std::iter::empty()`. `pending()` provides a stream that never becomes ready. --- tokio/src/stream/empty.rs | 48 ++++++++++++++++++++++++++++++++ tokio/src/stream/mod.rs | 6 ++++ tokio/src/stream/pending.rs | 52 +++++++++++++++++++++++++++++++++++ tokio/tests/stream_empty.rs | 11 ++++++++ tokio/tests/stream_pending.rs | 14 ++++++++++ 5 files changed, 131 insertions(+) create mode 100644 tokio/src/stream/empty.rs create mode 100644 tokio/src/stream/pending.rs create mode 100644 tokio/tests/stream_empty.rs create mode 100644 tokio/tests/stream_pending.rs diff --git a/tokio/src/stream/empty.rs b/tokio/src/stream/empty.rs new file mode 100644 index 000000000..a320d44b1 --- /dev/null +++ b/tokio/src/stream/empty.rs @@ -0,0 +1,48 @@ +use crate::stream::Stream; + +use core::marker::PhantomData; +use core::pin::Pin; +use core::task::{Context, Poll}; + +/// Stream for the [`empty`] function. +#[derive(Debug)] +#[must_use = "streams do nothing unless polled"] +pub struct Empty(PhantomData); + +impl Unpin for Empty {} + +/// Creates a stream that yields nothing. +/// +/// The returned stream is immediately ready and returns `None`. Use +/// [`stream::pending()`](super::pending()) to obtain a stream that is never +/// ready. +/// +/// # Examples +/// +/// Basic usage: +/// +/// ``` +/// use tokio::stream::{self, StreamExt}; +/// +/// #[tokio::main] +/// async fn main() { +/// let mut none = stream::empty::(); +/// +/// assert_eq!(None, none.next().await); +/// } +/// ``` +pub const fn empty() -> Empty { + Empty(PhantomData) +} + +impl Stream for Empty { + type Item = T; + + fn poll_next(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll> { + Poll::Ready(None) + } + + fn size_hint(&self) -> (usize, Option) { + (0, Some(0)) + } +} diff --git a/tokio/src/stream/mod.rs b/tokio/src/stream/mod.rs index 2bbb68020..b7b02d02f 100644 --- a/tokio/src/stream/mod.rs +++ b/tokio/src/stream/mod.rs @@ -10,6 +10,9 @@ use all::AllFuture; mod any; use any::AnyFuture; +mod empty; +pub use empty::{empty, Empty}; + mod filter; use filter::Filter; @@ -31,6 +34,9 @@ use merge::Merge; mod next; use next::Next; +mod pending; +pub use pending::{pending, Pending}; + mod try_next; use try_next::TryNext; diff --git a/tokio/src/stream/pending.rs b/tokio/src/stream/pending.rs new file mode 100644 index 000000000..8d954a98a --- /dev/null +++ b/tokio/src/stream/pending.rs @@ -0,0 +1,52 @@ +use crate::stream::Stream; + +use core::marker::PhantomData; +use core::pin::Pin; +use core::task::{Context, Poll}; + +/// Stream for the [`pending`] function. +#[derive(Debug)] +#[must_use = "streams do nothing unless polled"] +pub struct Pending(PhantomData); + +impl Unpin for Pending {} + +/// Creates a stream that is never ready +/// +/// The returned stream is never ready. Attempting to call +/// [`next()`](crate::stream::StreamExt::next) will never complete. Use +/// [`stream::empty()`](super::empty()) to obtain a stream that is is +/// immediately empty but returns no values. +/// +/// # Examples +/// +/// Basic usage: +/// +/// ```no_run +/// use tokio::stream::{self, StreamExt}; +/// +/// #[tokio::main] +/// async fn main() { +/// let mut never = stream::empty::(); +/// +/// // This will never complete +/// never.next().await; +/// +/// unreachable!(); +/// } +/// ``` +pub const fn pending() -> Pending { + Pending(PhantomData) +} + +impl Stream for Pending { + type Item = T; + + fn poll_next(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll> { + Poll::Pending + } + + fn size_hint(&self) -> (usize, Option) { + (0, None) + } +} diff --git a/tokio/tests/stream_empty.rs b/tokio/tests/stream_empty.rs new file mode 100644 index 000000000..f278076d1 --- /dev/null +++ b/tokio/tests/stream_empty.rs @@ -0,0 +1,11 @@ +use tokio::stream::{self, Stream, StreamExt}; + +#[tokio::test] +async fn basic_usage() { + let mut stream = stream::empty::(); + + for _ in 0..2 { + assert_eq!(stream.size_hint(), (0, Some(0))); + assert_eq!(None, stream.next().await); + } +} diff --git a/tokio/tests/stream_pending.rs b/tokio/tests/stream_pending.rs new file mode 100644 index 000000000..f4d3080de --- /dev/null +++ b/tokio/tests/stream_pending.rs @@ -0,0 +1,14 @@ +use tokio::stream::{self, Stream, StreamExt}; +use tokio_test::{assert_pending, task}; + +#[tokio::test] +async fn basic_usage() { + let mut stream = stream::pending::(); + + for _ in 0..2 { + assert_eq!(stream.size_hint(), (0, None)); + + let mut next = task::spawn(async { stream.next().await }); + assert_pending!(next.poll()); + } +}