Просмотр исходного кода

event-graph: Improve parent fetching, and improve tests.

parazyd 2 лет назад
Родитель
Сommit
57f0896fb2
2 измененных файлов с 115 добавлено и 34 удалено
  1. 63 32
      src/event_graph2/proto.rs
  2. 52 2
      src/event_graph2/tests.rs

+ 63 - 32
src/event_graph2/proto.rs

@@ -187,68 +187,99 @@ impl ProtocolEventGraph {
                 "Event {} is new", event_id,
             );
 
-            let mut missing_parents = HashSet::new();
+            let mut missing_parents = vec![];
             for parent_id in event.parents.iter() {
                 // `event.validate()` should have already made sure that
-                // not all parents are NULL.
+                // not all parents are NULL, and that there are no duplicates.
                 if parent_id == &NULL_ID {
                     continue
                 }
 
                 if !self.event_graph.dag.contains_key(parent_id.as_bytes()).unwrap() {
-                    missing_parents.insert(*parent_id);
+                    missing_parents.push(*parent_id);
                 }
             }
 
             // If we have missing parents, then we have to attempt to
-            // fetch them from this peer.
+            // fetch them from this peer. Do this recursively until we
+            // find all of them.
             if !missing_parents.is_empty() {
+                // We track the received events in a vec. If/when we get all
+                // of them, we need to insert them in reverse so the DAG state
+                // stays correct and unreferenced tips represent the actual thing
+                // they should. If we insert them out of order, then we might have
+                // wrong unreferenced tips.
+                // TODO: What should we do if at some point the events become too old?
+                let mut received_events = vec![];
+
                 debug!(
                     target: "event_graph::protocol::handle_event_put()",
                     "Event has {} missing parents. Requesting...", missing_parents.len(),
                 );
-                let mut received_events = HashMap::new();
-                for parent_id in missing_parents.iter() {
-                    debug!(
-                        target: "event_graph::protocol::handle_event_put()",
-                        "Requesting {}...", parent_id,
-                    );
-                    self.channel.send(&EventReq(*parent_id)).await?;
-                    let parent = match timeout(REPLY_TIMEOUT, self.ev_rep_sub.receive()).await {
-                        Ok(parent) => parent?,
-                        Err(_) => {
+
+                while !missing_parents.is_empty() {
+                    for parent_id in missing_parents.clone().iter() {
+                        debug!(
+                            target: "event_graph::protocol::handle_event_put()",
+                            "Requesting {}...", parent_id,
+                        );
+
+                        self.channel.send(&EventReq(*parent_id)).await?;
+                        let parent = match timeout(REPLY_TIMEOUT, self.ev_rep_sub.receive()).await {
+                            Ok(parent) => parent?,
+                            Err(_) => {
+                                error!(
+                                    target: "event_graph::protocol::handle_event_put()",
+                                    "[EVENTGRAPH] Timeout while waiting for parent {} from {}",
+                                    parent_id, self.channel.address(),
+                                );
+                                self.channel.stop().await;
+                                return Err(Error::ChannelStopped)
+                            }
+                        };
+                        let parent = parent.0.clone();
+
+                        if &parent.id() != parent_id {
                             error!(
                                 target: "event_graph::protocol::handle_event_put()",
-                                "[EVENTGRAPH] Timeout while waiting for parent {} from {}",
-                                parent_id, self.channel.address(),
+                                "[EVENTGRAPH] Peer {} replied with a wrong event: {}",
+                                self.channel.address(), parent.id(),
                             );
                             self.channel.stop().await;
                             return Err(Error::ChannelStopped)
                         }
-                    };
-                    let parent = parent.0.clone();
 
-                    if &parent.id() != parent_id {
-                        error!(
+                        debug!(
                             target: "event_graph::protocol::handle_event_put()",
-                            "[EVENTGRAPH] Peer {} replied with a wrong event: {}",
-                            self.channel.address(), parent.id(),
+                            "Got correct parent event {}", parent.id(),
                         );
-                        self.channel.stop().await;
-                        return Err(Error::ChannelStopped)
-                    }
 
-                    debug!(
-                        target: "event_graph::protocol::handle_event_put()",
-                        "Got correct parent event {}", parent.id(),
-                    );
+                        received_events.push(parent.clone());
+                        let pos = missing_parents.iter().position(|id| id == &parent.id()).unwrap();
+                        missing_parents.remove(pos);
+
+                        // See if we have the upper parents
+                        for upper_parent in parent.parents.iter() {
+                            if upper_parent == &NULL_ID {
+                                continue
+                            }
+
+                            if !self.event_graph.dag.contains_key(upper_parent.as_bytes()).unwrap()
+                            {
+                                debug!(
+                                    target: "event_graph::protocol::handle_event_put()",
+                                    "Found upper missing parent event{}", upper_parent,
+                                );
+                                missing_parents.push(*upper_parent);
+                            }
+                        }
+                    }
+                } // <-- while !missing_parents.is_empty()
 
-                    received_events.insert(parent.id(), parent);
-                }
                 // At this point we should've got all the events.
                 // We should add them to the DAG.
                 // TODO: FIXME: Also validate these events.
-                for event in received_events.values() {
+                for event in received_events.iter().rev() {
                     self.event_graph.dag_insert(event).await.unwrap();
                 }
             } // <-- !missing_parents.is_empty()

+ 52 - 2
src/event_graph2/tests.rs

@@ -48,14 +48,18 @@ fn eventgraph_propagation() {
     cfg.add_filter_ignore("net::protocol_ping".to_string());
     cfg.add_filter_ignore("net::channel::subscribe_stop()".to_string());
     cfg.add_filter_ignore("net::hosts".to_string());
+    cfg.add_filter_ignore("net::session".to_string());
     cfg.add_filter_ignore("net::message_subscriber".to_string());
     cfg.add_filter_ignore("net::protocol_address".to_string());
     cfg.add_filter_ignore("net::protocol_version".to_string());
+    cfg.add_filter_ignore("net::protocol_registry".to_string());
     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,
+        //simplelog::LevelFilter::Info,
+        simplelog::LevelFilter::Debug,
         //simplelog::LevelFilter::Trace,
         cfg.build(),
         simplelog::TerminalMode::Mixed,
@@ -176,6 +180,52 @@ async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
         assert!(tips.get(&event_id).is_some(), "Node {}", i);
     }
 
+    // ==============================================================
+    // 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 event0_id = random_node.dag_insert(&event0).await.unwrap();
+    let event1 = Event::new(vec![1, 2, 3, 4, 1], random_node.clone()).await;
+    let event1_id = random_node.dag_insert(&event1).await.unwrap();
+    let event2 = Event::new(vec![1, 2, 3, 4, 2], random_node.clone()).await;
+    let event2_id = random_node.dag_insert(&event2).await.unwrap();
+    // 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);
+
+    let event_chain =
+        vec![(event0_id, event0.parents), (event1_id, event1.parents), (event2_id, event2.parents)];
+
+    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;
+
+    // ==========================================
+    // 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);
+    }
+
+    // ===========================================
+    // 6. Create multiple events on multiple nodes
+    // ===========================================
+
+    // ==========================================
+    // 7. Assert that everyone has all the events
+    // ==========================================
+
     // Stop the P2P network
     for eg in eg_instances.iter() {
         eg.p2p.clone().stop().await;