mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-22 00:00:11 +02:00
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.
This commit is contained in:
+23
-13
@@ -357,27 +357,37 @@ impl<T: Io, C: Codec> Sink for Framed<T, C> {
|
||||
type SinkError = io::Error;
|
||||
|
||||
fn start_send(&mut self, item: C::Out) -> StartSend<C::Out, io::Error> {
|
||||
// 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(()));
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user