extern crate rmp_serde as rmps;
use std::convert::{From, TryFrom};
use std::fmt::Display;
use std::io::Write;
pub const RPC_HEADER_LEN: usize = 48;
const MAX_RPCNAME_LEN: usize = 36;
pub type Bytes = Vec<u8>;
#[derive(Debug, Clone)]
pub struct Error(String);
impl Error {
pub fn new(s: &'static str) -> Self {
Error(String::from(s))
}
}
impl std::fmt::Display for Error {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
impl<E: Display + std::error::Error> From<E> for Error {
#[inline]
fn from(e: E) -> Self {
Error(format!("{}", e))
}
}
#[derive(Clone, Copy, Debug)]
pub enum Mode {
Request = 1,
Respond = 2,
Subcribe = 3,
Publish = 4,
NotFound = 5,
NotMatch = 6,
HeartBeat = 7,
NoAccess = 8,
}
impl Default for Mode {
fn default() -> Self {
Mode::NoAccess
}
}
impl TryFrom<u8> for Mode {
type Error = Error;
fn try_from(v: u8) -> Result<Self, Self::Error> {
match v {
1 => Ok(Mode::Request),
2 => Ok(Mode::Respond),
3 => Ok(Mode::Subcribe),
4 => Ok(Mode::Publish),
5 => Ok(Mode::NotFound),
6 => Ok(Mode::NotMatch),
7 => Ok(Mode::HeartBeat),
8 => Ok(Mode::NoAccess),
_ => Err(Error::new("未知的消息类型")),
}
}
}
impl TryFrom<Mode> for Error {
type Error = ();
#[inline]
fn try_from(m: Mode) -> Result<Self, Self::Error> {
match m {
Mode::NotFound => Ok(Error::new("没有找到相应的函数")),
Mode::NotMatch => Ok(Error::new("函数参数不匹配")),
Mode::NoAccess => Ok(Error::new("没有权限")),
_ => Err(()),
}
}
}
#[derive(Debug)]
pub struct Msg {
ver: u8, id: u32, name: [u8; MAX_RPCNAME_LEN], mode: Mode, crc32: u32, buf: Vec<u8>, }
impl Default for Msg {
fn default() -> Self {
Msg { ver: 0, id: 0, name: [0u8; MAX_RPCNAME_LEN], mode: Default::default(), crc32: 0, buf: Default::default() }
}
}
impl Msg {
pub fn new(id: u32, name: &str) -> Self {
if name.len() >= MAX_RPCNAME_LEN {
panic!("the length of name must less than {}", MAX_RPCNAME_LEN);
}
let mut msg = Msg::default();
msg.id = id;
name.as_bytes().iter().enumerate().for_each(|(i, v)| {
msg.name[i] = *v;
});
msg
}
#[inline]
fn u32_to_array(v: u32) -> [u8; 4] {
[v as u8, (v >> 8) as u8, (v >> 16) as u8, (v >> 24) as u8]
}
#[inline]
fn slice_to_u32(slice: &[u8]) -> u32 {
let p = &slice[0] as *const u8 as *const u32;
unsafe { *p }
}
#[allow(unused_must_use)]
fn encode_header(&mut self, len: usize) {
let mode = self.mode as u32;
let lvm: u32 = ((len & 0x3ffffff) as u32) | ((self.ver as u32) << 26) | (mode << 28);
self.buf.write(&Msg::u32_to_array(lvm)[..]);
self.buf.write(&Msg::u32_to_array(self.id)[..]);
self.buf.write(&self.name[..]);
self.buf.write(&Msg::u32_to_array(self.crc32)[..]);
}
pub fn encode_without_body(mut self, mode: Mode) -> Bytes {
self.buf.clear();
self.buf.reserve(RPC_HEADER_LEN);
self.mode = mode;
self.encode_header(RPC_HEADER_LEN);
self.buf
}
#[allow(unused_must_use)]
pub fn encode_with_bytes(mut self, mode: Mode, data: &Bytes) -> Bytes {
self.buf.clear();
self.buf.reserve(RPC_HEADER_LEN + data.len());
self.mode = mode;
self.encode_header(RPC_HEADER_LEN + data.len());
self.buf.write(&data[..]);
self.buf
}
#[allow(unused_must_use)]
pub fn encode<Body: serde::ser::Serialize>(mut self, mode: Mode, body: &Body) -> Bytes {
self.buf.clear();
self.buf.reserve(128);
self.mode = mode;
self.encode_header(0);
rmps::encode::write(&mut self.buf, body);
let lvm = Msg::slice_to_u32(&self.buf[0..4]) | ((self.buf.len() & 0x3ffffff) as u32);
self.buf[0..4].copy_from_slice(&Msg::u32_to_array(lvm)[..]);
self.buf
}
pub fn decode(header: &[u8]) -> Result<Self, Error> {
let lvm = Msg::slice_to_u32(&header[0..4]);
let len = (lvm & 0x3ffffff) as usize;
let mut msg = Msg {
ver: ((lvm >> 26) & 0x3) as u8,
mode: Mode::try_from((lvm >> 28) as u8)?,
id: Msg::slice_to_u32(&header[4..8]),
crc32: Msg::slice_to_u32(&header[8 + MAX_RPCNAME_LEN..]),
buf: if len > RPC_HEADER_LEN { vec![0u8; len - RPC_HEADER_LEN] } else { Vec::new() },
..Default::default()
};
msg.name.copy_from_slice(&header[8..8 + MAX_RPCNAME_LEN]);
Ok(msg)
}
pub fn id(&self) -> u32 {
self.id
}
pub fn name(&self) -> &[u8] {
for i in 0..MAX_RPCNAME_LEN {
if self.name[i] == 0 {
return &self.name[..i];
}
}
&self.name[..]
}
pub fn mode(&self) -> Mode {
self.mode
}
pub fn body(&mut self) -> Option<&mut [u8]> {
if self.buf.is_empty() {
None
} else {
Some(&mut self.buf[..])
}
}
#[inline]
pub fn headeronly(&self) -> bool {
return self.buf.is_empty();
}
#[inline]
pub fn parse<Args: for<'a> serde::de::Deserialize<'a>>(&self) -> Result<Args, Error> {
let body = rmps::decode::from_slice::<Args>(&self.buf[..])?;
Ok(body)
}
}