diff --git a/src/net/tcp.rs b/src/net/tcp.rs index 4ae431247..f5b3ccfb3 100644 --- a/src/net/tcp.rs +++ b/src/net/tcp.rs @@ -5,10 +5,10 @@ use std::net::{self, SocketAddr, Shutdown}; use futures::stream::Stream; use futures::sync::oneshot; -use futures::{self, Future, failed, Poll, Async}; +use futures::{Future, Poll, Async}; use mio; -use io::Io; +use io::{Io, IoFuture}; use reactor::{Handle, PollEvented}; /// An I/O object representing a TCP socket listening for incoming connections. @@ -224,42 +224,12 @@ pub struct TcpStream { /// Future returned by `TcpStream::connect` which will resolve to a `TcpStream` /// when the stream is connected. pub struct TcpStreamNew { - inner: TcpStreamNewFuture, + inner: TcpStreamNewState, } -type TcpStreamNewConnected = ::futures::AndThen< - ::futures::future::FutureResult< - ::reactor::PollEvented<::mio::tcp::TcpStream>, - ::std::io::Error, - >, - TcpStreamConnect, - fn (::reactor::PollEvented<::mio::tcp::TcpStream>) -> ::net::tcp::TcpStreamConnect ->; - -enum TcpStreamNewFuture { - Connected(TcpStreamNewConnected), - Error(::futures::future::Err), -} - -impl Future for TcpStreamNewFuture { - type Item = TcpStream; - type Error = io::Error; - fn poll(&mut self) -> Result, ::std::io::Error> { - match *self { - TcpStreamNewFuture::Connected(ref mut stream) => stream.poll(), - TcpStreamNewFuture::Error(ref mut error) => error.poll(), - } - } -} - -impl TcpStreamConnect { - fn from_stream(io: ::reactor::PollEvented<::mio::tcp::TcpStream>) -> Self { - TcpStreamConnect::Waiting(TcpStream { io: io }) - } -} - -enum TcpStreamConnect { +enum TcpStreamNewState { Waiting(TcpStream), + Error(io::Error), Empty, } @@ -272,20 +242,19 @@ impl TcpStream { /// connection or during the socket creation, that error will be returned to /// the future instead. pub fn connect(addr: &SocketAddr, handle: &Handle) -> TcpStreamNew { - let future = match mio::tcp::TcpStream::connect(addr) { + let inner = match mio::tcp::TcpStream::connect(addr) { Ok(tcp) => TcpStream::new(tcp, handle), - Err(e) => TcpStreamNewFuture::Error(failed(e)), + Err(e) => TcpStreamNewState::Error(e), }; - TcpStreamNew { inner: future } + TcpStreamNew { inner: inner } } fn new(connected_stream: mio::tcp::TcpStream, handle: &Handle) - -> TcpStreamNewFuture - { - let tcp = PollEvented::new(connected_stream, handle); - TcpStreamNewFuture::Connected( - futures::done(tcp).and_then(TcpStreamConnect::from_stream) - ) + -> TcpStreamNewState { + match PollEvented::new(connected_stream, handle) { + Ok(io) => TcpStreamNewState::Waiting(TcpStream { io: io }), + Err(e) => TcpStreamNewState::Error(e), + } } /// Creates a new `TcpStream` from the pending socket inside the given @@ -308,12 +277,12 @@ impl TcpStream { /// (perhaps to `INADDR_ANY`) before this method is called. pub fn connect_stream(stream: net::TcpStream, addr: &SocketAddr, - handle: &Handle) -> TcpStreamNew { - let future = match mio::tcp::TcpStream::connect_stream(stream, addr) { + handle: &Handle) -> IoFuture { + let state = match mio::tcp::TcpStream::connect_stream(stream, addr) { Ok(tcp) => TcpStream::new(tcp, handle), - Err(e) => TcpStreamNewFuture::Error(failed(e)), + Err(e) => TcpStreamNewState::Error(e), }; - TcpStreamNew { inner: future } + state.boxed() } /// Test whether this socket is ready to be read or not. @@ -485,15 +454,22 @@ impl Future for TcpStreamNew { } } -impl Future for TcpStreamConnect { +impl Future for TcpStreamNewState { type Item = TcpStream; type Error = io::Error; fn poll(&mut self) -> Poll { { let stream = match *self { - TcpStreamConnect::Waiting(ref s) => s, - TcpStreamConnect::Empty => panic!("can't poll TCP stream twice"), + TcpStreamNewState::Waiting(ref s) => s, + TcpStreamNewState::Error(_) => { + let e = match mem::replace(self, TcpStreamNewState::Empty) { + TcpStreamNewState::Error(e) => e, + _ => panic!(), + }; + return Err(e) + } + TcpStreamNewState::Empty => panic!("can't poll TCP stream twice"), }; // Once we've connected, wait for the stream to be writable as @@ -509,9 +485,9 @@ impl Future for TcpStreamConnect { return Err(e) } } - match mem::replace(self, TcpStreamConnect::Empty) { - TcpStreamConnect::Waiting(stream) => Ok(Async::Ready(stream)), - TcpStreamConnect::Empty => panic!(), + match mem::replace(self, TcpStreamNewState::Empty) { + TcpStreamNewState::Waiting(stream) => Ok(Async::Ready(stream)), + _ => panic!(), } } }