#![warn(missing_docs)]
extern crate atomic_immut;
extern crate bytecodec;
extern crate byteorder;
extern crate factory;
extern crate fibers;
extern crate fibers_tasque;
extern crate futures;
extern crate prometrics;
#[macro_use]
extern crate slog;
#[macro_use]
extern crate trackable;
pub use error::{Error, ErrorKind};
pub mod client {
pub use client_service::{ClientService, ClientServiceBuilder, ClientServiceHandle};
pub use client_side_handlers::Response;
pub use rpc_client::{CallClient, CastClient, Options};
}
pub mod channel;
pub mod metrics;
pub mod server {
pub use rpc_server::{Server, ServerBuilder};
pub use server_side_handlers::{HandleCall, HandleCast, NoReply, Reply};
}
use client::{CallClient, CastClient, ClientServiceHandle};
mod client_service;
mod client_side_channel;
mod client_side_handlers;
mod error;
mod message;
mod message_stream;
mod packet;
mod rpc_client;
mod rpc_server;
mod server_side_channel;
mod server_side_handlers;
pub type Result<T> = std::result::Result<T, Error>;
#[derive(Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct ProcedureId(pub u32);
impl std::fmt::Debug for ProcedureId {
fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
write!(f, "ProcedureId(0x{:08x})", self.0)
}
}
pub trait Call: Sized + Send + Sync + 'static {
const ID: ProcedureId;
const NAME: &'static str;
type Req: Send + 'static;
type ReqEncoder: bytecodec::Encode<Item = Self::Req> + Send + 'static;
type ReqDecoder: bytecodec::Decode<Item = Self::Req> + Send + 'static;
type Res: Send + 'static;
type ResEncoder: bytecodec::Encode<Item = Self::Res> + Send + 'static;
type ResDecoder: bytecodec::Decode<Item = Self::Res> + Send + 'static;
#[allow(unused_variables)]
fn enable_async_request(request: &Self::Req) -> bool {
false
}
#[allow(unused_variables)]
fn enable_async_response(response: &Self::Res) -> bool {
false
}
fn client(service: &ClientServiceHandle) -> CallClient<Self>
where
Self::ReqEncoder: Default,
Self::ResDecoder: Default,
{
Self::client_with_codec(service, Default::default(), Default::default())
}
fn client_with_decoder(
service: &ClientServiceHandle,
decoder: Self::ResDecoder,
) -> CallClient<Self>
where
Self::ReqEncoder: Default,
{
Self::client_with_codec(service, decoder, Default::default())
}
fn client_with_encoder(
service: &ClientServiceHandle,
encoder: Self::ReqEncoder,
) -> CallClient<Self>
where
Self::ResDecoder: Default,
{
Self::client_with_codec(service, Default::default(), encoder)
}
fn client_with_codec(
service: &ClientServiceHandle,
decoder: Self::ResDecoder,
encoder: Self::ReqEncoder,
) -> CallClient<Self> {
CallClient::new(service, decoder, encoder)
}
}
pub trait Cast: Sized + Sync + Send + 'static {
const ID: ProcedureId;
const NAME: &'static str;
type Notification: Send + 'static;
type Encoder: bytecodec::Encode<Item = Self::Notification> + Send + 'static;
type Decoder: bytecodec::Decode<Item = Self::Notification> + Send + 'static;
#[allow(unused_variables)]
fn enable_async(notification: &Self::Notification) -> bool {
false
}
fn client(service: &ClientServiceHandle) -> CastClient<Self>
where
Self::Encoder: Default,
{
Self::client_with_encoder(service, Default::default())
}
fn client_with_encoder(
service: &ClientServiceHandle,
encoder: Self::Encoder,
) -> CastClient<Self> {
CastClient::new(service, encoder)
}
}
#[cfg(test)]
mod test {
use bytecodec::bytes::{BytesEncoder, RemainingBytesDecoder};
use fibers::{Executor, InPlaceExecutor, Spawn};
use futures::Future;
use client::ClientServiceBuilder;
use server::{HandleCall, Reply, ServerBuilder};
use {Call, ProcedureId};
struct EchoRpc;
impl Call for EchoRpc {
const ID: ProcedureId = ProcedureId(0);
const NAME: &'static str = "echo";
type Req = Vec<u8>;
type ReqEncoder = BytesEncoder<Vec<u8>>;
type ReqDecoder = RemainingBytesDecoder;
type Res = Vec<u8>;
type ResEncoder = BytesEncoder<Vec<u8>>;
type ResDecoder = RemainingBytesDecoder;
fn enable_async_request(x: &Self::Req) -> bool {
x == b"async"
}
fn enable_async_response(x: &Self::Res) -> bool {
x == b"async"
}
}
struct EchoHandler;
impl HandleCall<EchoRpc> for EchoHandler {
fn handle_call(&self, request: <EchoRpc as Call>::Req) -> Reply<EchoRpc> {
Reply::done(request)
}
}
#[test]
fn it_works() {
let mut executor = track_try_unwrap!(track_any_err!(InPlaceExecutor::new()));
let server_addr = "127.0.0.1:1920".parse().unwrap();
let server = ServerBuilder::new(server_addr)
.add_call_handler(EchoHandler)
.finish(executor.handle());
executor.spawn(server.map_err(|e| panic!("{}", e)));
let service = ClientServiceBuilder::new().finish(executor.handle());
let service_handle = service.handle();
let request = Vec::from(&b"hello"[..]);
let response = EchoRpc::client(&service_handle).call(server_addr, request.clone());
executor.spawn(service.map_err(|e| panic!("{}", e)));
let result = track_try_unwrap!(track_any_err!(executor.run_future(response)));
assert_eq!(result.ok(), Some(request));
let metrics = service_handle
.metrics()
.channels()
.as_map()
.load()
.get(&server_addr)
.cloned()
.unwrap();
assert_eq!(metrics.async_outgoing_messages(), 0);
assert_eq!(metrics.async_incoming_messages(), 0);
}
#[test]
fn large_message_works() {
let mut executor = track_try_unwrap!(track_any_err!(InPlaceExecutor::new()));
let server_addr = "127.0.0.1:1921".parse().unwrap();
let server = ServerBuilder::new(server_addr)
.add_call_handler(EchoHandler)
.finish(executor.handle());
executor.spawn(server.map_err(|e| panic!("{}", e)));
let service = ClientServiceBuilder::new().finish(executor.handle());
let request = vec![0; 10 * 1024 * 1024];
let response = EchoRpc::client(&service.handle()).call(server_addr, request.clone());
executor.spawn(service.map_err(|e| panic!("{}", e)));
let result = track_try_unwrap!(track_any_err!(executor.run_future(response)));
assert_eq!(result.ok(), Some(request));
}
#[test]
fn async_works() {
let mut executor = track_try_unwrap!(track_any_err!(InPlaceExecutor::new()));
let server_addr = "127.0.0.1:1922".parse().unwrap();
let server = ServerBuilder::new(server_addr)
.add_call_handler(EchoHandler)
.finish(executor.handle());
executor.spawn(server.map_err(|e| panic!("{}", e)));
let service = ClientServiceBuilder::new().finish(executor.handle());
let service_handle = service.handle();
let request = Vec::from(&b"async"[..]);
let response = EchoRpc::client(&service_handle).call(server_addr, request.clone());
executor.spawn(service.map_err(|e| panic!("{}", e)));
let result = track_try_unwrap!(track_any_err!(executor.run_future(response)));
assert_eq!(result.ok(), Some(request));
let metrics = service_handle
.metrics()
.channels()
.as_map()
.load()
.get(&server_addr)
.cloned()
.unwrap();
assert_eq!(metrics.async_outgoing_messages(), 1);
assert_eq!(metrics.async_incoming_messages(), 1);
}
}