|
|
@@ -20,15 +20,15 @@
|
|
|
|
|
|
use std::sync::Arc;
|
|
|
|
|
|
-use log::info;
|
|
|
-use rand::{prelude::SliceRandom, Rng};
|
|
|
+use log::{info, warn};
|
|
|
+use rand::{prelude::SliceRandom, rngs::ThreadRng};
|
|
|
use smol::{channel, future, Executor};
|
|
|
use url::Url;
|
|
|
|
|
|
use crate::{
|
|
|
event_graph::{
|
|
|
proto::{EventPut, ProtocolEventGraph},
|
|
|
- Event, EventGraph, NULL_ID,
|
|
|
+ Event, EventGraph,
|
|
|
},
|
|
|
net::{P2p, Settings, SESSION_ALL},
|
|
|
system::sleep,
|
|
|
@@ -40,9 +40,7 @@ const N_CONNS: usize = 2;
|
|
|
//const N_NODES: usize = 50;
|
|
|
//const N_CONNS: usize = N_NODES / 3;
|
|
|
|
|
|
-#[test]
|
|
|
-#[ignore]
|
|
|
-fn eventgraph_propagation() {
|
|
|
+fn init_logger() {
|
|
|
let mut cfg = simplelog::ConfigBuilder::new();
|
|
|
cfg.add_filter_ignore("sled".to_string());
|
|
|
cfg.add_filter_ignore("net::protocol_ping".to_string());
|
|
|
@@ -56,81 +54,87 @@ fn eventgraph_propagation() {
|
|
|
cfg.add_filter_ignore("net::channel::send()".to_string());
|
|
|
cfg.add_filter_ignore("net::channel::start()".to_string());
|
|
|
cfg.add_filter_ignore("net::channel::subscribe_msg()".to_string());
|
|
|
-
|
|
|
- simplelog::TermLogger::init(
|
|
|
- //simplelog::LevelFilter::Info,
|
|
|
- simplelog::LevelFilter::Debug,
|
|
|
+ cfg.add_filter_ignore("net::channel::main_receive_loop()".to_string());
|
|
|
+ cfg.add_filter_ignore("net::tcp".to_string());
|
|
|
+
|
|
|
+ // We check this error so we can execute same file tests in parallel,
|
|
|
+ // otherwise second one fails to init logger here.
|
|
|
+ if simplelog::TermLogger::init(
|
|
|
+ simplelog::LevelFilter::Info,
|
|
|
+ //simplelog::LevelFilter::Debug,
|
|
|
//simplelog::LevelFilter::Trace,
|
|
|
cfg.build(),
|
|
|
simplelog::TerminalMode::Mixed,
|
|
|
simplelog::ColorChoice::Auto,
|
|
|
)
|
|
|
- .unwrap();
|
|
|
-
|
|
|
- let ex = Arc::new(Executor::new());
|
|
|
- let ex_ = ex.clone();
|
|
|
- let (signal, shutdown) = channel::unbounded::<()>();
|
|
|
-
|
|
|
- // Run a thread for each node.
|
|
|
- easy_parallel::Parallel::new()
|
|
|
- .each(0..N_NODES, |_| future::block_on(ex.run(shutdown.recv())))
|
|
|
- .finish(|| {
|
|
|
- future::block_on(async {
|
|
|
- eventgraph_propagation_real(ex_).await;
|
|
|
- drop(signal);
|
|
|
- })
|
|
|
- });
|
|
|
+ .is_err()
|
|
|
+ {
|
|
|
+ warn!(target: "test_harness", "Logger already initialized");
|
|
|
+ }
|
|
|
}
|
|
|
|
|
|
-async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
|
|
|
- let mut eg_instances = vec![];
|
|
|
- let mut rng = rand::thread_rng();
|
|
|
+async fn spawn_node(
|
|
|
+ inbound_addrs: Vec<Url>,
|
|
|
+ peers: Vec<Url>,
|
|
|
+ ex: Arc<Executor<'static>>,
|
|
|
+) -> Arc<EventGraph> {
|
|
|
+ let settings = Settings {
|
|
|
+ localnet: true,
|
|
|
+ inbound_addrs,
|
|
|
+ outbound_connections: 0,
|
|
|
+ outbound_connect_timeout: 2,
|
|
|
+ inbound_connections: usize::MAX,
|
|
|
+ peers,
|
|
|
+ allowed_transports: vec!["tcp".to_string()],
|
|
|
+ ..Default::default()
|
|
|
+ };
|
|
|
+
|
|
|
+ let p2p = P2p::new(settings, ex.clone()).await;
|
|
|
+ let sled_db = sled::Config::new().temporary(true).open().unwrap();
|
|
|
+ let event_graph = EventGraph::new(p2p.clone(), sled_db, "dag", 1, ex.clone()).await.unwrap();
|
|
|
+ *event_graph.synced.write().await = true;
|
|
|
+ let event_graph_ = event_graph.clone();
|
|
|
+
|
|
|
+ // Register the P2P protocols
|
|
|
+ let registry = p2p.protocol_registry();
|
|
|
+ registry
|
|
|
+ .register(SESSION_ALL, move |channel, _| {
|
|
|
+ let event_graph_ = event_graph_.clone();
|
|
|
+ async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
|
|
|
+ })
|
|
|
+ .await;
|
|
|
+
|
|
|
+ event_graph
|
|
|
+}
|
|
|
|
|
|
- let mut genesis_event_id = NULL_ID;
|
|
|
+async fn bootstrap_nodes(
|
|
|
+ peer_indexes: &[usize],
|
|
|
+ starting_port: usize,
|
|
|
+ rng: &mut ThreadRng,
|
|
|
+ ex: Arc<Executor<'static>>,
|
|
|
+) -> Vec<Arc<EventGraph>> {
|
|
|
+ let mut eg_instances = vec![];
|
|
|
|
|
|
// Initialize the nodes
|
|
|
for i in 0..N_NODES {
|
|
|
// Everyone will connect to N_CONNS random peers.
|
|
|
+ let mut peer_indexes_copy = peer_indexes.to_owned();
|
|
|
+ peer_indexes_copy.remove(i);
|
|
|
+ let peer_indexes_to_connect: Vec<_> =
|
|
|
+ peer_indexes_copy.choose_multiple(rng, N_CONNS).collect();
|
|
|
+
|
|
|
let mut peers = vec![];
|
|
|
- for _ in 0..N_CONNS {
|
|
|
- let mut port = 13200 + i;
|
|
|
- while port == 13200 + i {
|
|
|
- port = 13200 + rng.gen_range(0..N_NODES);
|
|
|
- }
|
|
|
+ for peer_index in peer_indexes_to_connect {
|
|
|
+ let port = starting_port + peer_index;
|
|
|
peers.push(Url::parse(&format!("tcp://127.0.0.1:{}", port)).unwrap());
|
|
|
}
|
|
|
|
|
|
- let settings = Settings {
|
|
|
- localnet: true,
|
|
|
- inbound_addrs: vec![Url::parse(&format!("tcp://127.0.0.1:{}", 13200 + i)).unwrap()],
|
|
|
- outbound_connections: 0,
|
|
|
- outbound_connect_timeout: 2,
|
|
|
- inbound_connections: usize::MAX,
|
|
|
+ let event_graph = spawn_node(
|
|
|
+ vec![Url::parse(&format!("tcp://127.0.0.1:{}", starting_port + i)).unwrap()],
|
|
|
peers,
|
|
|
- allowed_transports: vec!["tcp".to_string()],
|
|
|
- ..Default::default()
|
|
|
- };
|
|
|
-
|
|
|
- let p2p = P2p::new(settings, ex.clone()).await;
|
|
|
- let sled_db = sled::Config::new().temporary(true).open().unwrap();
|
|
|
- let event_graph =
|
|
|
- EventGraph::new(p2p.clone(), sled_db, "dag", 1, ex.clone()).await.unwrap();
|
|
|
- let event_graph_ = event_graph.clone();
|
|
|
-
|
|
|
- // Take the last sled item since there's only 1
|
|
|
- if genesis_event_id == NULL_ID {
|
|
|
- let (id, _) = event_graph.dag.last().unwrap().unwrap();
|
|
|
- genesis_event_id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
|
|
|
- }
|
|
|
-
|
|
|
- // Register the P2P protocols
|
|
|
- let registry = p2p.protocol_registry();
|
|
|
- registry
|
|
|
- .register(SESSION_ALL, move |channel, _| {
|
|
|
- let event_graph_ = event_graph_.clone();
|
|
|
- async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
|
|
|
- })
|
|
|
- .await;
|
|
|
+ ex.clone(),
|
|
|
+ )
|
|
|
+ .await;
|
|
|
|
|
|
eg_instances.push(event_graph);
|
|
|
}
|
|
|
@@ -140,64 +144,121 @@ async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
|
|
|
eg.p2p.clone().start().await.unwrap();
|
|
|
}
|
|
|
|
|
|
- info!("Waiting 10s until all peers connect");
|
|
|
- sleep(10).await;
|
|
|
+ info!("Waiting 5s until all peers connect");
|
|
|
+ sleep(5).await;
|
|
|
+
|
|
|
+ eg_instances
|
|
|
+}
|
|
|
+
|
|
|
+async fn assert_dags(
|
|
|
+ eg_instances: &Vec<Arc<EventGraph>>,
|
|
|
+ expected_len: usize,
|
|
|
+ rng: &mut ThreadRng,
|
|
|
+) {
|
|
|
+ let random_node = eg_instances.choose(rng).unwrap();
|
|
|
+ let last_layer_tips =
|
|
|
+ random_node.unreferenced_tips.read().await.last_key_value().unwrap().1.clone();
|
|
|
+ for (i, eg) in eg_instances.iter().enumerate() {
|
|
|
+ let node_last_layer_tips =
|
|
|
+ eg.unreferenced_tips.read().await.last_key_value().unwrap().1.clone();
|
|
|
+ assert!(
|
|
|
+ eg.dag.len() == expected_len,
|
|
|
+ "Node {}, expected {} events, have {}",
|
|
|
+ i,
|
|
|
+ expected_len,
|
|
|
+ eg.dag.len()
|
|
|
+ );
|
|
|
+ assert_eq!(
|
|
|
+ node_last_layer_tips, last_layer_tips,
|
|
|
+ "Node {} contains malformed unreferenced tips",
|
|
|
+ i
|
|
|
+ );
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+macro_rules! test_body {
|
|
|
+ ($real_call:ident) => {
|
|
|
+ init_logger();
|
|
|
+
|
|
|
+ let ex = Arc::new(Executor::new());
|
|
|
+ let ex_ = ex.clone();
|
|
|
+ let (signal, shutdown) = channel::unbounded::<()>();
|
|
|
+
|
|
|
+ // Run a thread for each node.
|
|
|
+ easy_parallel::Parallel::new()
|
|
|
+ .each(0..N_NODES, |_| future::block_on(ex.run(shutdown.recv())))
|
|
|
+ .finish(|| {
|
|
|
+ future::block_on(async {
|
|
|
+ $real_call(ex_).await;
|
|
|
+ drop(signal);
|
|
|
+ })
|
|
|
+ });
|
|
|
+ };
|
|
|
+}
|
|
|
+
|
|
|
+#[test]
|
|
|
+fn eventgraph_propagation() {
|
|
|
+ test_body!(eventgraph_propagation_real);
|
|
|
+}
|
|
|
+
|
|
|
+async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
|
|
|
+ let mut rng = rand::thread_rng();
|
|
|
+ let peer_indexes: Vec<usize> = (0..N_NODES).collect();
|
|
|
+
|
|
|
+ // Bootstrap nodes
|
|
|
+ let mut eg_instances = bootstrap_nodes(&peer_indexes, 13200, &mut rng, ex.clone()).await;
|
|
|
+
|
|
|
+ // Grab genesis event
|
|
|
+ let random_node = eg_instances.choose(&mut rng).unwrap();
|
|
|
+ let (id, _) = random_node.dag.last().unwrap().unwrap();
|
|
|
+ let genesis_event_id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
|
|
|
|
|
|
// =========================================
|
|
|
// 1. Assert that everyone's DAG is the same
|
|
|
// =========================================
|
|
|
- for (i, eg) in eg_instances.iter().enumerate() {
|
|
|
- let tips = eg.unreferenced_tips.read().await;
|
|
|
- assert!(eg.dag.len() == 1, "Node {}", i);
|
|
|
- assert!(tips.len() == 1, "Node {}", i);
|
|
|
- assert!(tips.get(&genesis_event_id).is_some(), "Node {}", i);
|
|
|
- }
|
|
|
+ assert_dags(&eg_instances, 1, &mut rng).await;
|
|
|
|
|
|
// ==========================================
|
|
|
// 2. Create an event in one node and publish
|
|
|
// ==========================================
|
|
|
- let random_node = eg_instances.choose(&mut rand::thread_rng()).unwrap();
|
|
|
- let event = Event::new(vec![1, 2, 3, 4], random_node.clone()).await;
|
|
|
+ let random_node = eg_instances.choose(&mut rng).unwrap();
|
|
|
+ let event = Event::new(vec![1, 2, 3, 4], random_node).await;
|
|
|
assert!(event.parents.contains(&genesis_event_id));
|
|
|
- // The node adds it to their DAG.
|
|
|
+ // The node adds it to their DAG, on layer 1.
|
|
|
let event_id = random_node.dag_insert(&[event.clone()]).await.unwrap()[0];
|
|
|
- let tips = random_node.unreferenced_tips.read().await;
|
|
|
- assert!(tips.len() == 1);
|
|
|
- assert!(tips.get(&event_id).is_some());
|
|
|
- drop(tips);
|
|
|
+ let tips_layers = random_node.unreferenced_tips.read().await;
|
|
|
+ // Since genesis was referenced, its layer (0) have been removed
|
|
|
+ assert_eq!(tips_layers.len(), 1);
|
|
|
+ assert!(tips_layers.last_key_value().unwrap().1.get(&event_id).is_some());
|
|
|
+ drop(tips_layers);
|
|
|
info!("Broadcasting event {}", event_id);
|
|
|
random_node.p2p.broadcast(&EventPut(event)).await;
|
|
|
- info!("Waiting 10s for event propagation");
|
|
|
- sleep(10).await;
|
|
|
+ info!("Waiting 5s for event propagation");
|
|
|
+ sleep(5).await;
|
|
|
|
|
|
// ====================================================
|
|
|
// 3. Assert that everyone has the new event in the DAG
|
|
|
// ====================================================
|
|
|
- for (i, eg) in eg_instances.iter().enumerate() {
|
|
|
- let tips = eg.unreferenced_tips.read().await;
|
|
|
- assert!(eg.dag.len() == 2, "Node {}", i);
|
|
|
- assert!(tips.len() == 1, "Node {}", i);
|
|
|
- assert!(tips.get(&event_id).is_some(), "Node {}", i);
|
|
|
- }
|
|
|
+ assert_dags(&eg_instances, 2, &mut rng).await;
|
|
|
|
|
|
// ==============================================================
|
|
|
// 4. Create multiple events on a node and broadcast the last one
|
|
|
// The `EventPut` logic should manage to fetch all of them,
|
|
|
// provided that the last one references the earlier ones.
|
|
|
// ==============================================================
|
|
|
- let random_node = eg_instances.choose(&mut rand::thread_rng()).unwrap();
|
|
|
- let event0 = Event::new(vec![1, 2, 3, 4, 0], random_node.clone()).await;
|
|
|
+ let random_node = eg_instances.choose(&mut rng).unwrap();
|
|
|
+ let event0 = Event::new(vec![1, 2, 3, 4, 0], random_node).await;
|
|
|
let event0_id = random_node.dag_insert(&[event0.clone()]).await.unwrap()[0];
|
|
|
- let event1 = Event::new(vec![1, 2, 3, 4, 1], random_node.clone()).await;
|
|
|
+ let event1 = Event::new(vec![1, 2, 3, 4, 1], random_node).await;
|
|
|
let event1_id = random_node.dag_insert(&[event1.clone()]).await.unwrap()[0];
|
|
|
- let event2 = Event::new(vec![1, 2, 3, 4, 2], random_node.clone()).await;
|
|
|
+ let event2 = Event::new(vec![1, 2, 3, 4, 2], random_node).await;
|
|
|
let event2_id = random_node.dag_insert(&[event2.clone()]).await.unwrap()[0];
|
|
|
- // Genesis event + event from 2. + upper 3 events
|
|
|
- assert!(random_node.dag.len() == 5);
|
|
|
- let tips = random_node.unreferenced_tips.read().await;
|
|
|
- assert!(tips.len() == 1);
|
|
|
- assert!(tips.get(&event2_id).is_some());
|
|
|
- drop(tips);
|
|
|
+ // Genesis event + event from 2. + upper 3 events (layer 4)
|
|
|
+ assert_eq!(random_node.dag.len(), 5);
|
|
|
+ let tips_layers = random_node.unreferenced_tips.read().await;
|
|
|
+ assert_eq!(tips_layers.len(), 1);
|
|
|
+ assert!(tips_layers.get(&4).unwrap().get(&event2_id).is_some());
|
|
|
+ drop(tips_layers);
|
|
|
|
|
|
let event_chain =
|
|
|
vec![(event0_id, event0.parents), (event1_id, event1.parents), (event2_id, event2.parents)];
|
|
|
@@ -205,142 +266,186 @@ async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
|
|
|
info!("Broadcasting event {}", event2_id);
|
|
|
info!("Event chain: {:#?}", event_chain);
|
|
|
random_node.p2p.broadcast(&EventPut(event2)).await;
|
|
|
- info!("Waiting 10s for event propagation");
|
|
|
- sleep(10).await;
|
|
|
+ info!("Waiting 5s for event propagation");
|
|
|
+ sleep(5).await;
|
|
|
|
|
|
// ==========================================
|
|
|
// 5. Assert that everyone has all the events
|
|
|
// ==========================================
|
|
|
- for (i, eg) in eg_instances.iter().enumerate() {
|
|
|
- let tips = eg.unreferenced_tips.read().await;
|
|
|
- assert!(eg.dag.len() == 5, "Node {}, expected 5 events, have {}", i, eg.dag.len());
|
|
|
- assert!(tips.len() == 1, "Node {}, expected 1 tip, have {}", i, tips.len());
|
|
|
- assert!(tips.get(&event2_id).is_some(), "Node {}, expected tip to be {}", i, event2_id);
|
|
|
- }
|
|
|
+ assert_dags(&eg_instances, 5, &mut rng).await;
|
|
|
|
|
|
// ===========================================
|
|
|
// 6. Create multiple events on multiple nodes
|
|
|
// ===========================================
|
|
|
// node 1
|
|
|
// =======
|
|
|
- let node1 = eg_instances.choose(&mut rand::thread_rng()).unwrap();
|
|
|
- let event0_1 = Event::new(vec![1, 2, 3, 4, 3], node1.clone()).await;
|
|
|
- let _ = node1.dag_insert(&[event0_1.clone()]).await.unwrap()[0];
|
|
|
+ let node1 = eg_instances.choose(&mut rng).unwrap();
|
|
|
+ let event0_1 = Event::new(vec![1, 2, 3, 4, 3], node1).await;
|
|
|
+ node1.dag_insert(&[event0_1.clone()]).await.unwrap();
|
|
|
node1.p2p.broadcast(&EventPut(event0_1)).await;
|
|
|
|
|
|
- let event1_1 = Event::new(vec![1, 2, 3, 4, 4], node1.clone()).await;
|
|
|
- let _ = node1.dag_insert(&[event1_1.clone()]).await.unwrap()[0];
|
|
|
+ let event1_1 = Event::new(vec![1, 2, 3, 4, 4], node1).await;
|
|
|
+ node1.dag_insert(&[event1_1.clone()]).await.unwrap();
|
|
|
node1.p2p.broadcast(&EventPut(event1_1)).await;
|
|
|
|
|
|
- let event2_1 = Event::new(vec![1, 2, 3, 4, 5], node1.clone()).await;
|
|
|
- let _ = node1.dag_insert(&[event2_1.clone()]).await.unwrap()[0];
|
|
|
+ let event2_1 = Event::new(vec![1, 2, 3, 4, 5], node1).await;
|
|
|
+ node1.dag_insert(&[event2_1.clone()]).await.unwrap();
|
|
|
node1.p2p.broadcast(&EventPut(event2_1)).await;
|
|
|
|
|
|
// =======
|
|
|
// node 2
|
|
|
// =======
|
|
|
- let node2 = eg_instances.choose(&mut rand::thread_rng()).unwrap();
|
|
|
- let event0_2 = Event::new(vec![1, 2, 3, 4, 6], node2.clone()).await;
|
|
|
- let _ = node2.dag_insert(&[event0_2.clone()]).await.unwrap()[0];
|
|
|
+ let node2 = eg_instances.choose(&mut rng).unwrap();
|
|
|
+ let event0_2 = Event::new(vec![1, 2, 3, 4, 6], node2).await;
|
|
|
+ node2.dag_insert(&[event0_2.clone()]).await.unwrap();
|
|
|
node2.p2p.broadcast(&EventPut(event0_2)).await;
|
|
|
- let event1_2 = Event::new(vec![1, 2, 3, 4, 7], node2.clone()).await;
|
|
|
- let _ = node2.dag_insert(&[event1_2.clone()]).await.unwrap()[0];
|
|
|
+
|
|
|
+ let event1_2 = Event::new(vec![1, 2, 3, 4, 7], node2).await;
|
|
|
+ node2.dag_insert(&[event1_2.clone()]).await.unwrap();
|
|
|
node2.p2p.broadcast(&EventPut(event1_2)).await;
|
|
|
|
|
|
- let event2_2 = Event::new(vec![1, 2, 3, 4, 8], node2.clone()).await;
|
|
|
- let _ = node2.dag_insert(&[event2_2.clone()]).await.unwrap()[0];
|
|
|
+ let event2_2 = Event::new(vec![1, 2, 3, 4, 8], node2).await;
|
|
|
+ node2.dag_insert(&[event2_2.clone()]).await.unwrap();
|
|
|
node2.p2p.broadcast(&EventPut(event2_2)).await;
|
|
|
|
|
|
// =======
|
|
|
// node 3
|
|
|
// =======
|
|
|
- let node3 = eg_instances.choose(&mut rand::thread_rng()).unwrap();
|
|
|
- let event0_3 = Event::new(vec![1, 2, 3, 4, 9], node3.clone()).await;
|
|
|
- let _ = node3.dag_insert(&[event0_3.clone()]).await.unwrap()[0];
|
|
|
+ let node3 = eg_instances.choose(&mut rng).unwrap();
|
|
|
+ let event0_3 = Event::new(vec![1, 2, 3, 4, 9], node3).await;
|
|
|
+ node3.dag_insert(&[event0_3.clone()]).await.unwrap();
|
|
|
node2.p2p.broadcast(&EventPut(event0_3)).await;
|
|
|
|
|
|
- let event1_3 = Event::new(vec![1, 2, 3, 4, 10], node3.clone()).await;
|
|
|
- let _ = node3.dag_insert(&[event1_3.clone()]).await.unwrap()[0];
|
|
|
+ let event1_3 = Event::new(vec![1, 2, 3, 4, 10], node3).await;
|
|
|
+ node3.dag_insert(&[event1_3.clone()]).await.unwrap();
|
|
|
node2.p2p.broadcast(&EventPut(event1_3)).await;
|
|
|
|
|
|
- let event2_3 = Event::new(vec![1, 2, 3, 4, 11], node3.clone()).await;
|
|
|
- let event2_3_id = node3.dag_insert(&[event2_3.clone()]).await.unwrap()[0];
|
|
|
+ let event2_3 = Event::new(vec![1, 2, 3, 4, 11], node3).await;
|
|
|
+ node3.dag_insert(&[event2_3.clone()]).await.unwrap();
|
|
|
node3.p2p.broadcast(&EventPut(event2_3)).await;
|
|
|
|
|
|
- info!("Waiting 10s for events propagation");
|
|
|
- sleep(10).await;
|
|
|
+ info!("Waiting 5s for events propagation");
|
|
|
+ sleep(5).await;
|
|
|
|
|
|
// ==========================================
|
|
|
// 7. Assert that everyone has all the events
|
|
|
// ==========================================
|
|
|
- for (i, eg) in eg_instances.iter().enumerate() {
|
|
|
- let tips = eg.unreferenced_tips.read().await;
|
|
|
- assert!(eg.dag.len() == 14, "Node {}, expected 14 events, have {}", i, eg.dag.len());
|
|
|
- // 5 events from 2. and 4. + 9 events from 6. = ^
|
|
|
- assert!(tips.get(&event2_3_id).is_some(), "Node {}, expected tip to be {}", i, event2_3_id);
|
|
|
- }
|
|
|
+ // 5 events from 2. and 4. + 9 events from 6. = 14
|
|
|
+ assert_dags(&eg_instances, 14, &mut rng).await;
|
|
|
|
|
|
// ============================================================
|
|
|
// 8. Start a new node and try to sync the DAG from other peers
|
|
|
// ============================================================
|
|
|
{
|
|
|
// Connect to N_CONNS random peers.
|
|
|
+ let peer_indexes_to_connect: Vec<_> =
|
|
|
+ peer_indexes.choose_multiple(&mut rng, N_CONNS).collect();
|
|
|
+
|
|
|
let mut peers = vec![];
|
|
|
- for _ in 0..N_CONNS {
|
|
|
- let port = 13200 + rng.gen_range(0..N_NODES);
|
|
|
+ for peer_index in peer_indexes_to_connect {
|
|
|
+ let port = 13200 + peer_index;
|
|
|
peers.push(Url::parse(&format!("tcp://127.0.0.1:{}", port)).unwrap());
|
|
|
}
|
|
|
|
|
|
- let settings = Settings {
|
|
|
- localnet: true,
|
|
|
- inbound_addrs: vec![
|
|
|
- Url::parse(&format!("tcp://127.0.0.1:{}", 13200 + N_NODES + 1)).unwrap()
|
|
|
- ],
|
|
|
- outbound_connections: 0,
|
|
|
- outbound_connect_timeout: 2,
|
|
|
- inbound_connections: usize::MAX,
|
|
|
+ let event_graph = spawn_node(
|
|
|
+ vec![Url::parse(&format!("tcp://127.0.0.1:{}", 13200 + N_NODES + 1)).unwrap()],
|
|
|
peers,
|
|
|
- allowed_transports: vec!["tcp".to_string()],
|
|
|
- ..Default::default()
|
|
|
- };
|
|
|
-
|
|
|
- let p2p = P2p::new(settings, ex.clone()).await;
|
|
|
- let sled_db = sled::Config::new().temporary(true).open().unwrap();
|
|
|
- let event_graph =
|
|
|
- EventGraph::new(p2p.clone(), sled_db, "dag", 1, ex.clone()).await.unwrap();
|
|
|
- let event_graph_ = event_graph.clone();
|
|
|
-
|
|
|
- // Register the P2P protocols
|
|
|
- let registry = p2p.protocol_registry();
|
|
|
- registry
|
|
|
- .register(SESSION_ALL, move |channel, _| {
|
|
|
- let event_graph_ = event_graph_.clone();
|
|
|
- async move { ProtocolEventGraph::init(event_graph_, channel).await.unwrap() }
|
|
|
- })
|
|
|
- .await;
|
|
|
+ ex.clone(),
|
|
|
+ )
|
|
|
+ .await;
|
|
|
|
|
|
eg_instances.push(event_graph.clone());
|
|
|
|
|
|
event_graph.p2p.clone().start().await.unwrap();
|
|
|
|
|
|
- info!("Waiting 10s for new node connection");
|
|
|
- sleep(10).await;
|
|
|
+ info!("Waiting 5s for new node connection");
|
|
|
+ sleep(5).await;
|
|
|
|
|
|
- event_graph.dag_sync().await.unwrap();
|
|
|
+ event_graph.dag_sync().await.unwrap()
|
|
|
}
|
|
|
|
|
|
- info!("Waiting 10s for things to settle");
|
|
|
- sleep(10).await;
|
|
|
-
|
|
|
// ============================================================
|
|
|
// 9. Assert the new synced DAG has the same contents as others
|
|
|
// ============================================================
|
|
|
- for (i, eg) in eg_instances.iter().enumerate() {
|
|
|
- let tips = eg.unreferenced_tips.read().await;
|
|
|
- assert!(eg.dag.len() == 14, "Node {}, expected 14 events, have {}", i, eg.dag.len());
|
|
|
- // 5 events from 2. and 4. + 9 events from 6. = ^
|
|
|
- assert!(tips.get(&event2_3_id).is_some(), "Node {}, expected tip to be {}", i, event2_3_id);
|
|
|
+ // 5 events from 2. and 4. + 9 events from 6. = 14
|
|
|
+ assert_dags(&eg_instances, 14, &mut rng).await;
|
|
|
+
|
|
|
+ // Stop the P2P network
|
|
|
+ for eg in eg_instances.iter() {
|
|
|
+ eg.p2p.clone().stop().await;
|
|
|
+ }
|
|
|
+}
|
|
|
+
|
|
|
+#[test]
|
|
|
+#[ignore]
|
|
|
+fn eventgraph_chaotic_propagation() {
|
|
|
+ test_body!(eventgraph_chaotic_propagation_real);
|
|
|
+}
|
|
|
+
|
|
|
+async fn eventgraph_chaotic_propagation_real(ex: Arc<Executor<'static>>) {
|
|
|
+ let mut rng = rand::thread_rng();
|
|
|
+ let peer_indexes: Vec<usize> = (0..N_NODES).collect();
|
|
|
+ let n_events: usize = 100000;
|
|
|
+
|
|
|
+ // Bootstrap nodes
|
|
|
+ let mut eg_instances = bootstrap_nodes(&peer_indexes, 14200, &mut rng, ex.clone()).await;
|
|
|
+
|
|
|
+ // =========================================
|
|
|
+ // 1. Assert that everyone's DAG is the same
|
|
|
+ // =========================================
|
|
|
+ assert_dags(&eg_instances, 1, &mut rng).await;
|
|
|
+
|
|
|
+ // ===========================================
|
|
|
+ // 2. Create multiple events on multiple nodes
|
|
|
+ for i in 0..n_events {
|
|
|
+ let random_node = eg_instances.choose(&mut rng).unwrap();
|
|
|
+ let event = Event::new(i.to_be_bytes().to_vec(), random_node).await;
|
|
|
+ random_node.dag_insert(&[event.clone()]).await.unwrap();
|
|
|
+ random_node.p2p.broadcast(&EventPut(event)).await;
|
|
|
}
|
|
|
+ info!("Waiting 5s for events propagation");
|
|
|
+ sleep(5).await;
|
|
|
+
|
|
|
+ // ==========================================
|
|
|
+ // 3. Assert that everyone has all the events
|
|
|
+ // ==========================================
|
|
|
+ assert_dags(&eg_instances, n_events + 1, &mut rng).await;
|
|
|
+
|
|
|
+ // ============================================================
|
|
|
+ // 4. Start a new node and try to sync the DAG from other peers
|
|
|
+ // ============================================================
|
|
|
+ {
|
|
|
+ // Connect to N_CONNS random peers.
|
|
|
+ let peer_indexes_to_connect: Vec<_> =
|
|
|
+ peer_indexes.choose_multiple(&mut rng, N_CONNS).collect();
|
|
|
+
|
|
|
+ let mut peers = vec![];
|
|
|
+ for peer_index in peer_indexes_to_connect {
|
|
|
+ let port = 14200 + peer_index;
|
|
|
+ peers.push(Url::parse(&format!("tcp://127.0.0.1:{}", port)).unwrap());
|
|
|
+ }
|
|
|
+
|
|
|
+ let event_graph = spawn_node(
|
|
|
+ vec![Url::parse(&format!("tcp://127.0.0.1:{}", 14200 + N_NODES + 1)).unwrap()],
|
|
|
+ peers,
|
|
|
+ ex.clone(),
|
|
|
+ )
|
|
|
+ .await;
|
|
|
+
|
|
|
+ eg_instances.push(event_graph.clone());
|
|
|
+
|
|
|
+ event_graph.p2p.clone().start().await.unwrap();
|
|
|
+
|
|
|
+ info!("Waiting 5s for new node connection");
|
|
|
+ sleep(5).await;
|
|
|
+
|
|
|
+ event_graph.dag_sync().await.unwrap()
|
|
|
+ }
|
|
|
+
|
|
|
+ // ============================================================
|
|
|
+ // 5. Assert the new synced DAG has the same contents as others
|
|
|
+ // ============================================================
|
|
|
+ assert_dags(&eg_instances, n_events + 1, &mut rng).await;
|
|
|
|
|
|
// Stop the P2P network
|
|
|
for eg in eg_instances.iter() {
|