tests.rs 16 KB

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