use errors::*;
use io::Stream;
use lib_futures::failed;
use lib_futures::future::select_ok;
use lib_futures::future::SelectOk;
use lib_futures::Async;
use lib_futures::Async::Ready;
use lib_futures::Failed;
use lib_futures::Future;
use lib_futures::Poll;
use myc::packets::PacketParser;
use std::io;
use std::net::ToSocketAddrs;
use tokio::net::ConnectFuture;
use tokio::net::TcpStream;
steps! {
ConnectingStream {
WaitForStream(SelectOk<ConnectFuture>),
Fail(Failed<(), Error>),
}
}
pub struct ConnectingStream {
step: Step,
}
pub fn new<S>(addr: S) -> ConnectingStream
where
S: ToSocketAddrs,
{
match addr.to_socket_addrs() {
Ok(addresses) => {
let mut streams = Vec::new();
for address in addresses {
streams.push(TcpStream::connect(&address));
}
if streams.len() > 0 {
ConnectingStream {
step: Step::WaitForStream(select_ok(streams)),
}
} else {
let err = io::Error::new(
io::ErrorKind::InvalidInput,
"could not resolve to any address",
);
ConnectingStream {
step: Step::Fail(failed(err.into())),
}
}
}
Err(err) => ConnectingStream {
step: Step::Fail(failed(err.into())),
},
}
}
impl Future for ConnectingStream {
type Item = Stream;
type Error = Error;
fn poll(&mut self) -> Poll<Self::Item, Self::Error> {
match try_ready!(self.either_poll()) {
Out::WaitForStream((stream, _)) => Ok(Ready(Stream {
closed: false,
parser: Some(PacketParser::empty()),
packets: Default::default(),
buf: Vec::new(),
endpoint: Some(stream.into()),
})),
Out::Fail(_) => unreachable!(),
}
}
}