Procházet zdrojové kódy

event_graph: Fast fail on parent fetch stalls

x před 1 měsícem
rodič
revize
a9fe0d75a2
2 změnil soubory, kde provedl 134 přidání a 18 odebrání
  1. 56 17
      src/event_graph/proto.rs
  2. 78 1
      src/event_graph/tests.rs

+ 56 - 17
src/event_graph/proto.rs

@@ -41,6 +41,7 @@ use tracing::{error, warn};
 
 use super::{
     event::Header,
+    filter_requested_event_rep,
     rln::{self, create_slash_proof, sss_recover, RLNNode, SlashBlob},
     Event, EventGraphPtr, LayerUTips, NULL_ID, NULL_PARENTS,
 };
@@ -148,6 +149,39 @@ pub(crate) fn cap_layer_tips(tips: &LayerUTips, limit: usize) -> LayerUTips {
     out
 }
 
+/// Match one parent-fetch response against the exact request and update the
+/// outstanding frontier only if the response resolves at least one missing ID.
+pub(crate) fn filter_parent_event_rep(
+    requested: &[blake3::Hash],
+    missing: &mut HashSet<blake3::Hash>,
+    known: &mut HashSet<blake3::Hash>,
+    events: Vec<Event>,
+    blobs: Vec<Vec<u8>>,
+) -> Result<Vec<(Event, Vec<u8>)>> {
+    let (matched_events, matched_blobs, _) = filter_requested_event_rep(requested, events, blobs)?;
+    let resolves_missing = matched_events.iter().any(|event| {
+        let id = event.id();
+        missing.contains(&id) && !known.contains(&id)
+    });
+    if !resolves_missing {
+        return Err(Error::DagSyncFailed)
+    }
+
+    let mut resolved = Vec::with_capacity(matched_events.len());
+    for (event, blob) in matched_events.into_iter().zip(matched_blobs) {
+        let id = event.id();
+        if missing.remove(&id) && known.insert(id) {
+            resolved.push((event, blob));
+        }
+    }
+
+    if resolved.is_empty() {
+        return Err(Error::DagSyncFailed)
+    }
+
+    Ok(resolved)
+}
+
 struct MovingWindow {
     times: VecDeque<NanoTimestamp>,
     expiry_time: NanoTimestamp,
@@ -536,9 +570,8 @@ impl ProtocolEventGraph {
     ) -> bool {
         // received[layer] = Vec<(event, blob)> - keeping events
         // paired with their blobs through the layer ordering so we
-        // can re-verify proofs at insert time. An empty blob means
-        // the serving peer didn't have one; dag_insert_with_blobs
-        // treats that as the trust-the-quorum fallback.
+        // can re-verify proofs at insert time. Non-genesis parents
+        // without blobs are rejected by dag_insert_with_blobs.
         let mut received: BTreeMap<u64, Vec<(Event, Vec<u8>)>> = BTreeMap::new();
         let mut known = HashSet::new();
         let mut depth = 0usize;
@@ -555,7 +588,8 @@ impl ProtocolEventGraph {
                 return false
             }
 
-            if self.channel.send(&EventReq(missing.iter().cloned().collect())).await.is_err() {
+            let requested: Vec<_> = missing.iter().copied().collect();
+            if self.channel.send(&EventReq(requested.clone())).await.is_err() {
                 return false
             }
 
@@ -567,23 +601,28 @@ impl ProtocolEventGraph {
                 return false
             };
 
-            // Pair each returned event with its corresponding blob.
-            let parents = rep.0.clone();
-            let blobs_in = rep.1.clone();
-            let blobs_aligned = blobs_in.len() == parents.len();
-            for (i, parent) in parents.into_iter().enumerate() {
-                let pid = parent.id();
-                if !missing.contains(&pid) {
-                    // Peer sent an event we didn't ask for
-                    self.channel.stop().await;
+            let parents = match filter_parent_event_rep(
+                &requested,
+                missing,
+                &mut known,
+                rep.0.clone(),
+                rep.1.clone(),
+            ) {
+                Ok(parents) => parents,
+                Err(e) => {
+                    warn!(
+                        target: "event_graph::protocol",
+                        "[EVENTGRAPH] parent fetch made no progress or returned invalid data: {e}",
+                    );
+                    let _ = self.clone().strike().await;
                     return false
                 }
-                let blob = if blobs_aligned { blobs_in[i].clone() } else { Vec::new() };
+            };
+
+            for (parent, blob) in parents {
                 received.entry(parent.header.layer).or_default().push((parent.clone(), blob));
-                known.insert(pid);
-                missing.remove(&pid);
 
-                // Check for more unknown grandparents
+                // Check for more unknown grandparents.
                 let store = self.event_graph.dag_store.read().await;
                 if let Some(slot) = store.get_slot(&dag_ts) {
                     for gp in parent.header.parents.iter() {

+ 78 - 1
src/event_graph/tests.rs

@@ -33,7 +33,10 @@ use crate::{
         compute_unreferenced_tips,
         event::Header,
         filter_requested_event_rep, merge_static_sync_event_rep,
-        proto::{cap_layer_tips, count_layer_tips, EventPut, SyncDirection, MAX_RANGE_PAGE_SIZE},
+        proto::{
+            cap_layer_tips, count_layer_tips, filter_parent_event_rep, EventPut, SyncDirection,
+            MAX_RANGE_PAGE_SIZE,
+        },
         rln::epoch_of,
         test_helpers::{
             archive_config, bounded_dag_store_config, init_logger, make_eg, make_eg_with_config,
@@ -89,6 +92,80 @@ fn evgr_event_rep_filter_matches_only_requested_ids() {
     assert!(filter_requested_event_rep(&requested, vec![event_a], Vec::new()).is_err());
 }
 
+#[test]
+fn evgr_parent_event_rep_requires_progress() {
+    let parent_a = test_event(b"parent-a", 21);
+    let parent_b = test_event(b"parent-b", 22);
+    let unrelated = test_event(b"unrelated-parent", 23);
+    let requested = vec![parent_a.id(), parent_b.id()];
+
+    let mut missing: HashSet<_> = requested.iter().copied().collect();
+    let mut known = HashSet::new();
+    assert!(filter_parent_event_rep(&requested, &mut missing, &mut known, vec![], vec![]).is_err());
+    assert_eq!(missing.len(), 2);
+    assert!(known.is_empty());
+
+    assert!(filter_parent_event_rep(
+        &requested,
+        &mut missing,
+        &mut known,
+        vec![unrelated],
+        vec![b"blob".to_vec()],
+    )
+    .is_err());
+    assert_eq!(missing.len(), 2);
+    assert!(known.is_empty());
+
+    assert!(filter_parent_event_rep(
+        &requested,
+        &mut missing,
+        &mut known,
+        vec![parent_a.clone(), parent_a.clone()],
+        vec![b"blob-a".to_vec(), b"blob-a-duplicate".to_vec()],
+    )
+    .is_err());
+    assert_eq!(missing.len(), 2);
+    assert!(known.is_empty());
+
+    let resolved = filter_parent_event_rep(
+        &requested,
+        &mut missing,
+        &mut known,
+        vec![parent_a.clone()],
+        vec![b"blob-a".to_vec()],
+    )
+    .unwrap();
+    assert_eq!(
+        resolved.iter().map(|(event, _)| event.id()).collect::<Vec<_>>(),
+        vec![parent_a.id()]
+    );
+    assert!(!missing.contains(&parent_a.id()));
+    assert!(missing.contains(&parent_b.id()));
+    assert!(known.contains(&parent_a.id()));
+
+    let current_request = vec![parent_b.id()];
+    assert!(filter_parent_event_rep(
+        &current_request,
+        &mut missing,
+        &mut known,
+        vec![parent_a.clone()],
+        vec![b"stale-blob-a".to_vec()],
+    )
+    .is_err());
+    assert!(missing.contains(&parent_b.id()));
+
+    let mut empty_missing = HashSet::new();
+    let mut empty_known = HashSet::new();
+    assert!(filter_parent_event_rep(
+        &[parent_a.id()],
+        &mut empty_missing,
+        &mut empty_known,
+        vec![parent_a],
+        vec![b"blob-a".to_vec()],
+    )
+    .is_err());
+}
+
 #[test]
 fn evgr_static_sync_merge_tracks_partial_requested_batches() {
     let event_a = test_event(b"static-requested-a", 11);