protocol.rs 9.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270
  1. use async_executor::Executor;
  2. use async_std::sync::Arc;
  3. use async_trait::async_trait;
  4. use chrono::Utc;
  5. use log::{debug, error};
  6. use crate::{
  7. net::{
  8. ChannelPtr, MessageSubscription, P2pPtr, ProtocolBase, ProtocolBasePtr,
  9. ProtocolJobsManager, ProtocolJobsManagerPtr,
  10. },
  11. Result,
  12. };
  13. use super::{
  14. dht::DhtPtr,
  15. messages::{KeyRequest, KeyResponse, LookupMapRequest, LookupMapResponse, LookupRequest},
  16. };
  17. pub struct Protocol {
  18. channel: ChannelPtr,
  19. notify_queue_sender: async_channel::Sender<KeyResponse>,
  20. req_sub: MessageSubscription<KeyRequest>,
  21. resp_sub: MessageSubscription<KeyResponse>,
  22. lookup_sub: MessageSubscription<LookupRequest>,
  23. lookup_map_sub: MessageSubscription<LookupMapRequest>,
  24. jobsman: ProtocolJobsManagerPtr,
  25. dht: DhtPtr,
  26. p2p: P2pPtr,
  27. }
  28. impl Protocol {
  29. pub async fn init(
  30. channel: ChannelPtr,
  31. notify_queue_sender: async_channel::Sender<KeyResponse>,
  32. dht: DhtPtr,
  33. p2p: P2pPtr,
  34. ) -> Result<ProtocolBasePtr> {
  35. debug!("Adding Protocol to the protocol registry");
  36. let msg_subsystem = channel.get_message_subsystem();
  37. msg_subsystem.add_dispatch::<KeyRequest>().await;
  38. msg_subsystem.add_dispatch::<KeyResponse>().await;
  39. msg_subsystem.add_dispatch::<LookupRequest>().await;
  40. msg_subsystem.add_dispatch::<LookupMapRequest>().await;
  41. let req_sub = channel.subscribe_msg::<KeyRequest>().await?;
  42. let resp_sub = channel.subscribe_msg::<KeyResponse>().await?;
  43. let lookup_sub = channel.subscribe_msg::<LookupRequest>().await?;
  44. let lookup_map_sub = channel.subscribe_msg::<LookupMapRequest>().await?;
  45. Ok(Arc::new(Self {
  46. channel: channel.clone(),
  47. notify_queue_sender,
  48. req_sub,
  49. resp_sub,
  50. lookup_sub,
  51. lookup_map_sub,
  52. jobsman: ProtocolJobsManager::new("Protocol", channel),
  53. dht,
  54. p2p,
  55. }))
  56. }
  57. async fn handle_receive_request(self: Arc<Self>) -> Result<()> {
  58. debug!("Protocol::handle_receive_request() [START]");
  59. let exclude_list = vec![self.channel.address()];
  60. loop {
  61. let req = match self.req_sub.receive().await {
  62. Ok(v) => v,
  63. Err(e) => {
  64. error!("Protocol::handle_receive_request(): recv fail: {}", e);
  65. continue
  66. }
  67. };
  68. let req_copy = (*req).clone();
  69. debug!("Protocol::handle_receive_request(): req: {:?}", req_copy);
  70. {
  71. let dht = &mut self.dht.write().await;
  72. if dht.seen.contains_key(&req_copy.id) {
  73. debug!(
  74. "Protocol::handle_receive_request(): We have already seen this request."
  75. );
  76. continue
  77. }
  78. dht.seen.insert(req_copy.id.clone(), Utc::now().timestamp());
  79. }
  80. let daemon = self.dht.read().await.id;
  81. if daemon != req_copy.to {
  82. if let Err(e) =
  83. self.p2p.broadcast_with_exclude(req_copy.clone(), &exclude_list).await
  84. {
  85. error!("Protocol::handle_receive_response(): p2p broadcast fail: {}", e);
  86. };
  87. continue
  88. }
  89. match self.dht.read().await.map.get(&req_copy.key) {
  90. Some(value) => {
  91. let response =
  92. KeyResponse::new(daemon, req_copy.from, req_copy.key, value.clone());
  93. debug!("Protocol::handle_receive_request(): sending response: {:?}", response);
  94. if let Err(e) = self.channel.send(response).await {
  95. error!("Protocol::handle_receive_request(): p2p broadcast of response failed: {}", e);
  96. };
  97. }
  98. None => {
  99. error!("Protocol::handle_receive_request(): Requested key doesn't exist locally: {}", req_copy.key);
  100. }
  101. }
  102. }
  103. }
  104. async fn handle_receive_response(self: Arc<Self>) -> Result<()> {
  105. debug!("Protocol::handle_receive_response() [START]");
  106. let exclude_list = vec![self.channel.address()];
  107. loop {
  108. let resp = match self.resp_sub.receive().await {
  109. Ok(v) => v,
  110. Err(e) => {
  111. error!("Protocol::handle_receive_response(): recv fail: {}", e);
  112. continue
  113. }
  114. };
  115. let resp_copy = (*resp).clone();
  116. debug!("Protocol::handle_receive_response(): resp: {:?}", resp_copy);
  117. {
  118. let dht = &mut self.dht.write().await;
  119. if dht.seen.contains_key(&resp_copy.id) {
  120. debug!(
  121. "Protocol::handle_receive_request(): We have already seen this request."
  122. );
  123. continue
  124. }
  125. dht.seen.insert(resp_copy.id.clone(), Utc::now().timestamp());
  126. }
  127. if self.dht.read().await.id != resp_copy.to {
  128. if let Err(e) =
  129. self.p2p.broadcast_with_exclude(resp_copy.clone(), &exclude_list).await
  130. {
  131. error!("Protocol::handle_receive_response(): p2p broadcast fail: {}", e);
  132. };
  133. continue
  134. }
  135. self.notify_queue_sender.send(resp_copy.clone()).await?;
  136. }
  137. }
  138. async fn handle_receive_lookup_request(self: Arc<Self>) -> Result<()> {
  139. debug!("Protocol::handle_receive_lookup_request() [START]");
  140. let exclude_list = vec![self.channel.address()];
  141. loop {
  142. let req = match self.lookup_sub.receive().await {
  143. Ok(v) => v,
  144. Err(e) => {
  145. error!("Protocol::handle_receive_lookup_request(): recv fail: {}", e);
  146. continue
  147. }
  148. };
  149. let req_copy = (*req).clone();
  150. debug!("Protocol::handle_receive_lookup_request(): req: {:?}", req_copy);
  151. if !(0..=1).contains(&req_copy.req_type) {
  152. debug!("Protocol::handle_receive_lookup_request(): Unknown request type.");
  153. continue
  154. }
  155. {
  156. let dht = &mut self.dht.write().await;
  157. if dht.seen.contains_key(&req_copy.id) {
  158. debug!(
  159. "Protocol::handle_receive_request(): We have already seen this request."
  160. );
  161. continue
  162. }
  163. dht.seen.insert(req_copy.id.clone(), Utc::now().timestamp());
  164. }
  165. let result = match req_copy.req_type {
  166. 0 => self
  167. .dht
  168. .write()
  169. .await
  170. .lookup_insert(req_copy.key.clone(), req_copy.daemon.clone()),
  171. _ => self
  172. .dht
  173. .write()
  174. .await
  175. .lookup_remove(req_copy.key.clone(), req_copy.daemon.clone()),
  176. };
  177. if let Err(e) = result {
  178. error!("Protocol::handle_receive_lookup_request(): request action failed: {}", e);
  179. continue
  180. };
  181. if let Err(e) = self.p2p.broadcast_with_exclude(req_copy, &exclude_list).await {
  182. error!("Protocol::handle_receive_lookup_request(): p2p broadcast fail: {}", e);
  183. };
  184. }
  185. }
  186. async fn handle_receive_lookup_map_request(self: Arc<Self>) -> Result<()> {
  187. debug!("Protocol::handle_receive_lookup_map_request() [START]");
  188. loop {
  189. let req = match self.lookup_map_sub.receive().await {
  190. Ok(v) => v,
  191. Err(e) => {
  192. error!("Protocol::handle_receive_lookup_map_request(): recv fail: {}", e);
  193. continue
  194. }
  195. };
  196. debug!("Protocol::handle_receive_lookup_map_request(): req: {:?}", req);
  197. {
  198. let dht = &mut self.dht.write().await;
  199. if dht.seen.contains_key(&req.id) {
  200. debug!(
  201. "Protocol::handle_receive_lookup_map_request(): We have already seen this request."
  202. );
  203. continue
  204. }
  205. dht.seen.insert(req.id.clone(), Utc::now().timestamp());
  206. }
  207. // Extra validations can be added here.
  208. let lookup = self.dht.read().await.lookup.clone();
  209. let response = LookupMapResponse::new(lookup);
  210. if let Err(e) = self.channel.send(response).await {
  211. error!("Protocol::handle_receive_lookup_map_request() channel send fail: {}", e);
  212. };
  213. }
  214. }
  215. }
  216. #[async_trait]
  217. impl ProtocolBase for Protocol {
  218. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
  219. debug!("Protocol::start() [START]");
  220. self.jobsman.clone().start(executor.clone());
  221. self.jobsman.clone().spawn(self.clone().handle_receive_request(), executor.clone()).await;
  222. self.jobsman.clone().spawn(self.clone().handle_receive_response(), executor.clone()).await;
  223. self.jobsman
  224. .clone()
  225. .spawn(self.clone().handle_receive_lookup_request(), executor.clone())
  226. .await;
  227. self.jobsman
  228. .clone()
  229. .spawn(self.clone().handle_receive_lookup_map_request(), executor.clone())
  230. .await;
  231. debug!("Protocol::start() [END]");
  232. Ok(())
  233. }
  234. fn name(&self) -> &'static str {
  235. "Protocol"
  236. }
  237. }