use std::{io, thread};
use std::sync::mpsc;
use tokio_core;
use core::futures::{self, Future};
pub enum UninitializedRemote {
Shared(tokio_core::reactor::Remote),
Unspawned,
}
impl UninitializedRemote {
pub fn initialize(self) -> io::Result<Remote> {
self.init_with_name("event.loop")
}
pub fn init_with_name<T: Into<String>>(self, name: T) -> io::Result<Remote> {
match self {
UninitializedRemote::Shared(remote) => Ok(Remote::Shared(remote)),
UninitializedRemote::Unspawned => RpcEventLoop::with_name(Some(name.into())).map(Remote::Spawned),
}
}
}
pub enum Remote {
Shared(tokio_core::reactor::Remote),
Spawned(RpcEventLoop),
}
impl Remote {
pub fn remote(&self) -> tokio_core::reactor::Remote {
match *self {
Remote::Shared(ref remote) => remote.clone(),
Remote::Spawned(ref eloop) => eloop.remote(),
}
}
pub fn close(self) {
if let Remote::Spawned(eloop) = self {
eloop.close()
}
}
pub fn wait(self) {
if let Remote::Spawned(eloop) = self {
let _ = eloop.wait();
}
}
}
pub struct RpcEventLoop {
remote: tokio_core::reactor::Remote,
close: Option<futures::Complete<()>>,
handle: Option<thread::JoinHandle<()>>,
}
impl Drop for RpcEventLoop {
fn drop(&mut self) {
self.close.take().map(|v| v.send(()));
}
}
impl RpcEventLoop {
pub fn spawn() -> io::Result<Self> {
RpcEventLoop::with_name(None)
}
pub fn with_name(name: Option<String>) -> io::Result<Self> {
let (stop, stopped) = futures::oneshot();
let (tx, rx) = mpsc::channel();
let mut tb = thread::Builder::new();
if let Some(name) = name {
tb = tb.name(name);
}
let handle = tb.spawn(move || {
let el = tokio_core::reactor::Core::new();
match el {
Ok(mut el) => {
tx.send(Ok(el.remote())).expect("Rx is blocking upper thread.");
let _ = el.run(futures::empty().select(stopped));
},
Err(err) => {
tx.send(Err(err)).expect("Rx is blocking upper thread.");
}
}
}).expect("Couldn't spawn a thread.");
let remote = rx.recv().expect("tx is transfered to a newly spawned thread.");
remote.map(|remote| RpcEventLoop {
remote: remote,
close: Some(stop),
handle: Some(handle),
})
}
pub fn remote(&self) -> tokio_core::reactor::Remote {
self.remote.clone()
}
pub fn wait(mut self) -> thread::Result<()> {
self.handle.take().expect("Handle is always set before self is consumed.").join()
}
pub fn close(mut self) {
let _ = self.close.take().expect("Close is always set before self is consumed.").send(()).map_err(|e| {
warn!("Event Loop is already finished. {:?}", e);
});
}
}