lib

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

commit 08af0d4bb35715a5151940a089beb1c31fdb0788
parent 75e1483797adb413d2f815cebeeb790f2a77fe84
Author: triesap <tyson@radroots.org>
Date:   Mon, 27 Jul 2026 08:03:33 +0000

transport: seal relay outcome diagnostics

- hide relay outcome state behind validated constructors and accessors
- cap diagnostic messages at the shared complete-request budget
- reject oversized and structurally open outcome wire values
- propagate bounded outcome failures through publish and fetch adapters

Diffstat:
Mcrates/transport_nostr/src/error.rs | 7+++++++
Mcrates/transport_nostr/src/fetch.rs | 2+-
Mcrates/transport_nostr/src/outbox.rs | 113+++++++++++++++++++++++++++++++++++++++++--------------------------------------
Mcrates/transport_nostr/src/outcome.rs | 212++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------
Mcrates/transport_nostr/src/publish.rs | 67+++++++++++++++++++++++++++++++++++--------------------------------
Mcrates/transport_nostr/tests/phase1_outbox_publication.rs | 3++-
Mcrates/transport_nostr/tests/transport.rs | 159++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
7 files changed, 393 insertions(+), 170 deletions(-)

diff --git a/crates/transport_nostr/src/error.rs b/crates/transport_nostr/src/error.rs @@ -93,6 +93,13 @@ pub enum RadrootsRelayTransportError { actual: usize, }, + #[error("Relay transport {field} uses {actual} UTF-8 bytes; maximum is {max}")] + DiagnosticLimitExceeded { + field: &'static str, + max: usize, + actual: usize, + }, + #[error("Relay transport {field} cannot be negative: {value}")] InvalidTimestamp { field: &'static str, value: i64 }, diff --git a/crates/transport_nostr/src/fetch.rs b/crates/transport_nostr/src/fetch.rs @@ -1147,7 +1147,7 @@ fn process_relay_fetch_items( .push(RadrootsRelayFetchRelayOutcome { relay_url, kind: RadrootsRelayFetchOutcomeKind::Closed, - relay_outcome: Some(RadrootsRelayOutcome::classify(message.as_str())), + relay_outcome: Some(RadrootsRelayOutcome::classify(message.as_str())?), message: Some(message), }); } diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs @@ -195,7 +195,7 @@ where target_strings, 0, message, - ), + )?, Err(error) => return Err(error), }; let target_receipts = target_receipts_from_relay_receipts(&publishable, &publish.relays); @@ -207,7 +207,7 @@ where for relay in &publish.relays { if relay .outcome - .kind + .kind() .transport_outcome_kind() .target_status() .counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted) @@ -347,7 +347,7 @@ where .validate_for_request(&delivery_request) .map_err(transport_error_to_relay_error)?; let relay_receipts = relay_receipts_from_transport_receipts(&delivery)?; - let target_receipts = target_receipts_from_transport_receipts(&publishable, &delivery); + let target_receipts = target_receipts_from_transport_receipts(&publishable, &delivery)?; for target_receipt in &target_receipts { complete_outbox_delivery_target(outbox, claimed, target_receipt, now_ms).await?; @@ -356,7 +356,7 @@ where for relay in &relay_receipts { if relay .outcome - .kind + .kind() .transport_outcome_kind() .target_status() .counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted) @@ -418,17 +418,17 @@ fn adapter_transport_failure_receipt( relay_urls: Vec<String>, quorum: usize, message: String, -) -> RadrootsRelayPublishReceipt { +) -> Result<RadrootsRelayPublishReceipt, RadrootsRelayTransportError> { let relays = relay_urls .into_iter() .map(|relay_url| { - RadrootsRelayPublishRelayReceipt::attempted( + Ok(RadrootsRelayPublishRelayReceipt::attempted( relay_url, - RadrootsRelayOutcome::connection_failed(message.clone()), - ) + RadrootsRelayOutcome::connection_failed(message.clone())?, + )) }) - .collect::<Vec<_>>(); - RadrootsRelayPublishReceipt { + .collect::<Result<Vec<_>, RadrootsRelayTransportError>>()?; + Ok(RadrootsRelayPublishReceipt { event_id, attempted_count: relays.len(), accepted_count: 0, @@ -437,7 +437,7 @@ fn adapter_transport_failure_receipt( quorum, quorum_met: false, relays, - } + }) } struct PublishableRelays { @@ -522,7 +522,7 @@ fn target_receipts_from_relay_receipts( attempted: relay_receipt.attempted, transport_status: relay_receipt .outcome - .kind + .kind() .transport_outcome_kind() .target_status(), outcome: relay_receipt.outcome.clone(), @@ -535,27 +535,28 @@ fn target_receipts_from_relay_receipts( fn target_receipts_from_transport_receipts( publishable: &PublishableRelays, delivery: &RadrootsTransportDeliveryReceipt, -) -> Vec<RadrootsOutboxPublishTargetReceipt> { - delivery - .target_receipts() - .iter() - .filter_map(|receipt| { - publishable - .relays - .iter() - .find(|target| target.endpoint_fingerprint == *receipt.target().fingerprint()) - .map(|target| 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: receipt.was_attempted(), - transport_status: receipt.status(), - outcome: relay_outcome_from_transport_outcome(receipt.outcome()), - }) - }) - .collect() +) -> Result<Vec<RadrootsOutboxPublishTargetReceipt>, RadrootsRelayTransportError> { + let mut target_receipts = Vec::new(); + for receipt in delivery.target_receipts() { + let Some(target) = publishable + .relays + .iter() + .find(|target| target.endpoint_fingerprint == *receipt.target().fingerprint()) + else { + continue; + }; + 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: receipt.was_attempted(), + transport_status: receipt.status(), + outcome: relay_outcome_from_transport_outcome(receipt.outcome())?, + }); + } + Ok(target_receipts) } async fn complete_outbox_delivery_target( @@ -629,8 +630,7 @@ async fn complete_outbox_delivery_target( receipt.delivery_target_id, receipt .outcome - .message - .as_deref() + .message() .unwrap_or("relay publish deferred until implemented"), now_ms, ) @@ -644,8 +644,7 @@ async fn complete_outbox_delivery_target( receipt.delivery_target_id, receipt .outcome - .message - .as_deref() + .message() .unwrap_or("relay publish skipped by policy"), now_ms, ) @@ -659,8 +658,7 @@ async fn complete_outbox_delivery_target( receipt.delivery_target_id, receipt .outcome - .message - .as_deref() + .message() .unwrap_or("relay publish retryable"), now_ms, ) @@ -674,8 +672,7 @@ async fn complete_outbox_delivery_target( receipt.delivery_target_id, receipt .outcome - .message - .as_deref() + .message() .unwrap_or("relay publish terminal"), now_ms, ) @@ -690,7 +687,7 @@ fn relay_receipts_from_transport_receipts( ) -> Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError> { let mut relay_receipts: Vec<RadrootsRelayPublishRelayReceipt> = Vec::new(); for receipt in delivery.target_receipts() { - let outcome = relay_outcome_from_transport_outcome(receipt.outcome()); + let outcome = relay_outcome_from_transport_outcome(receipt.outcome())?; let relay_receipt = if receipt.was_attempted() { RadrootsRelayPublishRelayReceipt::attempted(receipt.target().uri().as_str(), outcome) } else { @@ -716,15 +713,12 @@ fn relay_receipts_from_transport_receipts( fn relay_outcome_from_transport_outcome( outcome: &RadrootsTransportOutcome, -) -> RadrootsRelayOutcome { +) -> Result<RadrootsRelayOutcome, RadrootsRelayTransportError> { let kind = outcome .code() .and_then(relay_outcome_kind_from_code) .unwrap_or_else(|| relay_outcome_kind_from_transport_outcome(outcome.kind())); - RadrootsRelayOutcome { - kind, - message: outcome.message().map(str::to_owned), - } + RadrootsRelayOutcome::try_new(kind, outcome.message().map(str::to_owned)) } fn relay_outcome_kind_from_code(code: &str) -> Option<crate::RadrootsRelayOutcomeKind> { @@ -1277,7 +1271,8 @@ mod tests { ) .expect("delivery receipt"); let delivered_transport_receipts = - target_receipts_from_transport_receipts(&publishable, &delivery); + target_receipts_from_transport_receipts(&publishable, &delivery) + .expect("bounded target receipts"); assert_eq!( delivered_transport_receipts[0].transport_status, RadrootsTransportDeliveryTargetStatus::Delivered @@ -1298,7 +1293,8 @@ mod tests { ], 2, "offline".to_owned(), - ); + ) + .expect("bounded adapter failure receipt"); assert_eq!(receipt.event_id, "event-1"); assert_eq!(receipt.attempted_count, 2); @@ -1377,9 +1373,10 @@ mod tests { let outcome = RadrootsTransportOutcome::new(transport_kind) .try_with_message(format!("{transport_kind:?}")) .expect("bounded test outcome message"); - let relay_outcome = relay_outcome_from_transport_outcome(&outcome); - assert_eq!(relay_outcome.kind, relay_kind); - assert_eq!(relay_outcome.message.as_deref(), outcome.message()); + let relay_outcome = + relay_outcome_from_transport_outcome(&outcome).expect("bounded relay outcome"); + assert_eq!(relay_outcome.kind(), relay_kind); + assert_eq!(relay_outcome.message(), outcome.message()); } let code_cases = [ @@ -1424,7 +1421,8 @@ mod tests { .try_with_code(code) .expect("bounded test outcome code") ) - .kind, + .expect("bounded relay outcome") + .kind(), relay_kind ); } @@ -1435,7 +1433,8 @@ mod tests { .try_with_code("unrecognized") .expect("bounded test outcome code") ) - .kind, + .expect("bounded relay outcome") + .kind(), RadrootsRelayOutcomeKind::Accepted ); @@ -1566,7 +1565,11 @@ mod tests { )], ) .expect("delivery receipt"); - assert!(target_receipts_from_transport_receipts(&invalid, &delivery).is_empty()); + assert!( + target_receipts_from_transport_receipts(&invalid, &delivery) + .expect("bounded target receipts") + .is_empty() + ); assert_eq!( relay_receipts_from_transport_receipts(&delivery) .expect("relay receipts") diff --git a/crates/transport_nostr/src/outcome.rs b/crates/transport_nostr/src/outcome.rs @@ -1,9 +1,11 @@ #![forbid(unsafe_code)] use radroots_transport::{ - RadrootsTransportError, RadrootsTransportOutcome, RadrootsTransportOutcomeKind, + RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES, RadrootsTransportError, RadrootsTransportOutcome, + RadrootsTransportOutcomeKind, }; -use serde::{Deserialize, Serialize}; +use serde::{Deserialize, Deserializer, Serialize, de}; +use std::fmt; #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)] pub enum RadrootsRelayOutcomeKind { @@ -104,10 +106,10 @@ impl RadrootsRelayOutcomeKind { } } -#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] pub struct RadrootsRelayOutcome { - pub kind: RadrootsRelayOutcomeKind, - pub message: Option<String>, + kind: RadrootsRelayOutcomeKind, + message: Option<String>, } impl RadrootsRelayOutcome { @@ -118,49 +120,67 @@ impl RadrootsRelayOutcome { } } - pub fn duplicate_accepted(message: impl Into<String>) -> Self { - Self { - kind: RadrootsRelayOutcomeKind::DuplicateAccepted, - message: Some(message.into()), - } + pub fn accepted_with_message( + message: impl Into<String>, + ) -> Result<Self, crate::RadrootsRelayTransportError> { + Self::try_new(RadrootsRelayOutcomeKind::Accepted, Some(message.into())) } - pub fn connection_failed(message: impl Into<String>) -> Self { - Self { - kind: RadrootsRelayOutcomeKind::ConnectionFailed, - message: Some(message.into()), - } + pub fn duplicate_accepted( + message: impl Into<String>, + ) -> Result<Self, crate::RadrootsRelayTransportError> { + Self::try_new( + RadrootsRelayOutcomeKind::DuplicateAccepted, + Some(message.into()), + ) } - pub fn unknown(message: impl Into<String>) -> Self { - Self { - kind: RadrootsRelayOutcomeKind::Unknown, - message: Some(message.into()), - } + pub fn connection_failed( + message: impl Into<String>, + ) -> Result<Self, crate::RadrootsRelayTransportError> { + Self::try_new( + RadrootsRelayOutcomeKind::ConnectionFailed, + Some(message.into()), + ) } - pub fn timeout(message: impl Into<String>) -> Self { - Self { - kind: RadrootsRelayOutcomeKind::Timeout, - message: Some(message.into()), - } + pub fn unknown(message: impl Into<String>) -> Result<Self, crate::RadrootsRelayTransportError> { + Self::try_new(RadrootsRelayOutcomeKind::Unknown, Some(message.into())) } - pub fn relay_url_rejected(message: impl Into<String>) -> Self { - Self { - kind: RadrootsRelayOutcomeKind::RelayUrlRejected, - message: Some(message.into()), - } + pub fn timeout(message: impl Into<String>) -> Result<Self, crate::RadrootsRelayTransportError> { + Self::try_new(RadrootsRelayOutcomeKind::Timeout, Some(message.into())) } - pub fn skipped_already_accepted(message: impl Into<String>) -> Self { - Self { - kind: RadrootsRelayOutcomeKind::SkippedAlreadyAccepted, - message: Some(message.into()), + pub fn relay_url_rejected( + message: impl Into<String>, + ) -> Result<Self, crate::RadrootsRelayTransportError> { + Self::try_new( + RadrootsRelayOutcomeKind::RelayUrlRejected, + Some(message.into()), + ) + } + + pub fn skipped_already_accepted( + message: impl Into<String>, + ) -> Result<Self, crate::RadrootsRelayTransportError> { + Self::try_new( + RadrootsRelayOutcomeKind::SkippedAlreadyAccepted, + Some(message.into()), + ) + } + + pub fn try_new( + kind: RadrootsRelayOutcomeKind, + message: Option<String>, + ) -> Result<Self, crate::RadrootsRelayTransportError> { + if let Some(message) = message.as_deref() { + ensure_relay_outcome_message(message)?; } + Ok(Self { kind, message }) } - pub fn classify(message: impl AsRef<str>) -> Self { + pub fn classify(message: impl AsRef<str>) -> Result<Self, crate::RadrootsRelayTransportError> { let message = message.as_ref().trim(); let lower = message.to_ascii_lowercase(); let kind = if lower.starts_with("duplicate:") { @@ -190,10 +210,15 @@ impl RadrootsRelayOutcome { } else { RadrootsRelayOutcomeKind::Unknown }; - Self { - kind, - message: Some(message.to_owned()), - } + Self::try_new(kind, Some(message.to_owned())) + } + + pub fn kind(&self) -> RadrootsRelayOutcomeKind { + self.kind + } + + pub fn message(&self) -> Option<&str> { + self.message.as_deref() } pub fn counts_toward_quorum(&self) -> bool { @@ -217,3 +242,112 @@ impl RadrootsRelayOutcome { Ok(outcome) } } + +fn ensure_relay_outcome_message(message: &str) -> Result<(), crate::RadrootsRelayTransportError> { + if message.len() > RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES { + return Err( + crate::RadrootsRelayTransportError::DiagnosticLimitExceeded { + field: "relay_outcome_message", + max: RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES, + actual: message.len(), + }, + ); + } + Ok(()) +} + +struct BoundedRelayOutcomeMessage; + +impl<'de> de::Visitor<'de> for BoundedRelayOutcomeMessage { + type Value = String; + + fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + write!( + formatter, + "a relay outcome message of at most {RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES} UTF-8 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, + { + ensure_relay_outcome_message(value) + .map_err(E::custom) + .map(|()| value.to_owned()) + } + + fn visit_string<E>(self, value: String) -> Result<Self::Value, E> + where + E: de::Error, + { + ensure_relay_outcome_message(value.as_str()) + .map_err(E::custom) + .map(|()| value) + } +} + +fn deserialize_relay_outcome_message<'de, D>(deserializer: D) -> Result<Option<String>, D::Error> +where + D: Deserializer<'de>, +{ + struct OptionalBoundedRelayOutcomeMessage; + + impl<'de> de::Visitor<'de> for OptionalBoundedRelayOutcomeMessage { + type Value = Option<String>; + + fn expecting(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("a bounded relay outcome message or null") + } + + fn visit_none<E>(self) -> Result<Self::Value, E> + where + E: de::Error, + { + Ok(None) + } + + fn visit_unit<E>(self) -> Result<Self::Value, E> + where + E: de::Error, + { + Ok(None) + } + + fn visit_some<D>(self, deserializer: D) -> Result<Self::Value, D::Error> + where + D: Deserializer<'de>, + { + deserializer + .deserialize_string(BoundedRelayOutcomeMessage) + .map(Some) + } + } + + deserializer.deserialize_option(OptionalBoundedRelayOutcomeMessage) +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct RadrootsRelayOutcomeWire { + kind: RadrootsRelayOutcomeKind, + #[serde(deserialize_with = "deserialize_relay_outcome_message")] + message: Option<String>, +} + +impl<'de> Deserialize<'de> for RadrootsRelayOutcome { + fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> + where + D: Deserializer<'de>, + { + let wire = RadrootsRelayOutcomeWire::deserialize(deserializer)?; + Self::try_new(wire.kind, wire.message).map_err(de::Error::custom) + } +} diff --git a/crates/transport_nostr/src/publish.rs b/crates/transport_nostr/src/publish.rs @@ -344,6 +344,9 @@ fn nostr_error_to_transport_error(error: RadrootsRelayTransportError) -> Radroot | RadrootsRelayTransportError::RequiredTargetNotRequested { .. } => { RadrootsTransportError::InvalidTransportKind } + RadrootsRelayTransportError::DiagnosticLimitExceeded { field, max, actual } => { + RadrootsTransportError::ResourceLimitExceeded { field, max, actual } + } #[cfg(feature = "storage")] RadrootsRelayTransportError::EventStore(_) | RadrootsRelayTransportError::Outbox(_) @@ -718,11 +721,11 @@ fn normalize_publish_receipts( receipt.relay_url.clone_from(&canonical); if (!receipt.attempted && matches!( - receipt.outcome.kind, + receipt.outcome.kind(), RadrootsRelayOutcomeKind::Accepted | RadrootsRelayOutcomeKind::DuplicateAccepted )) || (receipt.attempted - && receipt.outcome.kind == RadrootsRelayOutcomeKind::SkippedAlreadyAccepted) + && receipt.outcome.kind() == RadrootsRelayOutcomeKind::SkippedAlreadyAccepted) { return Err( RadrootsRelayTransportError::InvalidPublishReceiptAttemptState { url: canonical }, @@ -734,23 +737,25 @@ fn normalize_publish_receipts( ); } } - Ok(requested_relays + requested_relays .iter() .map(|relay_url| { - by_relay.remove(relay_url).unwrap_or_else(|| { - RadrootsRelayPublishRelayReceipt::skipped( + if let Some(receipt) = by_relay.remove(relay_url) { + Ok(receipt) + } else { + Ok(RadrootsRelayPublishRelayReceipt::skipped( relay_url, - RadrootsRelayOutcome::unknown("relay adapter omitted target receipt"), - ) - }) + RadrootsRelayOutcome::unknown("relay adapter omitted target receipt")?, + )) + } }) - .collect()) + .collect() } fn relay_receipt_counts_toward_quorum(receipt: &RadrootsRelayPublishRelayReceipt) -> bool { receipt.outcome.counts_toward_quorum() && (receipt.attempted - || receipt.outcome.kind == RadrootsRelayOutcomeKind::SkippedAlreadyAccepted) + || receipt.outcome.kind() == RadrootsRelayOutcomeKind::SkippedAlreadyAccepted) } fn relay_publish_satisfies_policy( @@ -768,7 +773,7 @@ fn relay_publish_satisfies_policy( relay_receipt_counts_toward_quorum(receipt) && receipt .outcome - .kind + .kind() .transport_outcome_kind() .target_status() .counts_as_satisfied(class) @@ -784,7 +789,7 @@ fn relay_publish_satisfies_policy( && relay_receipt_counts_toward_quorum(receipt) && receipt .outcome - .kind + .kind() .transport_outcome_kind() .target_status() .counts_as_satisfied(class) @@ -927,7 +932,7 @@ impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter { }) .collect::<BTreeMap<_, _>>(); if connected_strings.is_empty() { - return Ok(target_strings + return target_strings .into_iter() .map(|relay_url| { let target_url = relay_url.trim_end_matches('/'); @@ -935,26 +940,26 @@ impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter { .get(target_url) .cloned() .unwrap_or_else(|| "relay did not connect".to_owned()); - RadrootsRelayPublishRelayReceipt::attempted( + Ok(RadrootsRelayPublishRelayReceipt::attempted( relay_url, - RadrootsRelayOutcome::connection_failed(reason), - ) + RadrootsRelayOutcome::connection_failed(reason)?, + )) }) - .collect()); + .collect(); } let output = match self.client.send_event_to(connected_strings, &event).await { Ok(output) => output, Err(error) => { let message = error.to_string(); - return Ok(target_strings + return target_strings .into_iter() .map(|relay_url| { - RadrootsRelayPublishRelayReceipt::attempted( + Ok(RadrootsRelayPublishRelayReceipt::attempted( relay_url, - RadrootsRelayOutcome::connection_failed(message.clone()), - ) + RadrootsRelayOutcome::connection_failed(message.clone())?, + )) }) - .collect()); + .collect(); } }; let mut receipts = Vec::new(); @@ -967,19 +972,16 @@ impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter { if success { receipts.push(RadrootsRelayPublishRelayReceipt::attempted( relay_url, - RadrootsRelayOutcome { - kind: RadrootsRelayOutcomeKind::Accepted, - message: Some( - "nostr-relay-pool-success-ok-message-unavailable".to_owned(), - ), - }, + 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()), + RadrootsRelayOutcome::connection_failed(reason.clone())?, )); continue; } @@ -992,9 +994,10 @@ impl RadrootsRelayPublishAdapter for RadrootsNostrClientPublishAdapter { }); let outcome = failed .map(RadrootsRelayOutcome::classify) - .unwrap_or_else(|| { - RadrootsRelayOutcome::classify("error: relay output omitted target") - }); + .transpose()? + .unwrap_or(RadrootsRelayOutcome::classify( + "error: relay output omitted target", + )?); receipts.push(RadrootsRelayPublishRelayReceipt::attempted( relay_url, outcome, )); diff --git a/crates/transport_nostr/tests/phase1_outbox_publication.rs b/crates/transport_nostr/tests/phase1_outbox_publication.rs @@ -175,7 +175,8 @@ async fn outbox_publication_all_seven_leaves_reuse_exact_bytes_and_dispatch_iden let first_adapter = RadrootsMockRelayPublishAdapter::new().with_outcome( RELAY, - RadrootsRelayOutcome::connection_failed("relay unavailable"), + RadrootsRelayOutcome::connection_failed("relay unavailable") + .expect("bounded relay outcome"), ); let first_transport = RadrootsNostrTransport::new(first_adapter.clone()); let first_receipt = publish_claimed_phase1_publication_target_with_transport( diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -19,12 +19,13 @@ use radroots_outbox::{ RadrootsOutboxOperationStatus, }; use radroots_transport::{ - RADROOTS_TRANSPORT_ENDPOINT_URI_MAX_BYTES, RADROOTS_TRANSPORT_TARGET_MAX_COUNT, - RadrootsTransport, RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest, - RadrootsTransportDeliveryTargetStatus, RadrootsTransportError, RadrootsTransportFetchReceipt, - RadrootsTransportFetchRequest, RadrootsTransportFuture, RadrootsTransportImplementationState, - RadrootsTransportKind, RadrootsTransportMeshScopeId, RadrootsTransportOutcome, - RadrootsTransportOutcomeKind, RadrootsTransportPayload, RadrootsTransportSatisfactionClass, + RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES, RADROOTS_TRANSPORT_ENDPOINT_URI_MAX_BYTES, + RADROOTS_TRANSPORT_TARGET_MAX_COUNT, RadrootsTransport, RadrootsTransportDeliveryReceipt, + RadrootsTransportDeliveryRequest, RadrootsTransportDeliveryTargetStatus, + RadrootsTransportError, RadrootsTransportFetchReceipt, RadrootsTransportFetchRequest, + RadrootsTransportFuture, RadrootsTransportImplementationState, RadrootsTransportKind, + RadrootsTransportMeshScopeId, RadrootsTransportOutcome, RadrootsTransportOutcomeKind, + RadrootsTransportPayload, RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy, RadrootsTransportStatus, RadrootsTransportTarget, RadrootsTransportTargetLabel, RadrootsTransportTargetReceipt, RadrootsTransportTargetSet, }; @@ -55,6 +56,12 @@ const RELAY_PRIMARY_WSS: &str = "wss://relay.example.com"; const RELAY_SECONDARY_WSS: &str = "wss://relay-2.example.com"; const RELAY_TERTIARY_WSS: &str = "wss://relay-3.example.com"; +fn bounded_relay_outcome( + outcome: Result<RadrootsRelayOutcome, RadrootsRelayTransportError>, +) -> RadrootsRelayOutcome { + outcome.expect("bounded relay outcome") +} + struct TransportFailurePublishAdapter; impl RadrootsRelayPublishAdapter for TransportFailurePublishAdapter { @@ -238,7 +245,9 @@ impl RadrootsRelayPublishAdapter for AttemptedSkippedRelayReceiptPublishAdapter .as_str(); Ok(vec![RadrootsRelayPublishRelayReceipt::attempted( relay, - RadrootsRelayOutcome::skipped_already_accepted("already accepted"), + bounded_relay_outcome(RadrootsRelayOutcome::skipped_already_accepted( + "already accepted", + )), )]) }) } @@ -1058,8 +1067,8 @@ fn outcome_prefix_classification_covers_required_kinds() { ]; for (message, kind) in cases { - let outcome = RadrootsRelayOutcome::classify(message); - assert_eq!(outcome.kind, kind); + let outcome = bounded_relay_outcome(RadrootsRelayOutcome::classify(message)); + assert_eq!(outcome.kind(), kind); } let labels = [ (RadrootsRelayOutcomeKind::Accepted, "accepted"), @@ -1099,14 +1108,32 @@ fn outcome_prefix_classification_covers_required_kinds() { assert_eq!(kind.as_str(), label); } - assert!(RadrootsRelayOutcome::classify("duplicate: already have it").counts_toward_quorum()); assert!( - RadrootsRelayOutcome::skipped_already_accepted("already accepted").counts_toward_quorum() + bounded_relay_outcome(RadrootsRelayOutcome::classify("duplicate: already have it")) + .counts_toward_quorum() + ); + assert!( + bounded_relay_outcome(RadrootsRelayOutcome::skipped_already_accepted( + "already accepted" + )) + .counts_toward_quorum() + ); + assert!( + bounded_relay_outcome(RadrootsRelayOutcome::classify("auth-required: challenge")) + .is_retryable() + ); + assert!( + bounded_relay_outcome(RadrootsRelayOutcome::classify("restricted: denied")) + .is_terminal_failure() + ); + assert!( + bounded_relay_outcome(RadrootsRelayOutcome::relay_url_rejected("unsafe relay")) + .is_terminal_failure() + ); + assert!( + bounded_relay_outcome(RadrootsRelayOutcome::classify("mute: pubkey muted")) + .is_terminal_failure() ); - assert!(RadrootsRelayOutcome::classify("auth-required: challenge").is_retryable()); - assert!(RadrootsRelayOutcome::classify("restricted: denied").is_terminal_failure()); - assert!(RadrootsRelayOutcome::relay_url_rejected("unsafe relay").is_terminal_failure()); - assert!(RadrootsRelayOutcome::classify("mute: pubkey muted").is_terminal_failure()); assert_eq!( RadrootsRelayOutcome::accepted() .to_transport_outcome() @@ -1122,62 +1149,100 @@ fn outcome_prefix_classification_covers_required_kinds() { radroots_transport::RadrootsTransportDeliveryTargetStatus::Accepted ); assert_eq!( - RadrootsRelayOutcome::timeout("timeout: no OK") + bounded_relay_outcome(RadrootsRelayOutcome::timeout("timeout: no OK")) .to_transport_outcome() .expect("bounded outcome") .kind(), radroots_transport::RadrootsTransportOutcomeKind::Timeout ); assert_eq!( - RadrootsRelayOutcome::timeout("timeout: no OK") + bounded_relay_outcome(RadrootsRelayOutcome::timeout("timeout: no OK")) .to_transport_outcome() .expect("bounded outcome") .status(), radroots_transport::RadrootsTransportDeliveryTargetStatus::FailedRetryable ); assert_eq!( - RadrootsRelayOutcome::classify("restricted: denied") + bounded_relay_outcome(RadrootsRelayOutcome::classify("restricted: denied")) .to_transport_outcome() .expect("bounded outcome") .kind(), radroots_transport::RadrootsTransportOutcomeKind::Rejected ); assert_eq!( - RadrootsRelayOutcome::classify("restricted: denied") + bounded_relay_outcome(RadrootsRelayOutcome::classify("restricted: denied")) .to_transport_outcome() .expect("bounded outcome") .status(), radroots_transport::RadrootsTransportDeliveryTargetStatus::FailedTerminal ); assert_eq!( - RadrootsRelayOutcome::relay_url_rejected("unsafe") + bounded_relay_outcome(RadrootsRelayOutcome::relay_url_rejected("unsafe")) .to_transport_outcome() .expect("bounded outcome") .kind(), radroots_transport::RadrootsTransportOutcomeKind::RouteUnavailable ); assert_eq!( - RadrootsRelayOutcome::connection_failed("offline") - .kind + bounded_relay_outcome(RadrootsRelayOutcome::connection_failed("offline")) + .kind() .as_str(), "connection_failed" ); assert_eq!( - RadrootsRelayOutcome::unknown("adapter omitted receipt") + bounded_relay_outcome(RadrootsRelayOutcome::unknown("adapter omitted receipt")) .to_transport_outcome() .expect("bounded outcome") .kind(), radroots_transport::RadrootsTransportOutcomeKind::TransportUnavailable ); assert_eq!( - RadrootsRelayOutcome::relay_url_rejected("unsafe") - .kind + bounded_relay_outcome(RadrootsRelayOutcome::relay_url_rejected("unsafe")) + .kind() .as_str(), "relay_url_rejected" ); } #[test] +fn relay_outcome_messages_are_bounded_and_strictly_decoded() { + let exact_message = "x".repeat(RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES); + let outcome = RadrootsRelayOutcome::unknown(exact_message.clone()) + .expect("maximum relay outcome message"); + assert_eq!(outcome.kind(), RadrootsRelayOutcomeKind::Unknown); + assert_eq!(outcome.message(), Some(exact_message.as_str())); + + let wire = serde_json::to_value(&outcome).expect("relay outcome JSON"); + let decoded = + serde_json::from_value::<RadrootsRelayOutcome>(wire).expect("strict relay outcome reload"); + assert_eq!(decoded, outcome); + + let one_over = "x".repeat(RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES + 1); + assert!(matches!( + RadrootsRelayOutcome::unknown(one_over.clone()), + Err(RadrootsRelayTransportError::DiagnosticLimitExceeded { + field: "relay_outcome_message", + max: RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES, + actual, + }) if actual == RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES + 1 + )); + let error = serde_json::from_value::<RadrootsRelayOutcome>(serde_json::json!({ + "kind": "Unknown", + "message": one_over, + })) + .expect_err("oversized relay outcome message rejected"); + assert!(error.to_string().contains("relay_outcome_message")); + assert!( + serde_json::from_value::<RadrootsRelayOutcome>(serde_json::json!({ + "kind": "Accepted", + "message": null, + "extra": true, + })) + .is_err() + ); +} + +#[test] fn relay_transport_error_wraps_transport_contract_errors() { let error = RadrootsRelayTransportError::from( radroots_transport::RadrootsTransportError::EmptyTargetSet, @@ -1200,11 +1265,11 @@ async fn mock_publish_preserves_exact_raw_json_and_counts_outcomes() { let adapter = RadrootsMockRelayPublishAdapter::new() .with_outcome( RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::classify("duplicate: already have it"), + bounded_relay_outcome(RadrootsRelayOutcome::classify("duplicate: already have it")), ) .with_outcome( RELAY_TERTIARY_WSS, - RadrootsRelayOutcome::classify("auth-required: challenge"), + bounded_relay_outcome(RadrootsRelayOutcome::classify("auth-required: challenge")), ); let receipt = publish_signed_event( @@ -1637,7 +1702,9 @@ async fn publish_receipts_track_terminal_skipped_and_adapter_errors() { .expect("targets"); let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome( RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::classify("restricted: group write denied"), + bounded_relay_outcome(RadrootsRelayOutcome::classify( + "restricted: group write denied", + )), ); let receipt = publish_signed_event( @@ -1659,11 +1726,11 @@ async fn publish_receipts_track_terminal_skipped_and_adapter_errors() { let skipped = RadrootsRelayPublishRelayReceipt::skipped( RELAY_TERTIARY_WSS, - RadrootsRelayOutcome::timeout("timeout: no OK"), + 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.outcome.kind(), RadrootsRelayOutcomeKind::Timeout); let error = publish_signed_event( &TransportFailurePublishAdapter, @@ -1693,7 +1760,9 @@ async fn publish_required_target_policy_uses_relay_fingerprints() { let adapter = RadrootsMockRelayPublishAdapter::new() .with_outcome( RELAY_PRIMARY_WSS, - RadrootsRelayOutcome::classify("restricted: required relay rejected"), + bounded_relay_outcome(RadrootsRelayOutcome::classify( + "restricted: required relay rejected", + )), ) .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted()); @@ -1747,7 +1816,7 @@ async fn publish_all_policy_uses_requested_target_count() { assert_eq!(receipt.relays.len(), 2); assert!(!receipt.relays[1].attempted); assert_eq!( - receipt.relays[1].outcome.kind, + receipt.relays[1].outcome.kind(), RadrootsRelayOutcomeKind::Unknown ); @@ -2515,7 +2584,7 @@ async fn fetch_ingests_events_and_records_transport_observations() { .relay_outcome .as_ref() .expect("auth outcome") - .kind, + .kind(), RadrootsRelayOutcomeKind::AuthRequired ); assert_eq!(receipt.relay_outcomes[2].relay_url, RELAY_TERTIARY_WSS); @@ -2524,7 +2593,7 @@ async fn fetch_ingests_events_and_records_transport_observations() { .relay_outcome .as_ref() .expect("restricted outcome") - .kind, + .kind(), RadrootsRelayOutcomeKind::Restricted ); assert_eq!( @@ -3111,7 +3180,7 @@ async fn fetch_event_cap_counts_accepted_in_filter_events_and_preserves_later_co .relay_outcome .as_ref() .expect("closed outcome") - .kind, + .kind(), RadrootsRelayOutcomeKind::AuthRequired ); assert_eq!( @@ -3531,11 +3600,13 @@ async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() { .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()) .with_outcome( RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::timeout("timeout: no OK"), + bounded_relay_outcome(RadrootsRelayOutcome::timeout("timeout: no OK")), ) .with_outcome( RELAY_TERTIARY_WSS, - RadrootsRelayOutcome::duplicate_accepted("duplicate: already have it"), + bounded_relay_outcome(RadrootsRelayOutcome::duplicate_accepted( + "duplicate: already have it", + )), ); let first = publish_claimed_outbox_event( &outbox, @@ -3657,7 +3728,7 @@ async fn outbox_transport_facade_persists_partial_success_and_retryable_failures .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()) .with_outcome( RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::timeout("timeout: transport facade"), + bounded_relay_outcome(RadrootsRelayOutcome::timeout("timeout: transport facade")), ); let transport = RadrootsNostrTransport::new(adapter); let published = publish_claimed_outbox_event_with_transport( @@ -4358,7 +4429,7 @@ async fn outbox_publish_required_target_failure_is_not_satisfied_by_optional_suc let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome( RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::timeout("required relay timeout"), + bounded_relay_outcome(RadrootsRelayOutcome::timeout("required relay timeout")), ); let published = publish_claimed_outbox_event( &outbox, @@ -4669,7 +4740,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 @@ -4938,11 +5009,15 @@ async fn outbox_publish_marks_published_when_delivery_plan_satisfaction_is_met_w .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()) .with_outcome( RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::duplicate_accepted("duplicate: already have it"), + bounded_relay_outcome(RadrootsRelayOutcome::duplicate_accepted( + "duplicate: already have it", + )), ) .with_outcome( RELAY_TERTIARY_WSS, - RadrootsRelayOutcome::classify("restricted: group write denied"), + bounded_relay_outcome(RadrootsRelayOutcome::classify( + "restricted: group write denied", + )), ); let published = publish_claimed_outbox_event( &outbox,