reqrep.rs 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165
  1. use std::io;
  2. use crate::{Decodable, Encodable, Result};
  3. use async_zmq::zmq;
  4. use rand::Rng;
  5. pub struct ReqRepAPI;
  6. impl ReqRepAPI {
  7. pub async fn start() {
  8. let context = zmq::Context::new();
  9. let frontend = context.socket(zmq::ROUTER).unwrap();
  10. let backend = context.socket(zmq::DEALER).unwrap();
  11. frontend
  12. .bind("tcp://127.0.0.1:3333")
  13. .expect("failed binding frontend");
  14. backend
  15. .bind("tcp://127.0.0.1:4444")
  16. .expect("failed binding backend");
  17. loop {
  18. let mut items = [
  19. frontend.as_poll_item(zmq::POLLIN),
  20. backend.as_poll_item(zmq::POLLIN),
  21. ];
  22. zmq::poll(&mut items, -1).unwrap();
  23. if items[0].is_readable() {
  24. loop {
  25. let message = frontend.recv_msg(0).unwrap();
  26. let more = message.get_more();
  27. backend
  28. .send(message, if more { zmq::SNDMORE } else { 0 })
  29. .unwrap();
  30. if !more {
  31. break;
  32. }
  33. }
  34. }
  35. if items[1].is_readable() {
  36. loop {
  37. let message = backend.recv_msg(0).unwrap();
  38. let more = message.get_more();
  39. frontend
  40. .send(message, if more { zmq::SNDMORE } else { 0 })
  41. .unwrap();
  42. if !more {
  43. break;
  44. }
  45. }
  46. }
  47. }
  48. }
  49. }
  50. #[derive(Debug, PartialEq)]
  51. pub struct Request {
  52. command: u8,
  53. id: u32,
  54. payload: Vec<u8>,
  55. }
  56. impl Request {
  57. pub fn new(command: u8, payload: Vec<u8>) -> Request {
  58. let id = Self::gen_id();
  59. Request {
  60. command,
  61. id,
  62. payload,
  63. }
  64. }
  65. fn gen_id() -> u32 {
  66. let mut rng = rand::thread_rng();
  67. rng.gen()
  68. }
  69. pub fn get_id(&self) -> u32 {
  70. self.id
  71. }
  72. }
  73. #[derive(Debug, PartialEq)]
  74. pub struct Reply {
  75. id: u32,
  76. error: u32,
  77. payload: Vec<u8>,
  78. }
  79. impl Reply {
  80. pub fn from(request: &Request, error: u32, payload: Vec<u8>) -> Reply {
  81. Reply {
  82. id: request.get_id(),
  83. error,
  84. payload,
  85. }
  86. }
  87. }
  88. impl Encodable for Request {
  89. fn encode<S: io::Write>(&self, mut s: S) -> Result<usize> {
  90. let mut len = 0;
  91. len += self.command.encode(&mut s)?;
  92. len += self.id.encode(&mut s)?;
  93. len += self.payload.encode(&mut s)?;
  94. Ok(len)
  95. }
  96. }
  97. impl Encodable for Reply {
  98. fn encode<S: io::Write>(&self, mut s: S) -> Result<usize> {
  99. let mut len = 0;
  100. len += self.id.encode(&mut s)?;
  101. len += self.error.encode(&mut s)?;
  102. len += self.payload.encode(&mut s)?;
  103. Ok(len)
  104. }
  105. }
  106. impl Decodable for Request {
  107. fn decode<D: io::Read>(mut d: D) -> Result<Self> {
  108. Ok(Self {
  109. command: Decodable::decode(&mut d)?,
  110. id: Decodable::decode(&mut d)?,
  111. payload: Decodable::decode(&mut d)?,
  112. })
  113. }
  114. }
  115. impl Decodable for Reply {
  116. fn decode<D: io::Read>(mut d: D) -> Result<Self> {
  117. Ok(Self {
  118. id: Decodable::decode(&mut d)?,
  119. error: Decodable::decode(&mut d)?,
  120. payload: Decodable::decode(&mut d)?,
  121. })
  122. }
  123. }
  124. #[cfg(test)]
  125. mod tests {
  126. use super::{Reply, Request, Result};
  127. use crate::serial::{deserialize, serialize};
  128. #[test]
  129. fn serialize_and_deserialize_request_test() {
  130. let request = Request::new(2, vec![2, 3, 4, 6, 4]);
  131. let serialized_request = serialize(&request);
  132. assert!((deserialize(&serialized_request) as Result<bool>).is_err());
  133. let deserialized_request = deserialize(&serialized_request).ok();
  134. assert_eq!(deserialized_request, Some(request));
  135. }
  136. #[test]
  137. fn serialize_and_deserialize_reply_test() {
  138. let request = Request::new(2, vec![2, 3, 4, 6, 4]);
  139. let reply = Reply::from(&request, 0, vec![2, 3, 4, 6, 4]);
  140. let serialized_reply = serialize(&reply);
  141. assert!((deserialize(&serialized_reply) as Result<bool>).is_err());
  142. let deserialized_reply = deserialize(&serialized_reply).ok();
  143. assert_eq!(deserialized_reply, Some(reply));
  144. }
  145. }