From 46dd38b7d766995e983b395dceca74a7f746ce6a Mon Sep 17 00:00:00 2001 From: Dan Burkert Date: Tue, 15 Nov 2016 23:31:34 -0800 Subject: [PATCH] Clean up Sink implementation on Framed This commit makes a few changes to the Sink implementation on Framed: * Backpressure is implemented for `start_send`. If the write buffer is over 8KiB and can't be flushed, no new items are accepted. * 0 length writes to the upstream transport are translated into a `WriteZero` error, as with `io::Write::write_all`. `write_all` checks for and ignores `Interrupted` errors, but I do not think this is necessary for non-blocking writes. * In `poll_complete`, the upstream transport is not flushed until *after* writing the entire write buffer. --- src/io/frame.rs | 36 +++++++++++++++++++++++------------- 1 file changed, 23 insertions(+), 13 deletions(-) diff --git a/src/io/frame.rs b/src/io/frame.rs index 20fc1682d..0686614d4 100644 --- a/src/io/frame.rs +++ b/src/io/frame.rs @@ -357,27 +357,37 @@ impl Sink for Framed { type SinkError = io::Error; fn start_send(&mut self, item: C::Out) -> StartSend { + // If the buffer is already over 8KiB, then attempt to flush it. If after flushing it's + // *still* over 8KiB, then apply backpressure (reject the send). + if self.wr.len() > 8 * 1024 { + try!(self.poll_complete()); + if self.wr.len() > 8 * 1024 { + return Ok(AsyncSink::NotReady(item)); + } + } + self.codec.encode(item, &mut self.wr); Ok(AsyncSink::Ready) } fn poll_complete(&mut self) -> Poll<(), io::Error> { + trace!("flushing framed transport"); + + while !self.wr.is_empty() { + trace!("writing; remaining={}", self.wr.len()); + let n = try_nb!(self.upstream.write(&self.wr)); + if n == 0 { + return Err(io::Error::new(io::ErrorKind::WriteZero, + "failed to write frame to transport")); + } + self.wr.drain(..n); + } + // Try flushing the underlying IO try_nb!(self.upstream.flush()); - trace!("flushing framed transport"); - - loop { - if self.wr.len() == 0 { - trace!("framed transport flushed"); - return Ok(Async::Ready(())); - } - - trace!("writing; remaining={:?}", self.wr.len()); - - let n = try_nb!(self.upstream.write(&self.wr)); - self.wr.drain(..n); - } + trace!("framed transport flushed"); + return Ok(Async::Ready(())); } }