فهرست منبع

event_graph: use the genesis timestamp as a dag name, this will make it easier to sort dags

oars 10 ماه پیش
والد
کامیت
6c93c11ba1
4فایلهای تغییر یافته به همراه42 افزوده شده و 37 حذف شده
  1. 4 4
      src/event_graph/event.rs
  2. 20 20
      src/event_graph/mod.rs
  3. 10 5
      src/event_graph/proto.rs
  4. 8 8
      src/event_graph/tests.rs

+ 4 - 4
src/event_graph/event.rs

@@ -44,13 +44,13 @@ pub struct Header {
 impl Header {
     // Create a new Header given EventGraph to retrieve the correct layout
     pub async fn new(event_graph: &EventGraph) -> Self {
-        let current_dag_name = event_graph.current_genesis.read().await.id();
+        let current_dag_name = event_graph.current_genesis.read().await.header.timestamp;
         let (layer, parents) = event_graph.get_next_layer_with_parents(&current_dag_name).await;
         Self { timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64, parents, layer }
     }
 
     pub async fn with_timestamp(timestamp: u64, event_graph: &EventGraph) -> Self {
-        let current_dag_name = event_graph.current_genesis.read().await.id();
+        let current_dag_name = event_graph.current_genesis.read().await.header.timestamp;
         let (layer, parents) = event_graph.get_next_layer_with_parents(&current_dag_name).await;
         Self { timestamp, parents, layer }
     }
@@ -246,7 +246,7 @@ mod tests {
             // Generate a dummy event graph
             let event_graph = make_event_graph().await?;
 
-            let dag_name = event_graph.current_genesis.read().await.id().to_string();
+            let dag_name = event_graph.current_genesis.read().await.header.timestamp.to_string();
             let hdr_tree_name = format!("headers_{dag_name}");
             let header_dag = event_graph.dag_store.read().await.get_dag(&hdr_tree_name);
 
@@ -267,7 +267,7 @@ mod tests {
             // Generate a dummy event graph
             let event_graph = make_event_graph().await?;
 
-            let dag_name = event_graph.current_genesis.read().await.id().to_string();
+            let dag_name = event_graph.current_genesis.read().await.header.timestamp.to_string();
             let hdr_tree_name = format!("headers_{dag_name}");
             let header_dag = event_graph.dag_store.read().await.get_dag(&hdr_tree_name);
 

+ 20 - 20
src/event_graph/mod.rs

@@ -102,8 +102,8 @@ pub type LayerUTips = BTreeMap<u64, HashSet<blake3::Hash>>;
 #[derive(Clone)]
 pub struct DAGStore {
     db: sled::Db,
-    header_dags: HashMap<Hash, (sled::Tree, LayerUTips)>,
-    main_dags: HashMap<Hash, (sled::Tree, LayerUTips)>,
+    header_dags: HashMap<u64, (sled::Tree, LayerUTips)>,
+    main_dags: HashMap<u64, (sled::Tree, LayerUTips)>,
 }
 
 impl DAGStore {
@@ -121,7 +121,7 @@ impl DAGStore {
                 };
                 let genesis = Event { header, content: GENESIS_CONTENTS.to_vec() };
 
-                let tree_name = genesis.id().to_string();
+                let tree_name = genesis.header.timestamp.to_string();
                 let hdr_tree_name = format!("headers_{tree_name}");
                 let hdr_dag = sled_db.open_tree(hdr_tree_name).unwrap();
                 let dag = sled_db.open_tree(tree_name).unwrap();
@@ -161,12 +161,12 @@ impl DAGStore {
                     }
                 }
                 let utips = self.find_unreferenced_tips(&dag).await;
-                considered_header_trees.insert(genesis.id(), (hdr_dag, utips.clone()));
-                considered_trees.insert(genesis.id(), (dag, utips));
+                considered_header_trees.insert(genesis.header.timestamp, (hdr_dag, utips.clone()));
+                considered_trees.insert(genesis.header.timestamp, (dag, utips));
             }
         } else {
             let genesis = generate_genesis(0);
-            let tree_name = genesis.id().to_string();
+            let tree_name = genesis.header.timestamp.to_string();
             let hdr_tree_name = format!("headers_{tree_name}");
             let hdr_dag = sled_db.open_tree(hdr_tree_name).unwrap();
             let dag = sled_db.open_tree(tree_name).unwrap();
@@ -205,8 +205,8 @@ impl DAGStore {
                 }
             }
             let utips = self.find_unreferenced_tips(&dag).await;
-            considered_header_trees.insert(genesis.id(), (hdr_dag, utips.clone()));
-            considered_trees.insert(genesis.id(), (dag, utips));
+            considered_header_trees.insert(genesis.header.timestamp, (hdr_dag, utips.clone()));
+            considered_trees.insert(genesis.header.timestamp, (dag, utips));
         }
 
         Self { db: sled_db, header_dags: considered_header_trees, main_dags: considered_trees }
@@ -227,7 +227,7 @@ impl DAGStore {
                 // since dags are sorted in reverse
                 let oldest_tree = sorted_dags.last().unwrap().name();
                 let oldest_key = String::from_utf8_lossy(&oldest_tree);
-                let oldest_key = blake3::Hash::from_str(&oldest_key).unwrap();
+                let oldest_key = u64::from_str(&oldest_key).unwrap();
 
                 let oldest_hdr_tree = self.header_dags.remove(&oldest_key).unwrap();
                 let oldest_tree = self.main_dags.remove(&oldest_key).unwrap();
@@ -246,8 +246,8 @@ impl DAGStore {
         let dag = self.get_dag(dag_name);
         dag.insert(genesis_event.id().as_bytes(), serialize_async(genesis_event).await).unwrap();
         let utips = self.find_unreferenced_tips(&dag).await;
-        self.header_dags.insert(genesis_event.id(), (hdr_dag, utips.clone()));
-        self.main_dags.insert(genesis_event.id(), (dag, utips));
+        self.header_dags.insert(genesis_event.header.timestamp, (hdr_dag, utips.clone()));
+        self.main_dags.insert(genesis_event.header.timestamp, (dag, utips));
     }
 
     // Get a DAG providing its name.
@@ -404,7 +404,7 @@ impl EventGraph {
 
         // Create the current genesis event based on the `hours_rotation`
         let current_genesis = generate_genesis(hours_rotation);
-        let current_dag_tree_name = current_genesis.id().to_string();
+        let current_dag_tree_name = current_genesis.header.timestamp.to_string();
         let dag_store = DAGStore {
             db: sled_db.clone(),
             header_dags: HashMap::default(),
@@ -735,7 +735,7 @@ impl EventGraph {
         let mut broadcasted_ids = self.broadcasted_ids.write().await;
         let mut current_genesis = self.current_genesis.write().await;
 
-        let dag_name = genesis_event.id().to_string();
+        let dag_name = genesis_event.header.timestamp.to_string();
         self.dag_store.write().await.add_dag(&dag_name, &genesis_event).await;
 
         // Clear bcast ids
@@ -794,7 +794,7 @@ impl EventGraph {
         }
 
         // Acquire exclusive locks to `broadcasted_ids`
-        let dag_name_hash = blake3::Hash::from_str(dag_name).unwrap();
+        let dag_timestamp = u64::from_str(dag_name).unwrap();
         let mut broadcasted_ids = self.broadcasted_ids.write().await;
 
         let main_dag = self.dag_store.read().await.get_dag(dag_name);
@@ -859,7 +859,7 @@ impl EventGraph {
         }
 
         let mut dag_store = self.dag_store.write().await;
-        let (_, unreferenced_tips) = &mut dag_store.main_dags.get_mut(&dag_name_hash).unwrap();
+        let (_, unreferenced_tips) = &mut dag_store.main_dags.get_mut(&dag_timestamp).unwrap();
 
         // Iterate over given events to update references and
         // send out notifications about them
@@ -910,8 +910,8 @@ impl EventGraph {
             self.event_pub.notify(event.clone()).await;
         }
 
-        dag_store.header_dags.get_mut(&dag_name_hash).unwrap().1 =
-            dag_store.main_dags.get(&dag_name_hash).unwrap().1.clone();
+        dag_store.header_dags.get_mut(&dag_timestamp).unwrap().1 =
+            dag_store.main_dags.get(&dag_timestamp).unwrap().1.clone();
 
         // Drop the exclusive locks
         drop(dag_store);
@@ -987,7 +987,7 @@ impl EventGraph {
     /// parents.
     async fn get_next_layer_with_parents(
         &self,
-        dag_name: &Hash,
+        dag_name: &u64,
     ) -> (u64, [blake3::Hash; N_EVENT_PARENTS]) {
         let store = self.dag_store.read().await;
         let (_, unreferenced_tips) = store.header_dags.get(dag_name).unwrap();
@@ -1113,7 +1113,7 @@ impl EventGraph {
     #[cfg(feature = "rpc")]
     pub async fn eventgraph_info(&self, id: u16, _params: JsonValue) -> JsonResult {
         let current_genesis = self.current_genesis.read().await;
-        let dag_name = current_genesis.id().to_string();
+        let dag_name = current_genesis.header.timestamp.to_string();
         let mut graph = HashMap::new();
         for iter_elem in self.dag_store.read().await.get_dag(&dag_name).iter() {
             let (id, val) = iter_elem.unwrap();
@@ -1146,7 +1146,7 @@ impl EventGraph {
         );
 
         let current_genesis = self.current_genesis.read().await;
-        let dag_name = current_genesis.id().to_string();
+        let dag_name = current_genesis.header.timestamp.to_string();
         let mut graph = HashMap::new();
         for iter_elem in self.dag_store.read().await.get_dag(&dag_name).iter() {
             let (id, val) = iter_elem.unwrap();

+ 10 - 5
src/event_graph/proto.rs

@@ -279,7 +279,7 @@ impl ProtocolEventGraph {
 
             // If we have already seen the event, we'll stay quiet.
             let current_genesis = self.event_graph.current_genesis.read().await;
-            let dag_name = current_genesis.id().to_string();
+            let dag_name = current_genesis.header.timestamp.to_string();
             let hdr_tree_name = format!("headers_{dag_name}");
             let event_id = event.id();
             if self
@@ -382,7 +382,7 @@ impl ProtocolEventGraph {
                 );
 
                 let current_genesis = self.event_graph.current_genesis.read().await;
-                let dag_name = current_genesis.id().to_string();
+                let dag_name = current_genesis.header.timestamp.to_string();
                 let hdr_tree_name = format!("headers_{dag_name}");
 
                 while !missing_parents.is_empty() {
@@ -650,7 +650,12 @@ impl ProtocolEventGraph {
 
             // We received header request. Let's find them, add them to
             // our bcast ids list, and reply with them.
-            let main_dag = self.event_graph.dag_store.read().await.get_dag(&dag_name);
+            let dag_timestamp = u64::from_str(&dag_name)?;
+            let store = self.event_graph.dag_store.read().await;
+            if !store.header_dags.contains_key(&dag_timestamp) {
+                continue
+            }
+            let main_dag = store.get_dag(&dag_name);
             let mut headers = vec![];
             for item in main_dag.iter() {
                 let (_, event) = item.unwrap();
@@ -699,9 +704,9 @@ impl ProtocolEventGraph {
 
             // We received a tip request. Let's find them, add them to
             // our bcast ids list, and reply with them.
-            let dag_name_hash = blake3::Hash::from_str(&dag_name).unwrap();
+            let dag_timestamp = u64::from_str(&dag_name)?;
             let store = self.event_graph.dag_store.read().await;
-            let (_, layers) = match store.header_dags.get(&dag_name_hash) {
+            let (_, layers) = match store.header_dags.get(&dag_timestamp) {
                 Some(v) => v,
                 None => continue,
             };

+ 8 - 8
src/event_graph/tests.rs

@@ -162,13 +162,13 @@ async fn bootstrap_nodes(
 
 async fn assert_dags(eg_instances: &[Arc<EventGraph>], expected_len: usize, rng: &mut ThreadRng) {
     let random_node = eg_instances.choose(rng).unwrap();
-    let random_node_genesis = random_node.current_genesis.read().await.id();
+    let random_node_genesis = random_node.current_genesis.read().await.header.timestamp;
     let store = random_node.dag_store.read().await;
     let (_, unreferenced_tips) = store.main_dags.get(&random_node_genesis).unwrap();
     let last_layer_tips = unreferenced_tips.last_key_value().unwrap().1.clone();
     for (i, eg) in eg_instances.iter().enumerate() {
         let current_genesis = eg.current_genesis.read().await;
-        let dag_name = current_genesis.id().to_string();
+        let dag_name = current_genesis.header.timestamp.to_string();
         let dag = eg.dag_store.read().await.get_dag(&dag_name);
         let unreferenced_tips = eg.dag_store.read().await.find_unreferenced_tips(&dag).await;
         let node_last_layer_tips = unreferenced_tips.last_key_value().unwrap().1.clone();
@@ -219,7 +219,7 @@ async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
     // Grab genesis event
     let random_node = eg_instances.choose(&mut rng).unwrap();
     let current_genesis = random_node.current_genesis.read().await;
-    let dag_name = current_genesis.id().to_string();
+    let dag_name = current_genesis.header.timestamp.to_string();
     let (id, _) = random_node.dag_store.read().await.get_dag(&dag_name).last().unwrap().unwrap();
     let genesis_event_id = blake3::Hash::from_bytes((&id as &[u8]).try_into().unwrap());
 
@@ -235,7 +235,7 @@ async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
     // ==========================================
     let random_node = eg_instances.choose(&mut rng).unwrap();
     let current_genesis = random_node.current_genesis.read().await;
-    let dag_name = current_genesis.id().to_string();
+    let dag_name = current_genesis.header.timestamp.to_string();
     let event = Event::new(vec![1, 2, 3, 4], random_node).await;
     assert!(event.header.parents.contains(&genesis_event_id));
     // The node adds it to their DAG, on layer 1.
@@ -243,7 +243,7 @@ async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
     let event_id = random_node.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap()[0];
 
     let store = random_node.dag_store.read().await;
-    let (_, tips_layers) = store.header_dags.get(&current_genesis.id()).unwrap();
+    let (_, tips_layers) = store.header_dags.get(&current_genesis.header.timestamp).unwrap();
 
     // Since genesis was referenced, its layer (0) have been removed
     assert_eq!(tips_layers.len(), 1);
@@ -277,9 +277,9 @@ async fn eventgraph_propagation_real(ex: Arc<Executor<'static>>) {
     let event2_id = random_node.dag_insert(slice::from_ref(&event2), &dag_name).await.unwrap()[0];
     // Genesis event + event from 2. + upper 3 events (layer 4)
     let current_genesis = random_node.current_genesis.read().await;
-    let dag_name = current_genesis.id().to_string();
+    let dag_name = current_genesis.header.timestamp.to_string();
     assert_eq!(random_node.dag_store.read().await.get_dag(&dag_name).len(), 5);
-    let random_node_genesis = random_node.current_genesis.read().await.id();
+    let random_node_genesis = random_node.current_genesis.read().await.header.timestamp;
     let store = random_node.dag_store.read().await;
     let (_, tips_layers) = store.header_dags.get(&random_node_genesis).unwrap();
     assert_eq!(tips_layers.len(), 1);
@@ -464,7 +464,7 @@ async fn eventgraph_chaotic_propagation_real(ex: Arc<Executor<'static>>) {
         let random_node = eg_instances.choose(&mut rng).unwrap();
         let event = Event::new(i.to_be_bytes().to_vec(), random_node).await;
         let current_genesis = random_node.current_genesis.read().await;
-        let dag_name = current_genesis.id().to_string();
+        let dag_name = current_genesis.header.timestamp.to_string();
         random_node.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
         random_node.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
         random_node.p2p.broadcast(&EventPut(event)).await;