lib

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

commit 09f3eede085fa13c075906dc4d7ba2ab094e300a
parent 7a6f7d64b4d33642183185b6340c5acc2ce73d63
Author: triesap <tyson@radroots.org>
Date:   Sat, 11 Jul 2026 08:01:09 +0000

transport: add Nostr transport facade dispatch

- add RadrootsNostrTransport delivery over signed event payloads
- add transport-backed claimed outbox publish support
- preserve target metadata, URL policy, relay outcomes, and publish time
- cover facade payload rejection and scoped required targets

Diffstat:
Mcrates/transport/src/delivery.rs | 7+++++++
Mcrates/transport_nostr/src/lib.rs | 7++++---
Mcrates/transport_nostr/src/outbox.rs | 375++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/transport_nostr/src/publish.rs | 241++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/transport_nostr/tests/transport.rs | 116++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
5 files changed, 730 insertions(+), 16 deletions(-)

diff --git a/crates/transport/src/delivery.rs b/crates/transport/src/delivery.rs @@ -168,6 +168,7 @@ pub struct RadrootsTransportDeliveryRequest { pub payload: RadrootsTransportPayload, pub target_set: RadrootsTransportTargetSet, pub satisfaction_policy: RadrootsTransportSatisfactionPolicy, + pub now_ms: i64, } impl RadrootsTransportDeliveryRequest { @@ -182,8 +183,14 @@ impl RadrootsTransportDeliveryRequest { payload, target_set, satisfaction_policy, + now_ms: 0, } } + + pub fn with_now_ms(mut self, now_ms: i64) -> Self { + self.now_ms = now_ms; + self + } } #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] diff --git a/crates/transport_nostr/src/lib.rs b/crates/transport_nostr/src/lib.rs @@ -25,13 +25,14 @@ pub use fetch::{ #[cfg(feature = "storage")] pub use outbox::{ RadrootsOutboxPublishPolicy, RadrootsOutboxPublishReceipt, RadrootsOutboxPublishTargetReceipt, - publish_claimed_outbox_event, + publish_claimed_outbox_event, publish_claimed_outbox_event_with_transport, }; pub use outcome::{RadrootsRelayOutcome, RadrootsRelayOutcomeKind}; #[cfg(feature = "client")] pub use publish::RadrootsNostrClientPublishAdapter; pub use publish::{ - RadrootsMockRelayPublishAdapter, RadrootsRelayPublishAdapter, RadrootsRelayPublishReceipt, - RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, publish_signed_event, + RadrootsMockRelayPublishAdapter, RadrootsNostrTransport, RadrootsRelayPublishAdapter, + RadrootsRelayPublishReceipt, RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, + publish_signed_event, }; pub use relay::{RadrootsRelayTargetSet, RadrootsRelayUrl, RadrootsRelayUrlPolicy}; diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs @@ -16,8 +16,11 @@ use radroots_outbox::{ RadrootsOutboxDeliveryTargetStatus, RadrootsOutboxEventStoreIngestReceipt, }; use radroots_transport::{ - RadrootsTransportKind, RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy, - RadrootsTransportTargetFingerprint, + RadrootsTransport, RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest, + RadrootsTransportDeliveryTargetStatus, RadrootsTransportError, RadrootsTransportKind, + RadrootsTransportOutcome, RadrootsTransportOutcomeKind, RadrootsTransportPayload, + RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy, + RadrootsTransportTarget, RadrootsTransportTargetFingerprint, RadrootsTransportTargetSet, }; #[derive(Clone, Debug, PartialEq, Eq)] @@ -250,6 +253,182 @@ where }) } +pub async fn publish_claimed_outbox_event_with_transport<T>( + outbox: &RadrootsOutbox, + event_store: &RadrootsEventStore, + transport: &T, + claimed: &RadrootsOutboxClaimedEvent, + policy: RadrootsOutboxPublishPolicy, + now_ms: i64, +) -> Result<RadrootsOutboxPublishReceipt, RadrootsRelayTransportError> +where + T: RadrootsTransport + ?Sized, +{ + let signed_event = claimed.signed_event.clone().ok_or( + RadrootsRelayTransportError::MissingSignedOutboxEvent(claimed.outbox_event_id), + )?; + let local_ingest = outbox + .ingest_signed_event_local( + event_store, + claimed.outbox_event_id, + claimed.claim_token.as_str(), + now_ms, + ) + .await?; + let publishable = publishable_relays(outbox, claimed, policy.republish_accepted_relays).await?; + if publishable.relays.is_empty() { + outbox + .complete_publish_attempt( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + "relay publish incomplete", + "relay publish terminal", + policy.next_attempt_after_ms, + now_ms, + ) + .await?; + return Ok(RadrootsOutboxPublishReceipt { + local_ingest, + event_id: signed_event.id, + attempted_count: 0, + accepted_count: publishable.accepted_count, + retryable_count: 0, + terminal_count: 0, + quorum: publishable.satisfaction_required_count, + quorum_met: publishable.accepted_count >= publishable.satisfaction_required_count, + target_receipts: Vec::new(), + relay_receipts: Vec::new(), + }); + } + RadrootsRelayTargetSet::new( + publishable + .relays + .iter() + .map(|target| target.relay_url.as_str()), + policy.relay_url_policy, + )?; + let transport_targets = publishable_transport_targets(&publishable)?; + let target_set = RadrootsTransportTargetSet::new(transport_targets)?; + let satisfaction_policy = transport_satisfaction_policy_for_publishable(&publishable)?; + let request_id = outbox_publish_idempotency_key( + claimed.outbox_event_id, + claimed.attempt_count, + signed_event.id.as_str(), + publishable.active_delivery_plan_id, + ); + let payload = RadrootsTransportPayload::signed_event_json( + signed_event.id.clone(), + signed_event.raw_json.clone(), + ) + .map_err(transport_error_to_relay_error)?; + let delivery = transport + .deliver( + RadrootsTransportDeliveryRequest::new( + request_id, + payload, + target_set, + satisfaction_policy, + ) + .with_now_ms(now_ms), + ) + .await + .map_err(transport_error_to_relay_error)?; + let target_receipts = target_receipts_from_transport_receipts(&publishable, &delivery); + + for target_receipt in &target_receipts { + if target_receipt.outcome.counts_toward_quorum() { + outbox + .mark_delivery_target_accepted( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + target_receipt.delivery_target_id, + now_ms, + ) + .await?; + } else if target_receipt.outcome.is_retryable() { + outbox + .mark_delivery_target_failed_retryable( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + target_receipt.delivery_target_id, + target_receipt + .outcome + .message + .as_deref() + .unwrap_or("relay publish retryable"), + now_ms, + ) + .await?; + } else { + outbox + .mark_delivery_target_failed_terminal( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + target_receipt.delivery_target_id, + target_receipt + .outcome + .message + .as_deref() + .unwrap_or("relay publish terminal"), + now_ms, + ) + .await?; + } + } + + for target_receipt in &target_receipts { + if target_receipt.outcome.counts_toward_quorum() { + ingest_publish_observation( + event_store, + &signed_event, + target_receipt.endpoint_uri.as_str(), + target_receipt.outcome.message.as_deref(), + now_ms, + ) + .await?; + } + } + + outbox + .complete_publish_attempt( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + "relay publish incomplete", + "relay publish terminal", + policy.next_attempt_after_ms, + now_ms, + ) + .await?; + + let relay_receipts = relay_receipts_from_transport_receipts(&delivery); + + Ok(RadrootsOutboxPublishReceipt { + local_ingest, + event_id: signed_event.id, + attempted_count: target_receipts + .iter() + .filter(|receipt| receipt.attempted) + .count(), + accepted_count: target_receipts + .iter() + .filter(|receipt| receipt.outcome.counts_toward_quorum()) + .count(), + retryable_count: target_receipts + .iter() + .filter(|receipt| receipt.outcome.is_retryable()) + .count(), + terminal_count: target_receipts + .iter() + .filter(|receipt| receipt.outcome.is_terminal_failure()) + .count(), + quorum: publishable.required_accept_count, + quorum_met: publishable.satisfied_count_after_receipts(&target_receipts) + >= publishable.satisfaction_required_count, + target_receipts, + relay_receipts, + }) +} + fn adapter_transport_failure_receipt( event_id: String, relay_urls: Vec<String>, @@ -344,6 +523,198 @@ fn target_receipts_from_relay_receipts( target_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.status + != RadrootsTransportDeliveryTargetStatus::SkippedPolicyDenied, + outcome: relay_outcome_from_transport_outcome(&receipt.outcome), + }) + }) + .collect() +} + +fn relay_receipts_from_transport_receipts( + delivery: &RadrootsTransportDeliveryReceipt, +) -> Vec<RadrootsRelayPublishRelayReceipt> { + delivery + .target_receipts + .iter() + .map(|receipt| { + RadrootsRelayPublishRelayReceipt::attempted( + receipt.target.uri.as_str(), + relay_outcome_from_transport_outcome(&receipt.outcome), + ) + }) + .collect() +} + +fn relay_outcome_from_transport_outcome( + outcome: &RadrootsTransportOutcome, +) -> RadrootsRelayOutcome { + let kind = outcome + .code + .as_deref() + .and_then(relay_outcome_kind_from_code) + .unwrap_or_else(|| relay_outcome_kind_from_transport_outcome(outcome.kind)); + RadrootsRelayOutcome { + kind, + message: outcome.message.clone(), + } +} + +fn relay_outcome_kind_from_code(code: &str) -> Option<crate::RadrootsRelayOutcomeKind> { + Some(match code { + "accepted" => crate::RadrootsRelayOutcomeKind::Accepted, + "duplicate_accepted" => crate::RadrootsRelayOutcomeKind::DuplicateAccepted, + "blocked" => crate::RadrootsRelayOutcomeKind::Blocked, + "rate_limited" => crate::RadrootsRelayOutcomeKind::RateLimited, + "invalid" => crate::RadrootsRelayOutcomeKind::Invalid, + "pow_required" => crate::RadrootsRelayOutcomeKind::PowRequired, + "restricted" => crate::RadrootsRelayOutcomeKind::Restricted, + "auth_required" => crate::RadrootsRelayOutcomeKind::AuthRequired, + "muted" => crate::RadrootsRelayOutcomeKind::Muted, + "unsupported" => crate::RadrootsRelayOutcomeKind::Unsupported, + "payment_required" => crate::RadrootsRelayOutcomeKind::PaymentRequired, + "error" => crate::RadrootsRelayOutcomeKind::Error, + "timeout" => crate::RadrootsRelayOutcomeKind::Timeout, + "connection_failed" => crate::RadrootsRelayOutcomeKind::ConnectionFailed, + "relay_url_rejected" => crate::RadrootsRelayOutcomeKind::RelayUrlRejected, + "skipped_already_accepted" => crate::RadrootsRelayOutcomeKind::SkippedAlreadyAccepted, + "unknown" => crate::RadrootsRelayOutcomeKind::Unknown, + _ => return None, + }) +} + +fn relay_outcome_kind_from_transport_outcome( + kind: RadrootsTransportOutcomeKind, +) -> crate::RadrootsRelayOutcomeKind { + match kind { + RadrootsTransportOutcomeKind::Accepted => crate::RadrootsRelayOutcomeKind::Accepted, + RadrootsTransportOutcomeKind::DuplicateAccepted => { + crate::RadrootsRelayOutcomeKind::DuplicateAccepted + } + RadrootsTransportOutcomeKind::Rejected => crate::RadrootsRelayOutcomeKind::Invalid, + RadrootsTransportOutcomeKind::RouteUnavailable => { + crate::RadrootsRelayOutcomeKind::RelayUrlRejected + } + RadrootsTransportOutcomeKind::PolicyDenied => crate::RadrootsRelayOutcomeKind::Restricted, + RadrootsTransportOutcomeKind::Timeout => crate::RadrootsRelayOutcomeKind::Timeout, + RadrootsTransportOutcomeKind::ConnectionFailed => { + crate::RadrootsRelayOutcomeKind::ConnectionFailed + } + RadrootsTransportOutcomeKind::TransportUnavailable => { + crate::RadrootsRelayOutcomeKind::Error + } + RadrootsTransportOutcomeKind::PayloadTooLarge => crate::RadrootsRelayOutcomeKind::Invalid, + RadrootsTransportOutcomeKind::Delivered + | RadrootsTransportOutcomeKind::Forwarded + | RadrootsTransportOutcomeKind::StoredByGateway + | RadrootsTransportOutcomeKind::Seen => crate::RadrootsRelayOutcomeKind::Accepted, + RadrootsTransportOutcomeKind::DeferredUntilImplemented => { + crate::RadrootsRelayOutcomeKind::Unsupported + } + } +} + +fn publishable_transport_targets( + publishable: &PublishableRelays, +) -> Result<Vec<RadrootsTransportTarget>, RadrootsRelayTransportError> { + publishable + .relays + .iter() + .map(|relay| { + RadrootsTransportTarget::new_with_metadata( + RadrootsTransportKind::Nostr, + relay.relay_url.as_str(), + relay + .target_scope + .as_deref() + .map(|scope| radroots_transport::RadrootsTransportMeshScopeId::parse(scope)) + .transpose() + .map_err(transport_error_to_relay_error)?, + relay + .target_label + .as_deref() + .map(|label| radroots_transport::RadrootsTransportTargetLabel::parse(label)) + .transpose() + .map_err(transport_error_to_relay_error)?, + ) + .map_err(transport_error_to_relay_error) + }) + .collect() +} + +fn transport_satisfaction_policy_for_publishable( + publishable: &PublishableRelays, +) -> Result<RadrootsTransportSatisfactionPolicy, RadrootsRelayTransportError> { + if publishable.required_targets.is_some() { + let required_targets = publishable + .relays + .iter() + .map(|relay| relay.endpoint_fingerprint.clone()) + .collect::<Vec<_>>(); + return RadrootsTransportSatisfactionPolicy::required_targets( + RadrootsTransportSatisfactionClass::Accepted, + required_targets, + ) + .map_err(transport_error_to_relay_error); + } + satisfaction_policy_for_required_accept_count( + publishable.required_accept_count, + publishable.relays.len(), + false, + ) +} + +fn transport_error_to_relay_error(error: RadrootsTransportError) -> RadrootsRelayTransportError { + match error { + RadrootsTransportError::EmptyTargetUri + | RadrootsTransportError::InvalidTargetUri + | RadrootsTransportError::EmptyTargetSet + | RadrootsTransportError::DuplicateTargetFingerprint + | RadrootsTransportError::InvalidTargetFingerprint => { + RadrootsRelayTransportError::TransportContract(error.to_string()) + } + RadrootsTransportError::EmptyPayloadId + | RadrootsTransportError::InvalidPayloadId + | RadrootsTransportError::EmptyPayloadLabel + | RadrootsTransportError::InvalidPayloadLabel + | RadrootsTransportError::EmptyPayloadBytes + | RadrootsTransportError::InvalidPayloadBytes + | RadrootsTransportError::InvalidPayloadDigest + | RadrootsTransportError::PayloadDigestMismatch => { + RadrootsRelayTransportError::NostrEventJson(error.to_string()) + } + RadrootsTransportError::EmptyTransportKind + | RadrootsTransportError::InvalidTransportKind + | RadrootsTransportError::EmptyTargetScope + | RadrootsTransportError::InvalidTargetScope + | RadrootsTransportError::EmptyTargetLabel + | RadrootsTransportError::InvalidTargetLabel + | RadrootsTransportError::InvalidSatisfactionPolicy + | RadrootsTransportError::EmptyRequiredTargetSet + | RadrootsTransportError::DuplicateRequiredTargetFingerprint => { + RadrootsRelayTransportError::Transport(error.to_string()) + } + } +} + async fn publishable_relays( outbox: &RadrootsOutbox, claimed: &RadrootsOutboxClaimedEvent, diff --git a/crates/transport_nostr/src/publish.rs b/crates/transport_nostr/src/publish.rs @@ -4,9 +4,14 @@ use crate::{RadrootsRelayOutcome, RadrootsRelayTargetSet, RadrootsRelayTransport #[cfg(feature = "client")] use core::time::Duration; use futures::future::BoxFuture; -use radroots_events::draft::RadrootsSignedEvent; +use radroots_events::draft::{RadrootsSignedEvent, RadrootsSignedEventParts}; use radroots_transport::{ - RadrootsTransportKind, RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget, + RadrootsTransport, RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest, + RadrootsTransportError, RadrootsTransportFetchReceipt, RadrootsTransportFetchRequest, + RadrootsTransportFuture, RadrootsTransportImplementationState, RadrootsTransportKind, + RadrootsTransportOutcome, RadrootsTransportOutcomeKind, RadrootsTransportPayload, + RadrootsTransportSatisfactionPolicy, RadrootsTransportStatus, RadrootsTransportTarget, + RadrootsTransportTargetReceipt, }; use serde::{Deserialize, Serialize}; use std::collections::{BTreeMap, BTreeSet}; @@ -104,6 +109,238 @@ pub trait RadrootsRelayPublishAdapter: Send + Sync { ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>>; } +impl<A> RadrootsRelayPublishAdapter for &A +where + A: RadrootsRelayPublishAdapter + ?Sized, +{ + fn publish<'a>( + &'a self, + request: RadrootsRelayPublishRequest, + ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> + { + (*self).publish(request) + } +} + +#[derive(Clone)] +pub struct RadrootsNostrTransport<A> { + adapter: A, + status: RadrootsTransportStatus, +} + +impl<A> RadrootsNostrTransport<A> { + pub fn new(adapter: A) -> Self { + Self { + adapter, + status: RadrootsTransportStatus::new( + RadrootsTransportKind::Nostr, + true, + RadrootsTransportImplementationState::Real, + true, + "ready", + ), + } + } + + pub fn with_status(mut self, status: RadrootsTransportStatus) -> Self { + self.status = status; + self + } + + pub fn adapter(&self) -> &A { + &self.adapter + } +} + +impl<A> RadrootsTransport for RadrootsNostrTransport<A> +where + A: RadrootsRelayPublishAdapter, +{ + fn transport_kind(&self) -> RadrootsTransportKind { + RadrootsTransportKind::Nostr + } + + fn status<'a>(&'a self) -> RadrootsTransportFuture<'a, RadrootsTransportStatus> { + Box::pin(async move { Ok(self.status.clone()) }) + } + + fn deliver<'a>( + &'a self, + request: RadrootsTransportDeliveryRequest, + ) -> RadrootsTransportFuture<'a, RadrootsTransportDeliveryReceipt> { + Box::pin(async move { + let signed_event = signed_event_from_transport_payload(&request.payload)?; + let targets = relay_targets_from_transport_targets(request.target_set.targets())?; + let relay_receipts = match self + .adapter + .publish( + RadrootsRelayPublishRequest::new(signed_event, targets, request.now_ms) + .with_satisfaction_policy(request.satisfaction_policy.clone()) + .with_idempotency_key(request.request_id.clone()), + ) + .await + { + Ok(receipts) => receipts, + Err(RadrootsRelayTransportError::Transport(message)) => { + return Ok(RadrootsTransportDeliveryReceipt { + request_id: request.request_id, + target_receipts: transport_failure_target_receipts( + request.target_set.targets(), + message.as_str(), + ), + }); + } + Err(error) => return Err(nostr_error_to_transport_error(error)), + }; + Ok(RadrootsTransportDeliveryReceipt { + request_id: request.request_id, + target_receipts: target_receipts_from_relay_receipts( + request.target_set.targets(), + relay_receipts.as_slice(), + ), + }) + }) + } + + fn fetch<'a>( + &'a self, + _request: RadrootsTransportFetchRequest, + ) -> RadrootsTransportFuture<'a, RadrootsTransportFetchReceipt> { + Box::pin( + async move { Err(radroots_transport::RadrootsTransportError::InvalidTransportKind) }, + ) + } +} + +fn nostr_error_to_transport_error(error: RadrootsRelayTransportError) -> RadrootsTransportError { + match error { + RadrootsRelayTransportError::TransportContract(_) => { + RadrootsTransportError::InvalidPayloadBytes + } + RadrootsRelayTransportError::RelayUrlParse { .. } + | RadrootsRelayTransportError::WsRequiresLocalhostPolicy { .. } + | RadrootsRelayTransportError::UnsupportedRelayScheme { .. } + | RadrootsRelayTransportError::EmptyRelayHost { .. } + | RadrootsRelayTransportError::RelayUrlUserinfo { .. } + | RadrootsRelayTransportError::RelayUrlQueryOrFragment { .. } + | RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. } + | RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination { .. } + | RadrootsRelayTransportError::EmptyTargetSet => RadrootsTransportError::InvalidTargetUri, + RadrootsRelayTransportError::NostrEventJson(_) | RadrootsRelayTransportError::Json(_) => { + RadrootsTransportError::InvalidPayloadBytes + } + RadrootsRelayTransportError::Transport(_) => RadrootsTransportError::InvalidTransportKind, + RadrootsRelayTransportError::EmptyFetchFilters + | RadrootsRelayTransportError::InvalidFetchLimit { .. } => { + RadrootsTransportError::InvalidTransportKind + } + #[cfg(feature = "storage")] + RadrootsRelayTransportError::EventStore(_) + | RadrootsRelayTransportError::Outbox(_) + | RadrootsRelayTransportError::MissingSignedOutboxEvent(_) => { + RadrootsTransportError::InvalidTransportKind + } + } +} + +#[derive(Deserialize)] +struct SignedEventJsonWire { + id: String, + pubkey: String, + created_at: u32, + kind: u32, + tags: Vec<Vec<String>>, + content: String, + sig: String, +} + +fn signed_event_from_transport_payload( + payload: &RadrootsTransportPayload, +) -> Result<RadrootsSignedEvent, RadrootsTransportError> { + let RadrootsTransportPayload::SignedEventJson { + event_id, raw_json, .. + } = payload + else { + return Err(RadrootsTransportError::InvalidPayloadBytes); + }; + let wire: SignedEventJsonWire = + serde_json::from_str(raw_json).map_err(|_| RadrootsTransportError::InvalidPayloadBytes)?; + if wire.id != *event_id { + return Err(RadrootsTransportError::InvalidPayloadId); + } + RadrootsSignedEvent::new(RadrootsSignedEventParts { + id: wire.id, + pubkey: wire.pubkey, + created_at: wire.created_at, + kind: wire.kind, + tags: wire.tags, + content: wire.content, + sig: wire.sig, + raw_json: raw_json.clone(), + }) + .map_err(|_| RadrootsTransportError::InvalidPayloadBytes) +} + +fn relay_targets_from_transport_targets( + targets: &[RadrootsTransportTarget], +) -> Result<RadrootsRelayTargetSet, RadrootsTransportError> { + let relays = targets + .iter() + .map(|target| { + if target.kind != RadrootsTransportKind::Nostr { + return Err(RadrootsTransportError::InvalidTargetUri); + } + let policy = if target.uri.as_str().starts_with("ws://") { + crate::RadrootsRelayUrlPolicy::Localhost + } else { + crate::RadrootsRelayUrlPolicy::Public + }; + crate::RadrootsRelayUrl::parse(target.uri.as_str(), policy) + .map_err(nostr_error_to_transport_error) + }) + .collect::<Result<Vec<_>, _>>()?; + RadrootsRelayTargetSet::from_urls(relays).map_err(nostr_error_to_transport_error) +} + +fn target_receipts_from_relay_receipts( + targets: &[RadrootsTransportTarget], + relay_receipts: &[RadrootsRelayPublishRelayReceipt], +) -> Vec<RadrootsTransportTargetReceipt> { + targets + .iter() + .cloned() + .map(|target| { + let relay_url = target.uri.as_str().trim_end_matches('/'); + let outcome = relay_receipts + .iter() + .find(|receipt| receipt.relay_url.trim_end_matches('/') == relay_url) + .map(|receipt| receipt.outcome.to_transport_outcome()) + .unwrap_or_else(|| { + RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::RouteUnavailable) + .with_message("relay adapter omitted target receipt") + }); + RadrootsTransportTargetReceipt::new(target, outcome) + }) + .collect() +} + +fn transport_failure_target_receipts( + targets: &[RadrootsTransportTarget], + message: &str, +) -> Vec<RadrootsTransportTargetReceipt> { + targets + .iter() + .cloned() + .map(|target| { + RadrootsTransportTargetReceipt::new( + target, + RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::ConnectionFailed) + .with_message(message.to_owned()), + ) + }) + .collect() +} + pub async fn publish_signed_event<A>( adapter: &A, request: RadrootsRelayPublishRequest, diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -17,17 +17,20 @@ use radroots_outbox::{ RadrootsOutboxOperationStatus, }; use radroots_transport::{ - RadrootsTransportKind, RadrootsTransportMeshScopeId, RadrootsTransportSatisfactionClass, - RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget, RadrootsTransportTargetLabel, + RadrootsTransport, RadrootsTransportDeliveryRequest, RadrootsTransportError, + RadrootsTransportKind, RadrootsTransportMeshScopeId, RadrootsTransportPayload, + RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy, + RadrootsTransportTarget, RadrootsTransportTargetLabel, RadrootsTransportTargetSet, }; use radroots_transport_nostr::{ - RadrootsMockRelayFetchAdapter, RadrootsMockRelayPublishAdapter, RadrootsOutboxPublishPolicy, - RadrootsRelayFetchFilters, RadrootsRelayFetchItem, RadrootsRelayFetchMode, - RadrootsRelayFetchOutcomeKind, 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_signed_event, + RadrootsMockRelayFetchAdapter, RadrootsMockRelayPublishAdapter, RadrootsNostrTransport, + RadrootsOutboxPublishPolicy, RadrootsRelayFetchFilters, RadrootsRelayFetchItem, + RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, 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_signed_event, }; use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; @@ -701,6 +704,101 @@ async fn mock_publish_preserves_exact_raw_json_and_counts_outcomes() { } #[tokio::test] +async fn nostr_transport_facade_delivers_signed_event_payloads() { + let signed = signed_post("facade payload"); + let adapter = RadrootsMockRelayPublishAdapter::new(); + let transport = RadrootsNostrTransport::new(&adapter); + let target = nostr_target(RELAY_PRIMARY_WSS); + let request = RadrootsTransportDeliveryRequest::new( + "facade-request-1", + RadrootsTransportPayload::signed_event_json(signed.id.clone(), signed.raw_json.clone()) + .expect("payload"), + RadrootsTransportTargetSet::new(vec![target.clone()]).expect("targets"), + RadrootsTransportSatisfactionPolicy::all_accepted(), + ); + + let receipt = transport.deliver(request).await.expect("delivery"); + + assert_eq!(adapter.captured_raw_events(), vec![signed.raw_json]); + assert_eq!(receipt.request_id, "facade-request-1"); + assert_eq!(receipt.target_receipts.len(), 1); + assert_eq!(receipt.target_receipts[0].target, target); + assert_eq!( + receipt.target_receipts[0].outcome.kind, + radroots_transport::RadrootsTransportOutcomeKind::Accepted + ); + assert!( + receipt + .is_satisfied_by(&RadrootsTransportSatisfactionPolicy::all_accepted()) + .expect("satisfaction") + ); +} + +#[tokio::test] +async fn nostr_transport_facade_rejects_unsupported_payloads_and_targets() { + let signed = signed_post("facade rejected"); + let transport = RadrootsNostrTransport::new(RadrootsMockRelayPublishAdapter::new()); + let target_set = + RadrootsTransportTargetSet::new(vec![nostr_target(RELAY_PRIMARY_WSS)]).expect("targets"); + let payload_error = transport + .deliver(RadrootsTransportDeliveryRequest::new( + "facade-request-payload", + RadrootsTransportPayload::opaque_bytes("not-signed-event", [1, 2, 3]).expect("payload"), + target_set, + RadrootsTransportSatisfactionPolicy::all_accepted(), + )) + .await + .expect_err("payload rejected"); + assert_eq!(payload_error, RadrootsTransportError::InvalidPayloadBytes); + + let non_nostr_target = RadrootsTransportTarget::new( + RadrootsTransportKind::Reticulum, + radroots_transport::RADROOTS_RETICULUM_PREVIEW_ENDPOINT_URI, + ) + .expect("reticulum target"); + let target_error = transport + .deliver(RadrootsTransportDeliveryRequest::new( + "facade-request-target", + RadrootsTransportPayload::signed_event_json(signed.id.clone(), signed.raw_json.clone()) + .expect("payload"), + RadrootsTransportTargetSet::new(vec![non_nostr_target]).expect("targets"), + RadrootsTransportSatisfactionPolicy::all_accepted(), + )) + .await + .expect_err("target rejected"); + assert_eq!(target_error, RadrootsTransportError::InvalidTargetUri); +} + +#[tokio::test] +async fn nostr_transport_facade_preserves_scoped_duplicate_target_metadata() { + let signed = signed_post("facade scoped duplicate"); + let adapter = RadrootsMockRelayPublishAdapter::new(); + let transport = RadrootsNostrTransport::new(&adapter); + let first = scoped_nostr_target(RELAY_PRIMARY_WSS, "local_food_buyers", "buyers"); + let second = scoped_nostr_target(RELAY_PRIMARY_WSS, "local_food_farmers", "farmers"); + let policy = RadrootsTransportSatisfactionPolicy::required_targets( + RadrootsTransportSatisfactionClass::Accepted, + vec![first.fingerprint.clone(), second.fingerprint.clone()], + ) + .expect("required targets"); + let request = RadrootsTransportDeliveryRequest::new( + "facade-request-scoped", + RadrootsTransportPayload::signed_event_json(signed.id.clone(), signed.raw_json.clone()) + .expect("payload"), + RadrootsTransportTargetSet::new(vec![first.clone(), second.clone()]).expect("targets"), + policy.clone(), + ); + + let receipt = transport.deliver(request).await.expect("delivery"); + + assert_eq!(receipt.target_receipts.len(), 2); + assert_eq!(receipt.target_receipts[0].target, first); + assert_eq!(receipt.target_receipts[1].target, second); + assert!(receipt.is_satisfied_by(&policy).expect("satisfaction")); + assert_eq!(adapter.captured_raw_events().len(), 1); +} + +#[tokio::test] async fn publish_receipts_track_terminal_skipped_and_adapter_errors() { let signed = signed_post("terminal"); let targets = RadrootsRelayTargetSet::new(