protocol_event.rs 9.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2023 Dyne.org foundation
  4. *
  5. * This program is free software: you can redistribute it and/or modify
  6. * it under the terms of the GNU Affero General Public License as
  7. * published by the Free Software Foundation, either version 3 of the
  8. * License, or (at your option) any later version.
  9. *
  10. * This program is distributed in the hope that it will be useful,
  11. * but WITHOUT ANY WARRANTY; without even the implied warranty of
  12. * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
  13. * GNU Affero General Public License for more details.
  14. *
  15. * You should have received a copy of the GNU Affero General Public License
  16. * along with this program. If not, see <https://www.gnu.org/licenses/>.
  17. */
  18. use std::fmt::Debug;
  19. use async_std::sync::{Arc, Mutex};
  20. use async_trait::async_trait;
  21. use darkfi_serial::{Decodable, Encodable, SerialDecodable, SerialEncodable};
  22. use log::{debug, info};
  23. use super::EventMsg;
  24. use crate::{
  25. event_graph::model::{Event, EventId, ModelPtr},
  26. net,
  27. util::{async_util::sleep, ringbuffer::RingBuffer},
  28. Result,
  29. };
  30. const SIZE_OF_SEEN_BUFFER: usize = 65536;
  31. #[derive(SerialEncodable, SerialDecodable, Clone, Debug, PartialEq, Eq, Hash)]
  32. pub struct InvItem {
  33. pub hash: EventId,
  34. }
  35. #[derive(SerialDecodable, SerialEncodable, Clone, Debug)]
  36. pub struct Inv {
  37. pub invs: Vec<InvItem>,
  38. }
  39. #[derive(SerialDecodable, SerialEncodable, Clone, Debug)]
  40. struct SyncEvent {
  41. leaves: Vec<EventId>,
  42. }
  43. #[derive(SerialDecodable, SerialEncodable, Clone, Debug)]
  44. struct GetData {
  45. events: Vec<EventId>,
  46. }
  47. pub type SeenPtr<T> = Arc<Seen<T>>;
  48. pub struct Seen<T> {
  49. seen: Mutex<RingBuffer<T>>,
  50. }
  51. impl<T: Eq + PartialEq + Clone> Seen<T> {
  52. pub fn new() -> SeenPtr<T> {
  53. Arc::new(Self { seen: Mutex::new(RingBuffer::new(SIZE_OF_SEEN_BUFFER)) })
  54. }
  55. pub async fn push(&self, item: &T) -> bool {
  56. let seen = &mut self.seen.lock().await;
  57. if !seen.contains(item) {
  58. seen.push(item.clone());
  59. return true
  60. }
  61. false
  62. }
  63. }
  64. pub struct ProtocolEvent<T>
  65. where
  66. T: Send + Sync + Encodable + Decodable + Debug + 'static,
  67. {
  68. jobsman: net::ProtocolJobsManagerPtr,
  69. event_sub: net::MessageSubscription<Event<T>>,
  70. inv_sub: net::MessageSubscription<Inv>,
  71. getdata_sub: net::MessageSubscription<GetData>,
  72. syncevent_sub: net::MessageSubscription<SyncEvent>,
  73. p2p: net::P2pPtr,
  74. channel: net::ChannelPtr,
  75. model: ModelPtr<T>,
  76. seen_event: SeenPtr<EventId>,
  77. seen_inv: SeenPtr<EventId>,
  78. }
  79. impl<T> ProtocolEvent<T>
  80. where
  81. T: Send + Sync + Encodable + Decodable + Clone + EventMsg + Debug + 'static,
  82. {
  83. pub async fn init(
  84. channel: net::ChannelPtr,
  85. p2p: net::P2pPtr,
  86. model: ModelPtr<T>,
  87. seen_event: SeenPtr<EventId>,
  88. seen_inv: SeenPtr<EventId>,
  89. ) -> net::ProtocolBasePtr {
  90. let message_subsytem = channel.get_message_subsystem();
  91. message_subsytem.add_dispatch::<Event<T>>().await;
  92. message_subsytem.add_dispatch::<Inv>().await;
  93. message_subsytem.add_dispatch::<GetData>().await;
  94. message_subsytem.add_dispatch::<SyncEvent>().await;
  95. let event_sub =
  96. channel.clone().subscribe_msg::<Event<T>>().await.expect("Missing Event dispatcher!");
  97. let inv_sub = channel.subscribe_msg::<Inv>().await.expect("Missing Inv dispatcher!");
  98. let getdata_sub =
  99. channel.clone().subscribe_msg::<GetData>().await.expect("Missing GetData dispatcher!");
  100. let syncevent_sub = channel
  101. .clone()
  102. .subscribe_msg::<SyncEvent>()
  103. .await
  104. .expect("Missing SyncEvent dispatcher!");
  105. Arc::new(Self {
  106. jobsman: net::ProtocolJobsManager::new("ProtocolEvent", channel.clone()),
  107. event_sub,
  108. inv_sub,
  109. getdata_sub,
  110. syncevent_sub,
  111. p2p,
  112. channel,
  113. model,
  114. seen_event,
  115. seen_inv,
  116. })
  117. }
  118. async fn handle_receive_event(self: Arc<Self>) -> Result<()> {
  119. debug!(target: "event_graph", "ProtocolEvent::handle_receive_event() [START]");
  120. let exclude_list = vec![self.channel.address()];
  121. loop {
  122. let event = self.event_sub.receive().await?;
  123. let event = (*event).to_owned();
  124. if !self.seen_event.push(&event.hash()).await {
  125. continue
  126. }
  127. info!("[P2P] Received: {:?}", event.action);
  128. self.new_event(&event).await?;
  129. self.send_inv(&event).await?;
  130. // Broadcast the msg
  131. self.p2p.broadcast_with_exclude(event, &exclude_list).await?;
  132. }
  133. }
  134. async fn handle_receive_inv(self: Arc<Self>) -> Result<()> {
  135. debug!(target: "event_graph", "ProtocolEvent::handle_receive_inv() [START]");
  136. let exclude_list = vec![self.channel.address()];
  137. loop {
  138. let inv = self.inv_sub.receive().await?;
  139. let inv = (*inv).to_owned();
  140. let inv_item = inv.invs[0].clone();
  141. // for inv in inv.invs.iter() {
  142. if !self.seen_inv.push(&inv_item.hash).await {
  143. continue
  144. }
  145. if self.model.lock().await.get_event(&inv_item.hash).is_none() {
  146. self.send_getdata(vec![inv_item.hash]).await?;
  147. }
  148. // }
  149. // Broadcast the inv msg
  150. self.p2p.broadcast_with_exclude(inv, &exclude_list).await?;
  151. }
  152. }
  153. async fn handle_receive_getdata(self: Arc<Self>) -> Result<()> {
  154. debug!(target: "event_graph", "ProtocolEvent::handle_receive_getdata() [START]");
  155. loop {
  156. let getdata = self.getdata_sub.receive().await?;
  157. let events = (*getdata).to_owned().events;
  158. for event_id in events {
  159. let model_event = self.model.lock().await.get_event(&event_id);
  160. if let Some(event) = model_event {
  161. self.channel.send(event).await?;
  162. }
  163. }
  164. }
  165. }
  166. async fn handle_receive_syncevent(self: Arc<Self>) -> Result<()> {
  167. debug!(target: "event_graph", "ProtocolEvent::handle_receive_syncevent() [START]");
  168. loop {
  169. let syncevent = self.syncevent_sub.receive().await?;
  170. let model = self.model.lock().await;
  171. let leaves = model.find_leaves();
  172. if leaves == syncevent.leaves {
  173. continue
  174. }
  175. for leaf in syncevent.leaves.iter() {
  176. if leaves.contains(leaf) {
  177. continue
  178. }
  179. let children = model.get_offspring(leaf);
  180. for child in children {
  181. self.channel.send(child).await?;
  182. }
  183. }
  184. }
  185. }
  186. // every 6 seconds send a SyncEvent msg
  187. async fn send_sync_hash_loop(self: Arc<Self>) -> Result<()> {
  188. debug!(target: "event_graph", "ProtocolEvent::send_sync_hash_loop() [START]");
  189. loop {
  190. sleep(6).await;
  191. let leaves = self.model.lock().await.find_leaves();
  192. self.channel.send(SyncEvent { leaves }).await?;
  193. }
  194. }
  195. async fn new_event(&self, event: &Event<T>) -> Result<()> {
  196. debug!(target: "event_graph", "ProtocolEvent::new_event()");
  197. let mut model = self.model.lock().await;
  198. model.add(event.clone()).await;
  199. Ok(())
  200. }
  201. async fn send_inv(&self, event: &Event<T>) -> Result<()> {
  202. debug!(target: "event_graph", "ProtocolEvent::send_inv()");
  203. self.p2p.broadcast(Inv { invs: vec![InvItem { hash: event.hash() }] }).await?;
  204. Ok(())
  205. }
  206. async fn send_getdata(&self, events: Vec<EventId>) -> Result<()> {
  207. debug!(target: "event_graph", "ProtocolEvent::send_getdata()");
  208. self.channel.send(GetData { events }).await?;
  209. Ok(())
  210. }
  211. }
  212. #[async_trait]
  213. impl<T> net::ProtocolBase for ProtocolEvent<T>
  214. where
  215. T: Send + Sync + Encodable + Decodable + Clone + EventMsg + Debug,
  216. {
  217. async fn start(self: Arc<Self>, executor: Arc<smol::Executor<'_>>) -> Result<()> {
  218. debug!(target: "event_graph", "ProtocolEvent::start() [START]");
  219. self.jobsman.clone().start(executor.clone());
  220. self.jobsman.clone().spawn(self.clone().handle_receive_event(), executor.clone()).await;
  221. self.jobsman.clone().spawn(self.clone().handle_receive_inv(), executor.clone()).await;
  222. self.jobsman.clone().spawn(self.clone().handle_receive_getdata(), executor.clone()).await;
  223. self.jobsman.clone().spawn(self.clone().handle_receive_syncevent(), executor.clone()).await;
  224. self.jobsman.clone().spawn(self.clone().send_sync_hash_loop(), executor.clone()).await;
  225. debug!(target: "event_graph", "ProtocolEvent::start() [END]");
  226. Ok(())
  227. }
  228. fn name(&self) -> &'static str {
  229. "ProtocolEvent"
  230. }
  231. }
  232. impl<T> net::Message for Event<T>
  233. where
  234. T: Send + Sync + Decodable + Encodable + 'static,
  235. {
  236. fn name() -> &'static str {
  237. "event"
  238. }
  239. }
  240. impl net::Message for Inv {
  241. fn name() -> &'static str {
  242. "inv"
  243. }
  244. }
  245. impl net::Message for SyncEvent {
  246. fn name() -> &'static str {
  247. "syncevent"
  248. }
  249. }
  250. impl net::Message for GetData {
  251. fn name() -> &'static str {
  252. "getdata"
  253. }
  254. }