protocol.rs 8.0 KB

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