Преглед изворни кода

event_graph: Make rotating event construction fallible

x пре 1 месец
родитељ
комит
2480a905c3

+ 3 - 3
bin/darkirc/src/irc/client.rs

@@ -546,14 +546,14 @@ impl Client {
             if !args_queue.is_empty() {
                 for _ in 0..args_queue.len() {
                     let privmsg = args_queue.pop_front().unwrap();
-                    pending_events.push(self.privmsg_to_event(privmsg).await);
+                    pending_events.push(self.privmsg_to_event(privmsg).await?);
                 }
                 return Ok(Some(pending_events))
             }
 
             // If queue is empty, create an event and return it
             let privmsg = self.args_to_privmsg(args).await;
-            let event = self.privmsg_to_event(privmsg).await;
+            let event = self.privmsg_to_event(privmsg).await?;
 
             return Ok(Some(vec![event]))
         }
@@ -574,7 +574,7 @@ impl Client {
     }
 
     // Internal helper function that creates an Event from PRIVMSG arguments
-    async fn privmsg_to_event(&self, mut privmsg: Privmsg) -> Event {
+    async fn privmsg_to_event(&self, mut privmsg: Privmsg) -> Result<Event> {
         // Encrypt the Privmsg if an encryption method is available.
         self.server.try_encrypt(&mut privmsg).await;
 

+ 17 - 16
src/event_graph/event.rs

@@ -22,7 +22,8 @@ use darkfi_serial::{async_trait, deserialize_async, Encodable, SerialDecodable,
 use sled_overlay::{sled, SledTreeOverlay};
 
 use super::{
-    util::HOUR_MS, EventGraph, EventGraphConfig, EVENT_TIME_DRIFT, NULL_ID, N_EVENT_PARENTS,
+    util::{unix_timestamp_millis, HOUR_MS},
+    EventGraph, EventGraphConfig, EVENT_TIME_DRIFT, NULL_ID, N_EVENT_PARENTS,
 };
 use crate::Result;
 
@@ -45,31 +46,31 @@ pub struct Header {
 }
 
 impl Header {
-    pub async fn new(content: &[u8], eg: &EventGraph) -> Self {
+    pub async fn new(content: &[u8], eg: &EventGraph) -> Result<Self> {
         let dag_ts = eg.current_genesis.read().await.header.timestamp;
-        let (layer, parents) = eg.get_next_layer_with_parents(&dag_ts).await;
-        Self {
-            timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64,
+        let (layer, parents) = eg.get_next_layer_with_parents(&dag_ts).await?;
+        Ok(Self {
+            timestamp: unix_timestamp_millis()?,
             parents,
             layer,
             content_hash: blake3::hash(content),
-        }
+        })
     }
 
     pub async fn new_static(content: &[u8], eg: &EventGraph) -> Result<Self> {
         let (layer, parents) = eg.get_next_layer_with_parents_static().await?;
         Ok(Self {
-            timestamp: UNIX_EPOCH.elapsed().unwrap().as_millis() as u64,
+            timestamp: unix_timestamp_millis()?,
             parents,
             layer,
             content_hash: blake3::hash(content),
         })
     }
 
-    pub async fn with_timestamp(timestamp: u64, content: &[u8], eg: &EventGraph) -> Self {
+    pub async fn with_timestamp(timestamp: u64, content: &[u8], eg: &EventGraph) -> Result<Self> {
         let dag_ts = eg.current_genesis.read().await.header.timestamp;
-        let (layer, parents) = eg.get_next_layer_with_parents(&dag_ts).await;
-        Self { timestamp, parents, layer, content_hash: blake3::hash(content) }
+        let (layer, parents) = eg.get_next_layer_with_parents(&dag_ts).await?;
+        Ok(Self { timestamp, parents, layer, content_hash: blake3::hash(content) })
     }
 
     /// Blake3 hash of `(timestamp, parents, layer, content_hash)`.
@@ -155,9 +156,9 @@ pub struct Event {
 }
 
 impl Event {
-    pub async fn new(data: Vec<u8>, eg: &EventGraph) -> Self {
-        let header = Header::new(&data, eg).await;
-        Self { header, content: data }
+    pub async fn new(data: Vec<u8>, eg: &EventGraph) -> Result<Self> {
+        let header = Header::new(&data, eg).await?;
+        Ok(Self { header, content: data })
     }
 
     pub async fn new_static(data: Vec<u8>, eg: &EventGraph) -> Result<Self> {
@@ -169,9 +170,9 @@ impl Event {
         self.header.id()
     }
 
-    pub async fn with_timestamp(ts: u64, data: Vec<u8>, eg: &EventGraph) -> Self {
-        let header = Header::with_timestamp(ts, &data, eg).await;
-        Self { header, content: data }
+    pub async fn with_timestamp(ts: u64, data: Vec<u8>, eg: &EventGraph) -> Result<Self> {
+        let header = Header::with_timestamp(ts, &data, eg).await?;
+        Ok(Self { header, content: data })
     }
 
     pub fn content(&self) -> &[u8] {

+ 6 - 2
src/event_graph/mod.rs

@@ -1929,8 +1929,12 @@ impl EventGraph {
     pub(crate) async fn get_next_layer_with_parents(
         &self,
         dag_ts: &u64,
-    ) -> (u64, [blake3::Hash; N_EVENT_PARENTS]) {
-        select_parents_from_tips(&self.dag_store.read().await.get_slot(dag_ts).unwrap().tips)
+    ) -> Result<(u64, [blake3::Hash; N_EVENT_PARENTS])> {
+        let store = self.dag_store.read().await;
+        let slot = store
+            .get_slot(dag_ts)
+            .ok_or_else(|| Error::Custom(format!("event graph DAG slot {dag_ts} not found")))?;
+        Ok(select_parents_from_tips(&slot.tips))
     }
 
     pub(crate) async fn get_next_layer_with_parents_static(

+ 22 - 10
src/event_graph/tests.rs

@@ -420,6 +420,18 @@ fn evgr_dag_store_rejects_corrupt_event_tree_on_open() {
     })
 }
 
+#[test]
+fn evgr_rotating_event_creation_rejects_missing_current_dag_slot() {
+    smol::block_on(async {
+        let eg = make_eg().await;
+        let dag_ts = eg.current_genesis.read().await.header.timestamp;
+        eg.dag_store.write().await.dags.remove(&dag_ts);
+
+        let result = Event::new(b"missing-current-slot".to_vec(), &eg).await;
+        assert!(matches!(result, Err(crate::Error::Custom(_))));
+    })
+}
+
 #[test]
 fn evgr_static_event_creation_rejects_corrupt_static_dag() {
     smol::block_on(async {
@@ -502,7 +514,7 @@ fn evgr_dag_insert_valid_and_duplicate() {
         let dag_name = dag_ts.to_string();
         let sub = eg.event_pub.clone().subscribe().await;
 
-        let event = Event::new(b"hello".to_vec(), &eg).await;
+        let event = Event::new(b"hello".to_vec(), &eg).await.unwrap();
         eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
         let ids = eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
         assert_eq!(ids.len(), 1);
@@ -528,7 +540,7 @@ fn evgr_duplicate_header_insert_does_not_duplicate_time_index() {
         let eg = make_eg().await;
         let dag_ts = eg.current_genesis.read().await.header.timestamp;
         let dag_name = dag_ts.to_string();
-        let event = Event::new(b"time-index-dedup".to_vec(), &eg).await;
+        let event = Event::new(b"time-index-dedup".to_vec(), &eg).await.unwrap();
 
         eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
         let indexed_after_first = {
@@ -557,7 +569,7 @@ fn evgr_dag_insert_without_header_skipped() {
     smol::block_on(async {
         let eg = make_eg().await;
         let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
-        let event = Event::new(b"orphan".to_vec(), &eg).await;
+        let event = Event::new(b"orphan".to_vec(), &eg).await.unwrap();
         let ids = eg.dag_insert(slice::from_ref(&event), &dag_name).await.unwrap();
         assert!(ids.is_empty());
     })
@@ -820,7 +832,7 @@ fn evgr_fetch_page_both_directions() {
         let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
         let base = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;
         for i in 0..(MAX_RANGE_PAGE_SIZE as u64 + 10) {
-            let ev = Event::with_timestamp(base + i, vec![(i % 251) as u8], &eg).await;
+            let ev = Event::with_timestamp(base + i, vec![(i % 251) as u8], &eg).await.unwrap();
             eg.header_dag_insert(vec![ev.header.clone()], &dag_name).await.unwrap();
             eg.dag_insert(slice::from_ref(&ev), &dag_name).await.unwrap();
         }
@@ -938,7 +950,7 @@ async fn propagation_with_real_blob(ex: Arc<Executor<'static>>) {
 
     let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
     let dag_name = dag_ts.to_string();
-    let event = Event::new(b"hello-via-rln".to_vec(), &nodes[0]).await;
+    let event = Event::new(b"hello-via-rln".to_vec(), &nodes[0]).await.unwrap();
 
     let message_id =
         alice.next_message_id(event.header.timestamp).expect("budget available on first signal");
@@ -980,7 +992,7 @@ async fn empty_blob_rejected(ex: Arc<Executor<'static>>) {
     let nodes = make_network(ex).await;
 
     let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
-    let event = Event::new(b"unauthenticated".to_vec(), &nodes[0]).await;
+    let event = Event::new(b"unauthenticated".to_vec(), &nodes[0]).await.unwrap();
 
     nodes[0].p2p.broadcast(&EventPut(event.clone(), vec![])).await;
     sleep(5).await;
@@ -1011,7 +1023,7 @@ async fn malformed_event_rejected_before_rln(ex: Arc<Executor<'static>>) {
     }
 
     let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
-    let event = Event::new(b"preflight-live".to_vec(), &nodes[0]).await;
+    let event = Event::new(b"preflight-live".to_vec(), &nodes[0]).await.unwrap();
     let message_id = alice.next_message_id(event.header.timestamp).expect("budget");
     let blob = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
     let internal_nullifier = blob.internal_nullifier;
@@ -1102,7 +1114,7 @@ async fn dag_sync_with_blob(ex: Arc<Executor<'static>>) {
     let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
     let dag_name = dag_ts.to_string();
 
-    let event = Event::new(b"synced-message".to_vec(), &nodes[0]).await;
+    let event = Event::new(b"synced-message".to_vec(), &nodes[0]).await.unwrap();
     let message_id = alice.next_message_id(event.header.timestamp).expect("budget");
     let blob_struct = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
     let blob = serialize_async(&blob_struct).await;
@@ -1156,7 +1168,7 @@ async fn dag_sync_rejects_bad_blob(ex: Arc<Executor<'static>>) {
 
     let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
     let dag_name = dag_ts.to_string();
-    let event = Event::new(b"crafted-injection".to_vec(), &nodes[0]).await;
+    let event = Event::new(b"crafted-injection".to_vec(), &nodes[0]).await.unwrap();
 
     // Garbage bytes - won't deserialize as a real RLN signal, won't
     // verify. We don't need a real (failing) proof to exercise the
@@ -1215,7 +1227,7 @@ async fn dormant_user_can_post_after_long_silence(ex: Arc<Executor<'static>>) {
         }
     }
 
-    let event = Event::new(b"long-silent-but-still-registered".to_vec(), &nodes[0]).await;
+    let event = Event::new(b"long-silent-but-still-registered".to_vec(), &nodes[0]).await.unwrap();
     let message_id = alice.next_message_id(event.header.timestamp).expect("budget");
     let blob_struct = alice.create_signal(&event, message_id, &nodes[0]).await.unwrap();
     let blob = serialize_async(&blob_struct).await;

+ 6 - 6
src/event_graph/tests_rln.rs

@@ -1254,7 +1254,7 @@ async fn dag_injection_rejected(ex: Arc<Executor<'static>>) {
     // node 0's tip set) but has a garbage blob. We pre-insert
     // its header so the structural validation passes on the
     // recipient.
-    let injected = Event::new(b"injected by malicious peer".to_vec(), &nodes[0]).await;
+    let injected = Event::new(b"injected by malicious peer".to_vec(), &nodes[0]).await.unwrap();
     let bad_blob = b"not-a-real-rln-blob".to_vec();
 
     // Node 0 records the bad event in its own DAG and stashes
@@ -1371,7 +1371,7 @@ fn rln_dag_insert_with_blobs_already_known_skips_verification() {
         let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
 
         // Build a real event so it passes structural validation.
-        let event = Event::new(b"already-known".to_vec(), &eg).await;
+        let event = Event::new(b"already-known".to_vec(), &eg).await.unwrap();
         eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
 
         // First insert via dag_insert (no blob -> trust-the-quorum
@@ -1414,7 +1414,7 @@ fn rln_dag_insert_with_blobs_rejects_missing_blob_on_non_genesis() {
     smol::block_on(async {
         let eg = make_eg().await;
         let dag_name = eg.current_genesis.read().await.header.timestamp.to_string();
-        let event = Event::new(b"missing-blob".to_vec(), &eg).await;
+        let event = Event::new(b"missing-blob".to_vec(), &eg).await.unwrap();
         eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
 
         // Empty blobs slice -> empty blob for every event -> reject.
@@ -1448,7 +1448,7 @@ fn rln_dag_insert_with_blobs_rejects_bad_content_before_verification() {
 
         let dag_ts = eg.current_genesis.read().await.header.timestamp;
         let dag_name = dag_ts.to_string();
-        let event = Event::new(b"preflight-sync".to_vec(), &eg).await;
+        let event = Event::new(b"preflight-sync".to_vec(), &eg).await.unwrap();
         let message_id = alice.next_message_id(event.header.timestamp).expect("budget");
         let blob = alice.create_signal(&event, message_id, &eg).await.unwrap();
         let internal_nullifier = blob.internal_nullifier;
@@ -1487,7 +1487,7 @@ fn rln_insert_signal_with_blob_rejects_missing_blob_on_non_genesis() {
         let eg = make_eg().await;
         let dag_ts = eg.current_genesis.read().await.header.timestamp;
         let dag_name = dag_ts.to_string();
-        let event = Event::new(b"public-missing-blob".to_vec(), &eg).await;
+        let event = Event::new(b"public-missing-blob".to_vec(), &eg).await.unwrap();
         assert_ne!(event.header.parents, NULL_PARENTS);
 
         let result = eg.insert_signal_with_blob(&event, &[], &dag_name).await;
@@ -1560,7 +1560,7 @@ fn rln_dag_blobs_pruned_with_dag_rotation() {
         let dag_name = original_dag_ts.to_string();
 
         // Insert a real event in the current DAG.
-        let event = Event::new(b"to-be-evicted".to_vec(), &eg).await;
+        let event = Event::new(b"to-be-evicted".to_vec(), &eg).await.unwrap();
         eg.header_dag_insert(vec![event.header.clone()], &dag_name).await.unwrap();
         eg.dag_insert(std::slice::from_ref(&event), &dag_name).await.unwrap();
 

+ 9 - 0
src/event_graph/util.rs

@@ -48,6 +48,15 @@ use crate::rpc::{
 /// Milliseconds in one hour.
 pub(super) const HOUR_MS: u64 = 3_600_000;
 
+/// Current UNIX timestamp in milliseconds.
+pub(super) fn unix_timestamp_millis() -> Result<u64> {
+    let elapsed = UNIX_EPOCH
+        .elapsed()
+        .map_err(|_| Error::Custom("system clock is before UNIX epoch".into()))?;
+    u64::try_from(elapsed.as_millis())
+        .map_err(|_| Error::Custom("system clock milliseconds exceed u64".into()))
+}
+
 /// Timestamp (millis) for the start of the hour `hours` offsets from now.
 pub(super) fn next_hour_timestamp(hours: i64) -> u64 {
     let now = UNIX_EPOCH.elapsed().unwrap().as_millis() as u64;