lib

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

commit 451da16aa8f02eadb821a075427cddc0f17ec564
parent d3438a0072c87d011f4587c3c0a6540509c93883
Author: triesap <tyson@radroots.org>
Date:   Mon, 27 Jul 2026 08:12:08 +0000

transport: seal aggregate publish receipts

- hide aggregate publish receipt fields behind stable accessors
- validate event identity relay cardinality and duplicate endpoints
- derive receipt counters and reject incoherent wire claims
- prove exact and one-over relay collections under strict reload

Diffstat:
Mcrates/transport_nostr/src/error.rs | 3+++
Mcrates/transport_nostr/src/outbox.rs | 37+++++++++++++++----------------------
Mcrates/transport_nostr/src/publish.rs | 273+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------
Mcrates/transport_nostr/tests/transport.rs | 132++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
4 files changed, 366 insertions(+), 79 deletions(-)

diff --git a/crates/transport_nostr/src/error.rs b/crates/transport_nostr/src/error.rs @@ -62,6 +62,9 @@ pub enum RadrootsRelayTransportError { #[error("Relay publish receipt contains invalid relay URL `{url}`: {reason}")] InvalidPublishReceiptRelayUrl { url: String, reason: String }, + #[error("Relay publish receipt has invalid {field}: {reason}")] + InvalidPublishReceipt { field: &'static str, reason: String }, + #[error("Relay publish receipt came from unrequested relay URL `{url}`")] UnexpectedPublishReceiptRelayUrl { url: String }, diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs @@ -198,13 +198,13 @@ where )?, Err(error) => return Err(error), }; - let target_receipts = target_receipts_from_relay_receipts(&publishable, &publish.relays); + let target_receipts = target_receipts_from_relay_receipts(&publishable, publish.relays()); for target_receipt in &target_receipts { complete_outbox_delivery_target(outbox, claimed, target_receipt, now_ms).await?; } - for relay in &publish.relays { + for relay in publish.relays() { if relay .outcome() .kind() @@ -232,9 +232,10 @@ where ) .await?; + let (event_id, relay_receipts) = publish.into_event_id_and_relays(); Ok(RadrootsOutboxPublishReceipt { local_ingest, - event_id: publish.event_id, + event_id, attempted_count: target_receipts .iter() .filter(|receipt| receipt.attempted) @@ -255,7 +256,7 @@ where quorum_met: publishable.satisfied_count_after_receipts(&target_receipts) >= publishable.satisfaction_required_count, target_receipts, - relay_receipts: publish.relays, + relay_receipts, }) } @@ -418,16 +419,7 @@ fn adapter_transport_failure_receipt( ) }) .collect::<Result<Vec<_>, RadrootsRelayTransportError>>()?; - Ok(RadrootsRelayPublishReceipt { - event_id, - attempted_count: relays.len(), - accepted_count: 0, - retryable_count: relays.len(), - terminal_count: 0, - quorum, - quorum_met: false, - relays, - }) + RadrootsRelayPublishReceipt::new(event_id, quorum, false, relays) } struct PublishableRelays { @@ -1276,8 +1268,9 @@ mod tests { #[test] fn adapter_transport_failure_receipts_preserve_each_target() { + let event_id = "a".repeat(64); let receipt = adapter_transport_failure_receipt( - "event-1".to_owned(), + event_id.clone(), vec![ "wss://relay-a.example".to_owned(), "wss://relay-b.example".to_owned(), @@ -1287,13 +1280,13 @@ mod tests { ) .expect("bounded adapter failure receipt"); - assert_eq!(receipt.event_id, "event-1"); - assert_eq!(receipt.attempted_count, 2); - assert_eq!(receipt.retryable_count, 2); - assert_eq!(receipt.terminal_count, 0); - assert_eq!(receipt.quorum, 2); - assert!(!receipt.quorum_met); - assert!(receipt.relays.iter().all(|relay| relay.was_attempted())); + assert_eq!(receipt.event_id(), event_id); + assert_eq!(receipt.attempted_count(), 2); + assert_eq!(receipt.retryable_count(), 2); + assert_eq!(receipt.terminal_count(), 0); + assert_eq!(receipt.quorum(), 2); + assert!(!receipt.quorum_met()); + 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 @@ -7,6 +7,7 @@ use core::time::Duration; use futures::future::BoxFuture; use radroots_event::{ draft::{RadrootsSignedEvent, RadrootsVerifiedSignedEvent}, + ids::RadrootsEventId, wire::RadrootsNip01EventWire, }; use radroots_transport::{ @@ -275,16 +276,245 @@ impl<'de> Deserialize<'de> for RadrootsRelayPublishRelayReceipt { } } -#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] pub struct RadrootsRelayPublishReceipt { - pub event_id: String, - pub attempted_count: usize, - pub accepted_count: usize, - pub retryable_count: usize, - pub terminal_count: usize, - pub quorum: usize, - pub quorum_met: bool, - pub relays: Vec<RadrootsRelayPublishRelayReceipt>, + event_id: RadrootsEventId, + attempted_count: usize, + accepted_count: usize, + retryable_count: usize, + terminal_count: usize, + quorum: usize, + quorum_met: bool, + relays: Vec<RadrootsRelayPublishRelayReceipt>, +} + +impl RadrootsRelayPublishReceipt { + pub(crate) fn new( + event_id: impl AsRef<str>, + quorum: usize, + quorum_met: bool, + relays: Vec<RadrootsRelayPublishRelayReceipt>, + ) -> Result<Self, RadrootsRelayTransportError> { + let event_id = RadrootsEventId::parse(event_id).map_err(|error| { + RadrootsRelayTransportError::InvalidPublishReceipt { + field: "event_id", + reason: error.to_string(), + } + })?; + Self::from_validated_event_id(event_id, quorum, quorum_met, relays) + } + + fn from_validated_event_id( + event_id: RadrootsEventId, + quorum: usize, + quorum_met: bool, + relays: Vec<RadrootsRelayPublishRelayReceipt>, + ) -> Result<Self, RadrootsRelayTransportError> { + if relays.is_empty() { + return Err(RadrootsRelayTransportError::InvalidPublishReceipt { + field: "relays", + reason: "relay receipts must not be empty".to_owned(), + }); + } + if relays.len() > radroots_transport::RADROOTS_TRANSPORT_TARGET_MAX_COUNT { + return Err(RadrootsRelayTransportError::InvalidPublishReceipt { + field: "relays", + reason: format!( + "relay receipt count {} exceeds maximum {}", + relays.len(), + radroots_transport::RADROOTS_TRANSPORT_TARGET_MAX_COUNT + ), + }); + } + let mut canonical_relays = Vec::with_capacity(relays.len()); + for receipt in &relays { + let canonical = RadrootsTransportTarget::nostr_relay(receipt.relay_url())? + .uri() + .as_str() + .to_owned(); + if canonical_relays.contains(&canonical) { + return Err( + RadrootsRelayTransportError::DuplicatePublishReceiptRelayUrl { url: canonical }, + ); + } + canonical_relays.push(canonical); + } + if quorum > relays.len() { + return Err(RadrootsRelayTransportError::InvalidPublishReceipt { + field: "quorum", + reason: format!( + "quorum {quorum} exceeds relay receipt count {}", + relays.len() + ), + }); + } + let attempted_count = relays + .iter() + .filter(|receipt| receipt.was_attempted()) + .count(); + let accepted_count = relays + .iter() + .filter(|receipt| relay_receipt_counts_toward_quorum(receipt)) + .count(); + let retryable_count = relays + .iter() + .filter(|receipt| receipt.outcome().is_retryable()) + .count(); + let terminal_count = relays + .iter() + .filter(|receipt| receipt.outcome().is_terminal_failure()) + .count(); + if quorum_met && accepted_count < quorum { + return Err(RadrootsRelayTransportError::InvalidPublishReceipt { + field: "quorum_met", + reason: format!("accepted count {accepted_count} cannot satisfy quorum {quorum}"), + }); + } + Ok(Self { + event_id, + attempted_count, + accepted_count, + retryable_count, + terminal_count, + quorum, + quorum_met, + relays, + }) + } + + pub fn event_id(&self) -> &str { + self.event_id.as_str() + } + + pub fn attempted_count(&self) -> usize { + self.attempted_count + } + + pub fn accepted_count(&self) -> usize { + self.accepted_count + } + + pub fn retryable_count(&self) -> usize { + self.retryable_count + } + + pub fn terminal_count(&self) -> usize { + self.terminal_count + } + + pub fn quorum(&self) -> usize { + self.quorum + } + + pub fn quorum_met(&self) -> bool { + self.quorum_met + } + + pub fn relays(&self) -> &[RadrootsRelayPublishRelayReceipt] { + &self.relays + } + + #[cfg(feature = "storage")] + pub(crate) fn into_event_id_and_relays( + self, + ) -> (String, Vec<RadrootsRelayPublishRelayReceipt>) { + (self.event_id.into_string(), self.relays) + } +} + +struct BoundedPublishRelayReceipts; + +impl<'de> de::Visitor<'de> for BoundedPublishRelayReceipts { + type Value = Vec<RadrootsRelayPublishRelayReceipt>; + + fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + write!( + formatter, + "one to {} relay publish receipts", + radroots_transport::RADROOTS_TRANSPORT_TARGET_MAX_COUNT + ) + } + + fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error> + where + A: de::SeqAccess<'de>, + { + let mut relays = Vec::new(); + while let Some(receipt) = sequence.next_element()? { + if relays.len() == radroots_transport::RADROOTS_TRANSPORT_TARGET_MAX_COUNT { + return Err(de::Error::custom(format!( + "relay receipt count exceeds maximum {}", + radroots_transport::RADROOTS_TRANSPORT_TARGET_MAX_COUNT + ))); + } + relays.push(receipt); + } + Ok(relays) + } +} + +fn deserialize_publish_relay_receipts<'de, D>( + deserializer: D, +) -> Result<Vec<RadrootsRelayPublishRelayReceipt>, D::Error> +where + D: Deserializer<'de>, +{ + deserializer.deserialize_seq(BoundedPublishRelayReceipts) +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct RadrootsRelayPublishReceiptWire { + event_id: RadrootsEventId, + attempted_count: usize, + accepted_count: usize, + retryable_count: usize, + terminal_count: usize, + quorum: usize, + quorum_met: bool, + #[serde(deserialize_with = "deserialize_publish_relay_receipts")] + relays: Vec<RadrootsRelayPublishRelayReceipt>, +} + +impl<'de> Deserialize<'de> for RadrootsRelayPublishReceipt { + fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> + where + D: Deserializer<'de>, + { + let wire = RadrootsRelayPublishReceiptWire::deserialize(deserializer)?; + let receipt = + Self::from_validated_event_id(wire.event_id, wire.quorum, wire.quorum_met, wire.relays) + .map_err(de::Error::custom)?; + for (field, expected, actual) in [ + ( + "attempted_count", + receipt.attempted_count, + wire.attempted_count, + ), + ( + "accepted_count", + receipt.accepted_count, + wire.accepted_count, + ), + ( + "retryable_count", + receipt.retryable_count, + wire.retryable_count, + ), + ( + "terminal_count", + receipt.terminal_count, + wire.terminal_count, + ), + ] { + if expected != actual { + return Err(de::Error::custom(format!( + "relay publish receipt {field} {actual} does not match derived count {expected}" + ))); + } + } + Ok(receipt) + } } pub trait RadrootsRelayPublishAdapter: Send + Sync { @@ -443,6 +673,7 @@ fn nostr_error_to_transport_error(error: RadrootsRelayTransportError) -> Radroot RadrootsRelayTransportError::EmptyFetchFilters | RadrootsRelayTransportError::InvalidFetchLimit { .. } | RadrootsRelayTransportError::FetchLimitTooLarge { .. } + | RadrootsRelayTransportError::InvalidPublishReceipt { .. } | RadrootsRelayTransportError::InvalidTimestamp { .. } | RadrootsRelayTransportError::InvalidIdempotencyKey { .. } | RadrootsRelayTransportError::RequiredTargetNotRequested { .. } => { @@ -776,30 +1007,8 @@ where let quorum = satisfaction_policy.required_target_count(target_count)?; let relays = normalize_publish_receipts(requested_relays.as_slice(), adapter.publish(request).await?)?; - let attempted_count = relays.iter().filter(|receipt| receipt.attempted).count(); - let accepted_count = relays - .iter() - .filter(|receipt| relay_receipt_counts_toward_quorum(receipt)) - .count(); - let retryable_count = relays - .iter() - .filter(|receipt| receipt.outcome.is_retryable()) - .count(); - let terminal_count = relays - .iter() - .filter(|receipt| receipt.outcome.is_terminal_failure()) - .count(); let quorum_met = relay_publish_satisfies_policy(&satisfaction_policy, target_count, &relays)?; - Ok(RadrootsRelayPublishReceipt { - event_id, - attempted_count, - accepted_count, - retryable_count, - terminal_count, - quorum, - quorum_met, - relays, - }) + RadrootsRelayPublishReceipt::new(event_id, quorum, quorum_met, relays) } fn normalize_publish_receipts( diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -1347,14 +1347,96 @@ async fn mock_publish_preserves_exact_raw_json_and_counts_outcomes() { adapter.captured_raw_events(), vec![signed.raw_json().to_owned()] ); - assert_eq!(receipt.attempted_count, 3); - assert_eq!(receipt.accepted_count, 2); - assert_eq!(receipt.retryable_count, 1); - assert!(receipt.quorum_met); + assert_eq!(receipt.attempted_count(), 3); + assert_eq!(receipt.accepted_count(), 2); + assert_eq!(receipt.retryable_count(), 1); + assert!(receipt.quorum_met()); serde_json::to_string(&receipt).expect("receipt json"); } #[tokio::test] +async fn aggregate_publish_receipts_enforce_counts_cardinality_and_strict_wire() { + let signed = signed_post("maximum aggregate publish receipt"); + let relay_urls = (0..RADROOTS_TRANSPORT_TARGET_MAX_COUNT) + .map(|index| format!("wss://relay-{index}.example.com")) + .collect::<Vec<_>>(); + let targets = RadrootsRelayTargetSet::new( + relay_urls.iter().map(String::as_str), + RadrootsRelayUrlPolicy::Public, + ) + .expect("maximum relay target set"); + let receipt = publish_signed_event( + &RadrootsMockRelayPublishAdapter::new(), + RadrootsRelayPublishRequest::new(verified_signed_event(signed), targets, 1_001) + .expect("maximum publish request"), + ) + .await + .expect("maximum aggregate publish receipt"); + assert_eq!(receipt.relays().len(), RADROOTS_TRANSPORT_TARGET_MAX_COUNT); + assert_eq!( + receipt.attempted_count(), + RADROOTS_TRANSPORT_TARGET_MAX_COUNT + ); + + let wire = serde_json::to_value(&receipt).expect("aggregate publish receipt JSON"); + let decoded = serde_json::from_value::<radroots_transport_nostr::RadrootsRelayPublishReceipt>( + wire.clone(), + ) + .expect("strict aggregate publish receipt reload"); + assert_eq!(decoded, receipt); + + let mut wrong_count = wire.clone(); + wrong_count["attempted_count"] = serde_json::json!(RADROOTS_TRANSPORT_TARGET_MAX_COUNT - 1); + assert!( + serde_json::from_value::<radroots_transport_nostr::RadrootsRelayPublishReceipt>( + wrong_count + ) + .is_err() + ); + + let mut duplicate = wire.clone(); + duplicate["relays"][1] = duplicate["relays"][0].clone(); + assert!( + serde_json::from_value::<radroots_transport_nostr::RadrootsRelayPublishReceipt>(duplicate) + .is_err() + ); + + let mut empty = wire.clone(); + empty["relays"] = serde_json::json!([]); + empty["attempted_count"] = serde_json::json!(0); + empty["accepted_count"] = serde_json::json!(0); + assert!( + serde_json::from_value::<radroots_transport_nostr::RadrootsRelayPublishReceipt>(empty) + .is_err() + ); + + let relay_template = wire["relays"][0].clone(); + let mut one_over = wire.clone(); + one_over["relays"] = serde_json::Value::Array( + (0..=RADROOTS_TRANSPORT_TARGET_MAX_COUNT) + .map(|index| { + let mut relay = relay_template.clone(); + relay["relay_url"] = serde_json::json!(format!("wss://wire-{index}.example.com")); + relay + }) + .collect(), + ); + one_over["attempted_count"] = serde_json::json!(RADROOTS_TRANSPORT_TARGET_MAX_COUNT + 1); + one_over["accepted_count"] = serde_json::json!(RADROOTS_TRANSPORT_TARGET_MAX_COUNT + 1); + assert!( + serde_json::from_value::<radroots_transport_nostr::RadrootsRelayPublishReceipt>(one_over) + .is_err() + ); + + let mut unknown = wire; + unknown["extra"] = serde_json::json!(true); + assert!( + serde_json::from_value::<radroots_transport_nostr::RadrootsRelayPublishReceipt>(unknown) + .is_err() + ); +} + +#[tokio::test] async fn nostr_transport_facade_delivers_signed_event_payloads() { let signed = signed_post("facade payload"); let adapter = RadrootsMockRelayPublishAdapter::new(); @@ -1715,7 +1797,7 @@ async fn nostr_transport_facade_matches_canonical_equivalent_relay_receipts() { ) .await .expect("relay publish"); - assert!(relay_receipt.quorum_met); + assert!(relay_receipt.quorum_met()); } #[tokio::test] @@ -1772,13 +1854,13 @@ async fn publish_receipts_track_terminal_skipped_and_adapter_errors() { .await .expect("publish"); - assert_eq!(receipt.event_id, signed.id_str()); - assert_eq!(receipt.attempted_count, 2); - assert_eq!(receipt.accepted_count, 1); - assert_eq!(receipt.retryable_count, 0); - assert_eq!(receipt.terminal_count, 1); - assert_eq!(receipt.quorum, 2); - assert!(!receipt.quorum_met); + assert_eq!(receipt.event_id(), signed.id_str()); + assert_eq!(receipt.attempted_count(), 2); + assert_eq!(receipt.accepted_count(), 1); + assert_eq!(receipt.retryable_count(), 0); + assert_eq!(receipt.terminal_count(), 1); + assert_eq!(receipt.quorum(), 2); + assert!(!receipt.quorum_met()); let skipped = bounded_publish_relay_receipt(RadrootsRelayPublishRelayReceipt::skipped( RELAY_TERTIARY_WSS, @@ -1837,9 +1919,9 @@ async fn publish_required_target_policy_uses_relay_fingerprints() { .await .expect("publish"); - assert_eq!(receipt.accepted_count, 1); - assert_eq!(receipt.quorum, 1); - assert!(!receipt.quorum_met); + assert_eq!(receipt.accepted_count(), 1); + assert_eq!(receipt.quorum(), 1); + assert!(!receipt.quorum_met()); } #[tokio::test] @@ -1864,15 +1946,15 @@ async fn publish_all_policy_uses_requested_target_count() { .await .expect("publish"); - assert_eq!(receipt.attempted_count, 1); - assert_eq!(receipt.accepted_count, 1); - assert_eq!(receipt.retryable_count, 1); - assert_eq!(receipt.quorum, 2); - assert!(!receipt.quorum_met); - assert_eq!(receipt.relays.len(), 2); - assert!(!receipt.relays[1].was_attempted()); + assert_eq!(receipt.attempted_count(), 1); + assert_eq!(receipt.accepted_count(), 1); + assert_eq!(receipt.retryable_count(), 1); + assert_eq!(receipt.quorum(), 2); + assert!(!receipt.quorum_met()); + assert_eq!(receipt.relays().len(), 2); + assert!(!receipt.relays()[1].was_attempted()); assert_eq!( - receipt.relays[1].outcome().kind(), + receipt.relays()[1].outcome().kind(), RadrootsRelayOutcomeKind::Unknown ); @@ -1884,8 +1966,8 @@ async fn publish_all_policy_uses_requested_target_count() { ) .await .expect("no-wait publish"); - assert_eq!(no_wait.quorum, 0); - assert!(no_wait.quorum_met); + assert_eq!(no_wait.quorum(), 0); + assert!(no_wait.quorum_met()); } #[tokio::test]