Merge pull request #95 from aturon/sink

Refactor framing to use Streams and Sinks
This commit is contained in:
Alex Crichton
2016-11-08 16:53:21 -07:00
committed by GitHub
4 changed files with 223 additions and 210 deletions
+179 -141
View File
@@ -17,12 +17,14 @@
//! For more information see the `EasyFramed` and `EasyBuf` types. //! For more information see the `EasyFramed` and `EasyBuf` types.
use std::io; use std::io;
use std::marker::PhantomData;
use std::ops::{Deref, DerefMut}; use std::ops::{Deref, DerefMut};
use std::sync::Arc; use std::sync::Arc;
use futures::{Async, Poll}; use futures::{Async, Poll, Stream, Sink, StartSend, AsyncSink};
use futures::sync::BiLock;
use io::{Io, FramedIo}; use io::Io;
/// A reference counted buffer of bytes. /// A reference counted buffer of bytes.
/// ///
@@ -201,35 +203,9 @@ impl<'a> Drop for EasyBufMut<'a> {
} }
} }
/// An implementation of the `FramedIo` trait building on instances of the
/// `Decode` and `Encode` traits.
///
/// Many I/O streams are simply a framed protocol on both the inbound and
/// outbound halves. In essence the underlying stream of bytes can be converted
/// to a stream of *frames*. This way instead of reading or writing bytes a
/// stream deals with reading and writing frames.
///
/// This struct is essentially a convenience implementation of the `FramedIo`
/// which only requires knowledge of how to (de)encode types. It is
/// constructed with an arbitrary `Io` instance along with a encoder and
/// decoder for the frames that this `EasyFramed` will be yielding.
///
/// This implementation of `FramedIo` uses the `EasyBuf` type from the `bytes`
/// crate for the backing storage, which should allow for zero-copy
/// decoding where possible.
pub struct EasyFramed<T, D, S> {
upstream: T,
decode: D,
encode: S,
eof: bool,
is_readable: bool,
rd: EasyBuf,
wr: Vec<u8>,
}
/// Decoding of a frame from an internal buffer. /// Decoding of a frame from an internal buffer.
/// ///
/// This trait is used when constructing an instance of `EasyFramed`. It defines how /// This trait is used when constructing an instance of `Framed`. It defines how
/// to decode the incoming bytes on a stream to the specified type of frame for /// to decode the incoming bytes on a stream to the specified type of frame for
/// that framed I/O stream. /// that framed I/O stream.
/// ///
@@ -237,13 +213,7 @@ pub struct EasyFramed<T, D, S> {
/// frame from a buffer of bytes. It has the option of returning `NotReady`, /// frame from a buffer of bytes. It has the option of returning `NotReady`,
/// indicating that more bytes need to be read before decoding can /// indicating that more bytes need to be read before decoding can
/// continue. /// continue.
pub trait Decode { pub trait Decode: Sized {
/// The type of frame that this decoder produces.
///
/// This is typically a frame being parsed from an input stream, such as an
/// HTTP request, a Redis command, etc.
type Out;
/// Attempts to decode a frame from the provided buffer of bytes. /// Attempts to decode a frame from the provided buffer of bytes.
/// ///
/// This method is called by `EasyFramed` whenever bytes are ready to be parsed. /// This method is called by `EasyFramed` whenever bytes are ready to be parsed.
@@ -264,7 +234,7 @@ pub trait Decode {
/// Finally, if the bytes in the buffer are malformed then an error is /// Finally, if the bytes in the buffer are malformed then an error is
/// returned indicating why. This informs `EasyFramed` that the stream is now /// returned indicating why. This informs `EasyFramed` that the stream is now
/// corrupt and should be terminated. /// corrupt and should be terminated.
fn decode(&mut self, buf: &mut EasyBuf) -> Result<Option<Self::Out>, io::Error>; fn decode(buf: &mut EasyBuf) -> Result<Option<Self>, io::Error>;
/// A default method available to be called when there are no more bytes /// A default method available to be called when there are no more bytes
/// available to be read from the underlying I/O. /// available to be read from the underlying I/O.
@@ -272,8 +242,8 @@ pub trait Decode {
/// This method defaults to calling `decode` and returns an error if /// This method defaults to calling `decode` and returns an error if
/// `Ok(None)` is returned. Typically this doesn't need to be implemented /// `Ok(None)` is returned. Typically this doesn't need to be implemented
/// unless the framing protocol differs near the end of the stream. /// unless the framing protocol differs near the end of the stream.
fn done(&mut self, buf: &mut EasyBuf) -> io::Result<Self::Out> { fn done(buf: &mut EasyBuf) -> io::Result<Self> {
match try!(self.decode(buf)) { match try!(Self::decode(buf)) {
Some(frame) => Ok(frame), Some(frame) => Ok(frame),
None => Err(io::Error::new(io::ErrorKind::Other, None => Err(io::Error::new(io::ErrorKind::Other,
"bytes remaining on stream")), "bytes remaining on stream")),
@@ -289,97 +259,32 @@ pub trait Decode {
/// buffer. That buffer is then written out when possible to the underlying I/O /// buffer. That buffer is then written out when possible to the underlying I/O
/// stream. /// stream.
pub trait Encode { pub trait Encode {
/// The frame that's being encoded to a byte buffer.
///
/// This type is the type of frame that's also being written to a `EasyFramed`.
type In;
/// Encodes a frame into the buffer provided. /// Encodes a frame into the buffer provided.
/// ///
/// This method will encode `msg` into the byte buffer provided by `buf`. /// This method will encode `msg` into the byte buffer provided by `buf`.
/// The `buf` provided is an internal buffer of the `EasyFramed` instance and /// The `buf` provided is an internal buffer of the `EasyFramed` instance and
/// will be written out when possible. /// will be written out when possible.
fn encode(&mut self, msg: Self::In, buf: &mut Vec<u8>); fn encode(self, buf: &mut Vec<u8>);
} }
impl<T, D, S> EasyFramed<T, D, S> struct ReadState {
where T: Io, eof: bool,
D: Decode, is_readable: bool,
S: Encode, rd: EasyBuf,
{ }
/// Creates a new instance of `EasyFramed` from the given component pieces.
///
/// This method will create a new instance of `EasyFramed` which implements
/// `FramedIo` for reading and writing frames from an underlying I/O stream.
/// The `upstream` argument here is the byte-based I/O stream that it will
/// be operating on. Data will be read from this stream and decoded with
/// `decode` into frames. Frames written to this instance will be
/// encoded by `encode` and then written to `upstream`.
///
/// The `rd` and `wr` buffers provided are used for reading and writing
/// bytes and provide a small amount of control over how buffering happens.
pub fn new(upstream: T,
decode: D,
encode: S) -> EasyFramed<T, D, S> {
trace!("creating new framed transport"); impl ReadState {
EasyFramed { fn new() -> ReadState {
upstream: upstream, ReadState {
decode: decode,
encode: encode,
is_readable: false,
eof: false, eof: false,
is_readable: false,
rd: EasyBuf::new(), rd: EasyBuf::new(),
wr: Vec::with_capacity(8 * 1024),
} }
} }
/// Returns a reference to the underlying I/O stream wrapped by `EasyFramed`.
///
/// Note that care should be taken to not tamper with the underlying stream
/// of data coming in as it may corrupt the stream of frames otherwise being
/// worked with.
pub fn get_ref(&self) -> &T {
&self.upstream
}
/// Returns a mutable reference to the underlying I/O stream wrapped by
/// `EasyFramed`.
///
/// Note that care should be taken to not tamper with the underlying stream
/// of data coming in as it may corrupt the stream of frames otherwise being
/// worked with.
pub fn get_mut(&mut self) -> &mut T {
&mut self.upstream
}
/// Consumes the `EasyFramed`, returning its underlying I/O stream.
///
/// Note that care should be taken to not tamper with the underlying stream
/// of data coming in as it may corrupt the stream of frames otherwise being
/// worked with.
pub fn into_inner(self) -> T {
self.upstream
}
} }
impl<T, P, S> FramedIo for EasyFramed<T, P, S> impl ReadState {
where T: Io, fn poll<T: Io, D: Decode>(&mut self, upstream: &mut T) -> Poll<Option<D>, io::Error> {
P: Decode,
S: Encode,
{
type In = S::In;
type Out = Option<P::Out>;
fn poll_read(&mut self) -> Async<()> {
if self.is_readable || self.upstream.poll_read().is_ready() {
Async::Ready(())
} else {
Async::NotReady
}
}
fn read(&mut self) -> Poll<Self::Out, io::Error> {
loop { loop {
// If the read buffer has any pending data, then it could be // If the read buffer has any pending data, then it could be
// possible that `decode` will return a new frame. We leave it to // possible that `decode` will return a new frame. We leave it to
@@ -389,14 +294,14 @@ impl<T, P, S> FramedIo for EasyFramed<T, P, S>
if self.rd.len() == 0 { if self.rd.len() == 0 {
return Ok(None.into()) return Ok(None.into())
} else { } else {
let frame = try!(self.decode.done(&mut self.rd)); let frame = try!(Decode::done(&mut self.rd));
return Ok(Some(frame).into()) return Ok(Async::Ready(Some(frame)))
} }
} }
trace!("attempting to decode a frame"); trace!("attempting to decode a frame");
if let Some(frame) = try!(self.decode.decode(&mut self.rd)) { if let Some(frame) = try!(Decode::decode(&mut self.rd)) {
trace!("frame decoded from buffer"); trace!("frame decoded from buffer");
return Ok(Some(frame).into()); return Ok(Async::Ready(Some(frame)));
} }
self.is_readable = false; self.is_readable = false;
} }
@@ -407,7 +312,7 @@ impl<T, P, S> FramedIo for EasyFramed<T, P, S>
// //
// TODO: shouldn't read_to_end, that may read a lot // TODO: shouldn't read_to_end, that may read a lot
let before = self.rd.len(); let before = self.rd.len();
let ret = self.upstream.read_to_end(&mut self.rd.get_mut()); let ret = upstream.read_to_end(&mut self.rd.get_mut());
match ret { match ret {
Ok(_n) => self.eof = true, Ok(_n) => self.eof = true,
Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => { Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => {
@@ -420,32 +325,28 @@ impl<T, P, S> FramedIo for EasyFramed<T, P, S>
self.is_readable = true; self.is_readable = true;
} }
} }
}
fn poll_write(&mut self) -> Async<()> { struct WriteState {
// Always accept writes and let the write buffer grow wr: Vec<u8>,
// }
// TODO: This may not be the best option for robustness, but for now it
// makes the microbenchmarks happy.
Async::Ready(())
}
fn write(&mut self, msg: Self::In) -> Poll<(), io::Error> { impl WriteState {
if !self.poll_write().is_ready() { fn new() -> WriteState {
return Err(io::Error::new(io::ErrorKind::InvalidInput, WriteState {
"transport not currently writable")); wr: Vec::with_capacity(8 * 1024),
} }
}
}
// Encode the msg impl WriteState {
self.encode.encode(msg, &mut self.wr); fn write<E: Encode>(&mut self, data: E) {
data.encode(&mut self.wr)
// TODO: should provide some backpressure, such as when the buffer is
// too full this returns `NotReady` or something like that.
Ok(Async::Ready(()))
} }
fn flush(&mut self) -> Poll<(), io::Error> { fn poll_complete<T: Io>(&mut self, upstream: &mut T) -> Poll<(), io::Error> {
// Try flushing the underlying IO // Try flushing the underlying IO
try_nb!(self.upstream.flush()); try_nb!(upstream.flush());
trace!("flushing framed transport"); trace!("flushing framed transport");
@@ -457,8 +358,145 @@ impl<T, P, S> FramedIo for EasyFramed<T, P, S>
trace!("writing; remaining={:?}", self.wr.len()); trace!("writing; remaining={:?}", self.wr.len());
let n = try_nb!(self.upstream.write(&self.wr)); let n = try_nb!(upstream.write(&self.wr));
self.wr.drain(..n); self.wr.drain(..n);
} }
} }
} }
/// A `Stream` interface to an underlying `Io` object, using the `Decode` trait
/// to decode frames.
pub struct FramedRead<T, D> {
upstream: BiLock<T>,
read_state: ReadState,
_phantom: PhantomData<D>,
}
impl<T: Io, D: Decode> Stream for FramedRead<T, D> {
type Item = D;
type Error = io::Error;
fn poll(&mut self) -> Poll<Option<D>, io::Error> {
if let Async::Ready(mut guard) = self.upstream.poll_lock() {
self.read_state.poll(&mut *guard)
} else {
Ok(Async::NotReady)
}
}
}
/// A `Sink` interface to an underlying `Io` object, using the `Encode` trait
/// to encode frames.
pub struct FramedWrite<T, E> {
upstream: BiLock<T>,
write_state: WriteState,
_phantom: PhantomData<E>,
}
impl<T: Io, E: Encode> Sink for FramedWrite<T, E> {
type SinkItem = E;
type SinkError = io::Error;
fn start_send(&mut self, item: E) -> StartSend<E, io::Error> {
self.write_state.write(item);
Ok(AsyncSink::Ready)
}
fn poll_complete(&mut self) -> Poll<(), io::Error> {
if let Async::Ready(mut guard) = self.upstream.poll_lock() {
self.write_state.poll_complete(&mut *guard)
} else {
Ok(Async::NotReady)
}
}
}
/// A unified `Stream` and `Sink` interface to an underlying `Io` object, using
/// the `Encode` and `Decode` traits to encode and decode frames.
pub struct Framed<T, D, E> {
upstream: T,
read_state: ReadState,
write_state: WriteState,
_phantom: PhantomData<(D, E)>,
}
impl<T: Io, D: Decode, E: Encode> Stream for Framed<T, D, E> {
type Item = D;
type Error = io::Error;
fn poll(&mut self) -> Poll<Option<D>, io::Error> {
self.read_state.poll(&mut self.upstream)
}
}
impl<T: Io, D: Decode, E: Encode> Sink for Framed<T, D, E> {
type SinkItem = E;
type SinkError = io::Error;
fn start_send(&mut self, item: E) -> StartSend<E, io::Error> {
self.write_state.write(item);
Ok(AsyncSink::Ready)
}
fn poll_complete(&mut self) -> Poll<(), io::Error> {
self.write_state.poll_complete(&mut self.upstream)
}
}
pub fn framed<T, D, E>(io: T) -> Framed<T, D, E> {
Framed {
upstream: io,
read_state: ReadState::new(),
write_state: WriteState::new(),
_phantom: PhantomData,
}
}
impl<T, D, E> Framed<T, D, E> {
/// Splits this `Stream + Sink` object into separate `Stream` and `Sink`
/// objects, which can be useful when you want to split ownership between
/// tasks, or allow direct interaction between the two objects (e.g. via
/// `Sink::send_all`).
pub fn split(self) -> (FramedRead<T, D>, FramedWrite<T, E>) {
let (a, b) = BiLock::new(self.upstream);
let read = FramedRead {
upstream: a,
read_state: ReadState::new(),
_phantom: PhantomData,
};
let write = FramedWrite {
upstream: b,
write_state: WriteState::new(),
_phantom: PhantomData,
};
(read, write)
}
/// Returns a reference to the underlying I/O stream wrapped by `Framed`.
///
/// Note that care should be taken to not tamper with the underlying stream
/// of data coming in as it may corrupt the stream of frames otherwise being
/// worked with.
pub fn get_ref(&self) -> &T {
&self.upstream
}
/// Returns a mutable reference to the underlying I/O stream wrapped by
/// `Framed`.
///
/// Note that care should be taken to not tamper with the underlying stream
/// of data coming in as it may corrupt the stream of frames otherwise being
/// worked with.
pub fn get_mut(&mut self) -> &mut T {
&mut self.upstream
}
/// Consumes the `Framed`, returning its underlying I/O stream.
///
/// Note that care should be taken to not tamper with the underlying stream
/// of data coming in as it may corrupt the stream of frames otherwise being
/// worked with.
pub fn into_inner(self) -> T {
self.upstream
}
}
+31
View File
@@ -32,6 +32,7 @@ macro_rules! try_nb {
} }
mod copy; mod copy;
mod frame;
mod flush; mod flush;
mod read_exact; mod read_exact;
mod read_to_end; mod read_to_end;
@@ -41,6 +42,7 @@ mod split;
mod window; mod window;
mod write_all; mod write_all;
pub use self::copy::{copy, Copy}; pub use self::copy::{copy, Copy};
pub use self::frame::{EasyBuf, EasyBufMut, FramedRead, FramedWrite, Framed, Decode, Encode};
pub use self::flush::{flush, Flush}; pub use self::flush::{flush, Flush};
pub use self::read_exact::{read_exact, ReadExact}; pub use self::read_exact::{read_exact, ReadExact};
pub use self::read_to_end::{read_to_end, ReadToEnd}; pub use self::read_to_end::{read_to_end, ReadToEnd};
@@ -106,6 +108,33 @@ pub trait Io: io::Read + io::Write {
Async::Ready(()) Async::Ready(())
} }
/// Provides a `Stream` and `Sink` interface for reading and writing to this
/// `Io` object, using `Decode` and `Encode` to read and write the raw data.
///
/// Raw I/O objects work with byte sequences, but higher-level code usually
/// wants to batch these into meaningful chunks, called "frames". This
/// method layers framing on top of an I/O object, by using the `Encode` and
/// `Decode` traits:
///
/// - `Encode` interprets frames we want to send into bytes;
/// - `Decode` interprets incoming bytes into a stream of frames.
///
/// Note that the incoming and outgoing frame types may be distinct.
///
/// This function returns a *single* object that is both `Stream` and
/// `Sink`; grouping this into a single object is often useful for layering
/// things like gzip or TLS, which require both read and write access to the
/// underlying object.
///
/// If you want to work more directly with the streams and sink, consider
/// calling `split` on the `Framed` returned by this method, which will
/// break them into separate objects, allowing them to interact more easily.
fn framed<D: Decode, E: Encode>(self) -> Framed<Self, D, E>
where Self: Sized,
{
frame::framed(self)
}
/// Helper method for splitting this read/write object into two halves. /// Helper method for splitting this read/write object into two halves.
/// ///
/// The two halves returned implement the `Read` and `Write` traits, /// The two halves returned implement the `Read` and `Write` traits,
@@ -130,6 +159,8 @@ pub trait Io: io::Read + io::Write {
/// ///
/// For a sample implementation of `FramedIo` you can take a look at the /// For a sample implementation of `FramedIo` you can take a look at the
/// `EasyFramed` type in the `easy` module of htis crate. /// `EasyFramed` type in the `easy` module of htis crate.
#[doc(hidden)]
#[deprecated(since = "0.1.1", note = "replaced by Sink + Stream")]
pub trait FramedIo { pub trait FramedIo {
/// Messages written /// Messages written
type In; type In;
-1
View File
@@ -112,4 +112,3 @@ mod heap;
pub mod channel; pub mod channel;
pub mod net; pub mod net;
pub mod reactor; pub mod reactor;
pub mod easy;
+13 -68
View File
@@ -5,85 +5,30 @@ extern crate futures;
use std::io; use std::io;
use std::net::Shutdown; use std::net::Shutdown;
use futures::{Future, Async, Poll}; use futures::{Future, Stream, Sink};
use futures::stream::Stream; use tokio_core::io::{write_all, read, Encode, Decode, EasyBuf, Io};
use tokio_core::io::{FramedIo, write_all, read};
use tokio_core::easy::{Encode, Decode, EasyFramed, EasyBuf};
use tokio_core::net::{TcpListener, TcpStream}; use tokio_core::net::{TcpListener, TcpStream};
use tokio_core::reactor::Core; use tokio_core::reactor::Core;
pub struct LineDecoder; pub struct Line(EasyBuf);
pub struct LineEncoder;
impl Decode for LineDecoder { impl Decode for Line {
type Out = EasyBuf; fn decode(buf: &mut EasyBuf) -> Result<Option<Line>, io::Error> {
fn decode(&mut self, buf: &mut EasyBuf) -> Result<Option<EasyBuf>, io::Error> {
match buf.as_slice().iter().position(|&b| b == b'\n') { match buf.as_slice().iter().position(|&b| b == b'\n') {
Some(i) => Ok(Some(buf.drain_to(i + 1).into())), Some(i) => Ok(Some(Line(buf.drain_to(i + 1).into()))),
None => Ok(None), None => Ok(None),
} }
} }
fn done(&mut self, buf: &mut EasyBuf) -> io::Result<EasyBuf> { fn done(buf: &mut EasyBuf) -> io::Result<Line> {
let amt = buf.len(); let amt = buf.len();
Ok(buf.drain_to(amt)) Ok(Line(buf.drain_to(amt)))
} }
} }
impl Encode for LineEncoder { impl Encode for Line {
type In = EasyBuf; fn encode(self, into: &mut Vec<u8>) {
into.extend_from_slice(self.0.as_slice());
fn encode(&mut self, msg: EasyBuf, into: &mut Vec<u8>) {
into.extend_from_slice(msg.as_slice());
}
}
pub struct EchoFramed<T> {
inner: T,
eof: bool,
}
impl<T, U> Future for EchoFramed<T>
where T: FramedIo<In = U, Out = Option<U>>,
{
type Item = ();
type Error = ();
fn poll(&mut self) -> Poll<(), ()> {
// Try to write out any buffered messages if we have them.
if self.inner.flush().expect("flush error").is_not_ready() {
return Ok(Async::NotReady)
}
// Wait until we can simultaneously read and write a message
while !self.eof &&
self.inner.poll_read().is_ready() &&
self.inner.poll_write().is_ready() {
let frame = match self.inner.read() {
Ok(Async::Ready(Some(frame))) => frame,
Ok(Async::Ready(None)) => {
self.eof = true;
break
}
Ok(Async::NotReady) => break,
Err(e) => panic!("error in read: {}", e),
};
match self.inner.write(frame) {
Ok(Async::Ready(())) => {}
Ok(Async::NotReady) => break,
Err(e) => panic!("error in write: {}", e),
}
}
// If we wrote some frames try to flush again. Ignore whether this is
// ready to finish or not as we're going to continue to return NotReady
if self.inner.flush().expect("flush error").is_ready() && self.eof {
Ok(().into())
} else {
Ok(Async::NotReady)
}
} }
} }
@@ -97,8 +42,8 @@ fn echo() {
let listener = TcpListener::bind(&"127.0.0.1:0".parse().unwrap(), &handle).unwrap(); let listener = TcpListener::bind(&"127.0.0.1:0".parse().unwrap(), &handle).unwrap();
let addr = listener.local_addr().unwrap(); let addr = listener.local_addr().unwrap();
let srv = listener.incoming().for_each(move |(socket, _)| { let srv = listener.incoming().for_each(move |(socket, _)| {
let framed = EasyFramed::new(socket, LineDecoder, LineEncoder); let (stream, sink) = socket.framed::<Line, Line>().split();
handle.spawn(EchoFramed { inner: framed, eof: false }); handle.spawn(sink.send_all(stream).map(|_| ()).map_err(|_| ()));
Ok(()) Ok(())
}); });