diff --git a/tokio/Cargo.toml b/tokio/Cargo.toml index c25475f2f..f340326e5 100644 --- a/tokio/Cargo.toml +++ b/tokio/Cargo.toml @@ -78,7 +78,7 @@ tokio-sync = { version = "0.2.0", optional = true, path = "../tokio-sync", featu tokio-threadpool = { version = "0.2.0", optional = true, path = "../tokio-threadpool" } tokio-tcp = { version = "0.2.0", optional = true, path = "../tokio-tcp", features = ["async-traits"] } tokio-udp = { version = "0.2.0", optional = true, path = "../tokio-udp" } -tokio-timer = { version = "0.3.0", optional = true, path = "../tokio-timer" } +tokio-timer = { version = "0.3.0", optional = true, path = "../tokio-timer", features = ["async-traits"] } tracing-core = { version = "0.1", optional = true } memchr = { version = "2.2", optional = true } diff --git a/tokio/src/util/enumerate.rs b/tokio/src/util/enumerate.rs deleted file mode 100644 index a3b916288..000000000 --- a/tokio/src/util/enumerate.rs +++ /dev/null @@ -1,84 +0,0 @@ -use futures::{try_ready, Async, Poll, Sink, StartSend, Stream}; - -/// A stream combinator which combines the yields the current item -/// plus its count starting from 0. -/// -/// This structure is produced by the `Stream::enumerate` method. -#[derive(Debug)] -#[must_use = "Does nothing unless polled"] -pub struct Enumerate { - inner: T, - count: usize, -} - -impl Enumerate { - pub(crate) fn new(stream: T) -> Self { - Self { - inner: stream, - count: 0, - } - } - - /// Acquires a reference to the underlying stream that this combinator is - /// pulling from. - pub fn get_ref(&self) -> &T { - &self.inner - } - - /// Acquires a mutable reference to the underlying stream that this - /// combinator is pulling from. - /// - /// Note that care must be taken to avoid tampering with the state of the - /// stream which may otherwise confuse this combinator. - pub fn get_mut(&mut self) -> &mut T { - &mut self.inner - } - - /// Consumes this combinator, returning the underlying stream. - /// - /// Note that this may discard intermediate state of this combinator, so - /// care should be taken to avoid losing resources when this is called. - pub fn into_inner(self) -> T { - self.inner - } -} - -impl Stream for Enumerate -where - T: Stream, -{ - type Item = (usize, T::Item); - type Error = T::Error; - - fn poll(&mut self) -> Poll, T::Error> { - match try_ready!(self.inner.poll()) { - Some(item) => { - let ret = Some((self.count, item)); - self.count += 1; - Ok(Async::Ready(ret)) - } - None => return Ok(Async::Ready(None)), - } - } -} - -// Forwarding impl of Sink from the underlying stream -impl Sink for Enumerate -where - T: Sink, -{ - type SinkItem = T::SinkItem; - type SinkError = T::SinkError; - - fn start_send(&mut self, item: T::SinkItem) -> StartSend { - self.inner.start_send(item) - } - - fn poll_complete(&mut self) -> Poll<(), T::SinkError> { - self.inner.poll_complete() - } - - fn close(&mut self) -> Poll<(), T::SinkError> { - self.inner.close() - } -} diff --git a/tokio/src/util/mod.rs b/tokio/src/util/mod.rs index 3a9ec31f9..3ebd1fc70 100644 --- a/tokio/src/util/mod.rs +++ b/tokio/src/util/mod.rs @@ -7,9 +7,8 @@ //! [`FutureExt`]: trait.FutureExt.html //! [`StreamExt`]: trait.StreamExt.html -// mod enumerate; mod future; -// mod stream; +mod stream; pub use self::future::FutureExt; -// pub use self::stream::StreamExt; +pub use self::stream::StreamExt; diff --git a/tokio/src/util/stream.rs b/tokio/src/util/stream.rs index 001a4496b..679a2fc81 100644 --- a/tokio/src/util/stream.rs +++ b/tokio/src/util/stream.rs @@ -1,11 +1,11 @@ -pub use crate::util::enumerate::Enumerate; - #[cfg(feature = "timer")] use std::time::Duration; #[cfg(feature = "timer")] use tokio_timer::{throttle::Throttle, Timeout}; +use futures_core::Stream; + /// An extension trait for `Stream` that provides a variety of convenient /// combinator functions. /// @@ -31,25 +31,6 @@ pub trait StreamExt: Stream { Throttle::new(self, duration) } - /// Creates a new stream which gives the current iteration count as well - /// as the next value. - /// - /// The stream returned yields pairs `(i, val)`, where `i` is the - /// current index of iteration and `val` is the value returned by the - /// iterator. - /// - /// # Overflow Behavior - /// - /// The method does no guarding against overflows, so counting elements of - /// an iterator with more than [`std::usize::MAX`] elements either produces the - /// wrong result or panics. - fn enumerate(self) -> Enumerate - where - Self: Sized, - { - Enumerate::new(self) - } - /// Creates a new stream which allows `self` until `timeout`. /// /// This combinator creates a new stream which wraps the receiving stream diff --git a/tokio/tests/enumerate.rs b/tokio/tests/enumerate.rs deleted file mode 100644 index eaf2766b8..000000000 --- a/tokio/tests/enumerate.rs +++ /dev/null @@ -1,24 +0,0 @@ -#![cfg(feature = "broken")] -#![deny(warnings, rust_2018_idioms)] - -use futures::sync::mpsc; -use tokio::util::StreamExt; - -#[test] -fn enumerate() { - use futures::*; - - let (mut tx, rx) = mpsc::channel(1); - - std::thread::spawn(|| { - for i in 0..5 { - tx = tx.send(i * 2).wait().unwrap(); - } - }); - - let result = rx.enumerate().collect(); - assert_eq!( - result.wait(), - Ok(vec![(0, 0), (1, 2), (2, 4), (3, 6), (4, 8)]) - ); -}