tests.rs 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520
  1. /* This file is part of DarkFi (https://dark.fi)
  2. *
  3. * Copyright (C) 2020-2026 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 test --release --features=event-graph --lib eventgraph_propagation -- --include-ignored
  19. use std::{collections::HashMap, slice, sync::Arc};
  20. use rand::{prelude::SliceRandom, rngs::ThreadRng};
  21. use sled_overlay::sled;
  22. use smol::{channel, future, Executor};
  23. use tracing::{info, warn};
  24. use url::Url;
  25. use crate::{
  26. event_graph::{
  27. proto::{EventPut, ProtocolEventGraph},
  28. Event, EventGraph,
  29. },
  30. net::{session::SESSION_DEFAULT, settings::NetworkProfile, P2p, Settings},
  31. system::{msleep, sleep},
  32. util::logger::{setup_test_logger, Level},
  33. };
  34. // Number of nodes to spawn and number of peers each node connects to
  35. const N_NODES: usize = 5;
  36. const N_CONNS: usize = 2;
  37. //const N_NODES: usize = 50;
  38. //const N_CONNS: usize = N_NODES / 3;
  39. fn init_logger() {
  40. let ignored_targets = [
  41. "sled",
  42. "net::protocol_ping",
  43. "net::channel::subscribe_stop()",
  44. "net::hosts",
  45. "net::session",
  46. "net::message_subscriber",
  47. "net::protocol_address",
  48. "net::protocol_version",
  49. "net::protocol_registry",
  50. "net::channel::send()",
  51. "net::channel::start()",
  52. "net::channel::subscribe_msg()",
  53. "net::channel::main_receive_loop()",
  54. "net::tcp",
  55. ];
  56. // We check this error so we can execute same file tests in parallel,
  57. // otherwise second one fails to init logger here.
  58. if setup_test_logger(
  59. &ignored_targets,
  60. false,
  61. Level::Info,
  62. //Level::Verbose,
  63. //Level::Debug,
  64. //Level::Tracing,
  65. )
  66. .is_err()
  67. {
  68. warn!(target: "test_harness", "Logger already initialized");
  69. }
  70. }
  71. async fn spawn_node(
  72. inbound_addrs: Vec<Url>,
  73. peers: Vec<Url>,
  74. ex: Arc<Executor<'static>>,
  75. ) -> Arc<EventGraph> {
  76. let mut profiles = HashMap::new();
  77. profiles.insert(
  78. "tcp".to_string(),
  79. NetworkProfile { outbound_connect_timeout: 2, ..Default::default() },
  80. );
  81. let settings = Settings {
  82. localnet: true,
  83. inbound_addrs,
  84. outbound_connections: 0,
  85. inbound_connections: usize::MAX,
  86. peers,
  87. active_profiles: vec!["tcp".to_string()],
  88. profiles,
  89. ..Default::default()
  90. };
  91. let p2p = P2p::new(settings, ex.clone()).await.unwrap();
  92. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  93. let event_graph =
  94. EventGraph::new(p2p.clone(), sled_db, "/tmp".into(), false, false, 1, ex.clone())
  95. .await
  96. .unwrap();
  97. *event_graph.synced.write().await = true;
  98. let event_graph_ = event_graph.clone();
  99. // Register the P2P protocols
  100. let registry = p2p.protocol_registry();
  101. registry
  102. .register(SESSION_DEFAULT, move |channel, _| {
  103. let event_graph_ = event_graph_.clone();
  104. async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
  105. })
  106. .await;
  107. event_graph
  108. }
  109. async fn bootstrap_nodes(
  110. peer_indexes: &[usize],
  111. starting_port: usize,
  112. rng: &mut ThreadRng,
  113. ex: Arc<Executor<'static>>,
  114. ) -> Vec<Arc<EventGraph>> {
  115. let mut eg_instances = vec![];
  116. // Initialize the nodes
  117. for i in 0..N_NODES {
  118. // Everyone will connect to N_CONNS random peers.
  119. let mut peer_indexes_copy = peer_indexes.to_owned();
  120. peer_indexes_copy.remove(i);
  121. let peer_indexes_to_connect: Vec<_> =
  122. peer_indexes_copy.choose_multiple(rng, N_CONNS).collect();
  123. let mut peers = vec![];
  124. for peer_index in peer_indexes_to_connect {
  125. let port = starting_port + peer_index;
  126. peers.push(Url::parse(&format!("tcp://127.0.0.1:{port}")).unwrap());
  127. }
  128. let event_graph = spawn_node(
  129. vec![Url::parse(&format!("tcp://127.0.0.1:{}", starting_port + i)).unwrap()],
  130. peers,
  131. ex.clone(),
  132. )
  133. .await;
  134. eg_instances.push(event_graph);
  135. }
  136. // Start the P2P network
  137. for eg in eg_instances.iter() {
  138. eg.p2p.clone().start().await.unwrap();
  139. }
  140. info!("Waiting 5s until all peers connect");
  141. sleep(5).await;
  142. eg_instances
  143. }
  144. async fn assert_dags(eg_instances: &[Arc<EventGraph>], expected_len: usize, rng: &mut ThreadRng) {
  145. let random_node = eg_instances.choose(rng).unwrap();
  146. let random_node_genesis = random_node.current_genesis.read().await.id();
  147. let store = random_node.dag_store.read().await;
  148. let (_, unreferenced_tips) = store.main_dags.get(&random_node_genesis).unwrap();
  149. let last_layer_tips = unreferenced_tips.last_key_value().unwrap().1.clone();
  150. for (i, eg) in eg_instances.iter().enumerate() {
  151. let current_genesis = eg.current_genesis.read().await;
  152. let dag_name = current_genesis.id().to_string();
  153. let dag = eg.dag_store.read().await.get_dag(&dag_name);
  154. let unreferenced_tips = eg.dag_store.read().await.find_unreferenced_tips(&dag).await;
  155. let node_last_layer_tips = unreferenced_tips.last_key_value().unwrap().1.clone();
  156. assert!(
  157. dag.len() == expected_len,
  158. "Node {i}, expected {expected_len} events, have {}",
  159. dag.len()
  160. );
  161. assert_eq!(
  162. node_last_layer_tips, last_layer_tips,
  163. "Node {i} contains malformed unreferenced tips"
  164. );
  165. }
  166. }
  167. macro_rules! test_body {
  168. ($real_call:ident) => {
  169. init_logger();
  170. let ex = Arc::new(Executor::new());
  171. let ex_ = ex.clone();
  172. let (signal, shutdown) = channel::unbounded::<()>();
  173. // Run a thread for each node.
  174. easy_parallel::Parallel::new()
  175. .each(0..N_NODES, |_| future::block_on(ex.run(shutdown.recv())))
  176. .finish(|| {
  177. future::block_on(async {
  178. $real_call(ex_).await;
  179. drop(signal);
  180. })
  181. });
  182. };
  183. }
  184. #[test]
  185. fn eventgraph_propagation() {
  186. test_body!(eventgraph_propagation_real);
  187. }
  188. async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
  189. let mut rng = rand::thread_rng();
  190. let peer_indexes: Vec<usize> = (0..N_NODES).collect();
  191. // Bootstrap nodes
  192. let mut eg_instances = bootstrap_nodes(&peer_indexes, 13200, &mut rng, ex.clone()).await;
  193. // Grab genesis event
  194. let random_node = eg_instances.choose(&mut rng).unwrap();
  195. let current_genesis = random_node.current_genesis.read().await;
  196. let dag_name = current_genesis.id().to_string();
  197. let (id, _) = random_node.dag_store.read().await.get_dag(&dag_name).last().unwrap().unwrap();
  198. let genesis_event_id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
  199. drop(current_genesis);
  200. // =========================================
  201. // 1. Assert that everyone's DAG is the same
  202. // =========================================
  203. assert_dags(&eg_instances, 1, &mut rng).await;
  204. // ==========================================
  205. // 2. Create an event in one node and publish
  206. // ==========================================
  207. let random_node = eg_instances.choose(&mut rng).unwrap();
  208. let current_genesis = random_node.current_genesis.read().await;
  209. let dag_name = current_genesis.id().to_string();
  210. let event = Event::new(vec![1, 2, 3, 4], random_node).await;
  211. assert!(event.header.parents.contains(&genesis_event_id));
  212. // The node adds it to their DAG, on layer 1.
  213. random_node.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  214. let event_id = random_node.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap()[0];
  215. let store = random_node.dag_store.read().await;
  216. let (_, tips_layers) = store.header_dags.get(&current_genesis.id()).unwrap();
  217. // Since genesis was referenced, its layer (0) have been removed
  218. assert_eq!(tips_layers.len(), 1);
  219. assert!(tips_layers.last_key_value().unwrap().1.get(&event_id).is_some());
  220. drop(store);
  221. drop(current_genesis);
  222. info!("Broadcasting event {event_id}");
  223. random_node.p2p.broadcast(&EventPut(event)).await;
  224. info!("Waiting 5s for event propagation");
  225. sleep(5).await;
  226. // ====================================================
  227. // 3. Assert that everyone has the new event in the DAG
  228. // ====================================================
  229. assert_dags(&eg_instances, 2, &mut rng).await;
  230. // ==============================================================
  231. // 4. Create multiple events on a node and broadcast the last one
  232. // The `EventPut` logic should manage to fetch all of them,
  233. // provided that the last one references the earlier ones.
  234. // ==============================================================
  235. let random_node = eg_instances.choose(&mut rng).unwrap();
  236. let event0 = Event::new(vec![1, 2, 3, 4, 0], random_node).await;
  237. random_node.header_dag_insert(vec![event0.header.clone()], &dag_name).await.unwrap();
  238. let event0_id = random_node.dag_insert(slice::from_ref(&event0), &dag_name).await.unwrap()[0];
  239. let event1 = Event::new(vec![1, 2, 3, 4, 1], random_node).await;
  240. random_node.header_dag_insert(vec![event1.header.clone()], &dag_name).await.unwrap();
  241. let event1_id = random_node.dag_insert(slice::from_ref(&event1), &dag_name).await.unwrap()[0];
  242. let event2 = Event::new(vec![1, 2, 3, 4, 2], random_node).await;
  243. random_node.header_dag_insert(vec![event2.header.clone()], &dag_name).await.unwrap();
  244. let event2_id = random_node.dag_insert(slice::from_ref(&event2), &dag_name).await.unwrap()[0];
  245. // Genesis event + event from 2. + upper 3 events (layer 4)
  246. let current_genesis = random_node.current_genesis.read().await;
  247. let dag_name = current_genesis.id().to_string();
  248. assert_eq!(random_node.dag_store.read().await.get_dag(&dag_name).len(), 5);
  249. let random_node_genesis = random_node.current_genesis.read().await.id();
  250. let store = random_node.dag_store.read().await;
  251. let (_, tips_layers) = store.header_dags.get(&random_node_genesis).unwrap();
  252. assert_eq!(tips_layers.len(), 1);
  253. assert!(tips_layers.get(&4).unwrap().get(&event2_id).is_some());
  254. drop(current_genesis);
  255. drop(store);
  256. let event_chain = vec![
  257. (event0_id, event0.header.parents),
  258. (event1_id, event1.header.parents),
  259. (event2_id, event2.header.parents),
  260. ];
  261. info!("Broadcasting event {event2_id}");
  262. info!("Event chain: {event_chain:#?}");
  263. random_node.p2p.broadcast(&EventPut(event2)).await;
  264. info!("Waiting 5s for event propagation");
  265. sleep(5).await;
  266. // ==========================================
  267. // 5. Assert that everyone has all the events
  268. // ==========================================
  269. assert_dags(&eg_instances, 5, &mut rng).await;
  270. // ===========================================
  271. // 6. Create multiple events on multiple nodes
  272. // ===========================================
  273. // node 1
  274. // =======
  275. let node1 = eg_instances.choose(&mut rng).unwrap();
  276. let event0_1 = Event::new(vec![1, 2, 3, 4, 3], node1).await;
  277. node1.header_dag_insert(vec![event0_1.header.clone()], &dag_name).await.unwrap();
  278. node1.dag_insert(slice::from_ref(&event0_1), &dag_name).await.unwrap();
  279. node1.p2p.broadcast(&EventPut(event0_1)).await;
  280. msleep(300).await;
  281. let event1_1 = Event::new(vec![1, 2, 3, 4, 4], node1).await;
  282. node1.header_dag_insert(vec![event1_1.header.clone()], &dag_name).await.unwrap();
  283. node1.dag_insert(slice::from_ref(&event1_1), &dag_name).await.unwrap();
  284. node1.p2p.broadcast(&EventPut(event1_1)).await;
  285. msleep(300).await;
  286. let event2_1 = Event::new(vec![1, 2, 3, 4, 5], node1).await;
  287. node1.header_dag_insert(vec![event2_1.header.clone()], &dag_name).await.unwrap();
  288. node1.dag_insert(slice::from_ref(&event2_1), &dag_name).await.unwrap();
  289. node1.p2p.broadcast(&EventPut(event2_1)).await;
  290. msleep(300).await;
  291. // =======
  292. // node 2
  293. // =======
  294. let node2 = eg_instances.choose(&mut rng).unwrap();
  295. let event0_2 = Event::new(vec![1, 2, 3, 4, 6], node2).await;
  296. node2.header_dag_insert(vec![event0_2.header.clone()], &dag_name).await.unwrap();
  297. node2.dag_insert(slice::from_ref(&event0_2), &dag_name).await.unwrap();
  298. node2.p2p.broadcast(&EventPut(event0_2)).await;
  299. msleep(300).await;
  300. let event1_2 = Event::new(vec![1, 2, 3, 4, 7], node2).await;
  301. node2.header_dag_insert(vec![event1_2.header.clone()], &dag_name).await.unwrap();
  302. node2.dag_insert(slice::from_ref(&event1_2), &dag_name).await.unwrap();
  303. node2.p2p.broadcast(&EventPut(event1_2)).await;
  304. msleep(300).await;
  305. let event2_2 = Event::new(vec![1, 2, 3, 4, 8], node2).await;
  306. node2.header_dag_insert(vec![event2_2.header.clone()], &dag_name).await.unwrap();
  307. node2.dag_insert(slice::from_ref(&event2_2), &dag_name).await.unwrap();
  308. node2.p2p.broadcast(&EventPut(event2_2)).await;
  309. msleep(300).await;
  310. // =======
  311. // node 3
  312. // =======
  313. let node3 = eg_instances.choose(&mut rng).unwrap();
  314. let event0_3 = Event::new(vec![1, 2, 3, 4, 9], node3).await;
  315. node3.header_dag_insert(vec![event0_3.header.clone()], &dag_name).await.unwrap();
  316. node3.dag_insert(slice::from_ref(&event0_3), &dag_name).await.unwrap();
  317. node3.p2p.broadcast(&EventPut(event0_3)).await;
  318. msleep(300).await;
  319. let event1_3 = Event::new(vec![1, 2, 3, 4, 10], node3).await;
  320. node3.header_dag_insert(vec![event1_3.header.clone()], &dag_name).await.unwrap();
  321. node3.dag_insert(slice::from_ref(&event1_3), &dag_name).await.unwrap();
  322. node3.p2p.broadcast(&EventPut(event1_3)).await;
  323. msleep(300).await;
  324. let event2_3 = Event::new(vec![1, 2, 3, 4, 11], node3).await;
  325. node3.header_dag_insert(vec![event2_3.header.clone()], &dag_name).await.unwrap();
  326. node3.dag_insert(slice::from_ref(&event2_3), &dag_name).await.unwrap();
  327. node3.p2p.broadcast(&EventPut(event2_3)).await;
  328. msleep(300).await;
  329. // /////
  330. // //
  331. // let node4 = eg_instances.choose(&mut rng).unwrap();
  332. // let event0_4 = Event::new(vec![1, 2, 3, 4, 12], node4).await;
  333. // node4.dag_insert(&[event0_4.clone()]).await.unwrap();
  334. // node4.p2p.broadcast(&EventPut(event0_4)).await;
  335. // sleep(1).await;
  336. // let event1_4 = Event::new(vec![1, 2, 3, 4, 13], node4).await;
  337. // node4.dag_insert(&[event1_4.clone()]).await.unwrap();
  338. // node4.p2p.broadcast(&EventPut(event1_4)).await;
  339. // sleep(1).await;
  340. // let event2_4 = Event::new(vec![1, 2, 3, 4, 14], node4).await;
  341. // node4.dag_insert(&[event2_4.clone()]).await.unwrap();
  342. // node4.p2p.broadcast(&EventPut(event2_4)).await;
  343. // // sleep(1).await;
  344. // ==========================================
  345. // 7. Assert that everyone has all the events
  346. // ==========================================
  347. // 5 events from 2. and 4. + 9 events from 6. = 14
  348. assert_dags(&eg_instances, 14, &mut rng).await;
  349. // ============================================================
  350. // 8. Start a new node and try to sync the DAG from other peers
  351. // ============================================================
  352. {
  353. // Connect to N_CONNS random peers.
  354. let peer_indexes_to_connect: Vec<_> =
  355. peer_indexes.choose_multiple(&mut rng, N_CONNS).collect();
  356. let mut peers = vec![];
  357. for peer_index in peer_indexes_to_connect {
  358. let port = 13200 + peer_index;
  359. peers.push(Url::parse(&format!("tcp://127.0.0.1:{port}")).unwrap());
  360. }
  361. let event_graph = spawn_node(
  362. vec![Url::parse(&format!("tcp://127.0.0.1:{}", 13200 + N_NODES + 1)).unwrap()],
  363. peers,
  364. ex.clone(),
  365. )
  366. .await;
  367. eg_instances.push(event_graph.clone());
  368. event_graph.p2p.clone().start().await.unwrap();
  369. info!("Waiting 5s for new node connection");
  370. sleep(5).await;
  371. event_graph.sync_selected(1, false).await.unwrap();
  372. }
  373. // ============================================================
  374. // 9. Assert the new synced DAG has the same contents as others
  375. // ============================================================
  376. // 5 events from 2. and 4. + 9 events from 6. = 14
  377. assert_dags(&eg_instances, 14, &mut rng).await;
  378. // Stop the P2P network
  379. for eg in eg_instances.iter() {
  380. eg.p2p.clone().stop().await;
  381. }
  382. }
  383. #[test]
  384. #[ignore]
  385. fn eventgraph_chaotic_propagation() {
  386. test_body!(eventgraph_chaotic_propagation_real);
  387. }
  388. async fn eventgraph_chaotic_propagation_real(ex: Arc<Executor<'static>>) {
  389. let mut rng = rand::thread_rng();
  390. let peer_indexes: Vec<usize> = (0..N_NODES).collect();
  391. let n_events: usize = 100000;
  392. // Bootstrap nodes
  393. let mut eg_instances = bootstrap_nodes(&peer_indexes, 14200, &mut rng, ex.clone()).await;
  394. // =========================================
  395. // 1. Assert that everyone's DAG is the same
  396. // =========================================
  397. assert_dags(&eg_instances, 1, &mut rng).await;
  398. // ===========================================
  399. // 2. Create multiple events on multiple nodes
  400. for i in 0..n_events {
  401. let random_node = eg_instances.choose(&mut rng).unwrap();
  402. let event = Event::new(i.to_be_bytes().to_vec(), random_node).await;
  403. let current_genesis = random_node.current_genesis.read().await;
  404. let dag_name = current_genesis.id().to_string();
  405. random_node.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  406. random_node.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  407. random_node.p2p.broadcast(&EventPut(event)).await;
  408. }
  409. info!("Waiting 5s for events propagation");
  410. sleep(5).await;
  411. // ==========================================
  412. // 3. Assert that everyone has all the events
  413. // ==========================================
  414. assert_dags(&eg_instances, n_events + 1, &mut rng).await;
  415. // ============================================================
  416. // 4. Start a new node and try to sync the DAG from other peers
  417. // ============================================================
  418. {
  419. // Connect to N_CONNS random peers.
  420. let peer_indexes_to_connect: Vec<_> =
  421. peer_indexes.choose_multiple(&mut rng, N_CONNS).collect();
  422. let mut peers = vec![];
  423. for peer_index in peer_indexes_to_connect {
  424. let port = 14200 + peer_index;
  425. peers.push(Url::parse(&format!("tcp://127.0.0.1:{port}")).unwrap());
  426. }
  427. let event_graph = spawn_node(
  428. vec![Url::parse(&format!("tcp://127.0.0.1:{}", 14200 + N_NODES + 1)).unwrap()],
  429. peers,
  430. ex.clone(),
  431. )
  432. .await;
  433. eg_instances.push(event_graph.clone());
  434. event_graph.p2p.clone().start().await.unwrap();
  435. info!("Waiting 5s for new node connection");
  436. sleep(5).await;
  437. event_graph.sync_selected(2, false).await.unwrap()
  438. }
  439. // ============================================================
  440. // 5. Assert the new synced DAG has the same contents as others
  441. // ============================================================
  442. assert_dags(&eg_instances, n_events + 1, &mut rng).await;
  443. // Stop the P2P network
  444. for eg in eg_instances.iter() {
  445. eg.p2p.clone().stop().await;
  446. }
  447. }