lib

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

commit bd7e16c21befc30890f0e99c14448d415cbbe83e
parent 3d299ea5caa6203207c6c821884fb87b7d9039d9
Author: triesap <tyson@radroots.org>
Date:   Wed,  1 Jul 2026 21:11:26 +0000

relay_transport: expose shared fetch receipts

Add a non-ingesting shared relay fetch receipt with accepted events, relay outcomes, failure evidence, and a runtime-tokio blocking runner so downstream CLI read paths can consume fail-closed relay fetch semantics without owning a direct Nostr fetch helper.

Validation: cargo extbuild run -- cargo fmt --all; cargo extbuild run -- cargo check -p radroots_relay_transport --all-features; cargo extbuild run -- cargo test -p radroots_relay_transport --all-features

Diffstat:
Mcrates/relay_transport/Cargo.toml | 2++
Mcrates/relay_transport/src/fetch.rs | 442+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------
Mcrates/relay_transport/src/lib.rs | 10+++++++---
Mcrates/relay_transport/tests/transport.rs | 83++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
4 files changed, 434 insertions(+), 103 deletions(-)

diff --git a/crates/relay_transport/Cargo.toml b/crates/relay_transport/Cargo.toml @@ -22,6 +22,7 @@ client = [ ] storage = ["dep:radroots_event_store", "dep:radroots_outbox", "client"] runtime-tokio = [ + "dep:tokio", "storage", "radroots_event_store/runtime-tokio", "radroots_outbox/runtime-tokio", @@ -50,6 +51,7 @@ nostr = { workspace = true } serde = { workspace = true, features = ["derive", "std"] } serde_json = { workspace = true, features = ["std"] } thiserror = { workspace = true } +tokio = { workspace = true, optional = true, features = ["rt"] } url = { workspace = true } [dev-dependencies] diff --git a/crates/relay_transport/src/fetch.rs b/crates/relay_transport/src/fetch.rs @@ -220,6 +220,36 @@ pub struct RadrootsRelayFetchEventReceipt { pub message: Option<String>, } +#[derive(Clone, Debug)] +pub struct RadrootsRelayFetchedEvent { + pub relay_url: String, + pub event: RadrootsNostrEvent, + pub raw_json: String, + pub observed_at_ms: i64, +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct RadrootsRelayFetchFailure { + pub relay_url: String, + pub reason: String, +} + +#[derive(Clone, Debug)] +pub struct RadrootsRelayFetchedEventsReceipt { + pub target_relays: Vec<String>, + pub connected_relays: Vec<String>, + pub failed_relays: Vec<RadrootsRelayFetchFailure>, + pub events: Vec<RadrootsRelayFetchedEvent>, + pub event_receipts: Vec<RadrootsRelayFetchEventReceipt>, + pub malformed_count: usize, + pub out_of_filter_count: usize, + pub skipped_over_limit_count: usize, + pub eose_count: usize, + pub closed_count: usize, + pub notice_count: usize, + pub relay_outcomes: Vec<RadrootsRelayFetchRelayOutcome>, +} + #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] pub struct RadrootsRelayFetchReceipt { pub inserted_count: usize, @@ -242,6 +272,42 @@ pub trait RadrootsRelayFetchAdapter: Send + Sync { ) -> BoxFuture<'a, Result<Vec<RadrootsRelayFetchItem>, RadrootsRelayTransportError>>; } +pub async fn fetch_relay_events<A>( + adapter: &A, + request: RadrootsRelayFetchRequest, +) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError> +where + A: RadrootsRelayFetchAdapter, +{ + let target_relays = request.relay_urls.clone(); + 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); + } + let filters = request.filters.as_slice().to_vec(); + let items = adapter.fetch(request).await?; + Ok( + process_relay_fetch_items(target_relays, filters, max_events, max_raw_events, items) + .into_fetched_events_receipt(), + ) +} + +#[cfg(feature = "runtime-tokio")] +pub fn fetch_relay_events_blocking<A>( + adapter: &A, + request: RadrootsRelayFetchRequest, +) -> Result<RadrootsRelayFetchedEventsReceipt, RadrootsRelayTransportError> +where + A: RadrootsRelayFetchAdapter, +{ + let runtime = tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .map_err(|error| RadrootsRelayTransportError::Transport(error.to_string()))?; + runtime.block_on(fetch_relay_events(adapter, request)) +} + pub async fn fetch_and_ingest_relay_events<A>( adapter: &A, event_store: &RadrootsEventStore, @@ -251,6 +317,7 @@ where A: RadrootsRelayFetchAdapter, { let mode = request.mode; + let target_relays = request.relay_urls.clone(); let max_events = request.max_events; let max_raw_events = request.max_raw_events; if request.filters.as_slice().is_empty() { @@ -258,86 +325,20 @@ where } let filters = request.filters.as_slice().to_vec(); let items = adapter.fetch(request).await?; - let mut receipt = RadrootsRelayFetchReceipt { - inserted_count: 0, - 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, - notice_count: 0, - events: Vec::new(), - relay_outcomes: Vec::new(), - }; - let mut scanned_raw_events = 0usize; - let mut accepted_events = 0usize; - for item in items { + let processed = + process_relay_fetch_items(target_relays, filters, max_events, max_raw_events, items); + let mut receipt = RadrootsRelayFetchReceipt::from_processed_counts(&processed); + for item in processed.items { match item { - RadrootsRelayFetchItem::Event { + RadrootsRelayProcessedFetchItem::Receipt(event_receipt) => { + receipt.events.push(event_receipt); + } + RadrootsRelayProcessedFetchItem::Accepted(RadrootsRelayFetchedEvent { relay_url, + event: raw_event, raw_json, observed_at_ms, - } => { - if scanned_raw_events >= max_raw_events { - receipt.skipped_over_limit_count += 1; - continue; - } - scanned_raw_events += 1; - let parsed = RadrootsNostrEvent::from_json(raw_json.as_str()); - let Ok(raw_event) = parsed else { - receipt.malformed_count += 1; - receipt.events.push(RadrootsRelayFetchEventReceipt { - relay_url, - event_id: None, - inserted: false, - duplicate: false, - 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()), - }); - continue; - }; - if !relay_fetch_event_matches_filters(&filters, &raw_event) { - receipt.out_of_filter_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: 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, @@ -398,36 +399,257 @@ where } } } + } + } + Ok(receipt) +} + +#[derive(Clone, Debug)] +enum RadrootsRelayProcessedFetchItem { + Accepted(RadrootsRelayFetchedEvent), + Receipt(RadrootsRelayFetchEventReceipt), +} + +#[derive(Clone, Debug)] +struct RadrootsRelayProcessedFetch { + target_relays: Vec<String>, + items: Vec<RadrootsRelayProcessedFetchItem>, + malformed_count: usize, + out_of_filter_count: usize, + skipped_over_limit_count: usize, + eose_count: usize, + closed_count: usize, + notice_count: usize, + relay_outcomes: Vec<RadrootsRelayFetchRelayOutcome>, +} + +impl RadrootsRelayProcessedFetch { + fn into_fetched_events_receipt(self) -> RadrootsRelayFetchedEventsReceipt { + let mut events = Vec::new(); + let mut event_receipts = Vec::new(); + for item in self.items { + match item { + RadrootsRelayProcessedFetchItem::Accepted(event) => { + event_receipts.push(accepted_fetch_event_receipt(&event)); + events.push(event); + } + RadrootsRelayProcessedFetchItem::Receipt(receipt) => event_receipts.push(receipt), + } + } + let connected_relays = self + .relay_outcomes + .iter() + .filter(|outcome| outcome.kind == RadrootsRelayFetchOutcomeKind::Eose) + .map(|outcome| outcome.relay_url.clone()) + .collect(); + let failed_relays = self + .relay_outcomes + .iter() + .filter(|outcome| outcome.kind == RadrootsRelayFetchOutcomeKind::Closed) + .map(|outcome| RadrootsRelayFetchFailure { + relay_url: outcome.relay_url.clone(), + reason: outcome.message.clone().unwrap_or_default(), + }) + .collect(); + RadrootsRelayFetchedEventsReceipt { + target_relays: self.target_relays, + connected_relays, + failed_relays, + events, + event_receipts, + malformed_count: self.malformed_count, + out_of_filter_count: self.out_of_filter_count, + skipped_over_limit_count: self.skipped_over_limit_count, + eose_count: self.eose_count, + closed_count: self.closed_count, + notice_count: self.notice_count, + relay_outcomes: self.relay_outcomes, + } + } +} + +impl RadrootsRelayFetchReceipt { + fn from_processed_counts(processed: &RadrootsRelayProcessedFetch) -> Self { + Self { + inserted_count: 0, + duplicate_count: 0, + malformed_count: processed.malformed_count, + out_of_filter_count: processed.out_of_filter_count, + skipped_over_limit_count: processed.skipped_over_limit_count, + unsupported_count: 0, + eose_count: processed.eose_count, + closed_count: processed.closed_count, + notice_count: processed.notice_count, + events: Vec::new(), + relay_outcomes: processed.relay_outcomes.clone(), + } + } +} + +fn process_relay_fetch_items( + target_relays: Vec<String>, + filters: Vec<RadrootsNostrFilter>, + max_events: usize, + max_raw_events: usize, + items: Vec<RadrootsRelayFetchItem>, +) -> RadrootsRelayProcessedFetch { + let mut processed = RadrootsRelayProcessedFetch { + target_relays, + items: Vec::new(), + malformed_count: 0, + out_of_filter_count: 0, + skipped_over_limit_count: 0, + eose_count: 0, + closed_count: 0, + notice_count: 0, + relay_outcomes: Vec::new(), + }; + let mut scanned_raw_events = 0usize; + let mut accepted_events = 0usize; + for item in items { + match item { + RadrootsRelayFetchItem::Event { + relay_url, + raw_json, + observed_at_ms, + } => { + if scanned_raw_events >= max_raw_events { + processed.skipped_over_limit_count += 1; + continue; + } + scanned_raw_events += 1; + let parsed = RadrootsNostrEvent::from_json(raw_json.as_str()); + let Ok(raw_event) = parsed else { + processed.malformed_count += 1; + processed + .items + .push(RadrootsRelayProcessedFetchItem::Receipt( + RadrootsRelayFetchEventReceipt { + relay_url, + event_id: None, + inserted: false, + duplicate: false, + 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()), + }, + )); + continue; + }; + if !relay_fetch_event_matches_filters(&filters, &raw_event) { + processed.out_of_filter_count += 1; + processed + .items + .push(RadrootsRelayProcessedFetchItem::Receipt( + RadrootsRelayFetchEventReceipt { + relay_url, + event_id: Some(raw_event.id.to_hex()), + inserted: false, + duplicate: false, + 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 { + processed.skipped_over_limit_count += 1; + processed + .items + .push(RadrootsRelayProcessedFetchItem::Receipt( + 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; + processed + .items + .push(RadrootsRelayProcessedFetchItem::Accepted( + RadrootsRelayFetchedEvent { + relay_url, + event: raw_event, + raw_json, + observed_at_ms, + }, + )); + } RadrootsRelayFetchItem::Eose { relay_url } => { - receipt.eose_count += 1; - receipt.relay_outcomes.push(RadrootsRelayFetchRelayOutcome { - relay_url, - kind: RadrootsRelayFetchOutcomeKind::Eose, - relay_outcome: None, - message: None, - }); + processed.eose_count += 1; + processed + .relay_outcomes + .push(RadrootsRelayFetchRelayOutcome { + relay_url, + kind: RadrootsRelayFetchOutcomeKind::Eose, + relay_outcome: None, + message: None, + }); } RadrootsRelayFetchItem::Closed { relay_url, message } => { - receipt.closed_count += 1; - receipt.relay_outcomes.push(RadrootsRelayFetchRelayOutcome { - relay_url, - kind: RadrootsRelayFetchOutcomeKind::Closed, - relay_outcome: Some(RadrootsRelayOutcome::classify(message.as_str())), - message: Some(message), - }); + processed.closed_count += 1; + processed + .relay_outcomes + .push(RadrootsRelayFetchRelayOutcome { + relay_url, + kind: RadrootsRelayFetchOutcomeKind::Closed, + relay_outcome: Some(RadrootsRelayOutcome::classify(message.as_str())), + message: Some(message), + }); } RadrootsRelayFetchItem::Notice { relay_url, message } => { - receipt.notice_count += 1; - receipt.relay_outcomes.push(RadrootsRelayFetchRelayOutcome { - relay_url, - kind: RadrootsRelayFetchOutcomeKind::Notice, - relay_outcome: None, - message: Some(message), - }); + processed.notice_count += 1; + processed + .relay_outcomes + .push(RadrootsRelayFetchRelayOutcome { + relay_url, + kind: RadrootsRelayFetchOutcomeKind::Notice, + relay_outcome: None, + message: Some(message), + }); } } } - Ok(receipt) + processed +} + +fn accepted_fetch_event_receipt( + event: &RadrootsRelayFetchedEvent, +) -> RadrootsRelayFetchEventReceipt { + RadrootsRelayFetchEventReceipt { + relay_url: event.relay_url.clone(), + event_id: Some(event.event.id.to_hex()), + inserted: false, + duplicate: false, + unsupported: false, + malformed: false, + out_of_filter: false, + skipped_over_limit: false, + projection_eligible: false, + verification_status: None, + message: Some("event accepted by relay fetch filters".to_owned()), + } } fn relay_fetch_event_matches_filters( @@ -474,7 +696,14 @@ async fn fetch_from_nostr_relays( }); continue; } - client.connect().await; + let connection_output = client.try_connect(timeout).await; + if connection_output.success.is_empty() { + items.push(RadrootsRelayFetchItem::Closed { + relay_url, + message: summarize_nostr_output_failures(&connection_output.failed), + }); + continue; + } let mut closed = false; for filter in filters.iter().cloned() { match client.fetch_events(filter, timeout).await { @@ -504,6 +733,21 @@ async fn fetch_from_nostr_relays( Ok(items) } +fn summarize_nostr_output_failures<K, E>(failed: &std::collections::HashMap<K, E>) -> String +where + K: std::fmt::Display + Eq + std::hash::Hash, + E: std::fmt::Display, +{ + if failed.is_empty() { + return "no relay acknowledged the operation".to_owned(); + } + failed + .iter() + .map(|(relay, error)| format!("{relay}: {error}")) + .collect::<Vec<_>>() + .join("; ") +} + #[derive(Clone, Default)] pub struct RadrootsMockRelayFetchAdapter { items: Arc<Mutex<Vec<RadrootsRelayFetchItem>>>, diff --git a/crates/relay_transport/src/lib.rs b/crates/relay_transport/src/lib.rs @@ -11,12 +11,16 @@ mod publish; mod relay; pub use error::RadrootsRelayTransportError; +#[cfg(all(feature = "storage", feature = "runtime-tokio"))] +pub use fetch::fetch_relay_events_blocking; #[cfg(feature = "storage")] pub use fetch::{ RadrootsMockRelayFetchAdapter, RadrootsNostrClientFetchAdapter, RadrootsRelayFetchAdapter, - RadrootsRelayFetchEventReceipt, RadrootsRelayFetchFilters, RadrootsRelayFetchItem, - RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchReceipt, - RadrootsRelayFetchRelayOutcome, RadrootsRelayFetchRequest, fetch_and_ingest_relay_events, + RadrootsRelayFetchEventReceipt, RadrootsRelayFetchFailure, RadrootsRelayFetchFilters, + RadrootsRelayFetchItem, RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, + RadrootsRelayFetchReceipt, RadrootsRelayFetchRelayOutcome, RadrootsRelayFetchRequest, + RadrootsRelayFetchedEvent, RadrootsRelayFetchedEventsReceipt, fetch_and_ingest_relay_events, + fetch_relay_events, }; #[cfg(feature = "storage")] pub use outbox::{ diff --git a/crates/relay_transport/tests/transport.rs b/crates/relay_transport/tests/transport.rs @@ -18,7 +18,8 @@ use radroots_relay_transport::{ RadrootsRelayOutcome, RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, RadrootsRelayTargetSet, RadrootsRelayTransportError, RadrootsRelayUrl, RadrootsRelayUrlPolicy, - fetch_and_ingest_relay_events, publish_claimed_outbox_event, publish_signed_event, + fetch_and_ingest_relay_events, fetch_relay_events, publish_claimed_outbox_event, + publish_signed_event, }; use std::net::{IpAddr, Ipv4Addr}; @@ -786,6 +787,86 @@ async fn fetch_event_cap_counts_accepted_in_filter_events_and_preserves_later_co } #[tokio::test] +async fn fetch_relay_events_applies_shared_filter_limit_and_outcome_evidence() { + let accepted = signed_event_with_kind_and_hashtag("shared fetch accepted", KIND_POST, "soil"); + let skipped = signed_event_with_kind_and_hashtag("shared fetch skipped", KIND_POST, "soil"); + let wrong_tag = + signed_event_with_kind_and_hashtag("shared fetch wrong tag", KIND_POST, "compost"); + let filter = radroots_nostr_filter_tag( + RadrootsNostrFilter::new() + .kind(RadrootsNostrKind::Custom(KIND_POST as u16)) + .limit(10), + "t", + vec!["soil".to_owned()], + ) + .expect("filter"); + let accepted_id = accepted.id.clone(); + let adapter = RadrootsMockRelayFetchAdapter::new(vec![ + RadrootsRelayFetchItem::Event { + relay_url: RELAY_PRIMARY_WSS.to_owned(), + raw_json: "{not json".to_owned(), + observed_at_ms: 2_100, + }, + RadrootsRelayFetchItem::Event { + relay_url: RELAY_PRIMARY_WSS.to_owned(), + raw_json: wrong_tag.raw_json, + observed_at_ms: 2_101, + }, + RadrootsRelayFetchItem::Event { + relay_url: RELAY_PRIMARY_WSS.to_owned(), + raw_json: accepted.raw_json.clone(), + observed_at_ms: 2_102, + }, + RadrootsRelayFetchItem::Event { + relay_url: RELAY_PRIMARY_WSS.to_owned(), + raw_json: skipped.raw_json, + observed_at_ms: 2_103, + }, + RadrootsRelayFetchItem::Eose { + relay_url: RELAY_PRIMARY_WSS.to_owned(), + }, + RadrootsRelayFetchItem::Closed { + relay_url: RELAY_SECONDARY_WSS.to_owned(), + message: "auth-required: challenge".to_owned(), + }, + RadrootsRelayFetchItem::Notice { + relay_url: RELAY_TERTIARY_WSS.to_owned(), + message: "notice: still visible".to_owned(), + }, + ]); + + let receipt = fetch_relay_events( + &adapter, + RadrootsRelayFetchRequest::fetch(2_100, 1, [filter]) + .expect("fetch request") + .with_relay_urls([RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS]), + ) + .await + .expect("fetch events"); + + assert_eq!( + receipt.target_relays, + vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS] + ); + assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]); + assert_eq!(receipt.failed_relays.len(), 1); + assert_eq!(receipt.failed_relays[0].relay_url, RELAY_SECONDARY_WSS); + assert_eq!(receipt.events.len(), 1); + assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id); + 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.eose_count, 1); + assert_eq!(receipt.closed_count, 1); + assert_eq!(receipt.notice_count, 1); + assert_eq!(receipt.event_receipts.len(), 4); + assert!(receipt.event_receipts[0].malformed); + assert!(receipt.event_receipts[1].out_of_filter); + assert!(!receipt.event_receipts[2].malformed); + assert!(receipt.event_receipts[3].skipped_over_limit); +} + +#[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");