lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit dbf7353f86f30b64ecb67f7a4983bb5bc720bb54
parent 7bebe57b59438518372ecdeed67ef905c4e3a286
Author: triesap <tyson@radroots.org>
Date:   Wed,  1 Jul 2026 20:11:37 +0000

relay_transport: require fetch filters at construction

Diffstat:
Mcrates/relay_transport/src/error.rs | 3+++
Mcrates/relay_transport/src/fetch.rs | 146++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------
Mcrates/relay_transport/src/lib.rs | 6+++---
Mcrates/relay_transport/tests/transport.rs | 93++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
4 files changed, 186 insertions(+), 62 deletions(-)

diff --git a/crates/relay_transport/src/error.rs b/crates/relay_transport/src/error.rs @@ -35,6 +35,9 @@ pub enum RadrootsRelayTransportError { #[error("Relay target set must not be empty")] EmptyTargetSet, + #[error("Relay fetch filters must not be empty")] + EmptyFetchFilters, + #[error("JSON error: {0}")] Json(#[from] serde_json::Error), diff --git a/crates/relay_transport/src/fetch.rs b/crates/relay_transport/src/fetch.rs @@ -23,36 +23,93 @@ pub enum RadrootsRelayFetchMode { } #[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsRelayFetchFilters { + filters: Vec<RadrootsNostrFilter>, +} + +impl RadrootsRelayFetchFilters { + pub fn new<I>(filters: I) -> Result<Self, RadrootsRelayTransportError> + where + I: IntoIterator<Item = RadrootsNostrFilter>, + { + let filters = filters.into_iter().collect::<Vec<_>>(); + if filters.is_empty() { + return Err(RadrootsRelayTransportError::EmptyFetchFilters); + } + Ok(Self { filters }) + } + + pub fn as_slice(&self) -> &[RadrootsNostrFilter] { + &self.filters + } +} + +impl AsRef<[RadrootsNostrFilter]> for RadrootsRelayFetchFilters { + fn as_ref(&self) -> &[RadrootsNostrFilter] { + self.as_slice() + } +} + +#[derive(Clone, Debug, PartialEq, Eq)] pub struct RadrootsRelayFetchRequest { - pub mode: RadrootsRelayFetchMode, - pub observed_at_ms: i64, - pub max_events: usize, - pub relay_urls: Vec<String>, - pub filters: Vec<RadrootsNostrFilter>, - pub timeout_ms: u64, + mode: RadrootsRelayFetchMode, + observed_at_ms: i64, + max_events: usize, + relay_urls: Vec<String>, + filters: RadrootsRelayFetchFilters, + timeout_ms: u64, } impl RadrootsRelayFetchRequest { - pub fn fetch(observed_at_ms: i64, max_events: usize) -> Self { - Self { - mode: RadrootsRelayFetchMode::Fetch, + pub fn fetch<I>( + observed_at_ms: i64, + max_events: usize, + filters: I, + ) -> Result<Self, RadrootsRelayTransportError> + where + I: IntoIterator<Item = RadrootsNostrFilter>, + { + Self::new( + RadrootsRelayFetchMode::Fetch, observed_at_ms, max_events, - relay_urls: Vec::new(), - filters: Vec::new(), - timeout_ms: DEFAULT_RELAY_FETCH_TIMEOUT_MS, - } + filters, + ) } - pub fn subscription(observed_at_ms: i64, max_events: usize) -> Self { - Self { - mode: RadrootsRelayFetchMode::Subscription, + pub fn subscription<I>( + observed_at_ms: i64, + max_events: usize, + filters: I, + ) -> Result<Self, RadrootsRelayTransportError> + where + I: IntoIterator<Item = RadrootsNostrFilter>, + { + Self::new( + RadrootsRelayFetchMode::Subscription, + observed_at_ms, + max_events, + filters, + ) + } + + fn new<I>( + mode: RadrootsRelayFetchMode, + observed_at_ms: i64, + max_events: usize, + filters: I, + ) -> Result<Self, RadrootsRelayTransportError> + where + I: IntoIterator<Item = RadrootsNostrFilter>, + { + Ok(Self { + mode, observed_at_ms, max_events, relay_urls: Vec::new(), - filters: Vec::new(), + filters: RadrootsRelayFetchFilters::new(filters)?, timeout_ms: DEFAULT_RELAY_FETCH_TIMEOUT_MS, - } + }) } pub fn with_relay_urls<I, S>(mut self, relay_urls: I) -> Self @@ -64,18 +121,34 @@ impl RadrootsRelayFetchRequest { self } - pub fn with_filters<I>(mut self, filters: I) -> Self - where - I: IntoIterator<Item = RadrootsNostrFilter>, - { - self.filters = filters.into_iter().collect(); - self - } - pub fn with_timeout_ms(mut self, timeout_ms: u64) -> Self { self.timeout_ms = timeout_ms; self } + + pub fn mode(&self) -> RadrootsRelayFetchMode { + self.mode + } + + pub fn observed_at_ms(&self) -> i64 { + self.observed_at_ms + } + + pub fn max_events(&self) -> usize { + self.max_events + } + + pub fn relay_urls(&self) -> &[String] { + &self.relay_urls + } + + pub fn filters(&self) -> &[RadrootsNostrFilter] { + self.filters.as_slice() + } + + pub fn timeout_ms(&self) -> u64 { + self.timeout_ms + } } #[derive(Clone, Debug, PartialEq, Eq)] @@ -158,7 +231,10 @@ where { let mode = request.mode; let max_events = request.max_events; - let filters = request.filters.clone(); + if request.filters.as_slice().is_empty() { + return Err(RadrootsRelayTransportError::EmptyFetchFilters); + } + let filters = request.filters.as_slice().to_vec(); let items = adapter.fetch(request).await?; let mut receipt = RadrootsRelayFetchReceipt { inserted_count: 0, @@ -311,8 +387,8 @@ fn relay_fetch_event_matches_filters( filters: &[RadrootsNostrFilter], event: &RadrootsNostrEvent, ) -> bool { - filters.is_empty() - || filters + !filters.is_empty() + && filters .iter() .any(|filter| filter.match_event(event, MatchEventOptions::new())) } @@ -335,12 +411,12 @@ async fn fetch_from_nostr_relays( if request.relay_urls.is_empty() { return Err(RadrootsRelayTransportError::EmptyTargetSet); } - if request.filters.is_empty() { - return Err(RadrootsRelayTransportError::Transport( - "relay fetch filters must not be empty".to_owned(), - )); + if request.filters.as_slice().is_empty() { + return Err(RadrootsRelayTransportError::EmptyFetchFilters); } let timeout = Duration::from_millis(request.timeout_ms); + let filters = request.filters.as_slice().to_vec(); + let observed_at_ms = request.observed_at_ms; let mut items = Vec::new(); for relay_url in request.relay_urls { let client = RadrootsNostrClient::new_signerless(); @@ -353,14 +429,14 @@ async fn fetch_from_nostr_relays( } client.connect().await; let mut closed = false; - for filter in request.filters.iter().cloned() { + for filter in filters.iter().cloned() { match client.fetch_events(filter, timeout).await { Ok(events) => { for event in events { items.push(RadrootsRelayFetchItem::Event { relay_url: relay_url.clone(), raw_json: event.as_json(), - observed_at_ms: request.observed_at_ms, + observed_at_ms, }); } } diff --git a/crates/relay_transport/src/lib.rs b/crates/relay_transport/src/lib.rs @@ -14,9 +14,9 @@ pub use error::RadrootsRelayTransportError; #[cfg(feature = "storage")] pub use fetch::{ RadrootsMockRelayFetchAdapter, RadrootsNostrClientFetchAdapter, RadrootsRelayFetchAdapter, - RadrootsRelayFetchEventReceipt, RadrootsRelayFetchItem, RadrootsRelayFetchMode, - RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchReceipt, RadrootsRelayFetchRelayOutcome, - RadrootsRelayFetchRequest, fetch_and_ingest_relay_events, + RadrootsRelayFetchEventReceipt, RadrootsRelayFetchFilters, RadrootsRelayFetchItem, + RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchReceipt, + RadrootsRelayFetchRelayOutcome, RadrootsRelayFetchRequest, fetch_and_ingest_relay_events, }; #[cfg(feature = "storage")] pub use outbox::{ diff --git a/crates/relay_transport/tests/transport.rs b/crates/relay_transport/tests/transport.rs @@ -131,6 +131,47 @@ fn unsupported_raw_event() -> String { event.as_json() } +fn post_relay_fetch_filter(limit: usize) -> RadrootsNostrFilter { + radroots_nostr_filter_tag( + RadrootsNostrFilter::new() + .kind(RadrootsNostrKind::Custom(KIND_POST as u16)) + .limit(limit), + "t", + vec!["soil".to_owned()], + ) + .expect("post relay fetch filter") +} + +fn unsupported_relay_fetch_filter(limit: usize) -> RadrootsNostrFilter { + RadrootsNostrFilter::new() + .kind(RadrootsNostrKind::Custom(999)) + .limit(limit) +} + +fn fixture_relay_fetch_request( + observed_at_ms: i64, + max_events: usize, +) -> RadrootsRelayFetchRequest { + RadrootsRelayFetchRequest::fetch( + observed_at_ms, + max_events, + [ + post_relay_fetch_filter(max_events), + unsupported_relay_fetch_filter(max_events), + ], + ) + .expect("fixture relay fetch request") +} + +fn post_relay_fetch_request(observed_at_ms: i64, max_events: usize) -> RadrootsRelayFetchRequest { + RadrootsRelayFetchRequest::fetch( + observed_at_ms, + max_events, + [post_relay_fetch_filter(max_events)], + ) + .expect("post relay fetch request") +} + fn tampered_raw_event() -> String { let signed = signed_post("trusted"); let mut value = @@ -434,6 +475,18 @@ async fn publish_receipts_track_terminal_skipped_and_adapter_errors() { assert!(matches!(error, RadrootsRelayTransportError::Transport(_))); } +#[test] +fn fetch_requests_reject_empty_filter_sets() { + assert!(matches!( + RadrootsRelayFetchRequest::fetch(1_000, 10, Vec::<RadrootsNostrFilter>::new()), + Err(RadrootsRelayTransportError::EmptyFetchFilters) + )); + assert!(matches!( + RadrootsRelayFetchRequest::subscription(1_000, 10, Vec::<RadrootsNostrFilter>::new()), + Err(RadrootsRelayTransportError::EmptyFetchFilters) + )); +} + #[tokio::test] async fn fetch_ingests_events_and_records_relay_observations() { let signed = signed_post("hello"); @@ -481,13 +534,10 @@ async fn fetch_ingests_events_and_records_relay_observations() { }, ]); - let receipt = fetch_and_ingest_relay_events( - &adapter, - &store, - RadrootsRelayFetchRequest::fetch(1_000, 10), - ) - .await - .expect("fetch ingest"); + let receipt = + fetch_and_ingest_relay_events(&adapter, &store, fixture_relay_fetch_request(1_000, 10)) + .await + .expect("fetch ingest"); assert_eq!(receipt.inserted_count, 3); assert_eq!(receipt.duplicate_count, 1); @@ -597,7 +647,7 @@ async fn fetch_rejects_out_of_filter_events_before_store_mutation() { let receipt = fetch_and_ingest_relay_events( &adapter, &store, - RadrootsRelayFetchRequest::fetch(1_005, 10).with_filters([filter]), + RadrootsRelayFetchRequest::fetch(1_005, 10, [filter]).expect("fetch request"), ) .await .expect("fetch ingest"); @@ -673,7 +723,7 @@ async fn fetch_event_cap_preserves_later_control_outcomes() { ]); let receipt = - fetch_and_ingest_relay_events(&adapter, &store, RadrootsRelayFetchRequest::fetch(1_100, 1)) + fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(1_100, 1)) .await .expect("fetch ingest"); @@ -717,7 +767,8 @@ async fn fetch_subscription_mode_and_store_errors_are_reported() { let receipt = fetch_and_ingest_relay_events( &adapter, &store, - RadrootsRelayFetchRequest::subscription(1_200, 10), + RadrootsRelayFetchRequest::subscription(1_200, 10, [post_relay_fetch_filter(10)]) + .expect("subscription request"), ) .await .expect("fetch ingest"); @@ -737,13 +788,10 @@ async fn fetch_subscription_mode_and_store_errors_are_reported() { raw_json: signed.raw_json, observed_at_ms: 1_210, }]); - let receipt = fetch_and_ingest_relay_events( - &adapter, - &closed_store, - RadrootsRelayFetchRequest::fetch(1_210, 10), - ) - .await - .expect("fetch ingest"); + let receipt = + fetch_and_ingest_relay_events(&adapter, &closed_store, post_relay_fetch_request(1_210, 10)) + .await + .expect("fetch ingest"); assert_eq!(receipt.inserted_count, 0); assert_eq!(receipt.malformed_count, 1); @@ -1499,13 +1547,10 @@ async fn smoke_relay_fetch_processes_one_thousand_event_receipts() { }); } let adapter = RadrootsMockRelayFetchAdapter::new(items); - let receipt = fetch_and_ingest_relay_events( - &adapter, - &store, - RadrootsRelayFetchRequest::fetch(10_000, 1_000), - ) - .await - .expect("fetch"); + let receipt = + fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(10_000, 1_000)) + .await + .expect("fetch"); assert_eq!(receipt.inserted_count, 1_000); assert_eq!(receipt.duplicate_count, 0);