reqrep.rs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415
  1. use async_std::sync::Arc;
  2. use std::convert::TryFrom;
  3. use std::io;
  4. use std::net::SocketAddr;
  5. use crate::serial::{deserialize, serialize};
  6. use crate::{Decodable, Encodable, Result};
  7. use async_executor::Executor;
  8. use bytes::Bytes;
  9. use futures::FutureExt;
  10. use log::*;
  11. use rand::Rng;
  12. use signal_hook::{consts::SIGINT, iterator::Signals};
  13. use zeromq::*;
  14. pub type PeerId = Vec<u8>;
  15. enum NetEvent {
  16. Receive(zeromq::ZmqMessage),
  17. Send((PeerId, Reply)),
  18. Stop,
  19. }
  20. pub fn addr_to_string(addr: SocketAddr) -> String {
  21. format!("tcp://{}", addr.to_string())
  22. }
  23. pub struct RepProtocol {
  24. addr: SocketAddr,
  25. socket: zeromq::RouterSocket,
  26. recv_queue: async_channel::Receiver<(PeerId, Reply)>,
  27. send_queue: async_channel::Sender<(PeerId, Request)>,
  28. channels: (
  29. async_channel::Sender<(PeerId, Reply)>,
  30. async_channel::Receiver<(PeerId, Request)>,
  31. ),
  32. service_name: String,
  33. }
  34. impl RepProtocol {
  35. pub fn new(addr: SocketAddr, service_name: String) -> RepProtocol {
  36. let socket = zeromq::RouterSocket::new();
  37. let (send_queue, recv_channel) = async_channel::unbounded::<(PeerId, Request)>();
  38. let (send_channel, recv_queue) = async_channel::unbounded::<(PeerId, Reply)>();
  39. let channels = (send_channel.clone(), recv_channel.clone());
  40. RepProtocol {
  41. addr,
  42. socket,
  43. recv_queue,
  44. send_queue,
  45. channels,
  46. service_name,
  47. }
  48. }
  49. pub async fn start(
  50. &mut self,
  51. ) -> Result<(
  52. async_channel::Sender<(PeerId, Reply)>,
  53. async_channel::Receiver<(PeerId, Request)>,
  54. )> {
  55. let addr = addr_to_string(self.addr);
  56. self.socket.bind(addr.as_str()).await?;
  57. info!("{} SERVICE: Bound To {}", self.service_name, addr);
  58. Ok(self.channels.clone())
  59. }
  60. pub async fn run(&mut self, executor: Arc<Executor<'_>>) -> Result<()> {
  61. info!("{} SERVICE: Running", self.service_name);
  62. let (stop_s, stop_r) = async_channel::unbounded::<()>();
  63. let mut signals = Signals::new(&[SIGINT])?;
  64. let stop_task = executor.spawn(async move {
  65. for _ in signals.forever() {
  66. stop_s.send(()).await?;
  67. break;
  68. }
  69. Ok::<(), crate::Error>(())
  70. });
  71. loop {
  72. let event = futures::select! {
  73. msg = self.socket.recv().fuse() => NetEvent::Receive(msg?),
  74. msg = self.recv_queue.recv().fuse() => NetEvent::Send(msg?),
  75. _ = stop_r.recv().fuse() => NetEvent::Stop
  76. };
  77. match event {
  78. NetEvent::Receive(msg) => {
  79. if let Some(peer) = msg.get(0) {
  80. if let Some(request) = msg.get(1) {
  81. let request: Vec<u8> = request.to_vec();
  82. let request: Request = deserialize(&request)?;
  83. self.send_queue.send((peer.to_vec(), request)).await?;
  84. }
  85. }
  86. }
  87. NetEvent::Send((peer, reply)) => {
  88. let peer = Bytes::from(peer);
  89. let mut msg: Vec<Bytes> = vec![peer];
  90. let reply: Vec<u8> = serialize(&reply);
  91. let reply = Bytes::from(reply);
  92. msg.push(reply);
  93. let reply = zeromq::ZmqMessage::try_from(msg)
  94. .map_err(|_| crate::Error::TryFromError)?;
  95. self.socket.send(reply).await?;
  96. }
  97. NetEvent::Stop => break,
  98. }
  99. }
  100. let _ = stop_task.cancel().await;
  101. warn!("{} SERVICE: Stopped", self.service_name);
  102. Ok(())
  103. }
  104. }
  105. pub struct ReqProtocol {
  106. addr: SocketAddr,
  107. socket: zeromq::DealerSocket,
  108. service_name: String,
  109. }
  110. impl ReqProtocol {
  111. pub fn new(addr: SocketAddr, service_name: String) -> ReqProtocol {
  112. let socket = zeromq::DealerSocket::new();
  113. ReqProtocol {
  114. addr,
  115. socket,
  116. service_name,
  117. }
  118. }
  119. pub async fn start(&mut self) -> Result<()> {
  120. let addr = addr_to_string(self.addr);
  121. self.socket.connect(addr.as_str()).await?;
  122. info!("{} SERVICE: Connected To {}", self.service_name, self.addr);
  123. Ok(())
  124. }
  125. pub async fn request(
  126. &mut self,
  127. command: u8,
  128. data: Vec<u8>,
  129. handle_error: &dyn Fn(u32),
  130. ) -> Result<Option<Vec<u8>>> {
  131. let request = Request::new(command, data);
  132. let req = serialize(&request);
  133. let req = bytes::Bytes::from(req);
  134. let req: zeromq::ZmqMessage = req.into();
  135. self.socket.send(req).await?;
  136. info!(
  137. "{} SERVICE: Sent Request {{ command: {} }}",
  138. self.service_name, command
  139. );
  140. let rep: zeromq::ZmqMessage = self.socket.recv().await?;
  141. if let Some(reply) = rep.get(0) {
  142. let reply: Vec<u8> = reply.to_vec();
  143. let reply: Reply = deserialize(&reply)?;
  144. info!(
  145. "{} SERVICE: Received Reply {{ error: {} }}",
  146. self.service_name,
  147. reply.has_error()
  148. );
  149. if reply.has_error() {
  150. // TODO return error status code instead of None
  151. // this is temporary
  152. handle_error(reply.get_error());
  153. return Ok(None);
  154. }
  155. assert!(reply.get_id() == request.get_id());
  156. Ok(Some(reply.get_payload()))
  157. } else {
  158. Err(crate::Error::ZmqError(
  159. "Couldn't parse ZmqMessage".to_string(),
  160. ))
  161. }
  162. }
  163. }
  164. pub struct Publisher {
  165. addr: SocketAddr,
  166. socket: zeromq::PubSocket,
  167. service_name: String,
  168. }
  169. impl Publisher {
  170. pub fn new(addr: SocketAddr, service_name: String) -> Publisher {
  171. let socket = zeromq::PubSocket::new();
  172. Publisher {
  173. addr,
  174. socket,
  175. service_name,
  176. }
  177. }
  178. pub async fn start(&mut self, recv_queue: async_channel::Receiver<Vec<u8>>) -> Result<()> {
  179. let addr = addr_to_string(self.addr);
  180. self.socket.bind(addr.as_str()).await?;
  181. info!(
  182. "{} PUBLISHER SERVICE : Bound To {}",
  183. self.service_name, addr
  184. );
  185. loop {
  186. let msg = recv_queue.recv().await?;
  187. self.publish(msg).await?;
  188. }
  189. }
  190. async fn publish(&mut self, data: Vec<u8>) -> Result<()> {
  191. let data = Bytes::from(data);
  192. self.socket.send(data.into()).await?;
  193. Ok(())
  194. }
  195. }
  196. pub struct Subscriber {
  197. addr: SocketAddr,
  198. socket: zeromq::SubSocket,
  199. service_name: String,
  200. }
  201. impl Subscriber {
  202. pub fn new(addr: SocketAddr, service_name: String) -> Subscriber {
  203. let socket = zeromq::SubSocket::new();
  204. Subscriber {
  205. addr,
  206. socket,
  207. service_name,
  208. }
  209. }
  210. pub async fn start(&mut self) -> Result<()> {
  211. let addr = addr_to_string(self.addr);
  212. self.socket.connect(addr.as_str()).await?;
  213. self.socket.subscribe("").await?;
  214. info!(
  215. "{} SUBSCRIBER SERVICE : Connected To {}",
  216. self.service_name, addr
  217. );
  218. Ok(())
  219. }
  220. pub async fn fetch<T: Decodable>(&mut self) -> Result<T> {
  221. let data = self.socket.recv().await?;
  222. match data.get(0) {
  223. Some(d) => {
  224. let data = d.to_vec();
  225. let data: T = deserialize(&data)?;
  226. Ok(data)
  227. }
  228. None => Err(crate::Error::ZmqError(
  229. "Couldn't parse ZmqMessage".to_string(),
  230. )),
  231. }
  232. }
  233. }
  234. #[derive(Debug, PartialEq)]
  235. pub struct Request {
  236. command: u8,
  237. id: u32,
  238. payload: Vec<u8>,
  239. }
  240. impl Request {
  241. pub fn new(command: u8, payload: Vec<u8>) -> Request {
  242. let id = Self::gen_id();
  243. Request {
  244. command,
  245. id,
  246. payload,
  247. }
  248. }
  249. fn gen_id() -> u32 {
  250. let mut rng = rand::thread_rng();
  251. rng.gen()
  252. }
  253. pub fn get_id(&self) -> u32 {
  254. self.id
  255. }
  256. pub fn get_command(&self) -> u8 {
  257. self.command
  258. }
  259. pub fn get_payload(&self) -> Vec<u8> {
  260. self.payload.clone()
  261. }
  262. }
  263. #[derive(Debug, PartialEq)]
  264. pub struct Reply {
  265. id: u32,
  266. error: u32,
  267. payload: Vec<u8>,
  268. }
  269. impl Reply {
  270. pub fn from(request: &Request, error: u32, payload: Vec<u8>) -> Reply {
  271. Reply {
  272. id: request.get_id(),
  273. error,
  274. payload,
  275. }
  276. }
  277. pub fn has_error(&self) -> bool {
  278. if self.error == 0 {
  279. false
  280. } else {
  281. true
  282. }
  283. }
  284. pub fn get_error(&self) -> u32 {
  285. self.error
  286. }
  287. pub fn get_payload(&self) -> Vec<u8> {
  288. self.payload.clone()
  289. }
  290. pub fn set_payload(&mut self, payload: Vec<u8>) {
  291. self.payload = payload;
  292. }
  293. pub fn set_error(&mut self, error: u32) {
  294. self.error = error;
  295. }
  296. pub fn get_id(&self) -> u32 {
  297. self.id
  298. }
  299. }
  300. impl Encodable for Request {
  301. fn encode<S: io::Write>(&self, mut s: S) -> Result<usize> {
  302. let mut len = 0;
  303. len += self.command.encode(&mut s)?;
  304. len += self.id.encode(&mut s)?;
  305. len += self.payload.encode(&mut s)?;
  306. Ok(len)
  307. }
  308. }
  309. impl Encodable for Reply {
  310. fn encode<S: io::Write>(&self, mut s: S) -> Result<usize> {
  311. let mut len = 0;
  312. len += self.id.encode(&mut s)?;
  313. len += self.error.encode(&mut s)?;
  314. len += self.payload.encode(&mut s)?;
  315. Ok(len)
  316. }
  317. }
  318. impl Decodable for Request {
  319. fn decode<D: io::Read>(mut d: D) -> Result<Self> {
  320. Ok(Self {
  321. command: Decodable::decode(&mut d)?,
  322. id: Decodable::decode(&mut d)?,
  323. payload: Decodable::decode(&mut d)?,
  324. })
  325. }
  326. }
  327. impl Decodable for Reply {
  328. fn decode<D: io::Read>(mut d: D) -> Result<Self> {
  329. Ok(Self {
  330. id: Decodable::decode(&mut d)?,
  331. error: Decodable::decode(&mut d)?,
  332. payload: Decodable::decode(&mut d)?,
  333. })
  334. }
  335. }
  336. #[cfg(test)]
  337. mod tests {
  338. use super::{Reply, Request, Result};
  339. use crate::serial::{deserialize, serialize};
  340. #[test]
  341. fn serialize_and_deserialize_request_test() {
  342. let request = Request::new(2, vec![2, 3, 4, 6, 4]);
  343. let serialized_request = serialize(&request);
  344. assert!((deserialize(&serialized_request) as Result<bool>).is_err());
  345. let deserialized_request = deserialize(&serialized_request).ok();
  346. assert_eq!(deserialized_request, Some(request));
  347. }
  348. #[test]
  349. fn serialize_and_deserialize_reply_test() {
  350. let request = Request::new(2, vec![2, 3, 4, 6, 4]);
  351. let reply = Reply::from(&request, 0, vec![2, 3, 4, 6, 4]);
  352. let serialized_reply = serialize(&reply);
  353. assert!((deserialize(&serialized_reply) as Result<bool>).is_err());
  354. let deserialized_reply = deserialize(&serialized_reply).ok();
  355. assert_eq!(deserialized_reply, Some(reply));
  356. }
  357. }