protocol_privmsg2.rs 9.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331
  1. use std::collections::VecDeque;
  2. use async_executor::Executor;
  3. use async_std::sync::{Arc, Mutex};
  4. use async_trait::async_trait;
  5. use fxhash::FxHashMap;
  6. use log::debug;
  7. use rand::{rngs::OsRng, RngCore};
  8. use darkfi::{
  9. net,
  10. serial::{SerialDecodable, SerialEncodable},
  11. util::async_util::sleep,
  12. Result,
  13. };
  14. use crate::mvc::{get_current_time, Event, EventId};
  15. const UNREAD_EVENT_EXPIRE_TIME: u64 = 3600; // in seconds
  16. const SIZE_OF_SEEN_BUFFER: usize = 65536;
  17. const MAX_CONFIRM: u8 = 4;
  18. type InvId = u64;
  19. #[derive(SerialEncodable, SerialDecodable, Clone, Debug, PartialEq, Eq, Hash)]
  20. struct InvItem {
  21. id: InvId,
  22. hash: EventId,
  23. }
  24. #[derive(Clone)]
  25. pub struct RingBuffer<T> {
  26. pub items: VecDeque<T>,
  27. }
  28. impl<T: Eq + PartialEq + Clone> RingBuffer<T> {
  29. pub fn new(capacity: usize) -> Self {
  30. let items = VecDeque::with_capacity(capacity);
  31. Self { items }
  32. }
  33. pub fn push(&mut self, val: T) {
  34. if self.items.len() == self.items.capacity() {
  35. self.items.pop_front();
  36. }
  37. self.items.push_back(val);
  38. }
  39. pub fn contains(&self, val: &T) -> bool {
  40. self.items.contains(val)
  41. }
  42. }
  43. pub struct Seen<T> {
  44. seen: Mutex<RingBuffer<T>>,
  45. }
  46. impl<T: Eq + PartialEq + Clone> Seen<T> {
  47. pub fn new() -> Self {
  48. Self { seen: Mutex::new(RingBuffer::new(SIZE_OF_SEEN_BUFFER)) }
  49. }
  50. pub async fn push(&self, item: &T) -> bool {
  51. let seen = &mut self.seen.lock().await;
  52. if !seen.contains(item) {
  53. seen.push(item.clone());
  54. return true
  55. }
  56. false
  57. }
  58. }
  59. pub struct UnreadEvents {
  60. events: Mutex<FxHashMap<EventId, Event>>,
  61. }
  62. impl UnreadEvents {
  63. pub fn new() -> Self {
  64. Self { events: Mutex::new(FxHashMap::default()) }
  65. }
  66. pub async fn contains(&self, key: &EventId) -> bool {
  67. self.events.lock().await.contains_key(key)
  68. }
  69. // Increase the read_confirms for an event, if it has exceeded the MAX_CONFIRM
  70. // then remove it from the hash_map and return Some(event), otherwise return None
  71. pub async fn inc_read_confirms(&self, key: &EventId) -> Option<Event> {
  72. let events = &mut self.events.lock().await;
  73. let mut result = None;
  74. if let Some(event) = events.get_mut(key) {
  75. event.read_confirms += 1;
  76. if event.read_confirms >= MAX_CONFIRM {
  77. result = Some(event.clone())
  78. }
  79. }
  80. if result.is_some() {
  81. events.remove(key);
  82. }
  83. result
  84. }
  85. pub async fn insert(&self, event: &Event) {
  86. let events = &mut self.events.lock().await;
  87. // prune expired events
  88. let mut prune_ids = vec![];
  89. for (id, e) in events.iter() {
  90. if e.timestamp + (UNREAD_EVENT_EXPIRE_TIME * 1000) < get_current_time() {
  91. prune_ids.push(id.clone());
  92. }
  93. }
  94. for id in prune_ids {
  95. events.remove(&id);
  96. }
  97. events.insert(event.hash().clone(), event.clone());
  98. }
  99. }
  100. #[derive(SerialDecodable, SerialEncodable, Clone, Debug)]
  101. struct Inv {
  102. invs: Vec<InvItem>,
  103. }
  104. #[derive(SerialDecodable, SerialEncodable, Clone, Debug)]
  105. struct SyncEvent {
  106. leaves: Vec<EventId>,
  107. }
  108. #[derive(SerialDecodable, SerialEncodable, Clone, Debug)]
  109. struct GetData {
  110. invs: Vec<InvItem>,
  111. }
  112. pub struct ProtocolEvent {
  113. jobsman: net::ProtocolJobsManagerPtr,
  114. event_sub: net::MessageSubscription<Event>,
  115. inv_sub: net::MessageSubscription<Inv>,
  116. getdata_sub: net::MessageSubscription<GetData>,
  117. syncevent_sub: net::MessageSubscription<SyncEvent>,
  118. p2p: net::P2pPtr,
  119. channel: net::ChannelPtr,
  120. seen_event: Seen<EventId>,
  121. seen_inv: Seen<InvId>,
  122. unread_events: UnreadEvents,
  123. }
  124. impl ProtocolEvent {
  125. pub async fn init(
  126. channel: net::ChannelPtr,
  127. p2p: net::P2pPtr,
  128. seen_event: Seen<EventId>,
  129. seen_inv: Seen<InvId>,
  130. unread_events: UnreadEvents,
  131. ) -> net::ProtocolBasePtr {
  132. let message_subsytem = channel.get_message_subsystem();
  133. message_subsytem.add_dispatch::<Event>().await;
  134. message_subsytem.add_dispatch::<Inv>().await;
  135. message_subsytem.add_dispatch::<GetData>().await;
  136. message_subsytem.add_dispatch::<SyncEvent>().await;
  137. let event_sub =
  138. channel.clone().subscribe_msg::<Event>().await.expect("Missing Event dispatcher!");
  139. let inv_sub = channel.subscribe_msg::<Inv>().await.expect("Missing Inv dispatcher!");
  140. let getdata_sub =
  141. channel.clone().subscribe_msg::<GetData>().await.expect("Missing GetData dispatcher!");
  142. let syncevent_sub = channel
  143. .clone()
  144. .subscribe_msg::<SyncEvent>()
  145. .await
  146. .expect("Missing SyncEvent dispatcher!");
  147. Arc::new(Self {
  148. jobsman: net::ProtocolJobsManager::new("ProtocolEvent", channel.clone()),
  149. event_sub,
  150. inv_sub,
  151. getdata_sub,
  152. syncevent_sub,
  153. p2p,
  154. channel,
  155. seen_event,
  156. seen_inv,
  157. unread_events,
  158. })
  159. }
  160. async fn handle_receive_inv(self: Arc<Self>) -> Result<()> {
  161. debug!(target: "ircd", "ProtocolEvent::handle_receive_inv() [START]");
  162. let exclude_list = vec![self.channel.address()];
  163. loop {
  164. let inv = self.inv_sub.receive().await?;
  165. let inv = (*inv).to_owned();
  166. for inv in inv.invs.iter() {
  167. if !self.seen_inv.push(&inv.id).await {
  168. continue
  169. }
  170. // On receive inv message, if the unread_events buffer has the event's hash then increase
  171. // the read_confirms, if not then send GetMsgs contain the event's hash
  172. if !self.unread_events.contains(&inv.hash).await {
  173. self.send_getdata(vec![inv.clone()]).await?;
  174. } else if let Some(event) = self.unread_events.inc_read_confirms(&inv.hash).await {
  175. self.new_event(&event).await?;
  176. }
  177. }
  178. // Broadcast the inv msg
  179. self.p2p.broadcast_with_exclude(inv, &exclude_list).await?;
  180. }
  181. }
  182. async fn handle_receive_event(self: Arc<Self>) -> Result<()> {
  183. debug!(target: "ircd", "ProtocolEvent::handle_receive_event() [START]");
  184. let exclude_list = vec![self.channel.address()];
  185. loop {
  186. let event = self.event_sub.receive().await?;
  187. let mut event = (*event).to_owned();
  188. if !self.seen_event.push(&event.hash()).await {
  189. continue
  190. }
  191. // If the event has read_confirms greater or equal to MAX_CONFIRM, it will be added to
  192. // the model, otherwise increase the event's read_confirms, add it to unread_events, and
  193. // broadcast an Inv msg
  194. if event.read_confirms >= MAX_CONFIRM {
  195. self.new_event(&event).await?;
  196. } else {
  197. event.read_confirms += 1;
  198. self.unread_events.insert(&event).await;
  199. self.send_inv(&event).await?;
  200. }
  201. // Broadcast the msg
  202. self.p2p.broadcast_with_exclude(event, &exclude_list).await?;
  203. }
  204. }
  205. async fn handle_receive_getdata(self: Arc<Self>) -> Result<()> {
  206. debug!(target: "ircd", "ProtocolEvent::handle_receive_getdata() [START]");
  207. loop {
  208. let getdata = self.getdata_sub.receive().await?;
  209. let invs = (*getdata).to_owned().invs;
  210. for inv in invs {}
  211. }
  212. }
  213. async fn handle_receive_syncevent(self: Arc<Self>) -> Result<()> {
  214. debug!(target: "ircd", "ProtocolEvent::handle_receive_syncevent() [START]");
  215. loop {
  216. let syncevent = self.syncevent_sub.receive().await?;
  217. }
  218. }
  219. // every 2 seconds send a Sync msg
  220. async fn send_sync_hash_loop(self: Arc<Self>) -> Result<()> {
  221. loop {
  222. //let leaves = self.model.find_leaves();
  223. //self.channel.send(SyncEvent { leaves }).await;
  224. sleep(2).await;
  225. }
  226. }
  227. async fn new_event(&self, event: &Event) -> Result<()> {
  228. // self.model.add(event).await?;
  229. Ok(())
  230. }
  231. async fn send_inv(&self, event: &Event) -> Result<()> {
  232. let id = OsRng.next_u64();
  233. self.p2p.broadcast(Inv { invs: vec![InvItem { id, hash: event.hash() }] }).await?;
  234. Ok(())
  235. }
  236. async fn send_getdata(&self, invs: Vec<InvItem>) -> Result<()> {
  237. self.channel.send(GetData { invs }).await?;
  238. Ok(())
  239. }
  240. }
  241. #[async_trait]
  242. impl net::ProtocolBase for ProtocolEvent {
  243. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
  244. debug!(target: "ircd", "ProtocolEvent::start() [START]");
  245. self.jobsman.clone().start(executor.clone());
  246. self.jobsman.clone().spawn(self.clone().handle_receive_event(), executor.clone()).await;
  247. self.jobsman.clone().spawn(self.clone().handle_receive_inv(), executor.clone()).await;
  248. self.jobsman.clone().spawn(self.clone().handle_receive_getdata(), executor.clone()).await;
  249. self.jobsman.clone().spawn(self.clone().handle_receive_syncevent(), executor.clone()).await;
  250. self.jobsman.clone().spawn(self.clone().send_sync_hash_loop(), executor.clone()).await;
  251. debug!(target: "ircd", "ProtocolEvent::start() [END]");
  252. Ok(())
  253. }
  254. fn name(&self) -> &'static str {
  255. "ProtocolEvent"
  256. }
  257. }
  258. impl net::Message for Event {
  259. fn name() -> &'static str {
  260. "event"
  261. }
  262. }
  263. impl net::Message for Inv {
  264. fn name() -> &'static str {
  265. "inv"
  266. }
  267. }
  268. impl net::Message for SyncEvent {
  269. fn name() -> &'static str {
  270. "syncevent"
  271. }
  272. }
  273. impl net::Message for GetData {
  274. fn name() -> &'static str {
  275. "getdata"
  276. }
  277. }