extern crate capnp;
extern crate capnp_futures;
extern crate futures;
use futures::{Future};
use futures::sync::oneshot;
use capnp::Error;
use capnp::capability::Promise;
use capnp::private::capability::{ClientHook, ServerHook};
use std::cell::{RefCell};
use std::rc::{Rc};
use task_set::TaskSet;
pub mod rpc_capnp {
include!(concat!(env!("OUT_DIR"), "/rpc_capnp.rs"));
}
pub mod rpc_twoparty_capnp {
include!(concat!(env!("OUT_DIR"), "/rpc_twoparty_capnp.rs"));
}
#[macro_export]
macro_rules! pry {
($expr:expr) => (
match $expr {
::std::result::Result::Ok(val) => val,
::std::result::Result::Err(err) => {
return ::capnp::capability::Promise::err(::std::convert::From::from(err))
}
})
}
mod broken;
mod local;
mod queued;
mod rpc;
mod attach;
mod forked_promise;
mod sender_queue;
mod split;
mod task_set;
pub mod twoparty;
pub trait OutgoingMessage {
fn get_body<'a>(&'a mut self) -> ::capnp::Result<::capnp::any_pointer::Builder<'a>>;
fn get_body_as_reader<'a>(&'a self) -> ::capnp::Result<::capnp::any_pointer::Reader<'a>>;
fn send(self: Box<Self>)
-> (Promise<Rc<::capnp::message::Builder<::capnp::message::HeapAllocator>>, ::capnp::Error>,
Rc<::capnp::message::Builder<::capnp::message::HeapAllocator>>);
fn take(self: Box<Self>) -> ::capnp::message::Builder<::capnp::message::HeapAllocator>;
}
pub trait IncomingMessage {
fn get_body<'a>(&'a self) -> ::capnp::Result<::capnp::any_pointer::Reader<'a>>;
}
pub trait Connection<VatId> {
fn get_peer_vat_id(&self) -> VatId;
fn new_outgoing_message(&mut self, first_segment_word_size: u32) -> Box<OutgoingMessage>;
fn receive_incoming_message(&mut self) -> Promise<Option<Box<IncomingMessage>>, Error>;
fn shutdown(&mut self, result: ::capnp::Result<()>) -> Promise<(), Error>;
}
pub trait VatNetwork<VatId> {
fn connect(&mut self, host_id: VatId) -> Option<Box<Connection<VatId>>>;
fn accept(&mut self) -> Promise<Box<Connection<VatId>>, ::capnp::Error>;
fn drive_until_shutdown(&mut self) -> Promise<(), Error>;
}
#[must_use = "futures do nothing unless polled"]
pub struct RpcSystem<VatId> where VatId: 'static {
network: Box<::VatNetwork<VatId>>,
bootstrap_cap: Box<ClientHook>,
connection_state: Rc<RefCell<Option<Rc<rpc::ConnectionState<VatId>>>>>,
tasks: TaskSet<(), Error>,
handle: ::task_set::TaskSetHandle<(), Error>
}
impl <VatId> RpcSystem <VatId> {
pub fn new(
mut network: Box<::VatNetwork<VatId>>,
bootstrap: Option<::capnp::capability::Client>) -> RpcSystem<VatId>
{
let bootstrap_cap = match bootstrap {
Some(cap) => cap.hook,
None => broken::new_cap(Error::failed("no bootstrap capabiity".to_string())),
};
let (mut handle, tasks) = TaskSet::new(Box::new(SystemTaskReaper));
let mut handle1 = handle.clone();
handle.add(network.drive_until_shutdown().then(move |r| {
let r = match r {
Ok(()) => Ok(()),
Err(e) => {
if e.kind != ::capnp::ErrorKind::Disconnected {
Err(e)
} else {
Ok(())
}
}
};
handle1.terminate(r);
Ok(())
}));
let mut result = RpcSystem {
network: network,
bootstrap_cap: bootstrap_cap,
connection_state: Rc::new(RefCell::new(None)),
tasks: tasks,
handle: handle.clone(),
};
let accept_loop = result.accept_loop();
handle.add(accept_loop);
result
}
pub fn bootstrap<T>(&mut self, vat_id: VatId) -> T
where T: ::capnp::capability::FromClientHook
{
let connection = match self.network.connect(vat_id) {
Some(connection) => connection,
None => {
return T::new(self.bootstrap_cap.clone());
}
};
let connection_state =
RpcSystem::get_connection_state(self.connection_state.clone(),
self.bootstrap_cap.clone(),
connection, self.handle.clone());
let hook = rpc::ConnectionState::bootstrap(connection_state.clone());
T::new(hook)
}
fn accept_loop(&mut self) -> Promise<(), Error> {
let connection_state_ref = self.connection_state.clone();
let bootstrap_cap = self.bootstrap_cap.clone();
let handle = self.handle.clone();
Promise::from_future(self.network.accept().map(move |connection| {
RpcSystem::get_connection_state(connection_state_ref,
bootstrap_cap,
connection,
handle);
}))
}
fn get_connection_state(connection_state_ref: Rc<RefCell<Option<Rc<rpc::ConnectionState<VatId>>>>>,
bootstrap_cap: Box<ClientHook>,
connection: Box<::Connection<VatId>>,
mut handle: ::task_set::TaskSetHandle<(), Error>)
-> Rc<rpc::ConnectionState<VatId>>
{
let (tasks, result) = match *connection_state_ref.borrow() {
Some(ref connection_state) => {
return connection_state.clone()
}
None => {
let (on_disconnect_fulfiller, on_disconnect_promise) =
oneshot::channel::<Promise<(), Error>>();
let connection_state_ref1 = connection_state_ref.clone();
handle.add(on_disconnect_promise.then(move |shutdown_promise| {
*connection_state_ref1.borrow_mut() = None;
match shutdown_promise {
Ok(s) => s,
Err(e) => Promise::err(Error::failed(format!("{}", e))),
}
}));
rpc::ConnectionState::new(bootstrap_cap, connection, on_disconnect_fulfiller)
}
};
*connection_state_ref.borrow_mut() = Some(result.clone());
handle.add(tasks);
result
}
pub fn get_disconnector(&self) -> rpc::Disconnector<VatId> {
rpc::Disconnector::new(self.connection_state.clone())
}
}
impl <VatId> Future for RpcSystem<VatId> where VatId: 'static {
type Item = ();
type Error = Error;
fn poll(&mut self) -> ::futures::Poll<Self::Item, Self::Error> {
self.tasks.poll()
}
}
pub struct Server;
impl ServerHook for Server {
fn new_client(server: Box<::capnp::capability::Server>) -> ::capnp::capability::Client {
::capnp::capability::Client::new(Box::new(local::Client::new(server)))
}
}
pub fn new_promise_client<T, F>(client_promise: F) -> T
where T: ::capnp::capability::FromClientHook,
F: ::futures::Future<Item=::capnp::capability::Client,Error=Error>,
F: 'static
{
let mut queued_client = ::queued::Client::new(None);
let weak_client = Rc::downgrade(&queued_client.inner);
queued_client.drive(client_promise.then(move |r| {
if let Some(queued_inner) = weak_client.upgrade() {
::queued::ClientInner::resolve(&queued_inner, r.map(|c| c.hook));
}
Ok(())
}));
T::new(Box::new(queued_client))
}
struct SystemTaskReaper;
impl ::task_set::TaskReaper<(), Error> for SystemTaskReaper {
fn task_failed(&mut self, error: Error) {
println!("ERROR: {}", error);
}
}