diff --git a/src/easy.rs b/src/io/frame.rs similarity index 72% rename from src/easy.rs rename to src/io/frame.rs index 4444d5ccf..32d9d838d 100644 --- a/src/easy.rs +++ b/src/io/frame.rs @@ -17,12 +17,14 @@ //! For more information see the `EasyFramed` and `EasyBuf` types. use std::io; +use std::marker::PhantomData; use std::ops::{Deref, DerefMut}; 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. /// @@ -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 { - upstream: T, - decode: D, - encode: S, - eof: bool, - is_readable: bool, - rd: EasyBuf, - wr: Vec, -} - /// 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 /// that framed I/O stream. /// @@ -237,13 +213,7 @@ pub struct EasyFramed { /// frame from a buffer of bytes. It has the option of returning `NotReady`, /// indicating that more bytes need to be read before decoding can /// continue. -pub trait Decode { - /// 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; - +pub trait Decode: Sized { /// Attempts to decode a frame from the provided buffer of bytes. /// /// 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 /// returned indicating why. This informs `EasyFramed` that the stream is now /// corrupt and should be terminated. - fn decode(&mut self, buf: &mut EasyBuf) -> Result, io::Error>; + fn decode(buf: &mut EasyBuf) -> Result, io::Error>; /// A default method available to be called when there are no more bytes /// 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 /// `Ok(None)` is returned. Typically this doesn't need to be implemented /// unless the framing protocol differs near the end of the stream. - fn done(&mut self, buf: &mut EasyBuf) -> io::Result { - match try!(self.decode(buf)) { + fn done(buf: &mut EasyBuf) -> io::Result { + match try!(Self::decode(buf)) { Some(frame) => Ok(frame), None => Err(io::Error::new(io::ErrorKind::Other, "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 /// stream. 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. /// /// This method will encode `msg` into the byte buffer provided by `buf`. /// The `buf` provided is an internal buffer of the `EasyFramed` instance and /// will be written out when possible. - fn encode(&mut self, msg: Self::In, buf: &mut Vec); + fn encode(self, buf: &mut Vec); } -impl EasyFramed - where T: Io, - D: Decode, - S: Encode, -{ - /// 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 { +struct ReadState { + eof: bool, + is_readable: bool, + rd: EasyBuf, +} - trace!("creating new framed transport"); - EasyFramed { - upstream: upstream, - decode: decode, - encode: encode, - is_readable: false, +impl ReadState { + fn new() -> ReadState { + ReadState { eof: false, + is_readable: false, 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 FramedIo for EasyFramed - where T: Io, - P: Decode, - S: Encode, -{ - type In = S::In; - type Out = Option; - - 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 { +impl ReadState { + fn poll(&mut self, upstream: &mut T) -> Poll, io::Error> { loop { // If the read buffer has any pending data, then it could be // possible that `decode` will return a new frame. We leave it to @@ -389,14 +294,14 @@ impl FramedIo for EasyFramed if self.rd.len() == 0 { return Ok(None.into()) } else { - let frame = try!(self.decode.done(&mut self.rd)); - return Ok(Some(frame).into()) + let frame = try!(Decode::done(&mut self.rd)); + return Ok(Async::Ready(Some(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"); - return Ok(Some(frame).into()); + return Ok(Async::Ready(Some(frame))); } self.is_readable = false; } @@ -407,7 +312,7 @@ impl FramedIo for EasyFramed // // TODO: shouldn't read_to_end, that may read a lot 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 { Ok(_n) => self.eof = true, Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => { @@ -420,32 +325,28 @@ impl FramedIo for EasyFramed self.is_readable = true; } } +} - fn poll_write(&mut self) -> Async<()> { - // Always accept writes and let the write buffer grow - // - // TODO: This may not be the best option for robustness, but for now it - // makes the microbenchmarks happy. - Async::Ready(()) - } +struct WriteState { + wr: Vec, +} - fn write(&mut self, msg: Self::In) -> Poll<(), io::Error> { - if !self.poll_write().is_ready() { - return Err(io::Error::new(io::ErrorKind::InvalidInput, - "transport not currently writable")); +impl WriteState { + fn new() -> WriteState { + WriteState { + wr: Vec::with_capacity(8 * 1024), } + } +} - // Encode the msg - self.encode.encode(msg, &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(())) +impl WriteState { + fn write(&mut self, data: E) { + data.encode(&mut self.wr) } - fn flush(&mut self) -> Poll<(), io::Error> { + fn poll_complete(&mut self, upstream: &mut T) -> Poll<(), io::Error> { // Try flushing the underlying IO - try_nb!(self.upstream.flush()); + try_nb!(upstream.flush()); trace!("flushing framed transport"); @@ -457,8 +358,145 @@ impl FramedIo for EasyFramed 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); } } } + +/// A `Stream` interface to an underlying `Io` object, using the `Decode` trait +/// to decode frames. +pub struct FramedRead { + upstream: BiLock, + read_state: ReadState, + _phantom: PhantomData, +} + +impl Stream for FramedRead { + type Item = D; + type Error = io::Error; + + fn poll(&mut self) -> Poll, 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 { + upstream: BiLock, + write_state: WriteState, + _phantom: PhantomData, +} + +impl Sink for FramedWrite { + type SinkItem = E; + type SinkError = io::Error; + + fn start_send(&mut self, item: E) -> StartSend { + 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 { + upstream: T, + read_state: ReadState, + write_state: WriteState, + _phantom: PhantomData<(D, E)>, +} + +impl Stream for Framed { + type Item = D; + type Error = io::Error; + + fn poll(&mut self) -> Poll, io::Error> { + self.read_state.poll(&mut self.upstream) + } +} + +impl Sink for Framed { + type SinkItem = E; + type SinkError = io::Error; + + fn start_send(&mut self, item: E) -> StartSend { + 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(io: T) -> Framed { + Framed { + upstream: io, + read_state: ReadState::new(), + write_state: WriteState::new(), + _phantom: PhantomData, + } +} + +impl Framed { + /// 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, FramedWrite) { + 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 + } +} diff --git a/src/io/mod.rs b/src/io/mod.rs index d013727d3..1b7f676e1 100644 --- a/src/io/mod.rs +++ b/src/io/mod.rs @@ -32,6 +32,7 @@ macro_rules! try_nb { } mod copy; +mod frame; mod flush; mod read_exact; mod read_to_end; @@ -41,6 +42,7 @@ mod split; mod window; mod write_all; 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::read_exact::{read_exact, ReadExact}; pub use self::read_to_end::{read_to_end, ReadToEnd}; @@ -106,6 +108,33 @@ pub trait Io: io::Read + io::Write { 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(self) -> Framed + where Self: Sized, + { + frame::framed(self) + } + /// Helper method for splitting this read/write object into two halves. /// /// 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 /// `EasyFramed` type in the `easy` module of htis crate. +#[doc(hidden)] +#[deprecated(since = "0.1.1", note = "replaced by Sink + Stream")] pub trait FramedIo { /// Messages written type In; diff --git a/src/lib.rs b/src/lib.rs index d6080164d..a2a5f0dc9 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -112,4 +112,3 @@ mod heap; pub mod channel; pub mod net; pub mod reactor; -pub mod easy; diff --git a/tests/line-frames.rs b/tests/line-frames.rs index 39301dc78..ee9e55f6c 100644 --- a/tests/line-frames.rs +++ b/tests/line-frames.rs @@ -5,85 +5,30 @@ extern crate futures; use std::io; use std::net::Shutdown; -use futures::{Future, Async, Poll}; -use futures::stream::Stream; -use tokio_core::io::{FramedIo, write_all, read}; -use tokio_core::easy::{Encode, Decode, EasyFramed, EasyBuf}; +use futures::{Future, Stream, Sink}; +use tokio_core::io::{write_all, read, Encode, Decode, EasyBuf, Io}; use tokio_core::net::{TcpListener, TcpStream}; use tokio_core::reactor::Core; -pub struct LineDecoder; -pub struct LineEncoder; +pub struct Line(EasyBuf); -impl Decode for LineDecoder { - type Out = EasyBuf; - - fn decode(&mut self, buf: &mut EasyBuf) -> Result, io::Error> { +impl Decode for Line { + fn decode(buf: &mut EasyBuf) -> Result, io::Error> { 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), } } - fn done(&mut self, buf: &mut EasyBuf) -> io::Result { + fn done(buf: &mut EasyBuf) -> io::Result { let amt = buf.len(); - Ok(buf.drain_to(amt)) + Ok(Line(buf.drain_to(amt))) } } -impl Encode for LineEncoder { - type In = EasyBuf; - - fn encode(&mut self, msg: EasyBuf, into: &mut Vec) { - into.extend_from_slice(msg.as_slice()); - } -} - -pub struct EchoFramed { - inner: T, - eof: bool, -} - -impl Future for EchoFramed - where T: FramedIo>, -{ - 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) - } +impl Encode for Line { + fn encode(self, into: &mut Vec) { + into.extend_from_slice(self.0.as_slice()); } } @@ -97,8 +42,8 @@ fn echo() { let listener = TcpListener::bind(&"127.0.0.1:0".parse().unwrap(), &handle).unwrap(); let addr = listener.local_addr().unwrap(); let srv = listener.incoming().for_each(move |(socket, _)| { - let framed = EasyFramed::new(socket, LineDecoder, LineEncoder); - handle.spawn(EchoFramed { inner: framed, eof: false }); + let (stream, sink) = socket.framed::().split(); + handle.spawn(sink.send_all(stream).map(|_| ()).map_err(|_| ())); Ok(()) });