From 89d969d518c9cf0b1e7a7b79ddcc755550191d20 Mon Sep 17 00:00:00 2001 From: Ben Boeckel Date: Fri, 7 Sep 2018 18:43:03 -0400 Subject: [PATCH] StreamExt: add a trait for additional Stream methods (#573) Primarily, it offers a `timeout` method for streams. --- src/lib.rs | 1 + src/util/mod.rs | 6 ++++- src/util/stream.rs | 62 ++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 68 insertions(+), 1 deletion(-) create mode 100644 src/util/stream.rs diff --git a/src/lib.rs b/src/lib.rs index 7bb343062..3d3ab483f 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -599,6 +599,7 @@ pub mod prelude { pub use util::{ FutureExt, + StreamExt, }; pub use ::std::io::{ diff --git a/src/util/mod.rs b/src/util/mod.rs index 06f0fe500..3ebd1fc70 100644 --- a/src/util/mod.rs +++ b/src/util/mod.rs @@ -1,10 +1,14 @@ //! Utilities for working with Tokio. //! //! This module contains utilities that are useful for working with Tokio. -//! Currently, this only includes [`FutureExt`], but this may grow over time. +//! Currently, this only includes [`FutureExt`] and [`StreamExt`], but this +//! may grow over time. //! //! [`FutureExt`]: trait.FutureExt.html +//! [`StreamExt`]: trait.StreamExt.html mod future; +mod stream; pub use self::future::FutureExt; +pub use self::stream::StreamExt; diff --git a/src/util/stream.rs b/src/util/stream.rs new file mode 100644 index 000000000..ef268483c --- /dev/null +++ b/src/util/stream.rs @@ -0,0 +1,62 @@ +use tokio_timer::Timeout; + +use futures::Stream; + +use std::time::Duration; + + +/// An extension trait for `Stream` that provides a variety of convenient +/// combinator functions. +/// +/// Currently, there only is a [`timeout`] function, but this will increase +/// over time. +/// +/// Users are not expected to implement this trait. All types that implement +/// `Stream` already implement `StreamExt`. +/// +/// This trait can be imported directly or via the Tokio prelude: `use +/// tokio::prelude::*`. +/// +/// [`timeout`]: #method.timeout +pub trait StreamExt: Stream { + + /// Creates a new stream which allows `self` until `timeout`. + /// + /// This combinator creates a new stream which wraps the receiving stream + /// with a timeout. For each item, the returned stream is allowed to execute + /// until it completes or `timeout` has elapsed, whichever happens first. + /// + /// If an item completes before `timeout` then the stream will yield + /// with that item. Otherwise the stream will yield to an error. + /// + /// # Examples + /// + /// ``` + /// # extern crate tokio; + /// # extern crate futures; + /// use tokio::prelude::*; + /// use std::time::Duration; + /// # use futures::future::{self, FutureResult}; + /// + /// # fn long_future() -> FutureResult<(), ()> { + /// # future::ok(()) + /// # } + /// # + /// # fn main() { + /// let stream = long_future() + /// .into_stream() + /// .timeout(Duration::from_secs(1)) + /// .for_each(|i| future::ok(println!("item = {:?}", i))) + /// .map_err(|e| println!("error = {:?}", e)); + /// + /// tokio::run(stream); + /// # } + /// ``` + fn timeout(self, timeout: Duration) -> Timeout + where Self: Sized, + { + Timeout::new(self, timeout) + } +} + +impl StreamExt for T where T: Stream {}