use futures::sync::mpsc;
use futures::{Future, Poll, Async, Stream, Sink, AsyncSink, StartSend};
use std::{fmt, io};
use streaming::{Message, Body};
use super::{Frame, Transport};
use buffer_one::BufferOne;
pub struct Pipeline<T> where T: Dispatch {
run: bool,
dispatch: BufferOne<DispatchSink<T>>,
out_body: Option<BodySender<T::BodyOut, T::Error>>,
in_body: Option<T::Stream>,
is_flushed: bool,
}
pub type PipelineMessage<T, B, E> = Result<Message<T, B>, E>;
pub trait Dispatch {
type Io;
type In;
type BodyIn;
type Out;
type BodyOut;
type Error: From<io::Error>;
type Stream: Stream<Item = Self::BodyIn, Error = Self::Error>;
type Transport: Transport<Item = Frame<Self::Out, Self::BodyOut, Self::Error>,
SinkItem = Frame<Self::In, Self::BodyIn, Self::Error>>;
fn transport(&mut self) -> &mut Self::Transport;
fn dispatch(&mut self, message: PipelineMessage<Self::Out, Body<Self::BodyOut, Self::Error>, Self::Error>) -> io::Result<()>;
fn poll(&mut self) -> Poll<Option<PipelineMessage<Self::In, Self::Stream, Self::Error>>, io::Error>;
fn has_in_flight(&self) -> bool;
}
struct DispatchSink<T> {
inner: T,
}
type BodySender<B, E> = BufferOne<mpsc::Sender<Result<B, E>>>;
impl<T> Pipeline<T> where T: Dispatch {
pub fn new(dispatch: T) -> Pipeline<T> {
let dispatch = DispatchSink { inner: dispatch };
let dispatch = BufferOne::new(dispatch);
Pipeline {
run: true,
dispatch: dispatch,
out_body: None,
in_body: None,
is_flushed: true,
}
}
fn is_done(&self) -> bool {
!self.run && self.is_flushed && !self.has_in_flight()
}
fn read_out_frames(&mut self) -> io::Result<()> {
while self.run {
if !self.check_out_body_stream() {
break;
}
if let Async::Ready(frame) = try!(self.dispatch.get_mut().inner.transport().poll()) {
try!(self.process_out_frame(frame));
} else {
break;
}
}
Ok(())
}
fn check_out_body_stream(&mut self) -> bool {
let body = match self.out_body {
Some(ref mut body) => body,
None => return true,
};
body.poll_ready().is_ready()
}
fn process_out_frame(&mut self,
frame: Option<Frame<T::Out, T::BodyOut, T::Error>>)
-> io::Result<()> {
trace!("process_out_frame");
match frame {
Some(Frame::Message { message, body }) => {
if body {
trace!("read out message with body");
let (tx, rx) = Body::pair();
let message = Message::WithBody(message, rx);
self.out_body = Some(BufferOne::new(tx));
if let Err(e) = self.dispatch.get_mut().inner.dispatch(Ok(message)) {
panic!("unimplemented error handling: {:?}", e);
}
} else {
trace!("read out message");
let message = Message::WithoutBody(message);
self.out_body = None;
if let Err(e) = self.dispatch.get_mut().inner.dispatch(Ok(message)) {
panic!("unimplemented error handling: {:?}", e);
}
}
}
Some(Frame::Body { chunk }) => {
match chunk {
Some(chunk) => {
trace!("read out body chunk");
try!(self.process_out_body_chunk(chunk));
}
None => {
trace!("read out body EOF");
let _ = self.out_body.take();
}
}
}
None => {
trace!("read None");
self.run = false;
}
Some(Frame::Error { .. }) => {
return Err(io::Error::new(io::ErrorKind::BrokenPipe, "An error occurred."));
}
}
Ok(())
}
fn process_out_body_chunk(&mut self, chunk: T::BodyOut) -> io::Result<()> {
trace!("process_out_body_chunk");
let mut reset = false;
match self.out_body {
Some(ref mut body) => {
debug!("sending a chunk");
match body.start_send(Ok(chunk)) {
Ok(AsyncSink::Ready) => debug!("immediately done"),
Err(_e) => reset = true, Ok(AsyncSink::NotReady(_)) => {
unreachable!();
}
}
}
None => {
debug!("interest canceled");
}
}
if reset {
self.out_body = None;
}
Ok(())
}
fn write_in_frames(&mut self) -> io::Result<()> {
trace!("write_in_frames");
while self.dispatch.poll_ready().is_ready() {
if !try!(self.write_in_body()) {
debug!("write in body not done");
break;
}
debug!("write in body done");
match try!(self.dispatch.get_mut().inner.poll()) {
Async::Ready(Some(Ok(message))) => {
trace!(" --> got message");
try!(self.write_in_message(Ok(message)));
}
Async::Ready(Some(Err(error))) => {
trace!(" --> got error");
try!(self.write_in_message(Err(error)));
}
Async::Ready(None) => {
trace!(" --> got None");
break;
}
Async::NotReady => break,
}
}
Ok(())
}
fn write_in_message(&mut self, message: Result<Message<T::In, T::Stream>, T::Error>) -> io::Result<()> {
trace!("write_in_message");
match message {
Ok(Message::WithoutBody(val)) => {
trace!("got in_flight value without body");
let msg = Frame::Message { message: val, body: false };
try!(assert_send(&mut self.dispatch, msg));
assert!(self.in_body.is_none());
self.in_body = None;
}
Ok(Message::WithBody(val, body)) => {
trace!("got in_flight value with body");
let msg = Frame::Message { message: val, body: true };
try!(assert_send(&mut self.dispatch, msg));
assert!(self.in_body.is_none());
self.in_body = Some(body);
}
Err(e) => {
trace!("got in_flight error");
let msg = Frame::Error { error: e };
try!(assert_send(&mut self.dispatch, msg));
}
}
Ok(())
}
fn write_in_body(&mut self) -> io::Result<bool> {
trace!("write_in_body");
if self.in_body.is_some() {
loop {
if !self.dispatch.poll_ready().is_ready() {
return Ok(false);
}
match self.in_body.as_mut().unwrap().poll() {
Ok(Async::Ready(Some(chunk))) => {
try!(assert_send(&mut self.dispatch,
Frame::Body { chunk: Some(chunk) }));
}
Ok(Async::Ready(None)) => {
try!(assert_send(&mut self.dispatch,
Frame::Body { chunk: None }));
break;
}
Err(_) => {
unimplemented!();
}
Ok(Async::NotReady) => {
debug!("not ready");
return Ok(false);
}
}
}
}
self.in_body = None;
Ok(true)
}
fn flush(&mut self) -> io::Result<()> {
self.is_flushed = try!(self.dispatch.poll_complete()).is_ready();
if let Some(ref mut out_body) = self.out_body {
if out_body.poll_complete().is_ok() {
return Ok(());
}
} else {
return Ok(());
}
self.out_body = None;
Ok(())
}
fn has_in_flight(&self) -> bool {
self.dispatch.get_ref().inner.has_in_flight()
}
}
impl<T> Future for Pipeline<T> where T: Dispatch {
type Item = ();
type Error = io::Error;
fn poll(&mut self) -> Poll<(), io::Error> {
trace!("Pipeline::tick");
self.dispatch.get_mut().inner.transport().tick();
try!(self.read_out_frames());
try!(self.write_in_frames());
try!(self.flush());
if self.is_done() {
return Ok(().into())
}
Ok(Async::NotReady)
}
}
impl<T> fmt::Debug for Pipeline<T>
where T: Dispatch + fmt::Debug,
T::In: fmt::Debug,
T::BodyIn: fmt::Debug,
T::BodyOut: fmt::Debug,
T::Error: fmt::Debug,
T::Stream: fmt::Debug,
{
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
f.debug_struct("Pipeline")
.field("run", &self.run)
.field("dispatch", &self.dispatch)
.field("out_body", &"Sender { ... }")
.field("in_body", &self.in_body)
.field("is_flushed", &self.is_flushed)
.finish()
}
}
impl<T: Dispatch> Sink for DispatchSink<T> {
type SinkItem = <T::Transport as Sink>::SinkItem;
type SinkError = io::Error;
fn start_send(&mut self, item: Self::SinkItem)
-> StartSend<Self::SinkItem, io::Error>
{
self.inner.transport().start_send(item)
}
fn poll_complete(&mut self) -> Poll<(), io::Error> {
self.inner.transport().poll_complete()
}
fn close(&mut self) -> Poll<(), io::Error> {
self.inner.transport().close()
}
}
impl<T: fmt::Debug> fmt::Debug for DispatchSink<T> {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
f.debug_struct("DispatchSink")
.field("inner", &self.inner)
.finish()
}
}
fn assert_send<S: Sink>(s: &mut S, item: S::SinkItem) -> Result<(), S::SinkError> {
match try!(s.start_send(item)) {
AsyncSink::Ready => Ok(()),
AsyncSink::NotReady(_) => {
panic!("sink reported itself as ready after `poll_ready` but was \
then unable to accept a message")
}
}
}