commit a7bdfa4bcbc9320c60e0d59d1f650829e6356d83
parent dbf7353f86f30b64ecb67f7a4983bb5bc720bb54
Author: triesap <tyson@radroots.org>
Date: Wed, 1 Jul 2026 20:16:17 +0000
relay_transport: count fetch limits after filtering
Diffstat:
2 files changed, 144 insertions(+), 16 deletions(-)
diff --git a/crates/relay_transport/src/fetch.rs b/crates/relay_transport/src/fetch.rs
@@ -15,6 +15,7 @@ use serde::{Deserialize, Serialize};
use std::sync::{Arc, Mutex, PoisonError};
const DEFAULT_RELAY_FETCH_TIMEOUT_MS: u64 = 10_000;
+const DEFAULT_RELAY_FETCH_RAW_SCAN_MULTIPLIER: usize = 64;
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum RadrootsRelayFetchMode {
@@ -55,6 +56,7 @@ pub struct RadrootsRelayFetchRequest {
mode: RadrootsRelayFetchMode,
observed_at_ms: i64,
max_events: usize,
+ max_raw_events: usize,
relay_urls: Vec<String>,
filters: RadrootsRelayFetchFilters,
timeout_ms: u64,
@@ -106,6 +108,7 @@ impl RadrootsRelayFetchRequest {
mode,
observed_at_ms,
max_events,
+ max_raw_events: default_raw_event_scan_limit(max_events),
relay_urls: Vec::new(),
filters: RadrootsRelayFetchFilters::new(filters)?,
timeout_ms: DEFAULT_RELAY_FETCH_TIMEOUT_MS,
@@ -126,6 +129,11 @@ impl RadrootsRelayFetchRequest {
self
}
+ pub fn with_raw_event_scan_limit(mut self, max_raw_events: usize) -> Self {
+ self.max_raw_events = max_raw_events;
+ self
+ }
+
pub fn mode(&self) -> RadrootsRelayFetchMode {
self.mode
}
@@ -138,6 +146,10 @@ impl RadrootsRelayFetchRequest {
self.max_events
}
+ pub fn max_raw_events(&self) -> usize {
+ self.max_raw_events
+ }
+
pub fn relay_urls(&self) -> &[String] {
&self.relay_urls
}
@@ -151,6 +163,13 @@ impl RadrootsRelayFetchRequest {
}
}
+fn default_raw_event_scan_limit(max_events: usize) -> usize {
+ max_events
+ .saturating_mul(DEFAULT_RELAY_FETCH_RAW_SCAN_MULTIPLIER)
+ .max(max_events)
+ .max(1)
+}
+
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum RadrootsRelayFetchItem {
Event {
@@ -195,6 +214,7 @@ pub struct RadrootsRelayFetchEventReceipt {
pub unsupported: bool,
pub malformed: bool,
pub out_of_filter: bool,
+ pub skipped_over_limit: bool,
pub projection_eligible: bool,
pub verification_status: Option<String>,
pub message: Option<String>,
@@ -206,6 +226,7 @@ pub struct RadrootsRelayFetchReceipt {
pub duplicate_count: usize,
pub malformed_count: usize,
pub out_of_filter_count: usize,
+ pub skipped_over_limit_count: usize,
pub unsupported_count: usize,
pub eose_count: usize,
pub closed_count: usize,
@@ -231,6 +252,7 @@ where
{
let mode = request.mode;
let max_events = request.max_events;
+ let max_raw_events = request.max_raw_events;
if request.filters.as_slice().is_empty() {
return Err(RadrootsRelayTransportError::EmptyFetchFilters);
}
@@ -241,6 +263,7 @@ where
duplicate_count: 0,
malformed_count: 0,
out_of_filter_count: 0,
+ skipped_over_limit_count: 0,
unsupported_count: 0,
eose_count: 0,
closed_count: 0,
@@ -248,7 +271,8 @@ where
events: Vec::new(),
relay_outcomes: Vec::new(),
};
- let mut processed_events = 0usize;
+ let mut scanned_raw_events = 0usize;
+ let mut accepted_events = 0usize;
for item in items {
match item {
RadrootsRelayFetchItem::Event {
@@ -256,10 +280,11 @@ where
raw_json,
observed_at_ms,
} => {
- if processed_events >= max_events {
+ if scanned_raw_events >= max_raw_events {
+ receipt.skipped_over_limit_count += 1;
continue;
}
- processed_events += 1;
+ scanned_raw_events += 1;
let parsed = RadrootsNostrEvent::from_json(raw_json.as_str());
let Ok(raw_event) = parsed else {
receipt.malformed_count += 1;
@@ -271,6 +296,7 @@ where
unsupported: false,
malformed: true,
out_of_filter: false,
+ skipped_over_limit: false,
projection_eligible: false,
verification_status: None,
message: Some("event JSON parse failed".to_owned()),
@@ -287,12 +313,31 @@ where
unsupported: false,
malformed: false,
out_of_filter: true,
+ skipped_over_limit: false,
projection_eligible: false,
verification_status: None,
message: Some("event did not match relay fetch filters".to_owned()),
});
continue;
}
+ if accepted_events >= max_events {
+ receipt.skipped_over_limit_count += 1;
+ receipt.events.push(RadrootsRelayFetchEventReceipt {
+ relay_url,
+ event_id: Some(raw_event.id.to_hex()),
+ inserted: false,
+ duplicate: false,
+ unsupported: false,
+ malformed: false,
+ out_of_filter: false,
+ skipped_over_limit: true,
+ projection_eligible: false,
+ verification_status: None,
+ message: Some("accepted relay fetch event limit reached".to_owned()),
+ });
+ continue;
+ }
+ accepted_events += 1;
let event = radroots_event_from_nostr(&raw_event);
let observation_type = match mode {
RadrootsRelayFetchMode::Fetch => RadrootsRelayObservationType::Fetch,
@@ -327,6 +372,7 @@ where
unsupported,
malformed: false,
out_of_filter: false,
+ skipped_over_limit: false,
projection_eligible: store_receipt.projection_eligible,
verification_status: Some(
store_receipt.verification_status.as_str().to_owned(),
@@ -344,6 +390,7 @@ where
unsupported: false,
malformed: true,
out_of_filter: false,
+ skipped_over_limit: false,
projection_eligible: false,
verification_status: None,
message: Some(error.to_string()),
diff --git a/crates/relay_transport/tests/transport.rs b/crates/relay_transport/tests/transport.rs
@@ -684,31 +684,35 @@ async fn fetch_rejects_out_of_filter_events_before_store_mutation() {
}
#[tokio::test]
-async fn fetch_event_cap_preserves_later_control_outcomes() {
- let first = signed_post("first capped event");
+async fn fetch_event_cap_counts_accepted_in_filter_events_and_preserves_later_control_outcomes() {
+ let accepted = signed_post("accepted capped event");
let skipped = signed_post("skipped capped event");
+ let wrong_tag = signed_event_with_kind_and_hashtag("wrong capped tag", KIND_POST, "compost");
+ let accepted_id = accepted.id.clone();
+ let skipped_id = skipped.id.clone();
+ let wrong_tag_id = wrong_tag.id.clone();
let store = RadrootsEventStore::open_memory().await.expect("store");
let adapter = RadrootsMockRelayFetchAdapter::new(vec![
RadrootsRelayFetchItem::Event {
relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: first.raw_json.clone(),
+ raw_json: "{not json".to_owned(),
+ observed_at_ms: 1_099,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: wrong_tag.raw_json,
observed_at_ms: 1_100,
},
RadrootsRelayFetchItem::Event {
relay_url: RELAY_PRIMARY_WSS.to_owned(),
- raw_json: skipped.raw_json,
+ raw_json: accepted.raw_json.clone(),
observed_at_ms: 1_101,
},
RadrootsRelayFetchItem::Event {
- relay_url: RELAY_SECONDARY_WSS.to_owned(),
- raw_json: "{not json".to_owned(),
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: skipped.raw_json,
observed_at_ms: 1_102,
},
- RadrootsRelayFetchItem::Event {
- relay_url: RELAY_SECONDARY_WSS.to_owned(),
- raw_json: unsupported_raw_event(),
- observed_at_ms: 1_103,
- },
RadrootsRelayFetchItem::Eose {
relay_url: RELAY_PRIMARY_WSS.to_owned(),
},
@@ -730,8 +734,14 @@ async fn fetch_event_cap_preserves_later_control_outcomes() {
assert_eq!(receipt.inserted_count, 1);
assert_eq!(receipt.duplicate_count, 0);
assert_eq!(receipt.unsupported_count, 0);
- assert_eq!(receipt.malformed_count, 0);
- assert_eq!(receipt.events.len(), 1);
+ assert_eq!(receipt.malformed_count, 1);
+ assert_eq!(receipt.out_of_filter_count, 1);
+ assert_eq!(receipt.skipped_over_limit_count, 1);
+ assert_eq!(receipt.events.len(), 4);
+ assert!(receipt.events[0].malformed);
+ assert!(receipt.events[1].out_of_filter);
+ assert!(receipt.events[2].inserted);
+ assert!(receipt.events[3].skipped_over_limit);
assert_eq!(receipt.eose_count, 1);
assert_eq!(receipt.closed_count, 1);
assert_eq!(receipt.notice_count, 1);
@@ -752,6 +762,77 @@ async fn fetch_event_cap_preserves_later_control_outcomes() {
receipt.relay_outcomes[2].kind,
RadrootsRelayFetchOutcomeKind::Notice
);
+ assert!(
+ store
+ .get_event(accepted_id.as_str())
+ .await
+ .expect("accepted lookup")
+ .is_some()
+ );
+ assert!(
+ store
+ .get_event(skipped_id.as_str())
+ .await
+ .expect("skipped lookup")
+ .is_none()
+ );
+ assert!(
+ store
+ .get_event(wrong_tag_id.as_str())
+ .await
+ .expect("wrong tag lookup")
+ .is_none()
+ );
+}
+
+#[tokio::test]
+async fn fetch_raw_scan_limit_bounds_noisy_adapter_output() {
+ let accepted = signed_post("raw scan accepted event");
+ let wrong_tag = signed_event_with_kind_and_hashtag("raw scan wrong tag", KIND_POST, "compost");
+ let accepted_id = accepted.id.clone();
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let adapter = RadrootsMockRelayFetchAdapter::new(vec![
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: "{not json".to_owned(),
+ observed_at_ms: 1_130,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: wrong_tag.raw_json,
+ observed_at_ms: 1_131,
+ },
+ RadrootsRelayFetchItem::Event {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ raw_json: accepted.raw_json,
+ observed_at_ms: 1_132,
+ },
+ RadrootsRelayFetchItem::Eose {
+ relay_url: RELAY_PRIMARY_WSS.to_owned(),
+ },
+ ]);
+
+ let receipt = fetch_and_ingest_relay_events(
+ &adapter,
+ &store,
+ post_relay_fetch_request(1_130, 1).with_raw_event_scan_limit(2),
+ )
+ .await
+ .expect("fetch ingest");
+
+ assert_eq!(receipt.inserted_count, 0);
+ assert_eq!(receipt.malformed_count, 1);
+ assert_eq!(receipt.out_of_filter_count, 1);
+ assert_eq!(receipt.skipped_over_limit_count, 1);
+ assert_eq!(receipt.events.len(), 2);
+ assert_eq!(receipt.eose_count, 1);
+ assert!(
+ store
+ .get_event(accepted_id.as_str())
+ .await
+ .expect("accepted lookup")
+ .is_none()
+ );
}
#[tokio::test]