protocol_event.rs 9.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294
  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, sync::Arc};
  19. use async_trait::async_trait;
  20. use darkfi_serial::{Decodable, Encodable, SerialDecodable, SerialEncodable};
  21. use log::debug;
  22. use smol::lock::Mutex;
  23. use super::EventMsg;
  24. use crate::{
  25. event_graph::model::{Event, EventId, ModelPtr},
  26. impl_p2p_message, net,
  27. net::Message,
  28. system::sleep,
  29. util::ringbuffer::RingBuffer,
  30. Result,
  31. };
  32. const SIZE_OF_SEEN_BUFFER: usize = 65536;
  33. #[derive(SerialEncodable, SerialDecodable, Clone, Debug, PartialEq, Eq, Hash)]
  34. pub struct InvItem {
  35. pub hash: EventId,
  36. }
  37. #[derive(SerialDecodable, SerialEncodable, Clone, Debug)]
  38. pub struct Inv {
  39. pub invs: Vec<InvItem>,
  40. }
  41. impl_p2p_message!(Inv, "inv");
  42. #[derive(SerialDecodable, SerialEncodable, Clone, Debug)]
  43. struct SyncEvent {
  44. leaves: Vec<EventId>,
  45. }
  46. impl_p2p_message!(SyncEvent, "syncevent");
  47. #[derive(SerialDecodable, SerialEncodable, Clone, Debug)]
  48. struct GetData {
  49. events: Vec<EventId>,
  50. }
  51. impl_p2p_message!(GetData, "getdata");
  52. pub type SeenPtr<T> = Arc<Seen<T>>;
  53. pub struct Seen<T> {
  54. seen: Mutex<RingBuffer<T, SIZE_OF_SEEN_BUFFER>>,
  55. }
  56. impl<T: Send + Sync + Eq + PartialEq + Clone> Seen<T> {
  57. pub fn new() -> SeenPtr<T> {
  58. Arc::new(Self { seen: Mutex::new(RingBuffer::new()) })
  59. }
  60. pub async fn push(&self, item: &T) -> bool {
  61. let seen = &mut self.seen.lock().await;
  62. if !seen.contains(item) {
  63. seen.push(item.clone());
  64. return true
  65. }
  66. false
  67. }
  68. }
  69. pub struct ProtocolEvent<T>
  70. where
  71. T: Send + Sync + Encodable + Decodable + Debug + 'static,
  72. {
  73. jobsman: net::ProtocolJobsManagerPtr,
  74. event_sub: net::MessageSubscription<Event<T>>,
  75. inv_sub: net::MessageSubscription<Inv>,
  76. getdata_sub: net::MessageSubscription<GetData>,
  77. syncevent_sub: net::MessageSubscription<SyncEvent>,
  78. p2p: net::P2pPtr,
  79. channel: net::ChannelPtr,
  80. model: ModelPtr<T>,
  81. seen_event: SeenPtr<EventId>,
  82. seen_inv: SeenPtr<EventId>,
  83. }
  84. impl<T> ProtocolEvent<T>
  85. where
  86. T: Send + Sync + Encodable + Decodable + Clone + EventMsg + Debug + 'static,
  87. {
  88. pub async fn init(
  89. channel: net::ChannelPtr,
  90. p2p: net::P2pPtr,
  91. model: ModelPtr<T>,
  92. seen_event: SeenPtr<EventId>,
  93. seen_inv: SeenPtr<EventId>,
  94. ) -> net::ProtocolBasePtr {
  95. let message_subsytem = channel.message_subsystem();
  96. message_subsytem.add_dispatch::<Event<T>>().await;
  97. message_subsytem.add_dispatch::<Inv>().await;
  98. message_subsytem.add_dispatch::<GetData>().await;
  99. message_subsytem.add_dispatch::<SyncEvent>().await;
  100. let event_sub =
  101. channel.clone().subscribe_msg::<Event<T>>().await.expect("Missing Event dispatcher!");
  102. let inv_sub = channel.subscribe_msg::<Inv>().await.expect("Missing Inv dispatcher!");
  103. let getdata_sub =
  104. channel.clone().subscribe_msg::<GetData>().await.expect("Missing GetData dispatcher!");
  105. let syncevent_sub = channel
  106. .clone()
  107. .subscribe_msg::<SyncEvent>()
  108. .await
  109. .expect("Missing SyncEvent dispatcher!");
  110. Arc::new(Self {
  111. jobsman: net::ProtocolJobsManager::new("ProtocolEvent", channel.clone()),
  112. event_sub,
  113. inv_sub,
  114. getdata_sub,
  115. syncevent_sub,
  116. p2p,
  117. channel,
  118. model,
  119. seen_event,
  120. seen_inv,
  121. })
  122. }
  123. // Receives an event, checks if we already have it, if not, add it
  124. // to tree, broadcast an inventory of the event and rebroadcast
  125. // the event itself.
  126. async fn handle_receive_event(self: Arc<Self>) -> Result<()> {
  127. debug!(target: "event_graph", "ProtocolEvent::handle_receive_event() [START]");
  128. let exclude_list = vec![self.channel.address().clone()];
  129. loop {
  130. let event = self.event_sub.receive().await?;
  131. let event = (*event).to_owned();
  132. if !self.seen_event.push(&event.hash()).await {
  133. continue
  134. }
  135. debug!("[P2P] Received: {:?}", event.action);
  136. self.new_event(&event).await?;
  137. self.send_inv(&event).await?;
  138. // Broadcast the msg
  139. self.p2p.broadcast_with_exclude(&event, &exclude_list).await;
  140. }
  141. }
  142. // Receives an inventory msg, checks if we already have it, if not,
  143. // and if hash in inv is not in tree then ask for the event,
  144. // and then rebroadcast the inv anyways.
  145. async fn handle_receive_inv(self: Arc<Self>) -> Result<()> {
  146. debug!(target: "event_graph", "ProtocolEvent::handle_receive_inv() [START]");
  147. let exclude_list = vec![self.channel.address().clone()];
  148. loop {
  149. let inv = self.inv_sub.receive().await?;
  150. let inv = (*inv).to_owned();
  151. let inv_item = inv.invs[0].clone();
  152. if !self.seen_inv.push(&inv_item.hash).await {
  153. continue
  154. }
  155. if self.model.lock().await.get_event(&inv_item.hash).is_none() {
  156. self.send_getdata(vec![inv_item.hash]).await?;
  157. }
  158. self.p2p.broadcast_with_exclude(&inv, &exclude_list).await;
  159. }
  160. }
  161. // Receives getdata msg, retrieve event data of contained eventID
  162. // from our tree and sends it.
  163. async fn handle_receive_getdata(self: Arc<Self>) -> Result<()> {
  164. debug!(target: "event_graph", "ProtocolEvent::handle_receive_getdata() [START]");
  165. loop {
  166. let getdata = self.getdata_sub.receive().await?;
  167. let events = (*getdata).to_owned().events;
  168. for event_id in events {
  169. let model_event = self.model.lock().await.get_event(&event_id);
  170. if let Some(event) = model_event {
  171. self.channel.send(&event).await?;
  172. }
  173. }
  174. }
  175. }
  176. // Receives sencevent, gets the offspring of the contained eventID and sends them.
  177. async fn handle_receive_syncevent(self: Arc<Self>) -> Result<()> {
  178. debug!(target: "event_graph", "ProtocolEvent::handle_receive_syncevent() [START]");
  179. loop {
  180. let syncevent = self.syncevent_sub.receive().await?;
  181. let model = self.model.lock().await;
  182. let leaves = model.find_leaves();
  183. if leaves == syncevent.leaves {
  184. continue
  185. }
  186. for leaf in syncevent.leaves.iter() {
  187. if leaves.contains(leaf) {
  188. continue
  189. }
  190. let mut children = vec![];
  191. model.get_offspring(leaf, &mut children)?;
  192. for child in children {
  193. self.channel.send(&child).await?;
  194. }
  195. }
  196. }
  197. }
  198. // every 6 seconds send a SyncEvent msg
  199. async fn send_sync_hash_loop(self: Arc<Self>) -> Result<()> {
  200. debug!(target: "event_graph", "ProtocolEvent::send_sync_hash_loop() [START]");
  201. loop {
  202. sleep(6).await;
  203. let leaves = self.model.lock().await.find_leaves();
  204. self.channel.send(&SyncEvent { leaves }).await?;
  205. }
  206. }
  207. async fn new_event(&self, event: &Event<T>) -> Result<()> {
  208. debug!(target: "event_graph", "ProtocolEvent::new_event()");
  209. let mut model = self.model.lock().await;
  210. model.add(event.clone()).await?;
  211. Ok(())
  212. }
  213. async fn send_inv(&self, event: &Event<T>) -> Result<()> {
  214. debug!(target: "event_graph", "ProtocolEvent::send_inv()");
  215. self.p2p.broadcast(&Inv { invs: vec![InvItem { hash: event.hash() }] }).await;
  216. Ok(())
  217. }
  218. async fn send_getdata(&self, events: Vec<EventId>) -> Result<()> {
  219. debug!(target: "event_graph", "ProtocolEvent::send_getdata()");
  220. self.channel.send(&GetData { events }).await?;
  221. Ok(())
  222. }
  223. }
  224. #[async_trait]
  225. impl<T> net::ProtocolBase for ProtocolEvent<T>
  226. where
  227. T: Send + Sync + Encodable + Decodable + Clone + EventMsg + Debug,
  228. {
  229. async fn start(self: Arc<Self>, executor: Arc<smol::Executor<'_>>) -> Result<()> {
  230. debug!(target: "event_graph", "ProtocolEvent::start() [START]");
  231. self.jobsman.clone().start(executor.clone());
  232. self.jobsman.clone().spawn(self.clone().handle_receive_event(), executor.clone()).await;
  233. self.jobsman.clone().spawn(self.clone().handle_receive_inv(), executor.clone()).await;
  234. self.jobsman.clone().spawn(self.clone().handle_receive_getdata(), executor.clone()).await;
  235. self.jobsman.clone().spawn(self.clone().handle_receive_syncevent(), executor.clone()).await;
  236. self.jobsman.clone().spawn(self.clone().send_sync_hash_loop(), executor.clone()).await;
  237. debug!(target: "event_graph", "ProtocolEvent::start() [END]");
  238. Ok(())
  239. }
  240. fn name(&self) -> &'static str {
  241. "ProtocolEvent"
  242. }
  243. }
  244. impl<T> net::Message for Event<T>
  245. where
  246. T: Send + Sync + Decodable + Encodable + 'static,
  247. {
  248. const NAME: &'static str = "event";
  249. }