tests.rs 7.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  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. use std::{collections::HashMap, sync::Arc};
  19. use log::info;
  20. use rand::{prelude::SliceRandom, Rng};
  21. use smol::{channel, future, Executor};
  22. use url::Url;
  23. use crate::{
  24. event_graph2::{
  25. proto::{EventPut, ProtocolEventGraph},
  26. Event, EventGraph, NULL_ID,
  27. },
  28. net::{P2p, Settings, SESSION_ALL},
  29. system::sleep,
  30. };
  31. /// Number of nodes to spawn
  32. const N_NODES: usize = 50;
  33. //const N_NODES: usize = 2;
  34. /// Number of peers each node connects to
  35. const N_CONNS: usize = N_NODES / 3;
  36. //const N_CONNS: usize = 1;
  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::message_subscriber".to_string());
  46. cfg.add_filter_ignore("net::protocol_address".to_string());
  47. cfg.add_filter_ignore("net::protocol_version".to_string());
  48. cfg.add_filter_ignore("net::channel::send()".to_string());
  49. simplelog::TermLogger::init(
  50. simplelog::LevelFilter::Info,
  51. //simplelog::LevelFilter::Debug,
  52. //simplelog::LevelFilter::Trace,
  53. cfg.build(),
  54. simplelog::TerminalMode::Mixed,
  55. simplelog::ColorChoice::Auto,
  56. )
  57. .unwrap();
  58. let ex = Arc::new(Executor::new());
  59. let ex_ = ex.clone();
  60. let (signal, shutdown) = channel::unbounded::<()>();
  61. // Run a thread for each node.
  62. easy_parallel::Parallel::new()
  63. .each(0..N_NODES, |_| future::block_on(ex.run(shutdown.recv())))
  64. .finish(|| {
  65. future::block_on(async {
  66. eventgraph_propagation_real(ex_).await;
  67. drop(signal);
  68. })
  69. });
  70. }
  71. async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
  72. let mut eg_instances = vec![];
  73. let mut rng = rand::thread_rng();
  74. let mut genesis_event_id = NULL_ID;
  75. // Initialize the nodes
  76. for i in 0..N_NODES {
  77. // Everyone will connect to N_CONNS random peers.
  78. let mut peers = vec![];
  79. for _ in 0..N_CONNS {
  80. let mut port = 13200 + i;
  81. while port == 13200 + i {
  82. port = 13200 + rng.gen_range(0..N_NODES);
  83. }
  84. peers.push(Url::parse(&format!("tcp://127.0.0.1:{}", port)).unwrap());
  85. }
  86. let settings = Settings {
  87. localnet: true,
  88. inbound_addrs: vec![Url::parse(&format!("tcp://127.0.0.1:{}", 13200 + i)).unwrap()],
  89. outbound_connections: 0,
  90. outbound_connect_timeout: 2,
  91. inbound_connections: usize::MAX,
  92. peers,
  93. allowed_transports: vec!["tcp".to_string()],
  94. ..Default::default()
  95. };
  96. let p2p = P2p::new(settings, ex.clone()).await;
  97. let sled_db = sled::Config::new().temporary(true).open().unwrap();
  98. let event_graph =
  99. EventGraph::new(p2p.clone(), &sled_db, "dag", 1, ex.clone()).await.unwrap();
  100. let event_graph_ = event_graph.clone();
  101. // Take the last sled item since there's only 1
  102. if genesis_event_id == NULL_ID {
  103. let (id, _) = event_graph.dag.last().unwrap().unwrap();
  104. genesis_event_id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
  105. }
  106. // Register the P2P protocols
  107. let registry = p2p.protocol_registry();
  108. registry
  109. .register(SESSION_ALL, move |channel, _| {
  110. let event_graph_ = event_graph_.clone();
  111. async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
  112. })
  113. .await;
  114. eg_instances.push(event_graph);
  115. }
  116. // Start the P2P network
  117. for eg in eg_instances.iter() {
  118. eg.p2p.clone().start().await.unwrap();
  119. }
  120. info!("Waiting 10s until all peers connect");
  121. sleep(10).await;
  122. // Now we get to the logic.
  123. //for i in 1..1001_u16 {
  124. for i in 1..3_u16 {
  125. // A random node creates an event
  126. let random_node = eg_instances.choose(&mut rand::thread_rng()).unwrap();
  127. let event = Event::new(i.to_le_bytes().to_vec(), random_node.clone()).await;
  128. // The node adds it to their DAG
  129. let event_id = random_node.dag_insert(&event).await.unwrap();
  130. // The node broadcasts it
  131. info!("Broadcasting {}", event_id);
  132. random_node.p2p.broadcast(&EventPut(event)).await;
  133. // Another random node creates five events and sends them out of order.
  134. let random_node = eg_instances.choose(&mut rand::thread_rng()).unwrap();
  135. let event0 = Event::new(i.to_le_bytes().to_vec(), random_node.clone()).await;
  136. let event0_id = random_node.dag_insert(&event0).await.unwrap();
  137. let event1 = Event::new(i.to_le_bytes().to_vec(), random_node.clone()).await;
  138. let event1_id = random_node.dag_insert(&event1).await.unwrap();
  139. let event2 = Event::new(i.to_le_bytes().to_vec(), random_node.clone()).await;
  140. let event2_id = random_node.dag_insert(&event2).await.unwrap();
  141. let event3 = Event::new(i.to_le_bytes().to_vec(), random_node.clone()).await;
  142. let event3_id = random_node.dag_insert(&event3).await.unwrap();
  143. let event4 = Event::new(i.to_le_bytes().to_vec(), random_node.clone()).await;
  144. let event4_id = random_node.dag_insert(&event4).await.unwrap();
  145. info!("Broadcasting {}", event3_id);
  146. random_node.p2p.broadcast(&EventPut(event3)).await;
  147. info!("Broadcasting {}", event2_id);
  148. random_node.p2p.broadcast(&EventPut(event2)).await;
  149. info!("Broadcasting {}", event4_id);
  150. random_node.p2p.broadcast(&EventPut(event4)).await;
  151. info!("Broadcasting {}", event1_id);
  152. random_node.p2p.broadcast(&EventPut(event1)).await;
  153. info!("Broadcasting {}", event0_id);
  154. random_node.p2p.broadcast(&EventPut(event0)).await;
  155. }
  156. info!("Waiting 20s until the p2p broadcasts settle");
  157. sleep(20).await;
  158. // Assert that everyone has the same DAG
  159. let mut contents = HashMap::new();
  160. for (i, eg) in eg_instances.iter().enumerate() {
  161. let mut ids = vec![];
  162. for r in eg.dag.iter() {
  163. let (id, _) = r.unwrap();
  164. ids.push(blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap()));
  165. }
  166. contents.insert(i, ids);
  167. }
  168. let value = contents.values().next().unwrap();
  169. assert!(contents.values().all(|v| v == value));
  170. // Assert that everyone's DAG sorts the same.
  171. let mut orders = HashMap::new();
  172. for (i, eg) in eg_instances.iter().enumerate() {
  173. let order = eg.order_events().await;
  174. orders.insert(i, order);
  175. }
  176. let value = orders.values().next().unwrap();
  177. for (i, order) in orders.iter() {
  178. assert!(
  179. order == value,
  180. "Node {} has wrong order:\n{:#?}\nvs{:#?}\nGENESIS:{}",
  181. i,
  182. order,
  183. value,
  184. genesis_event_id,
  185. );
  186. }
  187. // Stop the P2P network
  188. for eg in eg_instances.iter() {
  189. eg.p2p.clone().stop().await;
  190. }
  191. }