tests.rs 33 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883
  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::{
  20. collections::{BTreeMap, HashMap, HashSet},
  21. slice,
  22. sync::Arc,
  23. time::{Duration, UNIX_EPOCH},
  24. };
  25. use darkfi_serial::{deserialize_async, serialize_async};
  26. use rand::{prelude::SliceRandom, rngs::ThreadRng};
  27. use sled_overlay::sled;
  28. use smol::{channel, future, Executor};
  29. use tracing::{info, warn};
  30. use url::Url;
  31. use crate::{
  32. error::Result,
  33. event_graph::{
  34. event::Header,
  35. proto::{EventPut, ProtocolEventGraph},
  36. util::next_rotation_timestamp,
  37. DAGStore, Event, EventGraph, EventGraphPtr, DAGS_MAX_NUMBER, GENESIS_CONTENTS,
  38. INITIAL_GENESIS, NULL_ID, N_EVENT_PARENTS,
  39. },
  40. net::{session::SESSION_DEFAULT, settings::NetworkProfile, P2p, Settings},
  41. system::{msleep, sleep, timeout::timeout},
  42. util::logger::{setup_test_logger, Level},
  43. Error,
  44. };
  45. // Number of nodes to spawn and number of peers each node connects to
  46. const N_NODES: usize = 5;
  47. const N_CONNS: usize = 2;
  48. //const N_NODES: usize = 50;
  49. //const N_CONNS: usize = N_NODES / 3;
  50. fn init_logger() {
  51. let ignored_targets = [
  52. "sled",
  53. "net::protocol_ping",
  54. "net::channel::subscribe_stop()",
  55. "net::hosts",
  56. "net::session",
  57. "net::message_subscriber",
  58. "net::protocol_address",
  59. "net::protocol_version",
  60. "net::protocol_registry",
  61. "net::channel::send()",
  62. "net::channel::start()",
  63. "net::channel::subscribe_msg()",
  64. "net::channel::main_receive_loop()",
  65. "net::tcp",
  66. ];
  67. // We check this error so we can execute same file tests in parallel,
  68. // otherwise second one fails to init logger here.
  69. if setup_test_logger(
  70. &ignored_targets,
  71. false,
  72. Level::Info,
  73. //Level::Verbose,
  74. //Level::Debug,
  75. //Level::Tracing,
  76. )
  77. .is_err()
  78. {
  79. warn!(target: "test_harness", "Logger already initialized");
  80. }
  81. }
  82. async fn spawn_node(
  83. inbound_addrs: Vec<Url>,
  84. peers: Vec<Url>,
  85. ex: Arc<Executor<'static>>,
  86. ) -> Arc<EventGraph> {
  87. let mut profiles = HashMap::new();
  88. profiles.insert(
  89. "tcp".to_string(),
  90. NetworkProfile { outbound_connect_timeout: 2, ..Default::default() },
  91. );
  92. let settings = Settings {
  93. localnet: true,
  94. inbound_addrs,
  95. outbound_connections: 0,
  96. inbound_connections: usize::MAX,
  97. peers,
  98. active_profiles: vec!["tcp".to_string()],
  99. profiles,
  100. ..Default::default()
  101. };
  102. let p2p = P2p::new(settings, ex.clone()).await.unwrap();
  103. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  104. let event_graph =
  105. EventGraph::new(p2p.clone(), sled_db, "/tmp".into(), false, false, 1, ex.clone())
  106. .await
  107. .unwrap();
  108. *event_graph.synced.write().await = true;
  109. let event_graph_ = event_graph.clone();
  110. // Register the P2P protocols
  111. let registry = p2p.protocol_registry();
  112. registry
  113. .register(SESSION_DEFAULT, move |channel, _| {
  114. let event_graph_ = event_graph_.clone();
  115. async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
  116. })
  117. .await;
  118. event_graph
  119. }
  120. async fn bootstrap_nodes(
  121. peer_indexes: &[usize],
  122. starting_port: usize,
  123. rng: &mut ThreadRng,
  124. ex: Arc<Executor<'static>>,
  125. ) -> Vec<Arc<EventGraph>> {
  126. let mut eg_instances = vec![];
  127. // Initialize the nodes
  128. for i in 0..N_NODES {
  129. // Everyone will connect to N_CONNS random peers.
  130. let mut peer_indexes_copy = peer_indexes.to_owned();
  131. peer_indexes_copy.remove(i);
  132. let peer_indexes_to_connect: Vec<_> =
  133. peer_indexes_copy.choose_multiple(rng, N_CONNS).collect();
  134. let mut peers = vec![];
  135. for peer_index in peer_indexes_to_connect {
  136. let port = starting_port + peer_index;
  137. peers.push(Url::parse(&format!("tcp://127.0.0.1:{port}")).unwrap());
  138. }
  139. let event_graph = spawn_node(
  140. vec![Url::parse(&format!("tcp://127.0.0.1:{}", starting_port + i)).unwrap()],
  141. peers,
  142. ex.clone(),
  143. )
  144. .await;
  145. eg_instances.push(event_graph);
  146. }
  147. // Start the P2P network
  148. for eg in eg_instances.iter() {
  149. eg.p2p.clone().start().await.unwrap();
  150. }
  151. info!("Waiting 5s until all peers connect");
  152. sleep(5).await;
  153. eg_instances
  154. }
  155. async fn assert_dags(eg_instances: &[Arc<EventGraph>], expected_len: usize, rng: &mut ThreadRng) {
  156. let random_node = eg_instances.choose(rng).unwrap();
  157. let random_node_genesis = random_node.current_genesis.read().await.header.timestamp;
  158. let store = random_node.dag_store.read().await;
  159. let (_, unreferenced_tips) = store.main_dags.get(&random_node_genesis).unwrap();
  160. let last_layer_tips = unreferenced_tips.last_key_value().unwrap().1.clone();
  161. for (i, eg) in eg_instances.iter().enumerate() {
  162. let current_genesis = eg.current_genesis.read().await;
  163. let dag_name = current_genesis.header.timestamp.to_string();
  164. let dag = eg.dag_store.read().await.get_dag(&dag_name);
  165. let unreferenced_tips = eg.dag_store.read().await.find_unreferenced_tips(&dag).await;
  166. let node_last_layer_tips = unreferenced_tips.last_key_value().unwrap().1.clone();
  167. assert!(
  168. dag.len() == expected_len,
  169. "Node {i}, expected {expected_len} events, have {}",
  170. dag.len()
  171. );
  172. assert_eq!(
  173. node_last_layer_tips, last_layer_tips,
  174. "Node {i} contains malformed unreferenced tips"
  175. );
  176. }
  177. }
  178. macro_rules! test_body {
  179. ($real_call:ident) => {
  180. init_logger();
  181. let ex = Arc::new(Executor::new());
  182. let ex_ = ex.clone();
  183. let (signal, shutdown) = channel::unbounded::<()>();
  184. // Run a thread for each node.
  185. easy_parallel::Parallel::new()
  186. .each(0..N_NODES, |_| future::block_on(ex.run(shutdown.recv())))
  187. .finish(|| {
  188. future::block_on(async {
  189. $real_call(ex_).await;
  190. drop(signal);
  191. })
  192. });
  193. };
  194. }
  195. #[test]
  196. fn eventgraph_propagation() {
  197. test_body!(eventgraph_propagation_real);
  198. }
  199. async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
  200. let mut rng = rand::thread_rng();
  201. let peer_indexes: Vec<usize> = (0..N_NODES).collect();
  202. // Bootstrap nodes
  203. let mut eg_instances = bootstrap_nodes(&peer_indexes, 13200, &mut rng, ex.clone()).await;
  204. // Grab genesis event
  205. let random_node = eg_instances.choose(&mut rng).unwrap();
  206. let current_genesis = random_node.current_genesis.read().await;
  207. let dag_name = current_genesis.header.timestamp.to_string();
  208. let (id, _) = random_node.dag_store.read().await.get_dag(&dag_name).last().unwrap().unwrap();
  209. let genesis_event_id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
  210. drop(current_genesis);
  211. // =========================================
  212. // 1. Assert that everyone's DAG is the same
  213. // =========================================
  214. assert_dags(&eg_instances, 1, &mut rng).await;
  215. // ==========================================
  216. // 2. Create an event in one node and publish
  217. // ==========================================
  218. let random_node = eg_instances.choose(&mut rng).unwrap();
  219. let current_genesis = random_node.current_genesis.read().await;
  220. let dag_name = current_genesis.header.timestamp.to_string();
  221. let event = Event::new(vec![1, 2, 3, 4], random_node).await;
  222. assert!(event.header.parents.contains(&genesis_event_id));
  223. // The node adds it to their DAG, on layer 1.
  224. random_node.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  225. let event_id = random_node.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap()[0];
  226. let store = random_node.dag_store.read().await;
  227. let (_, tips_layers) = store.header_dags.get(&current_genesis.header.timestamp).unwrap();
  228. // Since genesis was referenced, its layer (0) have been removed
  229. assert_eq!(tips_layers.len(), 1);
  230. assert!(tips_layers.last_key_value().unwrap().1.get(&event_id).is_some());
  231. drop(store);
  232. drop(current_genesis);
  233. info!("Broadcasting event {event_id}");
  234. random_node.p2p.broadcast(&EventPut(event, vec![])).await;
  235. info!("Waiting 5s for event propagation");
  236. sleep(5).await;
  237. // ====================================================
  238. // 3. Assert that everyone has the new event in the DAG
  239. // ====================================================
  240. assert_dags(&eg_instances, 2, &mut rng).await;
  241. // ==============================================================
  242. // 4. Create multiple events on a node and broadcast the last one
  243. // The `EventPut` logic should manage to fetch all of them,
  244. // provided that the last one references the earlier ones.
  245. // ==============================================================
  246. let random_node = eg_instances.choose(&mut rng).unwrap();
  247. let event0 = Event::new(vec![1, 2, 3, 4, 0], random_node).await;
  248. random_node.header_dag_insert(vec![event0.header.clone()], &dag_name).await.unwrap();
  249. let event0_id = random_node.dag_insert(slice::from_ref(&event0), &dag_name).await.unwrap()[0];
  250. let event1 = Event::new(vec![1, 2, 3, 4, 1], random_node).await;
  251. random_node.header_dag_insert(vec![event1.header.clone()], &dag_name).await.unwrap();
  252. let event1_id = random_node.dag_insert(slice::from_ref(&event1), &dag_name).await.unwrap()[0];
  253. let event2 = Event::new(vec![1, 2, 3, 4, 2], random_node).await;
  254. random_node.header_dag_insert(vec![event2.header.clone()], &dag_name).await.unwrap();
  255. let event2_id = random_node.dag_insert(slice::from_ref(&event2), &dag_name).await.unwrap()[0];
  256. // Genesis event + event from 2. + upper 3 events (layer 4)
  257. let current_genesis = random_node.current_genesis.read().await;
  258. let dag_name = current_genesis.header.timestamp.to_string();
  259. assert_eq!(random_node.dag_store.read().await.get_dag(&dag_name).len(), 5);
  260. let random_node_genesis = random_node.current_genesis.read().await.header.timestamp;
  261. let store = random_node.dag_store.read().await;
  262. let (_, tips_layers) = store.header_dags.get(&random_node_genesis).unwrap();
  263. assert_eq!(tips_layers.len(), 1);
  264. assert!(tips_layers.get(&4).unwrap().get(&event2_id).is_some());
  265. drop(current_genesis);
  266. drop(store);
  267. let event_chain = vec![
  268. (event0_id, event0.header.parents),
  269. (event1_id, event1.header.parents),
  270. (event2_id, event2.header.parents),
  271. ];
  272. info!("Broadcasting event {event2_id}");
  273. info!("Event chain: {event_chain:#?}");
  274. random_node.p2p.broadcast(&EventPut(event2, vec![])).await;
  275. info!("Waiting 5s for event propagation");
  276. sleep(5).await;
  277. // ==========================================
  278. // 5. Assert that everyone has all the events
  279. // ==========================================
  280. assert_dags(&eg_instances, 5, &mut rng).await;
  281. // ===========================================
  282. // 6. Create multiple events on multiple nodes
  283. // ===========================================
  284. // node 1
  285. // =======
  286. let node1 = eg_instances.choose(&mut rng).unwrap();
  287. let event0_1 = Event::new(vec![1, 2, 3, 4, 3], node1).await;
  288. node1.header_dag_insert(vec![event0_1.header.clone()], &dag_name).await.unwrap();
  289. node1.dag_insert(slice::from_ref(&event0_1), &dag_name).await.unwrap();
  290. node1.p2p.broadcast(&EventPut(event0_1, vec![])).await;
  291. msleep(300).await;
  292. let event1_1 = Event::new(vec![1, 2, 3, 4, 4], node1).await;
  293. node1.header_dag_insert(vec![event1_1.header.clone()], &dag_name).await.unwrap();
  294. node1.dag_insert(slice::from_ref(&event1_1), &dag_name).await.unwrap();
  295. node1.p2p.broadcast(&EventPut(event1_1, vec![])).await;
  296. msleep(300).await;
  297. let event2_1 = Event::new(vec![1, 2, 3, 4, 5], node1).await;
  298. node1.header_dag_insert(vec![event2_1.header.clone()], &dag_name).await.unwrap();
  299. node1.dag_insert(slice::from_ref(&event2_1), &dag_name).await.unwrap();
  300. node1.p2p.broadcast(&EventPut(event2_1, vec![])).await;
  301. msleep(300).await;
  302. // =======
  303. // node 2
  304. // =======
  305. let node2 = eg_instances.choose(&mut rng).unwrap();
  306. let event0_2 = Event::new(vec![1, 2, 3, 4, 6], node2).await;
  307. node2.header_dag_insert(vec![event0_2.header.clone()], &dag_name).await.unwrap();
  308. node2.dag_insert(slice::from_ref(&event0_2), &dag_name).await.unwrap();
  309. node2.p2p.broadcast(&EventPut(event0_2, vec![])).await;
  310. msleep(300).await;
  311. let event1_2 = Event::new(vec![1, 2, 3, 4, 7], node2).await;
  312. node2.header_dag_insert(vec![event1_2.header.clone()], &dag_name).await.unwrap();
  313. node2.dag_insert(slice::from_ref(&event1_2), &dag_name).await.unwrap();
  314. node2.p2p.broadcast(&EventPut(event1_2, vec![])).await;
  315. msleep(300).await;
  316. let event2_2 = Event::new(vec![1, 2, 3, 4, 8], node2).await;
  317. node2.header_dag_insert(vec![event2_2.header.clone()], &dag_name).await.unwrap();
  318. node2.dag_insert(slice::from_ref(&event2_2), &dag_name).await.unwrap();
  319. node2.p2p.broadcast(&EventPut(event2_2, vec![])).await;
  320. msleep(300).await;
  321. // =======
  322. // node 3
  323. // =======
  324. let node3 = eg_instances.choose(&mut rng).unwrap();
  325. let event0_3 = Event::new(vec![1, 2, 3, 4, 9], node3).await;
  326. node3.header_dag_insert(vec![event0_3.header.clone()], &dag_name).await.unwrap();
  327. node3.dag_insert(slice::from_ref(&event0_3), &dag_name).await.unwrap();
  328. node3.p2p.broadcast(&EventPut(event0_3, vec![])).await;
  329. msleep(300).await;
  330. let event1_3 = Event::new(vec![1, 2, 3, 4, 10], node3).await;
  331. node3.header_dag_insert(vec![event1_3.header.clone()], &dag_name).await.unwrap();
  332. node3.dag_insert(slice::from_ref(&event1_3), &dag_name).await.unwrap();
  333. node3.p2p.broadcast(&EventPut(event1_3, vec![])).await;
  334. msleep(300).await;
  335. let event2_3 = Event::new(vec![1, 2, 3, 4, 11], node3).await;
  336. node3.header_dag_insert(vec![event2_3.header.clone()], &dag_name).await.unwrap();
  337. node3.dag_insert(slice::from_ref(&event2_3), &dag_name).await.unwrap();
  338. node3.p2p.broadcast(&EventPut(event2_3, vec![])).await;
  339. msleep(300).await;
  340. // /////
  341. // //
  342. // let node4 = eg_instances.choose(&mut rng).unwrap();
  343. // let event0_4 = Event::new(vec![1, 2, 3, 4, 12], node4).await;
  344. // node4.dag_insert(&[event0_4.clone()]).await.unwrap();
  345. // node4.p2p.broadcast(&EventPut(event0_4)).await;
  346. // sleep(1).await;
  347. // let event1_4 = Event::new(vec![1, 2, 3, 4, 13], node4).await;
  348. // node4.dag_insert(&[event1_4.clone()]).await.unwrap();
  349. // node4.p2p.broadcast(&EventPut(event1_4)).await;
  350. // sleep(1).await;
  351. // let event2_4 = Event::new(vec![1, 2, 3, 4, 14], node4).await;
  352. // node4.dag_insert(&[event2_4.clone()]).await.unwrap();
  353. // node4.p2p.broadcast(&EventPut(event2_4)).await;
  354. // // sleep(1).await;
  355. // ==========================================
  356. // 7. Assert that everyone has all the events
  357. // ==========================================
  358. // 5 events from 2. and 4. + 9 events from 6. = 14
  359. assert_dags(&eg_instances, 14, &mut rng).await;
  360. // ============================================================
  361. // 8. Start a new node and try to sync the DAG from other peers
  362. // ============================================================
  363. {
  364. // Connect to N_CONNS random peers.
  365. let peer_indexes_to_connect: Vec<_> =
  366. peer_indexes.choose_multiple(&mut rng, N_CONNS).collect();
  367. let mut peers = vec![];
  368. for peer_index in peer_indexes_to_connect {
  369. let port = 13200 + peer_index;
  370. peers.push(Url::parse(&format!("tcp://127.0.0.1:{port}")).unwrap());
  371. }
  372. let event_graph = spawn_node(
  373. vec![Url::parse(&format!("tcp://127.0.0.1:{}", 13200 + N_NODES + 1)).unwrap()],
  374. peers,
  375. ex.clone(),
  376. )
  377. .await;
  378. eg_instances.push(event_graph.clone());
  379. event_graph.p2p.clone().start().await.unwrap();
  380. info!("Waiting 5s for new node connection");
  381. sleep(5).await;
  382. event_graph.sync_selected(1, false).await.unwrap();
  383. }
  384. // ============================================================
  385. // 9. Assert the new synced DAG has the same contents as others
  386. // ============================================================
  387. // 5 events from 2. and 4. + 9 events from 6. = 14
  388. assert_dags(&eg_instances, 14, &mut rng).await;
  389. // Stop the P2P network
  390. for eg in eg_instances.iter() {
  391. eg.p2p.clone().stop().await;
  392. }
  393. }
  394. #[test]
  395. #[ignore]
  396. fn eventgraph_chaotic_propagation() {
  397. test_body!(eventgraph_chaotic_propagation_real);
  398. }
  399. async fn eventgraph_chaotic_propagation_real(ex: Arc<Executor<'static>>) {
  400. let mut rng = rand::thread_rng();
  401. let peer_indexes: Vec<usize> = (0..N_NODES).collect();
  402. let n_events: usize = 100000;
  403. // Bootstrap nodes
  404. let mut eg_instances = bootstrap_nodes(&peer_indexes, 14200, &mut rng, ex.clone()).await;
  405. // =========================================
  406. // 1. Assert that everyone's DAG is the same
  407. // =========================================
  408. assert_dags(&eg_instances, 1, &mut rng).await;
  409. // ===========================================
  410. // 2. Create multiple events on multiple nodes
  411. for i in 0..n_events {
  412. let random_node = eg_instances.choose(&mut rng).unwrap();
  413. let event = Event::new(i.to_be_bytes().to_vec(), random_node).await;
  414. let current_genesis = random_node.current_genesis.read().await;
  415. let dag_name = current_genesis.header.timestamp.to_string();
  416. random_node.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
  417. random_node.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
  418. random_node.p2p.broadcast(&EventPut(event, vec![])).await;
  419. }
  420. info!("Waiting 5s for events propagation");
  421. sleep(5).await;
  422. // ==========================================
  423. // 3. Assert that everyone has all the events
  424. // ==========================================
  425. assert_dags(&eg_instances, n_events + 1, &mut rng).await;
  426. // ============================================================
  427. // 4. Start a new node and try to sync the DAG from other peers
  428. // ============================================================
  429. {
  430. // Connect to N_CONNS random peers.
  431. let peer_indexes_to_connect: Vec<_> =
  432. peer_indexes.choose_multiple(&mut rng, N_CONNS).collect();
  433. let mut peers = vec![];
  434. for peer_index in peer_indexes_to_connect {
  435. let port = 14200 + peer_index;
  436. peers.push(Url::parse(&format!("tcp://127.0.0.1:{port}")).unwrap());
  437. }
  438. let event_graph = spawn_node(
  439. vec![Url::parse(&format!("tcp://127.0.0.1:{}", 14200 + N_NODES + 1)).unwrap()],
  440. peers,
  441. ex.clone(),
  442. )
  443. .await;
  444. eg_instances.push(event_graph.clone());
  445. event_graph.p2p.clone().start().await.unwrap();
  446. info!("Waiting 5s for new node connection");
  447. sleep(5).await;
  448. event_graph.sync_selected(2, false).await.unwrap()
  449. }
  450. // ============================================================
  451. // 5. Assert the new synced DAG has the same contents as others
  452. // ============================================================
  453. assert_dags(&eg_instances, n_events + 1, &mut rng).await;
  454. // Stop the P2P network
  455. for eg in eg_instances.iter() {
  456. eg.p2p.clone().stop().await;
  457. }
  458. }
  459. // DAGStore tests
  460. async fn make_dag_store() -> Result<DAGStore> {
  461. let sled_db = sled::Config::new().temporary(true).open()?;
  462. let hours_rotation = 1;
  463. let dag_store = DAGStore {
  464. db: sled_db.clone(),
  465. header_dags: BTreeMap::default(),
  466. main_dags: BTreeMap::default(),
  467. }
  468. .new(sled_db.clone(), hours_rotation)
  469. .await;
  470. Ok(dag_store)
  471. }
  472. #[test]
  473. fn header_dags_and_main_dags_length_equals_dags_max_number() -> Result<()> {
  474. smol::block_on(async {
  475. let dag_store = make_dag_store().await?;
  476. assert_eq!(dag_store.header_dags.len() as i8, DAGS_MAX_NUMBER);
  477. assert_eq!(dag_store.main_dags.len() as i8, DAGS_MAX_NUMBER);
  478. Ok(())
  479. })
  480. }
  481. #[test]
  482. fn all_dag_trees_are_created_on_sled_after_dag_store_creation() -> Result<()> {
  483. smol::block_on(async {
  484. let dag_store = make_dag_store().await?;
  485. let dag_trees: Vec<String> =
  486. dag_store.db.tree_names().iter().map(|n| String::from_utf8_lossy(n).into()).collect();
  487. // Should have 2 * DAGS_MAX_NUMBER trees + 1 (the default tree)
  488. assert_eq!(dag_trees.len() as i8, DAGS_MAX_NUMBER * 2 + 1);
  489. for (dag_timestamp, _) in dag_store.header_dags {
  490. assert!(dag_trees.contains(&format!("headers_{dag_timestamp}")));
  491. }
  492. for (dag_timestamp, _) in dag_store.main_dags {
  493. assert!(dag_trees.contains(&dag_timestamp.to_string()));
  494. }
  495. Ok(())
  496. })
  497. }
  498. #[test]
  499. fn genesis_events_or_headers_are_added_to_all_trees_and_utips() -> Result<()> {
  500. smol::block_on(async {
  501. let dag_store = make_dag_store().await?;
  502. for (_, (tree, layer_utips)) in dag_store.header_dags {
  503. let genesis_header = tree.first()?;
  504. // A Genesis Header is found in sled tree
  505. assert!(genesis_header.is_some());
  506. let (genesis_hash, genesis_header) = genesis_header.unwrap();
  507. let genesis_header: Header = deserialize_async(&genesis_header).await?;
  508. let genesis_hash: blake3::Hash = deserialize_async(&genesis_hash).await?;
  509. assert_eq!(genesis_header.layer, 0);
  510. assert!(genesis_header.parents.iter().all(|p| *p == NULL_ID));
  511. // The Genesis Header hash is stored as Unreferenced tip
  512. assert!(layer_utips.contains_key(&0));
  513. assert!(layer_utips.get(&0).unwrap().contains(&genesis_hash));
  514. }
  515. for (_, (tree, layer_utips)) in dag_store.main_dags {
  516. let genesis_event = tree.first()?;
  517. // A Genesis Event is found in sled tree
  518. assert!(genesis_event.is_some());
  519. let (genesis_hash, genesis_event) = genesis_event.unwrap();
  520. let genesis_event: Event = deserialize_async(&genesis_event).await?;
  521. let genesis_hash: blake3::Hash = deserialize_async(&genesis_hash).await?;
  522. assert_eq!(genesis_event.header.layer, 0);
  523. assert!(genesis_event.header.parents.iter().all(|p| *p == NULL_ID));
  524. assert_eq!(genesis_event.content, GENESIS_CONTENTS);
  525. // The Genesis Header hash is stored as Unreferenced tip
  526. assert!(layer_utips.contains_key(&0));
  527. assert!(layer_utips.get(&0).unwrap().contains(&genesis_hash));
  528. }
  529. Ok(())
  530. })
  531. }
  532. #[test]
  533. fn adding_new_dag_removes_oldest_dag_tree() -> Result<()> {
  534. smol::block_on(async {
  535. let mut dag_store = make_dag_store().await?;
  536. let oldest_dag_timestamp = dag_store.main_dags.first_key_value().unwrap().0.to_owned();
  537. // Next dag to add
  538. let next_rotation = next_rotation_timestamp(INITIAL_GENESIS, 1);
  539. let header =
  540. Header { timestamp: next_rotation, parents: [NULL_ID; N_EVENT_PARENTS], layer: 0 };
  541. let next_genesis = Event { header, content: GENESIS_CONTENTS.to_vec() };
  542. dag_store.add_dag(&next_genesis.header.timestamp.to_string(), &next_genesis).await;
  543. // The length of the dags should stay the same after adding
  544. assert_eq!(dag_store.main_dags.len() as i8, DAGS_MAX_NUMBER);
  545. assert_eq!(dag_store.header_dags.len() as i8, DAGS_MAX_NUMBER);
  546. // We should have an entry with the new dag timestamp
  547. assert!(dag_store.main_dags.contains_key(&next_rotation));
  548. assert!(dag_store.header_dags.contains_key(&next_rotation));
  549. // The oldest dag entry should have been removed
  550. assert!(!dag_store.main_dags.contains_key(&oldest_dag_timestamp));
  551. assert!(!dag_store.header_dags.contains_key(&oldest_dag_timestamp));
  552. let dag_trees: Vec<String> =
  553. dag_store.db.tree_names().iter().map(|n| String::from_utf8_lossy(n).into()).collect();
  554. // The number of dag trees should stay the same after adding
  555. assert_eq!(dag_trees.len() as i8, 2 * DAGS_MAX_NUMBER + 1);
  556. // We should have a tree with the new dag timestamp value
  557. assert!(dag_trees.contains(&next_rotation.to_string()));
  558. assert!(dag_trees.contains(&format!("headers_{next_rotation}")));
  559. // The oldest dag sled tree should have been removed
  560. assert!(!dag_trees.contains(&oldest_dag_timestamp.to_string()));
  561. assert!(!dag_trees.contains(&format!("headers_{oldest_dag_timestamp}")));
  562. Ok(())
  563. })
  564. }
  565. #[test]
  566. fn sort_moves_current_dag_to_front() -> Result<()> {
  567. smol::block_on(async {
  568. let dag_store = make_dag_store().await?;
  569. let trees = dag_store.sort_dags().await;
  570. let first_tree_name: String =
  571. String::from_utf8_lossy(&trees.first().unwrap().name()).into();
  572. assert_eq!(
  573. first_tree_name,
  574. dag_store.main_dags.last_key_value().unwrap().0.to_owned().to_string()
  575. );
  576. Ok(())
  577. })
  578. }
  579. #[test]
  580. fn unreferenced_tips_are_found() -> Result<()> {
  581. smol::block_on(async {
  582. let dag_store = make_dag_store().await?;
  583. let current_dag_tree = dag_store.main_dags.last_key_value().unwrap().1 .0.clone();
  584. let current_dag_genesis_hash = *dag_store
  585. .main_dags
  586. .last_key_value()
  587. .unwrap()
  588. .1
  589. .1
  590. .get(&0)
  591. .unwrap()
  592. .iter()
  593. .next()
  594. .unwrap();
  595. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  596. parents[0] = current_dag_genesis_hash;
  597. let event2 = Event {
  598. header: Header {
  599. timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64,
  600. parents,
  601. layer: 1,
  602. },
  603. content: "event2".as_bytes().to_vec(),
  604. };
  605. let event2_hash = event2.id();
  606. current_dag_tree.insert(event2_hash.as_bytes(), serialize_async(&event2).await)?;
  607. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  608. parents[0] = event2_hash;
  609. let event3 = Event {
  610. header: Header {
  611. timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64,
  612. parents,
  613. layer: 2,
  614. },
  615. content: "event3".as_bytes().to_vec(),
  616. };
  617. let event3_hash = event3.id();
  618. current_dag_tree.insert(event3_hash.as_bytes(), serialize_async(&event3).await)?;
  619. let mut parents = [NULL_ID; N_EVENT_PARENTS];
  620. parents[0] = current_dag_genesis_hash;
  621. let event4 = Event {
  622. header: Header {
  623. timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64,
  624. parents,
  625. layer: 2,
  626. },
  627. content: "event4".as_bytes().to_vec(),
  628. };
  629. let event4_hash = event4.id();
  630. current_dag_tree.insert(event4_hash.as_bytes(), serialize_async(&event4).await)?;
  631. let layer_utips = dag_store.find_unreferenced_tips(&current_dag_tree).await;
  632. // We have unreferenced tips only on the 2nd layer
  633. assert_eq!(layer_utips.len(), 1);
  634. // We have two unreferenced tips
  635. let tip_hashes = layer_utips.get(&2).unwrap();
  636. assert_eq!(tip_hashes.len(), 2);
  637. // Event3 and Event4 are the only unreferenced tips
  638. assert!(tip_hashes.contains(&event3_hash));
  639. assert!(tip_hashes.contains(&event4_hash));
  640. Ok(())
  641. })
  642. }
  643. // EventGraph tests
  644. async fn make_event_graph() -> Result<EventGraphPtr> {
  645. let ex = Arc::new(Executor::new());
  646. let p2p = P2p::new(Settings::default(), ex.clone()).await?;
  647. let sled_db = sled::Config::new().temporary(true).open()?;
  648. EventGraph::new(p2p, sled_db, "/tmp".into(), false, false, 1, ex).await
  649. }
  650. #[test]
  651. fn dag_insert_on_invalid_dag_name() -> Result<()> {
  652. smol::block_on(async {
  653. let event_graph = make_event_graph().await?;
  654. let new_event = Event::new("new_event".as_bytes().to_vec(), &event_graph).await;
  655. // Using a dag name that is not a u64 timestamp gives an error
  656. let res = event_graph.dag_insert(&[new_event], "non_timestamp_dag_name").await;
  657. assert!(res.is_err());
  658. let err = res.unwrap_err();
  659. match err {
  660. Error::ParseIntError(_) => {}
  661. _ => panic!("expected parse error"),
  662. }
  663. Ok(())
  664. })
  665. }
  666. #[test]
  667. fn invalid_header_dag_insert() -> Result<()> {
  668. smol::block_on(async {
  669. let event_graph = make_event_graph().await?;
  670. let dag_name = event_graph
  671. .dag_store
  672. .read()
  673. .await
  674. .main_dags
  675. .last_key_value()
  676. .unwrap()
  677. .0
  678. .clone()
  679. .to_string();
  680. let new_event = Event::new("new_event".as_bytes().to_vec(), &event_graph).await;
  681. // Inserting an invalid event gives an error
  682. let mut event_timestamp_too_old = new_event.clone();
  683. event_timestamp_too_old.header.timestamp = 1000;
  684. let res =
  685. event_graph.header_dag_insert(vec![event_timestamp_too_old.header], &dag_name).await;
  686. assert!(res.is_err());
  687. let err = res.unwrap_err();
  688. match err {
  689. Error::HeaderIsInvalid => {}
  690. _ => panic!("expected invalid header error"),
  691. }
  692. Ok(())
  693. })
  694. }
  695. #[test]
  696. fn dag_insert_without_inserting_header() -> Result<()> {
  697. smol::block_on(async {
  698. let event_graph = make_event_graph().await?;
  699. let dag_name = event_graph
  700. .dag_store
  701. .read()
  702. .await
  703. .main_dags
  704. .last_key_value()
  705. .unwrap()
  706. .0
  707. .clone()
  708. .to_string();
  709. let new_event = Event::new("new_event".as_bytes().to_vec(), &event_graph).await;
  710. let res = event_graph.dag_insert(slice::from_ref(&new_event), &dag_name).await;
  711. // Inserting event without inserting its header first gets skipped
  712. assert!(res.is_ok() && res.unwrap().is_empty());
  713. Ok(())
  714. })
  715. }
  716. #[test]
  717. fn dag_insert_duplicate_event() -> Result<()> {
  718. smol::block_on(async {
  719. let event_graph = make_event_graph().await?;
  720. let dag_name = event_graph
  721. .dag_store
  722. .read()
  723. .await
  724. .main_dags
  725. .last_key_value()
  726. .unwrap()
  727. .0
  728. .clone()
  729. .to_string();
  730. let new_event = Event::new("new_event".as_bytes().to_vec(), &event_graph).await;
  731. event_graph.header_dag_insert(vec![new_event.header.clone()], &dag_name).await?;
  732. let res = event_graph.dag_insert(slice::from_ref(&new_event), &dag_name).await;
  733. // Proper insertion
  734. assert!(res.is_ok() && res.unwrap().len() == 1);
  735. // Inserting duplicate event gets skipped
  736. let res = event_graph.dag_insert(&[new_event], &dag_name).await;
  737. assert!(res.is_ok() && res.unwrap().is_empty());
  738. Ok(())
  739. })
  740. }
  741. #[test]
  742. fn dag_insert_valid_event() -> Result<()> {
  743. smol::block_on(async {
  744. let event_graph = make_event_graph().await?;
  745. let dag_name = *event_graph.dag_store.read().await.main_dags.last_key_value().unwrap().0;
  746. let new_event_sub = event_graph.event_pub.clone().subscribe().await;
  747. let new_event = Event::new("new_event".as_bytes().to_vec(), &event_graph).await;
  748. event_graph
  749. .header_dag_insert(vec![new_event.header.clone()], &dag_name.to_string())
  750. .await?;
  751. let res = event_graph.dag_insert(slice::from_ref(&new_event), &dag_name.to_string()).await;
  752. assert!(res.is_ok() && res.unwrap().len() == 1);
  753. // Unreferenced tips is updated
  754. let layer_utips =
  755. event_graph.dag_store.read().await.main_dags.get(&dag_name).unwrap().1.clone();
  756. assert!(layer_utips.get(&1).unwrap().contains(&new_event.id()));
  757. // The new event notification is sent to subscriber
  758. let dur = Duration::from_secs(1);
  759. let Ok(res) = timeout(dur, new_event_sub.receive()).await else {
  760. panic!("Event is not sent to subscriber")
  761. };
  762. assert_eq!(res.id(), new_event.id());
  763. Ok(())
  764. })
  765. }