use bytes::{BufMut, BytesMut};
use futures::{future::{self, Either}, prelude::*};
use holly::{actor, sink, stream, prelude::*};
use std::{io, net::SocketAddr};
use tokio::{
codec::{Encoder, Decoder, Framed},
net::{TcpListener, TcpStream},
runtime::Runtime
};
#[derive(Clone, Debug)]
enum Frame {
Ping(u8),
Pong(u8),
Finish
}
struct Codec;
impl Encoder for Codec {
type Item = Frame;
type Error = io::Error;
fn encode(&mut self, item: Self::Item, dst: &mut BytesMut) -> Result<(), Self::Error> {
dst.reserve(2);
match item {
Frame::Ping(n) => dst.put(&[1, n][..]),
Frame::Pong(n) => dst.put(&[2, n][..]),
Frame::Finish => dst.put(&[3, 0][..]),
}
Ok(())
}
}
impl Decoder for Codec {
type Item = Frame;
type Error = io::Error;
fn decode(&mut self, src: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
if src.len() < 2 {
return Ok(None)
}
match &src.split_to(2)[..] {
[1, n] => Ok(Some(Frame::Ping(*n))),
[2, n] => Ok(Some(Frame::Pong(*n))),
[3, 0] => Ok(Some(Frame::Finish)),
_ => Err(io::ErrorKind::InvalidData.into())
}
}
}
enum PlayerMessage<T> {
Listen(SocketAddr, Addr<T, ()>),
Connect(SocketAddr, Addr<T, ()>),
Data(Frame),
IoError(io::Error)
}
impl<T> From<stream::Event<(), Frame, io::Error>> for PlayerMessage<T> {
fn from(x: stream::Event<(), Frame, io::Error>) -> Self {
match x {
stream::Event::Item(_, f) => PlayerMessage::Data(f),
stream::Event::End(_) => PlayerMessage::IoError(io::ErrorKind::UnexpectedEof.into()),
stream::Event::Error(_, e) => PlayerMessage::IoError(e)
}
}
}
#[derive(Debug)]
enum Error {
Actor(holly::Error),
Io(io::Error)
}
impl From<io::Error> for Error {
fn from(e: io::Error) -> Self {
Error::Io(e)
}
}
impl From<holly::Error> for Error {
fn from(e: holly::Error) -> Self {
Error::Actor(e)
}
}
enum Player {
Init,
Running {
input: stream::KillCord,
output: Box<dyn Sink<SinkItem = Frame, SinkError = io::Error> + Send>
},
Closing {
output: Box<dyn Sink<SinkItem = Frame, SinkError = io::Error> + Send>
}
}
impl<T> Actor<PlayerMessage<T>, Error> for Player
where
T: From<()> + Send + 'static
{
type Result = Box<dyn Future<Item = State<Self, PlayerMessage<T>>, Error = Error> + Send>;
fn process(self, ctx: &mut Context<PlayerMessage<T>>, msg: Option<PlayerMessage<T>>) -> Self::Result {
match self {
Player::Init => match msg {
Some(PlayerMessage::Listen(a, reply)) => {
println!("Player {}: Listening at: {}", ctx.id(), a);
let f = future::result(TcpListener::bind(&a)).from_err()
.and_then(move |l| {
reply.send(()).from_err().map(|_| l)
})
.and_then(move |l| l.incoming().from_err()
.into_future()
.map_err(|e| e.0)
.and_then(move |(conn, _)| match conn {
None => Either::A(future::ok(State::Done)),
Some(c) => {
let (sink, stream) = Framed::new(c, Codec).split();
let (r, s) = stream::stoppable((), stream);
let output = Box::new(sink);
let state = Player::Running { input: r, output };
let boxed = Box::new(s.map(Into::into));
Either::B(future::ok(State::Stream(state, boxed)))
}
}));
Box::new(f)
}
Some(PlayerMessage::Connect(a, reply)) => {
println!("Player {}: Connecting to: {}", ctx.id(), a);
let f = TcpStream::connect(&a).from_err()
.and_then(move |c| {
let (sink, stream) = Framed::new(c, Codec).split();
let (r, s) = stream::stoppable((), stream);
let output = Box::new(sink);
let state = Player::Running { input: r, output };
reply.send(()).from_err().map(|_| {
State::Stream(state, Box::new(s.map(Into::into)))
})
});
Box::new(f)
}
_ => Box::new(future::ok(State::Done))
}
Player::Running { input, output } => match msg {
Some(PlayerMessage::IoError(e)) => {
println!("Player {}: i/o error: {}", ctx.id(), e);
Box::new(future::ok(State::Done))
}
Some(PlayerMessage::Data(Frame::Finish)) => {
println!("Player {}: finish", ctx.id());
Box::new(future::ok(State::Close(Player::Closing { output })))
}
Some(PlayerMessage::Data(Frame::Ping(0))) => {
let id = ctx.id();
println!("Player {}: Ping(0)", id);
let f = sink::send(output, Frame::Finish)
.and_then(sink::flush)
.map(move |output| State::Close(Player::Closing { output }))
.map_err(move |e| {
println!("Player {}: Error sending Ping(0): {}", id, e);
e.into()
});
Box::new(f)
}
Some(PlayerMessage::Data(Frame::Ping(n))) => {
let id = ctx.id();
println!("Player {}: Ping({})", id, n);
let f = sink::send(output, Frame::Pong(n - 1))
.and_then(sink::flush)
.map(move |os| State::Ready(Player::Running { input, output: os }))
.map_err(move |e| {
println!("Player {}: Error sending Ping({}): {}", id, n - 1, e);
e.into()
});
Box::new(f)
}
Some(PlayerMessage::Data(Frame::Pong(0))) => {
let id = ctx.id();
println!("Player {}: Pong(0)", id);
let f = sink::send(output, Frame::Finish)
.and_then(sink::flush)
.map(move |output| State::Close(Player::Closing { output }))
.map_err(move |e| {
println!("Player {}: Error sending Pong(0): {}", id, e);
e.into()
});
Box::new(f)
}
Some(PlayerMessage::Data(Frame::Pong(n))) => {
let id = ctx.id();
println!("Player {}: Pong({})", id, n);
let f = sink::send(output, Frame::Ping(n - 1))
.and_then(sink::flush)
.map(move |os| State::Ready(Player::Running { input, output: os }))
.map_err(move |e| {
println!("Player {}: Error sending Pong({}): {}", id, n - 1, e);
e.into()
});
Box::new(f)
}
Some(PlayerMessage::Listen(..)) | Some(PlayerMessage::Connect(..)) | None => {
Box::new(future::ok(State::Done))
}
}
Player::Closing { output } => match msg {
None => {
let i = ctx.id();
println!("Player {}: closing output", i);
let f = sink::close(output)
.map(move |_| {
println!("Player {}: done", i);
State::Done
})
.map_err(move |e| {
println!("Player {}: Error closing output: {}", i, e);
e.into()
});
Box::new(f)
}
_ => Box::new(future::ok(State::Ready(Player::Closing { output })))
}
}
}
}
enum Root {
Init,
ServerStarted,
ClientStarted(Addr<PlayerMessage<RootMessage>>),
}
enum RootMessage {
Start,
Reply(()),
PlayerFailed(actor::Fail<Error>)
}
impl From<()> for RootMessage {
fn from(r: ()) -> Self {
RootMessage::Reply(r)
}
}
impl From<actor::Fail<Error>> for RootMessage {
fn from(x: actor::Fail<Error>) -> Self {
RootMessage::PlayerFailed(x)
}
}
impl Actor<RootMessage, holly::Error> for Root {
type Result = Box<dyn Future<Item = State<Self, RootMessage>, Error = holly::Error> + Send>;
fn process(self, ctx: &mut Context<RootMessage>, msg: Option<RootMessage>) -> Self::Result {
match (self, msg) {
(Root::Init, Some(RootMessage::Start)) => {
println!("Root: Starting server.");
let listen_addr = "127.0.0.1:12345".parse().unwrap();
let addr = ctx.address();
let opts = actor::Options::default().supervisor(ctx.address());
let f = future::result(ctx.scheduler().spawn_ext(Player::Init, opts))
.and_then(move |player| {
player.send(PlayerMessage::Listen(listen_addr, addr.cast()))
})
.map(|_| State::Ready(Root::ServerStarted));
Box::new(f)
}
(Root::ServerStarted, Some(RootMessage::Reply(()))) => {
println!("Root: Starting client.");
let connect_addr = "127.0.0.1:12345".parse().unwrap();
let addr = ctx.address();
let opts = actor::Options::default().supervisor(ctx.address());
let f = future::result(ctx.scheduler().spawn_ext(Player::Init, opts))
.and_then(move |player| {
player.send(PlayerMessage::Connect(connect_addr, addr.cast()))
})
.map(move |player| {
State::Ready(Root::ClientStarted(player))
});
Box::new(f)
}
(Root::ClientStarted(player), Some(RootMessage::Reply(()))) => {
println!("Root: Player {} - Go!", player.id());
let f = player.send(PlayerMessage::Data(Frame::Ping(32)))
.map(|_| State::Done);
Box::new(f)
}
(_, Some(RootMessage::PlayerFailed(f))) => {
println!("Root: Player {} failed: {:?}", f.id(), f.error());
Box::new(future::ok(State::Done))
}
_ => Box::new(future::ok(State::Done))
}
}
}
#[test]
fn ping_pong_over_tcp() -> Result<(), Box<dyn std::error::Error>> {
let rt = Runtime::new()?;
let exec = Scheduler::new(rt.executor());
let mut r = exec.spawn(Root::Init)?;
r.send_now(RootMessage::Start)?;
rt.shutdown_on_idle().wait().unwrap();
Ok(())
}