lib

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

commit d3438a0072c87d011f4587c3c0a6540509c93883
parent 08af0d4bb35715a5151940a089beb1c31fdb0788
Author: triesap <tyson@radroots.org>
Date:   Mon, 27 Jul 2026 08:07:25 +0000

transport: seal relay publish receipts

- hide per-relay publish receipt state behind checked accessors
- validate relay receipt identities before adapter results are admitted
- bound strict wire decoding at the shared endpoint ceiling
- retain canonical duplicate and request provenance validation

Diffstat:
Mcrates/transport_nostr/src/outbox.rs | 51+++++++++++++++++++++------------------------------
Mcrates/transport_nostr/src/publish.rs | 156++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mcrates/transport_nostr/tests/transport.rs | 124+++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------------
3 files changed, 241 insertions(+), 90 deletions(-)

diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs @@ -206,23 +206,18 @@ where for relay in &publish.relays { if relay - .outcome + .outcome() .kind() .transport_outcome_kind() .target_status() .counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted) && publishable - .targets_for_relay(relay.relay_url.as_str()) + .targets_for_relay(relay.relay_url()) .next() .is_some() { - ingest_publish_observation( - event_store, - &signed_event, - relay.relay_url.as_str(), - now_ms, - ) - .await?; + ingest_publish_observation(event_store, &signed_event, relay.relay_url(), now_ms) + .await?; } } @@ -355,23 +350,18 @@ where for relay in &relay_receipts { if relay - .outcome + .outcome() .kind() .transport_outcome_kind() .target_status() .counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted) && publishable - .targets_for_relay(relay.relay_url.as_str()) + .targets_for_relay(relay.relay_url()) .next() .is_some() { - ingest_publish_observation( - event_store, - &signed_event, - relay.relay_url.as_str(), - now_ms, - ) - .await?; + ingest_publish_observation(event_store, &signed_event, relay.relay_url(), now_ms) + .await?; } } @@ -422,10 +412,10 @@ fn adapter_transport_failure_receipt( let relays = relay_urls .into_iter() .map(|relay_url| { - Ok(RadrootsRelayPublishRelayReceipt::attempted( + RadrootsRelayPublishRelayReceipt::attempted( relay_url, RadrootsRelayOutcome::connection_failed(message.clone())?, - )) + ) }) .collect::<Result<Vec<_>, RadrootsRelayTransportError>>()?; Ok(RadrootsRelayPublishReceipt { @@ -512,20 +502,20 @@ fn target_receipts_from_relay_receipts( ) -> Vec<RadrootsOutboxPublishTargetReceipt> { let mut target_receipts = Vec::new(); for relay_receipt in relay_receipts { - for target in publishable.targets_for_relay(relay_receipt.relay_url.as_str()) { + for target in publishable.targets_for_relay(relay_receipt.relay_url()) { target_receipts.push(RadrootsOutboxPublishTargetReceipt { delivery_target_id: target.delivery_target_id, endpoint_uri: target.relay_url.clone(), endpoint_fingerprint: target.endpoint_fingerprint.clone(), target_scope: target.target_scope.clone(), target_label: target.target_label.clone(), - attempted: relay_receipt.attempted, + attempted: relay_receipt.was_attempted(), transport_status: relay_receipt - .outcome + .outcome() .kind() .transport_outcome_kind() .target_status(), - outcome: relay_receipt.outcome.clone(), + outcome: relay_receipt.outcome().clone(), }); } } @@ -689,18 +679,18 @@ fn relay_receipts_from_transport_receipts( for receipt in delivery.target_receipts() { let outcome = relay_outcome_from_transport_outcome(receipt.outcome())?; let relay_receipt = if receipt.was_attempted() { - RadrootsRelayPublishRelayReceipt::attempted(receipt.target().uri().as_str(), outcome) + RadrootsRelayPublishRelayReceipt::attempted(receipt.target().uri().as_str(), outcome)? } else { - RadrootsRelayPublishRelayReceipt::skipped(receipt.target().uri().as_str(), outcome) + RadrootsRelayPublishRelayReceipt::skipped(receipt.target().uri().as_str(), outcome)? }; if let Some(existing) = relay_receipts .iter() - .find(|existing| existing.relay_url == relay_receipt.relay_url) + .find(|existing| existing.relay_url() == relay_receipt.relay_url()) { if existing != &relay_receipt { return Err( RadrootsRelayTransportError::ConflictingTransportReceiptRelayUrl { - url: relay_receipt.relay_url, + url: relay_receipt.relay_url().to_owned(), }, ); } @@ -1250,7 +1240,8 @@ mod tests { &[RadrootsRelayPublishRelayReceipt::attempted( target.uri().as_str(), RadrootsRelayOutcome::accepted(), - )], + ) + .expect("bounded relay receipt")], ); assert_eq!( accepted_relay_receipts[0].transport_status, @@ -1302,7 +1293,7 @@ mod tests { assert_eq!(receipt.terminal_count, 0); assert_eq!(receipt.quorum, 2); assert!(!receipt.quorum_met); - assert!(receipt.relays.iter().all(|relay| relay.attempted)); + assert!(receipt.relays.iter().all(|relay| relay.was_attempted())); } #[test] diff --git a/crates/transport_nostr/src/publish.rs b/crates/transport_nostr/src/publish.rs @@ -17,8 +17,9 @@ use radroots_transport::{ RadrootsTransportPayload, RadrootsTransportSatisfactionPolicy, RadrootsTransportStatus, RadrootsTransportTarget, RadrootsTransportTargetReceipt, }; -use serde::{Deserialize, Serialize}; +use serde::{Deserialize, Deserializer, Serialize, de}; use std::collections::{BTreeMap, BTreeSet}; +use std::fmt; use std::sync::{Arc, Mutex, PoisonError}; use crate::RadrootsRelayOutcomeKind; @@ -146,28 +147,131 @@ fn validate_publish_idempotency_key(value: &str) -> Result<(), RadrootsRelayTran Ok(()) } -#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] pub struct RadrootsRelayPublishRelayReceipt { - pub relay_url: String, - pub outcome: RadrootsRelayOutcome, - pub attempted: bool, + relay_url: String, + outcome: RadrootsRelayOutcome, + attempted: bool, } impl RadrootsRelayPublishRelayReceipt { - pub fn attempted(relay_url: impl Into<String>, outcome: RadrootsRelayOutcome) -> Self { - Self { - relay_url: relay_url.into(), - outcome, - attempted: true, - } + pub fn attempted( + relay_url: impl Into<String>, + outcome: RadrootsRelayOutcome, + ) -> Result<Self, RadrootsRelayTransportError> { + Self::try_new(relay_url.into(), outcome, true) } - pub fn skipped(relay_url: impl Into<String>, outcome: RadrootsRelayOutcome) -> Self { - Self { - relay_url: relay_url.into(), + pub fn skipped( + relay_url: impl Into<String>, + outcome: RadrootsRelayOutcome, + ) -> Result<Self, RadrootsRelayTransportError> { + Self::try_new(relay_url.into(), outcome, false) + } + + fn try_new( + relay_url: String, + outcome: RadrootsRelayOutcome, + attempted: bool, + ) -> Result<Self, RadrootsRelayTransportError> { + validate_publish_receipt_relay_url(relay_url.as_str())?; + Ok(Self { + relay_url, outcome, - attempted: false, + attempted, + }) + } + + pub fn relay_url(&self) -> &str { + self.relay_url.as_str() + } + + pub fn outcome(&self) -> &RadrootsRelayOutcome { + &self.outcome + } + + pub fn was_attempted(&self) -> bool { + self.attempted + } +} + +fn validate_publish_receipt_relay_url(relay_url: &str) -> Result<(), RadrootsRelayTransportError> { + RadrootsTransportTarget::nostr_relay(relay_url).map_err(|error| { + RadrootsRelayTransportError::InvalidPublishReceiptRelayUrl { + url: if relay_url.len() <= radroots_transport::RADROOTS_TRANSPORT_ENDPOINT_URI_MAX_BYTES + { + relay_url.to_owned() + } else { + "<oversized>".to_owned() + }, + reason: error.to_string(), } + })?; + Ok(()) +} + +struct BoundedPublishReceiptRelayUrl; + +impl<'de> de::Visitor<'de> for BoundedPublishReceiptRelayUrl { + type Value = String; + + fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + write!( + formatter, + "a relay URL of at most {} UTF-8 bytes", + radroots_transport::RADROOTS_TRANSPORT_ENDPOINT_URI_MAX_BYTES + ) + } + + fn visit_borrowed_str<E>(self, value: &'de str) -> Result<Self::Value, E> + where + E: de::Error, + { + self.visit_str(value) + } + + fn visit_str<E>(self, value: &str) -> Result<Self::Value, E> + where + E: de::Error, + { + validate_publish_receipt_relay_url(value) + .map_err(E::custom) + .map(|()| value.to_owned()) + } + + fn visit_string<E>(self, value: String) -> Result<Self::Value, E> + where + E: de::Error, + { + validate_publish_receipt_relay_url(value.as_str()) + .map_err(E::custom) + .map(|()| value) + } +} + +fn deserialize_publish_receipt_relay_url<'de, D>(deserializer: D) -> Result<String, D::Error> +where + D: Deserializer<'de>, +{ + deserializer.deserialize_string(BoundedPublishReceiptRelayUrl) +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct RadrootsRelayPublishRelayReceiptWire { + #[serde(deserialize_with = "deserialize_publish_receipt_relay_url")] + relay_url: String, + outcome: RadrootsRelayOutcome, + attempted: bool, +} + +impl<'de> Deserialize<'de> for RadrootsRelayPublishRelayReceipt { + fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> + where + D: Deserializer<'de>, + { + let wire = RadrootsRelayPublishRelayReceiptWire::deserialize(deserializer)?; + Self::try_new(wire.relay_url, wire.outcome, wire.attempted).map_err(de::Error::custom) } } @@ -743,10 +847,10 @@ fn normalize_publish_receipts( if let Some(receipt) = by_relay.remove(relay_url) { Ok(receipt) } else { - Ok(RadrootsRelayPublishRelayReceipt::skipped( + RadrootsRelayPublishRelayReceipt::skipped( relay_url, RadrootsRelayOutcome::unknown("relay adapter omitted target receipt")?, - )) + ) } }) .collect() @@ -849,7 +953,7 @@ impl RadrootsRelayPublishAdapter for RadrootsMockRelayPublishAdapter { .lock() .map_err(captured_raw_event_lock_error)? .push(request.signed_event.signed_event().raw_json().to_owned()); - Ok(request + request .targets .relays() .iter() @@ -861,7 +965,7 @@ impl RadrootsRelayPublishAdapter for RadrootsMockRelayPublishAdapter { .unwrap_or_else(RadrootsRelayOutcome::accepted); RadrootsRelayPublishRelayReceipt::attempted(relay.as_str(), outcome) }) - .collect()) + .collect() }) } } @@ -940,10 +1044,10 @@ impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter { .get(target_url) .cloned() .unwrap_or_else(|| "relay did not connect".to_owned()); - Ok(RadrootsRelayPublishRelayReceipt::attempted( + RadrootsRelayPublishRelayReceipt::attempted( relay_url, RadrootsRelayOutcome::connection_failed(reason)?, - )) + ) }) .collect(); } @@ -954,10 +1058,10 @@ impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter { return target_strings .into_iter() .map(|relay_url| { - Ok(RadrootsRelayPublishRelayReceipt::attempted( + RadrootsRelayPublishRelayReceipt::attempted( relay_url, RadrootsRelayOutcome::connection_failed(message.clone())?, - )) + ) }) .collect(); } @@ -975,14 +1079,14 @@ impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter { RadrootsRelayOutcome::accepted_with_message( "nostr-relay-pool-success-ok-message-unavailable", )?, - )); + )?); continue; } if let Some(reason) = connection_failures.get(target_url) { receipts.push(RadrootsRelayPublishRelayReceipt::attempted( relay_url, RadrootsRelayOutcome::connection_failed(reason.clone())?, - )); + )?); continue; } let failed = output.failed.iter().find_map(|(failed_url, message)| { @@ -1000,7 +1104,7 @@ impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter { )?); receipts.push(RadrootsRelayPublishRelayReceipt::attempted( relay_url, outcome, - )); + )?); } Ok(receipts) }) diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -62,6 +62,12 @@ fn bounded_relay_outcome( outcome.expect("bounded relay outcome") } +fn bounded_publish_relay_receipt( + receipt: Result<RadrootsRelayPublishRelayReceipt, RadrootsRelayTransportError>, +) -> RadrootsRelayPublishRelayReceipt { + receipt.expect("bounded relay publish receipt") +} + struct TransportFailurePublishAdapter; impl RadrootsRelayPublishAdapter for TransportFailurePublishAdapter { @@ -87,9 +93,11 @@ impl RadrootsRelayPublishAdapter for PartialPublishAdapter { ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> { Box::pin(async { - Ok(vec![RadrootsRelayPublishRelayReceipt::attempted( - RELAY_PRIMARY_WSS, - RadrootsRelayOutcome::accepted(), + Ok(vec![bounded_publish_relay_receipt( + RadrootsRelayPublishRelayReceipt::attempted( + RELAY_PRIMARY_WSS, + RadrootsRelayOutcome::accepted(), + ), )]) }) } @@ -104,9 +112,11 @@ impl RadrootsRelayPublishAdapter for SlashSpelledRelayReceiptPublishAdapter { ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> { Box::pin(async { - Ok(vec![RadrootsRelayPublishRelayReceipt::attempted( - format!("{RELAY_PRIMARY_WSS}/"), - RadrootsRelayOutcome::accepted(), + Ok(vec![bounded_publish_relay_receipt( + RadrootsRelayPublishRelayReceipt::attempted( + format!("{RELAY_PRIMARY_WSS}/"), + RadrootsRelayOutcome::accepted(), + ), )]) }) } @@ -145,14 +155,14 @@ impl RadrootsRelayPublishAdapter for UnknownRelayReceiptPublishAdapter { .as_str() .to_owned(); Ok(vec![ - RadrootsRelayPublishRelayReceipt::attempted( + bounded_publish_relay_receipt(RadrootsRelayPublishRelayReceipt::attempted( relay, RadrootsRelayOutcome::accepted(), - ), - RadrootsRelayPublishRelayReceipt::attempted( + )), + bounded_publish_relay_receipt(RadrootsRelayPublishRelayReceipt::attempted( RELAY_TERTIARY_WSS, RadrootsRelayOutcome::accepted(), - ), + )), ]) }) } @@ -175,14 +185,14 @@ impl RadrootsRelayPublishAdapter for DuplicateRelayReceiptPublishAdapter { .as_str() .to_owned(); Ok(vec![ - RadrootsRelayPublishRelayReceipt::attempted( + bounded_publish_relay_receipt(RadrootsRelayPublishRelayReceipt::attempted( relay.clone(), RadrootsRelayOutcome::accepted(), - ), - RadrootsRelayPublishRelayReceipt::attempted( + )), + bounded_publish_relay_receipt(RadrootsRelayPublishRelayReceipt::attempted( format!("{relay}/"), RadrootsRelayOutcome::accepted(), - ), + )), ]) }) } @@ -197,10 +207,11 @@ impl RadrootsRelayPublishAdapter for InvalidRelayReceiptPublishAdapter { ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> { Box::pin(async { - Ok(vec![RadrootsRelayPublishRelayReceipt::attempted( + RadrootsRelayPublishRelayReceipt::attempted( "not a relay URL", RadrootsRelayOutcome::accepted(), - )]) + ) + .map(|receipt| vec![receipt]) }) } } @@ -220,9 +231,8 @@ impl RadrootsRelayPublishAdapter for SkippedAcceptedRelayReceiptPublishAdapter { .first() .expect("fixture target") .as_str(); - Ok(vec![RadrootsRelayPublishRelayReceipt::skipped( - relay, - RadrootsRelayOutcome::accepted(), + Ok(vec![bounded_publish_relay_receipt( + RadrootsRelayPublishRelayReceipt::skipped(relay, RadrootsRelayOutcome::accepted()), )]) }) } @@ -243,11 +253,13 @@ impl RadrootsRelayPublishAdapter for AttemptedSkippedRelayReceiptPublishAdapter .first() .expect("fixture target") .as_str(); - Ok(vec![RadrootsRelayPublishRelayReceipt::attempted( - relay, - bounded_relay_outcome(RadrootsRelayOutcome::skipped_already_accepted( - "already accepted", - )), + Ok(vec![bounded_publish_relay_receipt( + RadrootsRelayPublishRelayReceipt::attempted( + relay, + bounded_relay_outcome(RadrootsRelayOutcome::skipped_already_accepted( + "already accepted", + )), + ), )]) }) } @@ -1243,6 +1255,50 @@ fn relay_outcome_messages_are_bounded_and_strictly_decoded() { } #[test] +fn relay_publish_receipts_validate_bounded_wire_identity() { + let prefix = "wss://relay.example.com/"; + let exact_url = format!( + "{prefix}{}", + "x".repeat(RADROOTS_TRANSPORT_ENDPOINT_URI_MAX_BYTES - prefix.len()) + ); + let receipt = RadrootsRelayPublishRelayReceipt::attempted( + exact_url.clone(), + RadrootsRelayOutcome::accepted(), + ) + .expect("maximum relay receipt URL"); + assert_eq!(receipt.relay_url(), exact_url); + assert!(receipt.was_attempted()); + assert_eq!(receipt.outcome().kind(), RadrootsRelayOutcomeKind::Accepted); + + let wire = serde_json::to_value(&receipt).expect("relay receipt JSON"); + let decoded = serde_json::from_value::<RadrootsRelayPublishRelayReceipt>(wire) + .expect("strict relay receipt reload"); + assert_eq!(decoded, receipt); + + let one_over = format!("{exact_url}x"); + assert!(matches!( + RadrootsRelayPublishRelayReceipt::attempted( + one_over, + RadrootsRelayOutcome::accepted(), + ), + Err(RadrootsRelayTransportError::InvalidPublishReceiptRelayUrl { url, .. }) + if url == "<oversized>" + )); + assert!( + serde_json::from_value::<RadrootsRelayPublishRelayReceipt>(serde_json::json!({ + "relay_url": RELAY_PRIMARY_WSS, + "outcome": { + "kind": "Accepted", + "message": null, + }, + "attempted": true, + "extra": true, + })) + .is_err() + ); +} + +#[test] fn relay_transport_error_wraps_transport_contract_errors() { let error = RadrootsRelayTransportError::from( radroots_transport::RadrootsTransportError::EmptyTargetSet, @@ -1724,13 +1780,13 @@ async fn publish_receipts_track_terminal_skipped_and_adapter_errors() { assert_eq!(receipt.quorum, 2); assert!(!receipt.quorum_met); - let skipped = RadrootsRelayPublishRelayReceipt::skipped( + let skipped = bounded_publish_relay_receipt(RadrootsRelayPublishRelayReceipt::skipped( RELAY_TERTIARY_WSS, bounded_relay_outcome(RadrootsRelayOutcome::timeout("timeout: no OK")), - ); - assert_eq!(skipped.relay_url, RELAY_TERTIARY_WSS); - assert!(!skipped.attempted); - assert_eq!(skipped.outcome.kind(), RadrootsRelayOutcomeKind::Timeout); + )); + assert_eq!(skipped.relay_url(), RELAY_TERTIARY_WSS); + assert!(!skipped.was_attempted()); + assert_eq!(skipped.outcome().kind(), RadrootsRelayOutcomeKind::Timeout); let error = publish_signed_event( &TransportFailurePublishAdapter, @@ -1814,9 +1870,9 @@ async fn publish_all_policy_uses_requested_target_count() { assert_eq!(receipt.quorum, 2); assert!(!receipt.quorum_met); assert_eq!(receipt.relays.len(), 2); - assert!(!receipt.relays[1].attempted); + assert!(!receipt.relays[1].was_attempted()); assert_eq!( - receipt.relays[1].outcome.kind(), + receipt.relays[1].outcome().kind(), RadrootsRelayOutcomeKind::Unknown ); @@ -4221,7 +4277,7 @@ async fn outbox_publish_fans_out_endpoint_receipts_to_scoped_logical_targets() { assert_eq!(published.quorum, 2); assert!(published.quorum_met); assert_eq!(published.relay_receipts.len(), 1); - assert_eq!(published.relay_receipts[0].relay_url, RELAY_PRIMARY_WSS); + assert_eq!(published.relay_receipts[0].relay_url(), RELAY_PRIMARY_WSS); assert_eq!(published.target_receipts.len(), 2); assert!( published @@ -4448,7 +4504,7 @@ async fn outbox_publish_required_target_failure_is_not_satisfied_by_optional_suc assert_eq!(published.quorum, 1); assert!(!published.quorum_met); assert_eq!(published.relay_receipts.len(), 1); - assert_eq!(published.relay_receipts[0].relay_url, RELAY_SECONDARY_WSS); + assert_eq!(published.relay_receipts[0].relay_url(), RELAY_SECONDARY_WSS); let event = outbox .get_event(receipt.outbox_event_id) .await @@ -4740,7 +4796,7 @@ async fn outbox_transport_publish_failure_releases_retryable_claim() { published .relay_receipts .iter() - .all(|relay| relay.outcome.kind() == RadrootsRelayOutcomeKind::ConnectionFailed) + .all(|relay| relay.outcome().kind() == RadrootsRelayOutcomeKind::ConnectionFailed) ); let event = outbox