StreamExt: add a trait for additional Stream methods (#573)

Primarily, it offers a `timeout` method for streams.
This commit is contained in:
Ben Boeckel
2018-09-07 15:43:03 -07:00
committed by Carl Lerche
parent 16664189c1
commit 89d969d518
3 changed files with 68 additions and 1 deletions
+1
View File
@@ -599,6 +599,7 @@ pub mod prelude {
pub use util::{
FutureExt,
StreamExt,
};
pub use ::std::io::{
+5 -1
View File
@@ -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;
+62
View File
@@ -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<Self>
where Self: Sized,
{
Timeout::new(self, timeout)
}
}
impl<T: ?Sized> StreamExt for T where T: Stream {}