tests.rs 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349
  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. // cargo +nightly test --release --features=event-graph --lib eventgraph_propagation -- --include-ignored
  19. use std::sync::Arc;
  20. use log::info;
  21. use rand::{prelude::SliceRandom, Rng};
  22. use smol::{channel, future, Executor};
  23. use url::Url;
  24. use crate::{
  25. event_graph::{
  26. proto::{EventPut, ProtocolEventGraph},
  27. Event, EventGraph, NULL_ID,
  28. },
  29. net::{P2p, Settings, SESSION_ALL},
  30. system::sleep,
  31. };
  32. // Number of nodes to spawn and number of peers each node connects to
  33. const N_NODES: usize = 5;
  34. const N_CONNS: usize = 2;
  35. //const N_NODES: usize = 50;
  36. //const N_CONNS: usize = N_NODES / 3;
  37. #[test]
  38. #[ignore]
  39. fn eventgraph_propagation() {
  40. let mut cfg = simplelog::ConfigBuilder::new();
  41. cfg.add_filter_ignore("sled".to_string());
  42. cfg.add_filter_ignore("net::protocol_ping".to_string());
  43. cfg.add_filter_ignore("net::channel::subscribe_stop()".to_string());
  44. cfg.add_filter_ignore("net::hosts".to_string());
  45. cfg.add_filter_ignore("net::session".to_string());
  46. cfg.add_filter_ignore("net::message_subscriber".to_string());
  47. cfg.add_filter_ignore("net::protocol_address".to_string());
  48. cfg.add_filter_ignore("net::protocol_version".to_string());
  49. cfg.add_filter_ignore("net::protocol_registry".to_string());
  50. cfg.add_filter_ignore("net::channel::send()".to_string());
  51. cfg.add_filter_ignore("net::channel::start()".to_string());
  52. cfg.add_filter_ignore("net::channel::subscribe_msg()".to_string());
  53. simplelog::TermLogger::init(
  54. //simplelog::LevelFilter::Info,
  55. simplelog::LevelFilter::Debug,
  56. //simplelog::LevelFilter::Trace,
  57. cfg.build(),
  58. simplelog::TerminalMode::Mixed,
  59. simplelog::ColorChoice::Auto,
  60. )
  61. .unwrap();
  62. let ex = Arc::new(Executor::new());
  63. let ex_ = ex.clone();
  64. let (signal, shutdown) = channel::unbounded::<()>();
  65. // Run a thread for each node.
  66. easy_parallel::Parallel::new()
  67. .each(0..N_NODES, |_| future::block_on(ex.run(shutdown.recv())))
  68. .finish(|| {
  69. future::block_on(async {
  70. eventgraph_propagation_real(ex_).await;
  71. drop(signal);
  72. })
  73. });
  74. }
  75. async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
  76. let mut eg_instances = vec![];
  77. let mut rng = rand::thread_rng();
  78. let mut genesis_event_id = NULL_ID;
  79. // Initialize the nodes
  80. for i in 0..N_NODES {
  81. // Everyone will connect to N_CONNS random peers.
  82. let mut peers = vec![];
  83. for _ in 0..N_CONNS {
  84. let mut port = 13200 + i;
  85. while port == 13200 + i {
  86. port = 13200 + rng.gen_range(0..N_NODES);
  87. }
  88. peers.push(Url::parse(&format!("tcp://127.0.0.1:{}", port)).unwrap());
  89. }
  90. let settings = Settings {
  91. localnet: true,
  92. inbound_addrs: vec![Url::parse(&format!("tcp://127.0.0.1:{}", 13200 + i)).unwrap()],
  93. outbound_connections: 0,
  94. outbound_connect_timeout: 2,
  95. inbound_connections: usize::MAX,
  96. peers,
  97. allowed_transports: vec!["tcp".to_string()],
  98. ..Default::default()
  99. };
  100. let p2p = P2p::new(settings, ex.clone()).await;
  101. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  102. let event_graph =
  103. EventGraph::new(p2p.clone(), sled_db, "dag", 1, ex.clone()).await.unwrap();
  104. let event_graph_ = event_graph.clone();
  105. // Take the last sled item since there's only 1
  106. if genesis_event_id == NULL_ID {
  107. let (id, _) = event_graph.dag.last().unwrap().unwrap();
  108. genesis_event_id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
  109. }
  110. // Register the P2P protocols
  111. let registry = p2p.protocol_registry();
  112. registry
  113. .register(SESSION_ALL, move |channel, _| {
  114. let event_graph_ = event_graph_.clone();
  115. async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
  116. })
  117. .await;
  118. eg_instances.push(event_graph);
  119. }
  120. // Start the P2P network
  121. for eg in eg_instances.iter() {
  122. eg.p2p.clone().start().await.unwrap();
  123. }
  124. info!("Waiting 10s until all peers connect");
  125. sleep(10).await;
  126. // =========================================
  127. // 1. Assert that everyone's DAG is the same
  128. // =========================================
  129. for (i, eg) in eg_instances.iter().enumerate() {
  130. let tips = eg.unreferenced_tips.read().await;
  131. assert!(eg.dag.len() == 1, "Node {}", i);
  132. assert!(tips.len() == 1, "Node {}", i);
  133. assert!(tips.get(&genesis_event_id).is_some(), "Node {}", i);
  134. }
  135. // ==========================================
  136. // 2. Create an event in one node and publish
  137. // ==========================================
  138. let random_node = eg_instances.choose(&mut rand::thread_rng()).unwrap();
  139. let event = Event::new(vec![1, 2, 3, 4], random_node.clone()).await;
  140. assert!(event.parents.contains(&genesis_event_id));
  141. // The node adds it to their DAG.
  142. let event_id = random_node.dag_insert(event.clone()).await.unwrap();
  143. let tips = random_node.unreferenced_tips.read().await;
  144. assert!(tips.len() == 1);
  145. assert!(tips.get(&event_id).is_some());
  146. drop(tips);
  147. info!("Broadcasting event {}", event_id);
  148. random_node.p2p.broadcast(&EventPut(event)).await;
  149. info!("Waiting 10s for event propagation");
  150. sleep(10).await;
  151. // ====================================================
  152. // 3. Assert that everyone has the new event in the DAG
  153. // ====================================================
  154. for (i, eg) in eg_instances.iter().enumerate() {
  155. let tips = eg.unreferenced_tips.read().await;
  156. assert!(eg.dag.len() == 2, "Node {}", i);
  157. assert!(tips.len() == 1, "Node {}", i);
  158. assert!(tips.get(&event_id).is_some(), "Node {}", i);
  159. }
  160. // ==============================================================
  161. // 4. Create multiple events on a node and broadcast the last one
  162. // The `EventPut` logic should manage to fetch all of them,
  163. // provided that the last one references the earlier ones.
  164. // ==============================================================
  165. let random_node = eg_instances.choose(&mut rand::thread_rng()).unwrap();
  166. let event0 = Event::new(vec![1, 2, 3, 4, 0], random_node.clone()).await;
  167. let event0_id = random_node.dag_insert(event0.clone()).await.unwrap();
  168. let event1 = Event::new(vec![1, 2, 3, 4, 1], random_node.clone()).await;
  169. let event1_id = random_node.dag_insert(event1.clone()).await.unwrap();
  170. let event2 = Event::new(vec![1, 2, 3, 4, 2], random_node.clone()).await;
  171. let event2_id = random_node.dag_insert(event2.clone()).await.unwrap();
  172. // Genesis event + event from 2. + upper 3 events
  173. assert!(random_node.dag.len() == 5);
  174. let tips = random_node.unreferenced_tips.read().await;
  175. assert!(tips.len() == 1);
  176. assert!(tips.get(&event2_id).is_some());
  177. drop(tips);
  178. let event_chain =
  179. vec![(event0_id, event0.parents), (event1_id, event1.parents), (event2_id, event2.parents)];
  180. info!("Broadcasting event {}", event2_id);
  181. info!("Event chain: {:#?}", event_chain);
  182. random_node.p2p.broadcast(&EventPut(event2)).await;
  183. info!("Waiting 10s for event propagation");
  184. sleep(10).await;
  185. // ==========================================
  186. // 5. Assert that everyone has all the events
  187. // ==========================================
  188. for (i, eg) in eg_instances.iter().enumerate() {
  189. let tips = eg.unreferenced_tips.read().await;
  190. assert!(eg.dag.len() == 5, "Node {}, expected 5 events, have {}", i, eg.dag.len());
  191. assert!(tips.len() == 1, "Node {}, expected 1 tip, have {}", i, tips.len());
  192. assert!(tips.get(&event2_id).is_some(), "Node {}, expected tip to be {}", i, event2_id);
  193. }
  194. // ===========================================
  195. // 6. Create multiple events on multiple nodes
  196. // ===========================================
  197. // node 1
  198. // =======
  199. let node1 = eg_instances.choose(&mut rand::thread_rng()).unwrap();
  200. let event0_1 = Event::new(vec![1, 2, 3, 4, 3], node1.clone()).await;
  201. let _ = node1.dag_insert(event0_1.clone()).await.unwrap();
  202. node1.p2p.broadcast(&EventPut(event0_1)).await;
  203. let event1_1 = Event::new(vec![1, 2, 3, 4, 4], node1.clone()).await;
  204. let _ = node1.dag_insert(event1_1.clone()).await.unwrap();
  205. node1.p2p.broadcast(&EventPut(event1_1)).await;
  206. let event2_1 = Event::new(vec![1, 2, 3, 4, 5], node1.clone()).await;
  207. let _ = node1.dag_insert(event2_1.clone()).await.unwrap();
  208. node1.p2p.broadcast(&EventPut(event2_1)).await;
  209. // =======
  210. // node 2
  211. // =======
  212. let node2 = eg_instances.choose(&mut rand::thread_rng()).unwrap();
  213. let event0_2 = Event::new(vec![1, 2, 3, 4, 6], node2.clone()).await;
  214. let _ = node2.dag_insert(event0_2.clone()).await.unwrap();
  215. node2.p2p.broadcast(&EventPut(event0_2)).await;
  216. let event1_2 = Event::new(vec![1, 2, 3, 4, 7], node2.clone()).await;
  217. let _ = node2.dag_insert(event1_2.clone()).await.unwrap();
  218. node2.p2p.broadcast(&EventPut(event1_2)).await;
  219. let event2_2 = Event::new(vec![1, 2, 3, 4, 8], node2.clone()).await;
  220. let _ = node2.dag_insert(event2_2.clone()).await.unwrap();
  221. node2.p2p.broadcast(&EventPut(event2_2)).await;
  222. // =======
  223. // node 3
  224. // =======
  225. let node3 = eg_instances.choose(&mut rand::thread_rng()).unwrap();
  226. let event0_3 = Event::new(vec![1, 2, 3, 4, 9], node3.clone()).await;
  227. let _ = node3.dag_insert(event0_3.clone()).await.unwrap();
  228. node2.p2p.broadcast(&EventPut(event0_3)).await;
  229. let event1_3 = Event::new(vec![1, 2, 3, 4, 10], node3.clone()).await;
  230. let _ = node3.dag_insert(event1_3.clone()).await.unwrap();
  231. node2.p2p.broadcast(&EventPut(event1_3)).await;
  232. let event2_3 = Event::new(vec![1, 2, 3, 4, 11], node3.clone()).await;
  233. let event2_3_id = node3.dag_insert(event2_3.clone()).await.unwrap();
  234. node3.p2p.broadcast(&EventPut(event2_3)).await;
  235. info!("Waiting 10s for events propagation");
  236. sleep(10).await;
  237. // ==========================================
  238. // 7. Assert that everyone has all the events
  239. // ==========================================
  240. for (i, eg) in eg_instances.iter().enumerate() {
  241. let tips = eg.unreferenced_tips.read().await;
  242. assert!(eg.dag.len() == 14, "Node {}, expected 14 events, have {}", i, eg.dag.len());
  243. // 5 events from 2. and 4. + 9 events from 6. = ^
  244. assert!(tips.get(&event2_3_id).is_some(), "Node {}, expected tip to be {}", i, event2_3_id);
  245. }
  246. // ============================================================
  247. // 8. Start a new node and try to sync the DAG from other peers
  248. // ============================================================
  249. {
  250. // Connect to N_CONNS random peers.
  251. let mut peers = vec![];
  252. for _ in 0..N_CONNS {
  253. let port = 13200 + rng.gen_range(0..N_NODES);
  254. peers.push(Url::parse(&format!("tcp://127.0.0.1:{}", port)).unwrap());
  255. }
  256. let settings = Settings {
  257. localnet: true,
  258. inbound_addrs: vec![
  259. Url::parse(&format!("tcp://127.0.0.1:{}", 13200 + N_NODES + 1)).unwrap()
  260. ],
  261. outbound_connections: 0,
  262. outbound_connect_timeout: 2,
  263. inbound_connections: usize::MAX,
  264. peers,
  265. allowed_transports: vec!["tcp".to_string()],
  266. ..Default::default()
  267. };
  268. let p2p = P2p::new(settings, ex.clone()).await;
  269. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  270. let event_graph =
  271. EventGraph::new(p2p.clone(), sled_db, "dag", 1, ex.clone()).await.unwrap();
  272. let event_graph_ = event_graph.clone();
  273. // Register the P2P protocols
  274. let registry = p2p.protocol_registry();
  275. registry
  276. .register(SESSION_ALL, move |channel, _| {
  277. let event_graph_ = event_graph_.clone();
  278. async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
  279. })
  280. .await;
  281. eg_instances.push(event_graph.clone());
  282. event_graph.p2p.clone().start().await.unwrap();
  283. info!("Waiting 10s for new node connection");
  284. sleep(10).await;
  285. event_graph.dag_sync().await.unwrap();
  286. }
  287. info!("Waiting 10s for things to settle");
  288. sleep(10).await;
  289. // ============================================================
  290. // 9. Assert the new synced DAG has the same contents as others
  291. // ============================================================
  292. for (i, eg) in eg_instances.iter().enumerate() {
  293. let tips = eg.unreferenced_tips.read().await;
  294. assert!(eg.dag.len() == 14, "Node {}, expected 14 events, have {}", i, eg.dag.len());
  295. // 5 events from 2. and 4. + 9 events from 6. = ^
  296. assert!(tips.get(&event2_3_id).is_some(), "Node {}, expected tip to be {}", i, event2_3_id);
  297. }
  298. // Stop the P2P network
  299. for eg in eg_instances.iter() {
  300. eg.p2p.clone().stop().await;
  301. }
  302. }