lib

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

commit cc2199e29ed5e42363dad506661c0464c7c7344f
parent 09caf2c6dfab061b1f7d63b30d0d883b7cce4b96
Author: triesap <tyson@radroots.org>
Date:   Mon, 27 Jul 2026 08:35:14 +0000

transport: bind fetched event envelopes

- make fetched event fields private behind read-only accessors
- bind each parsed event to its exact verified raw JSON bytes
- validate relay identity observation time and event byte limits
- prove mismatch endpoint timestamp and one-over rejection paths

Diffstat:
Mcrates/transport_nostr/src/fetch.rs | 114+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------
Mcrates/transport_nostr/tests/transport.rs | 83+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------
2 files changed, 155 insertions(+), 42 deletions(-)

diff --git a/crates/transport_nostr/src/fetch.rs b/crates/transport_nostr/src/fetch.rs @@ -859,10 +859,81 @@ impl<'de> Deserialize<'de> for RadrootsRelayFetchEventReceipt { #[derive(Clone, Debug)] pub struct RadrootsRelayFetchedEvent { - pub relay_url: String, - pub event: RadrootsNostrEvent, - pub raw_json: String, - pub observed_at_ms: i64, + relay_url: String, + event: RadrootsNostrEvent, + raw_json: String, + observed_at_ms: i64, +} + +impl RadrootsRelayFetchedEvent { + pub fn new( + relay_url: impl Into<String>, + event: RadrootsNostrEvent, + raw_json: impl Into<String>, + observed_at_ms: i64, + ) -> Result<Self, RadrootsRelayTransportError> { + let relay_url = relay_url.into(); + let raw_json = raw_json.into(); + let fetched = Self::from_verified(relay_url, event, raw_json, observed_at_ms)?; + RadrootsEventIngest::from_raw_json(fetched.raw_json.clone(), observed_at_ms)?; + Ok(fetched) + } + + fn from_verified( + relay_url: String, + event: RadrootsNostrEvent, + raw_json: String, + observed_at_ms: i64, + ) -> Result<Self, RadrootsRelayTransportError> { + let relay_url = canonical_fetch_receipt_relay_url(relay_url.as_str())?; + ensure_nonnegative_timestamp("observed_at_ms", observed_at_ms)?; + if raw_json.len() > DEFAULT_RAW_JSON_MAX_BYTES { + return Err(RadrootsRelayTransportError::FetchLimitTooLarge { + field: "fetched_event_raw_json_bytes", + max: DEFAULT_RAW_JSON_MAX_BYTES, + actual: raw_json.len(), + }); + } + let decoded = RadrootsNostrEvent::from_json(raw_json.as_str()) + .map_err(|error| RadrootsRelayTransportError::NostrEventJson(error.to_string()))?; + if decoded != event { + return Err(invalid_fetch_receipt( + "fetched_event", + "event object does not match its raw JSON bytes", + )); + } + Ok(Self { + relay_url, + event, + raw_json, + observed_at_ms, + }) + } + + pub fn relay_url(&self) -> &str { + self.relay_url.as_str() + } + + pub fn event(&self) -> &RadrootsNostrEvent { + &self.event + } + + pub fn raw_json(&self) -> &str { + self.raw_json.as_str() + } + + pub fn observed_at_ms(&self) -> i64 { + self.observed_at_ms + } + + fn into_parts(self) -> (String, RadrootsNostrEvent, String, i64) { + ( + self.relay_url, + self.event, + self.raw_json, + self.observed_at_ms, + ) + } } #[derive(Clone, Debug, PartialEq, Eq, Serialize)] @@ -1064,18 +1135,9 @@ where RadrootsRelayProcessedFetchItem::Receipt(event_receipt) => { receipt.events.push(event_receipt); } - RadrootsRelayProcessedFetchItem::Accepted(RadrootsRelayFetchedEvent { - relay_url, - event: raw_event, - raw_json, - observed_at_ms, - }) - | RadrootsRelayProcessedFetchItem::Duplicate(RadrootsRelayFetchedEvent { - relay_url, - event: raw_event, - raw_json, - observed_at_ms, - }) => { + RadrootsRelayProcessedFetchItem::Accepted(event) + | RadrootsRelayProcessedFetchItem::Duplicate(event) => { + let (relay_url, raw_event, raw_json, observed_at_ms) = event.into_parts(); let observation_type = match mode { RadrootsRelayFetchMode::Fetch => RadrootsTransportObservationType::Fetch, RadrootsRelayFetchMode::Subscription => { @@ -1623,12 +1685,12 @@ fn process_relay_fetch_items( processed .items .push(RadrootsRelayProcessedFetchItem::Duplicate( - RadrootsRelayFetchedEvent { + RadrootsRelayFetchedEvent::from_verified( relay_url, - event: raw_event, + raw_event, raw_json, observed_at_ms, - }, + )?, )); continue; } @@ -1663,12 +1725,12 @@ fn process_relay_fetch_items( processed .items .push(RadrootsRelayProcessedFetchItem::Accepted( - RadrootsRelayFetchedEvent { + RadrootsRelayFetchedEvent::from_verified( relay_url, - event: raw_event, + raw_event, raw_json, observed_at_ms, - }, + )?, )); } RadrootsRelayFetchItemBody::Eose { .. } => { @@ -1725,8 +1787,8 @@ fn accepted_fetch_event_receipt( event: &RadrootsRelayFetchedEvent, ) -> Result<RadrootsRelayFetchEventReceipt, RadrootsRelayTransportError> { RadrootsRelayFetchEventReceipt { - relay_url: event.relay_url.clone(), - event_id: Some(event.event.id.to_hex()), + relay_url: event.relay_url().to_owned(), + event_id: Some(event.event().id.to_hex()), inserted: false, duplicate: false, not_persisted: false, @@ -1747,8 +1809,8 @@ fn duplicate_fetch_event_receipt( event: &RadrootsRelayFetchedEvent, ) -> Result<RadrootsRelayFetchEventReceipt, RadrootsRelayTransportError> { RadrootsRelayFetchEventReceipt { - relay_url: event.relay_url.clone(), - event_id: Some(event.event.id.to_hex()), + relay_url: event.relay_url().to_owned(), + event_id: Some(event.event().id.to_hex()), inserted: false, duplicate: true, not_persisted: false, diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -9,9 +9,9 @@ use radroots_event_store::{ RadrootsEventStore, RadrootsTransportObservationRow, RadrootsTransportObservationType, }; use radroots_nostr::prelude::{ - RadrootsNostrFilter, RadrootsNostrKeys, RadrootsNostrKind, RadrootsNostrSecretKey, - RadrootsNostrTag, RadrootsNostrTagKind, RadrootsNostrTimestamp, radroots_nostr_filter_tag, - radroots_nostr_sign_frozen_draft, + RadrootsNostrEvent, RadrootsNostrFilter, RadrootsNostrKeys, RadrootsNostrKind, + RadrootsNostrSecretKey, RadrootsNostrTag, RadrootsNostrTagKind, RadrootsNostrTimestamp, + radroots_nostr_filter_tag, radroots_nostr_sign_frozen_draft, }; use radroots_outbox::{ RadrootsOutbox, RadrootsOutboxClaimedEvent, RadrootsOutboxDeliveryPlanInput, @@ -39,12 +39,13 @@ use radroots_transport_nostr::{ RadrootsRelayFetchEventVerification, RadrootsRelayFetchEventVisibility, RadrootsRelayFetchFailure, RadrootsRelayFetchFilters, RadrootsRelayFetchItem, RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchRelayOutcome, - RadrootsRelayFetchRequest, RadrootsRelayOutcome, RadrootsRelayOutcomeKind, - RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, - RadrootsRelayTargetSet, RadrootsRelayTransportError, RadrootsRelayUrl, RadrootsRelayUrlPolicy, - fetch_and_ingest_relay_events, fetch_relay_events, fetch_relay_events_blocking, - publish_claimed_outbox_event, publish_claimed_outbox_event_with_transport, - publish_signed_event, verified_signed_event_payload, + RadrootsRelayFetchRequest, RadrootsRelayFetchedEvent, RadrootsRelayOutcome, + RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt, + RadrootsRelayPublishRequest, RadrootsRelayTargetSet, RadrootsRelayTransportError, + RadrootsRelayUrl, RadrootsRelayUrlPolicy, fetch_and_ingest_relay_events, fetch_relay_events, + fetch_relay_events_blocking, publish_claimed_outbox_event, + publish_claimed_outbox_event_with_transport, publish_signed_event, + verified_signed_event_payload, }; use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; @@ -2395,8 +2396,8 @@ fn fetch_blocking_facade_runs_mock_adapter() { .expect("blocking fetch"); assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id); - assert_eq!(receipt.events[0].observed_at_ms, 1_090); + assert_eq!(receipt.events[0].event().id.to_hex(), accepted_id); + assert_eq!(receipt.events[0].observed_at_ms(), 1_090); assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]); } @@ -2416,8 +2417,8 @@ async fn fetch_canonicalizes_adapter_relay_spelling_and_uses_request_observation .expect("canonical fetch"); assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].relay_url, RELAY_PRIMARY_WSS); - assert_eq!(receipt.events[0].observed_at_ms, 1_091); + assert_eq!(receipt.events[0].relay_url(), RELAY_PRIMARY_WSS); + assert_eq!(receipt.events[0].observed_at_ms(), 1_091); assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]); } @@ -2435,7 +2436,7 @@ async fn fetch_verifies_events_before_acceptance_budgeting() { .expect("verified fetch"); assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id); + assert_eq!(receipt.events[0].event().id.to_hex(), accepted_id); assert_eq!(receipt.verification_failed_count, 1); assert_eq!(receipt.skipped_over_limit_count, 0); assert_eq!( @@ -2468,7 +2469,7 @@ async fn fetch_deduplicates_event_ids_without_starving_unique_events() { receipt .events .iter() - .map(|event| event.event.id.to_hex()) + .map(|event| event.event().id.to_hex()) .collect::<Vec<_>>(), vec![first_id.clone(), second_id] ); @@ -3311,7 +3312,7 @@ async fn fetch_relay_events_applies_shared_filter_limit_and_outcome_evidence() { 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.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); @@ -3834,6 +3835,56 @@ fn fetch_event_receipts_reject_oversized_and_incoherent_wire_state() { ); } +#[test] +fn fetched_events_bind_raw_bytes_identity_endpoint_and_observation_time() { + let signed = signed_post("fetched event envelope"); + let event = RadrootsNostrEvent::from_json(signed.raw_json()).expect("signed Nostr event"); + let fetched = + RadrootsRelayFetchedEvent::new(RELAY_PRIMARY_WSS, event.clone(), signed.raw_json(), 1_500) + .expect("bounded fetched event"); + assert_eq!(fetched.relay_url(), RELAY_PRIMARY_WSS); + assert_eq!(fetched.event(), &event); + assert_eq!(fetched.raw_json(), signed.raw_json()); + assert_eq!(fetched.observed_at_ms(), 1_500); + + assert!(matches!( + RadrootsRelayFetchedEvent::new(" ", event.clone(), signed.raw_json(), 1_500), + Err(RadrootsRelayTransportError::InvalidFetchReceipt { + field: "relay_url", + .. + }) + )); + assert!(matches!( + RadrootsRelayFetchedEvent::new(RELAY_PRIMARY_WSS, event.clone(), signed.raw_json(), -1,), + Err(RadrootsRelayTransportError::InvalidTimestamp { + field: "observed_at_ms", + value: -1, + }) + )); + assert!(matches!( + RadrootsRelayFetchedEvent::new( + RELAY_PRIMARY_WSS, + event.clone(), + "x".repeat(DEFAULT_RAW_JSON_MAX_BYTES + 1), + 1_500, + ), + Err(RadrootsRelayTransportError::FetchLimitTooLarge { + field: "fetched_event_raw_json_bytes", + max: DEFAULT_RAW_JSON_MAX_BYTES, + actual, + }) if actual == DEFAULT_RAW_JSON_MAX_BYTES + 1 + )); + + let other = signed_post("different fetched event envelope"); + assert!(matches!( + RadrootsRelayFetchedEvent::new(RELAY_PRIMARY_WSS, event, other.raw_json(), 1_500,), + Err(RadrootsRelayTransportError::InvalidFetchReceipt { + field: "fetched_event", + .. + }) + )); +} + #[tokio::test] async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() { let outbox = RadrootsOutbox::open_memory().await.expect("outbox");