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

event_graph: Enforce rotating DAG body closure

x пре 1 месец
родитељ
комит
48e2334ebf
3 измењених фајлова са 263 додато и 13 уклоњено
  1. 84 10
      src/event_graph/mod.rs
  2. 12 3
      src/event_graph/proto.rs
  3. 167 0
      src/event_graph/tests.rs

+ 84 - 10
src/event_graph/mod.rs

@@ -648,6 +648,20 @@ pub struct EventGraph {
     rln_app_id: rln::RlnAppId,
 }
 
+fn sort_event_indices(events: &[Event], indices: &mut [usize]) {
+    indices.sort_by(|a, b| {
+        let a_event = &events[*a];
+        let b_event = &events[*b];
+        let a_id = a_event.id();
+        let b_id = b_event.id();
+        a_event
+            .header
+            .layer
+            .cmp(&b_event.header.layer)
+            .then_with(|| a_id.as_bytes().cmp(b_id.as_bytes()))
+    });
+}
+
 impl EventGraph {
     /// Create a new Event Graph.
     pub async fn new(
@@ -1714,13 +1728,19 @@ impl EventGraph {
             (already_have, structurally_valid)
         };
 
+        let mut candidates: Vec<usize> = (0..events.len()).collect();
+        sort_event_indices(events, &mut candidates);
+
         let mut accepted: Vec<usize> = Vec::with_capacity(events.len());
-        for (i, ev) in events.iter().enumerate() {
+        let mut accepted_body_ids = HashSet::with_capacity(events.len());
+        for i in candidates {
+            let ev = &events[i];
+            let eid = ev.id();
             if !structurally_valid[i] {
                 error!(
                     target: "event_graph::dag_insert",
                     "[DAG_INSERT] event {} failed structural validation before RLN verification; skipping",
-                    ev.id(),
+                    eid,
                 );
                 continue
             }
@@ -1731,6 +1751,7 @@ impl EventGraph {
             // (epoch, internal_nullifier, x, y) tuple.
             if already_have[i] {
                 accepted.push(i);
+                accepted_body_ids.insert(eid);
                 continue
             }
             // Genesis-shaped events have no blob and no proof - they're
@@ -1739,29 +1760,45 @@ impl EventGraph {
                 accepted.push(i);
                 continue
             }
+
+            if !self.parents_have_bodies(ev, dag_ts, &accepted_body_ids).await? {
+                error!(
+                    target: "event_graph::dag_insert",
+                    "[DAG_INSERT] event {} has a missing parent body; skipping before RLN verification",
+                    eid,
+                );
+                continue
+            }
+
             let blob = blobs.get(i).cloned().unwrap_or_default();
             if blob.is_empty() {
                 if require_blobs {
                     error!(
                         target: "event_graph::dag_insert",
-                        "[DAG_INSERT] sync event {} arrived without an RLN blob; rejecting. \
-                         Every non-genesis rotating-DAG event must carry a blob.",
-                        ev.id(),
+                        concat!(
+                            "[DAG_INSERT] sync event {} arrived without an RLN blob; rejecting. ",
+                            "Every non-genesis rotating-DAG event must carry a blob.",
+                        ),
+                        eid,
                     );
                     continue
                 }
                 // Lenient path: caller pre-verified. Accept the event
                 // structurally without running the RLN verifier on it.
                 accepted.push(i);
+                accepted_body_ids.insert(eid);
                 continue
             }
             match self.rln_verify_signal(ev, &blob).await {
-                rln::SignalCheck::Accepted => accepted.push(i),
+                rln::SignalCheck::Accepted => {
+                    accepted.push(i);
+                    accepted_body_ids.insert(eid);
+                }
                 rln::SignalCheck::Rejected => {
                     error!(
                         target: "event_graph::dag_insert",
                         "[DAG_INSERT] sync event {} failed RLN re-verification; skipping",
-                        ev.id(),
+                        eid,
                     );
                 }
                 rln::SignalCheck::Slashable(_) => {
@@ -1774,7 +1811,7 @@ impl EventGraph {
                     error!(
                         target: "event_graph::dag_insert",
                         "[DAG_INSERT] sync event {} is slashable (slot reuse); skipping",
-                        ev.id(),
+                        eid,
                     );
                 }
             }
@@ -1784,16 +1821,22 @@ impl EventGraph {
         let mut store = self.dag_store.write().await;
         let slot = store.get_slot_mut(&dag_ts).ok_or(Error::DagSyncFailed)?;
 
+        let mut accepted = accepted;
+        sort_event_indices(events, &mut accepted);
+
         let mut ids = Vec::with_capacity(accepted.len());
+        let mut committed_indices = Vec::with_capacity(accepted.len());
         let mut overlay = SledTreeOverlay::new(&slot.main_tree);
+        let mut staged_body_ids = HashSet::with_capacity(accepted.len());
 
-        for &i in &accepted {
+        'commit: for &i in &accepted {
             let ev = &events[i];
             let eid = ev.id();
             if ev.header.parents == NULL_PARENTS {
                 continue
             }
             if slot.main_tree.contains_key(eid.as_bytes())? {
+                staged_body_ids.insert(eid);
                 continue
             }
             if !slot.header_tree.contains_key(eid.as_bytes())? {
@@ -1802,8 +1845,20 @@ impl EventGraph {
             if !ev.dag_validate(&slot.header_tree, &self.config, dag_ts).await? {
                 return Err(Error::EventIsInvalid)
             }
+            for pid in ev.header.parents.iter().filter(|pid| **pid != NULL_ID) {
+                if !staged_body_ids.contains(pid) && !slot.main_tree.contains_key(pid.as_bytes())? {
+                    error!(
+                        target: "event_graph::dag_insert",
+                        "[DAG_INSERT] event {} has parent header {} but no committed parent body; skipping",
+                        eid, pid,
+                    );
+                    continue 'commit
+                }
+            }
+
             let se = serialize_async(ev).await;
             overlay.insert(eid.as_bytes(), &se)?;
+            staged_body_ids.insert(eid);
             if self.replay_mode {
                 replayer_log(&self.datastore, "insert".into(), se)?;
             }
@@ -1817,6 +1872,7 @@ impl EventGraph {
             }
 
             ids.push(eid);
+            committed_indices.push(i);
         }
 
         if let Some(b) = overlay.aggregate() {
@@ -1825,7 +1881,7 @@ impl EventGraph {
             return Ok(vec![])
         }
 
-        for &i in &accepted {
+        for &i in &committed_indices {
             let ev = &events[i];
             let eid = ev.id();
             if ev.header.parents == NULL_PARENTS {
@@ -1849,6 +1905,24 @@ impl EventGraph {
         Ok(ids)
     }
 
+    async fn parents_have_bodies(
+        &self,
+        ev: &Event,
+        dag_ts: u64,
+        accepted_body_ids: &HashSet<blake3::Hash>,
+    ) -> Result<bool> {
+        let store = self.dag_store.read().await;
+        let Some(slot) = store.get_slot(&dag_ts) else { return Ok(false) };
+
+        for pid in ev.header.parents.iter().filter(|pid| **pid != NULL_ID) {
+            if !accepted_body_ids.contains(pid) && !slot.main_tree.contains_key(pid.as_bytes())? {
+                return Ok(false)
+            }
+        }
+
+        Ok(true)
+    }
+
     pub async fn header_dag_insert(&self, headers: Vec<Header>, dag_name: &str) -> Result<()> {
         let dag_ts = u64::from_str(dag_name)?;
 

+ 12 - 3
src/event_graph/proto.rs

@@ -649,13 +649,22 @@ impl ProtocolEventGraph {
             return false
         }
 
-        // dag_insert_with_blobs verifies each event's blob (when
-        // present) and skips events that fail RLN re-verification.
-        // This closes sync-time injection via fetch_parents.
+        // dag_insert_with_blobs verifies each event's blob and skips events
+        // that fail RLN re-verification or parent-body closure. Parent fetch is
+        // strict: if any fetched parent body failed to commit, the original
+        // child must not proceed to RLN verification.
         if self.event_graph.dag_insert_with_blobs(&events, &blobs, dag_name).await.is_err() {
             return false
         }
 
+        let store = self.event_graph.dag_store.read().await;
+        let Some(slot) = store.get_slot(&dag_ts) else { return false };
+        for event in &events {
+            if !slot.main_tree.contains_key(event.id().as_bytes()).unwrap_or(false) {
+                return false
+            }
+        }
+
         true
     }
 

+ 167 - 0
src/event_graph/tests.rs

@@ -1326,6 +1326,82 @@ fn evgr_fetch_missing_events_rejects_corrupt_header_record() {
     })
 }
 
+async fn make_parent_child_with_child_blob(
+    source: &EventGraphPtr,
+    alice: &mut TestIdentity,
+    parent_content: &[u8],
+    child_content: &[u8],
+) -> (Event, Event, Vec<u8>) {
+    let dag_ts = source.current_genesis.read().await.header.timestamp;
+    let dag_name = dag_ts.to_string();
+
+    let parent = Event::new(parent_content.to_vec(), source).await.unwrap();
+    source.header_dag_insert(vec![parent.header.clone()], &dag_name).await.unwrap();
+    source.dag_insert(slice::from_ref(&parent), &dag_name).await.unwrap();
+
+    let child = Event::new(child_content.to_vec(), source).await.unwrap();
+    assert!(child.header.parents.contains(&parent.id()));
+    let message_id = alice.next_message_id(child.header.timestamp).expect("budget");
+    let child_blob = alice.create_signal(&child, message_id, source).await.unwrap();
+
+    (parent, child, serialize_async(&child_blob).await)
+}
+
+async fn seed_rotating_event_unchecked(
+    eg: &EventGraphPtr,
+    event: &Event,
+    blob: &[u8],
+    dag_name: &str,
+) {
+    eg.header_dag_insert(vec![event.header.clone()], dag_name).await.unwrap();
+    eg.dag_insert(slice::from_ref(event), dag_name).await.unwrap();
+    eg.dag_blob_store(&event.id(), blob).unwrap();
+    eg.broadcasted_ids.write().await.insert(event.id());
+}
+
+#[test]
+fn evgr_dag_insert_rejects_child_when_parent_body_rejected() {
+    smol::block_on(async {
+        let source = make_eg().await;
+        let recipient = make_eg().await;
+        let mut alice = TestIdentity::new();
+        alice.register_directly(&source).await.unwrap();
+        alice.register_directly(&recipient).await.unwrap();
+
+        let dag_ts = source.current_genesis.read().await.header.timestamp;
+        let dag_name = dag_ts.to_string();
+        let (parent, child, child_blob) = make_parent_child_with_child_blob(
+            &source,
+            &mut alice,
+            b"parent-body-closure-direct",
+            b"child-body-closure-direct",
+        )
+        .await;
+        let bad_parent_blob = b"bad-parent-rln-blob".to_vec();
+
+        recipient
+            .header_dag_insert(vec![child.header.clone(), parent.header.clone()], &dag_name)
+            .await
+            .unwrap();
+        let inserted = recipient
+            .dag_insert_with_blobs(
+                &[child.clone(), parent.clone()],
+                &[child_blob.clone(), bad_parent_blob],
+                &dag_name,
+            )
+            .await
+            .unwrap();
+        assert!(inserted.is_empty());
+
+        let store = recipient.dag_store.read().await;
+        let slot = store.get_slot(&dag_ts).unwrap();
+        assert!(!slot.main_tree.contains_key(parent.id().as_bytes()).unwrap());
+        assert!(!slot.main_tree.contains_key(child.id().as_bytes()).unwrap());
+        drop(store);
+        assert!(recipient.dag_blob_fetch(&child.id()).unwrap().is_none());
+    })
+}
+
 #[test]
 fn evgr_multi_node_dag_sync_with_blob() {
     init_logger();
@@ -1382,6 +1458,52 @@ async fn dag_sync_with_blob(ex: Arc<Executor<'static>>) {
     shutdown_network(&nodes).await;
 }
 
+#[test]
+fn evgr_multi_node_dag_sync_rejects_child_when_parent_body_rejected() {
+    init_logger();
+    run_multi_node_test(dag_sync_rejects_child_when_parent_body_rejected);
+}
+async fn dag_sync_rejects_child_when_parent_body_rejected(ex: Arc<Executor<'static>>) {
+    let nodes = make_network(ex).await;
+
+    let mut alice = TestIdentity::new();
+    for eg in &nodes {
+        alice.register_directly(eg).await.unwrap();
+    }
+
+    let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
+    let dag_name = dag_ts.to_string();
+    let (parent, child, child_blob) = make_parent_child_with_child_blob(
+        &nodes[0],
+        &mut alice,
+        b"parent-body-closure-sync",
+        b"child-body-closure-sync",
+    )
+    .await;
+    let bad_parent_blob = b"bad-parent-rln-blob-sync".to_vec();
+
+    for eg in nodes.iter().take(4) {
+        seed_rotating_event_unchecked(eg, &parent, &bad_parent_blob, &dag_name).await;
+        seed_rotating_event_unchecked(eg, &child, &child_blob, &dag_name).await;
+    }
+
+    nodes[4].dag_sync(dag_ts).await.unwrap();
+    sleep(2).await;
+
+    let store = nodes[4].dag_store.read().await;
+    let slot = store.get_slot(&dag_ts).unwrap();
+    assert!(
+        !slot.main_tree.contains_key(parent.id().as_bytes()).unwrap(),
+        "sync accepted a parent whose blob failed RLN verification",
+    );
+    assert!(
+        !slot.main_tree.contains_key(child.id().as_bytes()).unwrap(),
+        "sync accepted a child body whose parent body was rejected",
+    );
+
+    shutdown_network(&nodes).await;
+}
+
 #[test]
 fn evgr_multi_node_dag_sync_rejects_bad_blob() {
     init_logger();
@@ -1427,6 +1549,51 @@ async fn dag_sync_rejects_bad_blob(ex: Arc<Executor<'static>>) {
     shutdown_network(&nodes).await;
 }
 
+#[test]
+fn evgr_multi_node_fetch_parents_rejects_child_when_parent_body_rejected() {
+    init_logger();
+    run_multi_node_test(fetch_parents_rejects_child_when_parent_body_rejected);
+}
+async fn fetch_parents_rejects_child_when_parent_body_rejected(ex: Arc<Executor<'static>>) {
+    let nodes = make_network(ex).await;
+
+    let mut alice = TestIdentity::new();
+    for eg in &nodes {
+        alice.register_directly(eg).await.unwrap();
+    }
+
+    let dag_ts = nodes[0].current_genesis.read().await.header.timestamp;
+    let dag_name = dag_ts.to_string();
+    let (parent, child, child_blob) = make_parent_child_with_child_blob(
+        &nodes[0],
+        &mut alice,
+        b"parent-body-closure-fetch",
+        b"child-body-closure-fetch",
+    )
+    .await;
+    let bad_parent_blob = b"bad-parent-rln-blob-fetch".to_vec();
+
+    for eg in nodes.iter().take(4) {
+        seed_rotating_event_unchecked(eg, &parent, &bad_parent_blob, &dag_name).await;
+    }
+
+    nodes[0].p2p.broadcast(&EventPut(child.clone(), child_blob)).await;
+    sleep(5).await;
+
+    let store = nodes[4].dag_store.read().await;
+    let slot = store.get_slot(&dag_ts).unwrap();
+    assert!(
+        !slot.main_tree.contains_key(parent.id().as_bytes()).unwrap(),
+        "fetch_parents accepted a parent whose blob failed RLN verification",
+    );
+    assert!(
+        !slot.main_tree.contains_key(child.id().as_bytes()).unwrap(),
+        "live EventPut accepted a child whose fetched parent body was rejected",
+    );
+
+    shutdown_network(&nodes).await;
+}
+
 #[test]
 fn evgr_multi_node_dormant_user_can_post_after_long_silence() {
     init_logger();