reqrep.rs 10 KB

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