krpc 0.2.0

A asynchronous RPC library(include client and server) which can use easly and communicate by tokio unix/tcp socket
Documentation
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,
	/// 未找到RPC服务
	NotFound = 5,
	/// RPC调用参数不匹配
	NotMatch = 6,
	/// 心跳消息
	HeartBeat = 7,
	/// 无权限调用RPC服务
	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,                     //版本,0:没有校验(本地连接使用),1:增加了CRC32校验(外部连接使用)
	id: u32,                     //消息唯一标识
	name: [u8; MAX_RPCNAME_LEN], //RPC调用函数名或订阅主题
	mode: Mode,                  //消息模式
	crc32: u32,                  //数据部分的CRC32校验值
	buf: Vec<u8>,                //接收消息的buffer
}

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 {
	/// 新建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);
		//重新填充头部4字节lvm的len
		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)[..]);
		//TODO 注意后续还需要根据ver把body的CRC32计算之后再次填充到相应位置
		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)
	}
	/// 获取msg的id
	pub fn id(&self) -> u32 {
		self.id
	}
	/// 获取msg的函数名
	pub fn name(&self) -> &[u8] {
		for i in 0..MAX_RPCNAME_LEN {
			if self.name[i] == 0 {
				return &self.name[..i];
			}
		}
		&self.name[..]
	}
	/// 获取msg的mode
	pub fn mode(&self) -> Mode {
		self.mode
	}
	/// 获取msg的body
	pub fn body(&mut self) -> Option<&mut [u8]> {
		if self.buf.is_empty() {
			None
		} else {
			Some(&mut self.buf[..])
		}
	}
	/// msg是否只有消息头
	#[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)
	}
}