use futures::{
future::{self, FutureResult},
prelude::*,
stream::futures_unordered::FuturesUnordered
};
use holly::{Error, prelude::*};
use std::{collections::HashMap, fmt::Display};
use tokio_threadpool::ThreadPool;
enum Request<T> {
Register(Addr<Message<T>>),
Broadcast(T)
}
enum Message<T> {
Id(u32),
Data(T)
}
struct Broadcast<T> {
id_counter: u32,
receivers: HashMap<u32, Addr<Message<T>>>
}
impl<T: Clone + Send + 'static> Actor<Request<T>, Error> for Broadcast<T> {
type Result = Box<dyn Future<Item = State<Self, Request<T>>, Error = Error> + Send>;
fn process(mut self, _ctx: &mut Context<Request<T>>, msg: Option<Request<T>>) -> Self::Result {
match msg {
Some(Request::Register(addr)) => {
let id = self.id_counter;
self.receivers.insert(id, addr.clone());
self.id_counter += 1;
let f = addr.send(Message::Id(id)).then(move |r| {
if r.is_err() {
self.receivers.remove(&id);
self.id_counter -= 1
}
Ok(State::Ready(self))
});
Box::new(f)
}
Some(Request::Broadcast(msg)) => {
let mut f = FuturesUnordered::new();
for a in self.receivers.values().cloned() {
f.push(a.send(Message::Data(msg.clone())))
}
let f = f.for_each(|_| Ok(())).map(move |_| State::Done);
Box::new(f)
}
None => {
Box::new(future::ok(State::Done))
}
}
}
}
struct Receiver(u32);
impl<T: Display> Actor<Message<T>, ()> for Receiver {
type Result = FutureResult<State<Self, Message<T>>, ()>;
fn process(mut self, _ctx: &mut Context<Message<T>>, msg: Option<Message<T>>) -> Self::Result {
match msg {
Some(Message::Id(id)) => {
self.0 = id;
future::ok(State::Ready(self))
}
Some(Message::Data(x)) => {
println!("Receiver {}: {}", self.0, x);
future::ok(State::Done)
}
None => future::ok(State::Done)
}
}
}
#[test]
fn broadcast() -> Result<(), Error> {
let pool = ThreadPool::new();
let exec = Scheduler::new(pool.sender().clone());
let mut broadcast = exec.spawn(Broadcast { id_counter: 0, receivers: HashMap::new() })?;
for _ in 0 .. 10 {
broadcast.send_now(Request::Register(exec.spawn(Receiver(0))?))?;
}
broadcast.send_now(Request::Broadcast("Hello receivers!"))?;
pool.shutdown_on_idle().wait().unwrap();
Ok(())
}