From 38df0d7f0ff072ba5ed676a34eaf8f6546f3e236 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Dawid=20Ci=C4=99=C5=BCarkiewicz?= Date: Wed, 16 Nov 2016 15:14:13 -0800 Subject: [PATCH 1/3] Implement `TcpListener::accept()` --- src/net/tcp.rs | 33 +++++++++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/src/net/tcp.rs b/src/net/tcp.rs index 8f7f0d247..f99833180 100644 --- a/src/net/tcp.rs +++ b/src/net/tcp.rs @@ -34,6 +34,39 @@ impl TcpListener { TcpListener::new(l, handle) } + /// Attempt to accept a connection and create a new connected `TcpStream` if successful. + /// + /// It is more idiomatic to treat incoming connection as a `Stream` of `TcpStream`s. + /// See `incoming()` for details. + pub fn accept(&self) -> io::Result<(TcpStream, SocketAddr)> { + if let Async::NotReady = self.io.poll_read() { + return Err(io::Error::new(io::ErrorKind::WouldBlock, "not ready")) + } + + let res = self.io.get_ref().accept(); + match res { + Err(e) => { + if e.kind() == io::ErrorKind::WouldBlock { + self.io.need_read(); + } + Err(e) + }, + Ok((sock, addr)) => { + let (tx, rx) = futures::oneshot(); + let remote = self.io.remote().clone(); + remote.spawn(move |handle| { + let res = PollEvented::new(sock, handle) + .map(move |io| { + (TcpStream { io: io }, addr) + }); + tx.complete(res); + Ok(()) + }); + rx.then(|r| r.expect("shouldn't be canceled")).wait() + } + } + } + /// Create a new TCP listener from the standard library's TCP listener. /// /// This method can be used when the `Handle::tcp_listen` method isn't From 936727fa528c527f48baeabb4d6b2d0195d274b3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Dawid=20Ci=C4=99=C5=BCarkiewicz?= Date: Fri, 2 Dec 2016 15:18:39 -0800 Subject: [PATCH 2/3] Don't block in `TcpListener::accept` Instead hold a pending Future of the next `TcpStream`. --- src/net/tcp.rs | 20 +++++++++++++++++--- 1 file changed, 17 insertions(+), 3 deletions(-) diff --git a/src/net/tcp.rs b/src/net/tcp.rs index f99833180..d9ec2748c 100644 --- a/src/net/tcp.rs +++ b/src/net/tcp.rs @@ -16,6 +16,7 @@ use reactor::{Handle, PollEvented}; /// various forms of processing. pub struct TcpListener { io: PollEvented, + pending_accept: Option>>, } /// Stream returned by the `TcpListener::incoming` function representing the @@ -38,7 +39,19 @@ impl TcpListener { /// /// It is more idiomatic to treat incoming connection as a `Stream` of `TcpStream`s. /// See `incoming()` for details. - pub fn accept(&self) -> io::Result<(TcpStream, SocketAddr)> { + pub fn accept(&mut self) -> io::Result<(TcpStream, SocketAddr)> { + if let Some(mut pending) = self.pending_accept.take() { + match pending.poll().expect("shouldn't be canceled") { + Async::NotReady => { + self.pending_accept = Some(pending); + return Err(io::Error::new(io::ErrorKind::WouldBlock, "not ready")) + }, + Async::Ready(r) => { + return r + } + } + } + if let Async::NotReady = self.io.poll_read() { return Err(io::Error::new(io::ErrorKind::WouldBlock, "not ready")) } @@ -62,7 +75,8 @@ impl TcpListener { tx.complete(res); Ok(()) }); - rx.then(|r| r.expect("shouldn't be canceled")).wait() + self.pending_accept = Some(rx); + return self.accept() } } } @@ -104,7 +118,7 @@ impl TcpListener { fn new(listener: mio::tcp::TcpListener, handle: &Handle) -> io::Result { let io = try!(PollEvented::new(listener, handle)); - Ok(TcpListener { io: io }) + Ok(TcpListener { io: io, pending_accept: None }) } /// Test whether this socket is ready to be read or not. From 1864ecc46f87b2daef023443a6987b1b9553a62e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Dawid=20Ci=C4=99=C5=BCarkiewicz?= Date: Fri, 16 Dec 2016 11:07:08 -0800 Subject: [PATCH 3/3] tcp: Express `incoming` in terms of `accept` --- src/net/tcp.rs | 27 +++------------------------ 1 file changed, 3 insertions(+), 24 deletions(-) diff --git a/src/net/tcp.rs b/src/net/tcp.rs index d9ec2748c..731f12ab5 100644 --- a/src/net/tcp.rs +++ b/src/net/tcp.rs @@ -145,38 +145,17 @@ impl TcpListener { } impl Stream for MyIncoming { - type Item = (mio::tcp::TcpStream, SocketAddr); + type Item = (TcpStream, SocketAddr); type Error = io::Error; fn poll(&mut self) -> Poll, io::Error> { - if let Async::NotReady = self.inner.io.poll_read() { - return Ok(Async::NotReady) - } - match self.inner.io.get_ref().accept() { - Ok(pair) => Ok(Async::Ready(Some(pair))), - Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => { - self.inner.io.need_read(); - Ok(Async::NotReady) - } - Err(e) => Err(e) - } + Ok(Async::Ready(Some(try_nb!(self.inner.accept())))) } } - let remote = self.io.remote().clone(); let stream = MyIncoming { inner: self }; Incoming { - inner: stream.and_then(move |(tcp, addr)| { - let (tx, rx) = futures::oneshot(); - remote.spawn(move |handle| { - let res = PollEvented::new(tcp, handle).map(move |io| { - (TcpStream { io: io }, addr) - }); - tx.complete(res); - Ok(()) - }); - rx.then(|r| r.expect("shouldn't be canceled")) - }).boxed(), + inner: stream.boxed(), } }