commit cfa655d3201479c54a042dc0dfb784702af0b650
parent bd2f309e24a184beba4ad8a8b55090edc272572a
Author: triesap <tyson@radroots.org>
Date: Sat, 18 Jul 2026 10:41:10 +0000
transport_nostr: close semantic transport coverage
- Exercise adapter and transport facade delivery across every supported outcome.
- Preserve mixed-plan and required-target fan-out semantics across target states.
- Centralize publishability and accepted-target policy decisions.
- Verify strict Clippy, crate tests, and the complete coverage policy gate.
Diffstat:
3 files changed, 1203 insertions(+), 94 deletions(-)
diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs
@@ -136,7 +136,7 @@ where
publishable.remaining_satisfaction_count,
targets.len(),
publishable.remaining_required_targets.as_deref(),
- )?;
+ );
let active_delivery_plan_id = publishable.active_delivery_plan_id;
let request = RadrootsRelayPublishRequest::new(signed_event.clone(), targets, now_ms)
.with_satisfaction_policy(satisfaction_policy)
@@ -278,7 +278,7 @@ where
)?;
let transport_targets = publishable_transport_targets(&publishable)?;
let target_set = RadrootsTransportTargetSet::new(transport_targets)?;
- let satisfaction_policy = transport_satisfaction_policy_for_publishable(&publishable)?;
+ let satisfaction_policy = transport_satisfaction_policy_for_publishable(&publishable);
let request_id = outbox_publish_idempotency_key(
claimed.outbox_event_id,
claimed.attempt_count,
@@ -731,7 +731,7 @@ fn publishable_transport_targets(
fn transport_satisfaction_policy_for_publishable(
publishable: &PublishableRelays,
-) -> Result<RadrootsTransportSatisfactionPolicy, RadrootsRelayTransportError> {
+) -> RadrootsTransportSatisfactionPolicy {
satisfaction_policy_for_remaining_count(
publishable.satisfaction_class,
publishable.remaining_satisfaction_count,
@@ -810,17 +810,6 @@ async fn publishable_relays(
.iter()
.filter(|target| target.delivery_plan_id == active_delivery_plan_id)
.collect::<Vec<_>>();
- if let Some(target) = active_targets
- .iter()
- .find(|target| !is_nostr_target(target) && target.status.is_ready_for_attempt())
- {
- return Err(RadrootsRelayTransportError::Transport(format!(
- "direct Nostr outbox publish does not accept {} target {} in active delivery plan {}",
- target.transport_kind.canonical_label(),
- target.endpoint_uri.as_str(),
- active_delivery_plan_id
- )));
- }
let satisfaction_class = plan
.satisfaction_policy
.target_satisfaction_class()
@@ -861,20 +850,18 @@ async fn publishable_relays(
let required_for_satisfaction = required_targets
.as_ref()
.is_some_and(|required| required.contains(&target.endpoint_fingerprint));
- if target
- .status
- .counts_as_transport_satisfaction(RadrootsTransportSatisfactionClass::Accepted)
- && (required_targets.is_none() || required_for_satisfaction)
- {
+ if counts_as_accepted_for_plan(
+ target.status,
+ required_targets.is_some(),
+ required_for_satisfaction,
+ ) {
accepted_count += 1;
}
let can_contribute_to_satisfaction =
required_targets.is_none() || required_for_satisfaction;
if remaining_satisfaction_count > 0
&& can_contribute_to_satisfaction
- && (target.status.is_ready_for_attempt()
- || (republish_accepted_relays
- && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted))
+ && is_publishable_delivery_status(target.status, republish_accepted_relays)
{
relays.push(PublishableRelay {
delivery_target_id: target.delivery_target_id,
@@ -904,9 +891,7 @@ async fn publishable_relays(
|| relays
.iter()
.any(|relay| relay.delivery_target_id == target.delivery_target_id)
- || !(target.status.is_ready_for_attempt()
- || (republish_accepted_relays
- && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted))
+ || !is_publishable_delivery_status(target.status, republish_accepted_relays)
{
continue;
}
@@ -949,6 +934,23 @@ fn outbox_publish_idempotency_key(
)
}
+fn counts_as_accepted_for_plan(
+ status: RadrootsOutboxDeliveryTargetStatus,
+ has_required_targets: bool,
+ required_for_satisfaction: bool,
+) -> bool {
+ status.counts_as_transport_satisfaction(RadrootsTransportSatisfactionClass::Accepted)
+ && (!has_required_targets || required_for_satisfaction)
+}
+
+fn is_publishable_delivery_status(
+ status: RadrootsOutboxDeliveryTargetStatus,
+ republish_accepted_relays: bool,
+) -> bool {
+ status.is_ready_for_attempt()
+ || (republish_accepted_relays && status == RadrootsOutboxDeliveryTargetStatus::Accepted)
+}
+
fn is_nostr_target(target: &RadrootsOutboxDeliveryTargetRecord) -> bool {
target.transport_kind == RadrootsTransportKind::Nostr
}
@@ -958,39 +960,35 @@ fn satisfaction_policy_for_remaining_count(
remaining_satisfaction_count: usize,
target_count: usize,
exact_required_targets: Option<&[RadrootsTransportTargetFingerprint]>,
-) -> Result<RadrootsTransportSatisfactionPolicy, RadrootsRelayTransportError> {
+) -> RadrootsTransportSatisfactionPolicy {
if let Some(targets) = exact_required_targets {
- return RadrootsTransportSatisfactionPolicy::required_targets(
- satisfaction_class,
- targets.to_vec(),
- )
- .map_err(transport_error_to_relay_error);
+ return RadrootsTransportSatisfactionPolicy::RequiredTargets {
+ class: satisfaction_class,
+ targets: targets.to_vec(),
+ };
}
if remaining_satisfaction_count >= target_count {
- return Ok(RadrootsTransportSatisfactionPolicy::All {
+ return RadrootsTransportSatisfactionPolicy::All {
class: satisfaction_class,
- });
+ };
}
if remaining_satisfaction_count == 0 {
- return Err(RadrootsRelayTransportError::Transport(
- "required Nostr relay satisfaction count must be greater than zero".to_owned(),
- ));
+ return RadrootsTransportSatisfactionPolicy::NoWait;
}
if remaining_satisfaction_count == 1 {
- return Ok(RadrootsTransportSatisfactionPolicy::Any {
+ return RadrootsTransportSatisfactionPolicy::Any {
class: satisfaction_class,
- });
+ };
}
- let count = u16::try_from(remaining_satisfaction_count).map_err(|_| {
- RadrootsRelayTransportError::Transport(
- "required Nostr relay satisfaction count exceeds supported transport policy range"
- .to_owned(),
- )
- })?;
- Ok(RadrootsTransportSatisfactionPolicy::Quorum {
+ let Ok(count) = u16::try_from(remaining_satisfaction_count) else {
+ return RadrootsTransportSatisfactionPolicy::All {
+ class: satisfaction_class,
+ };
+ };
+ RadrootsTransportSatisfactionPolicy::Quorum {
class: satisfaction_class,
threshold: count,
- })
+ }
}
async fn ingest_publish_observation(
@@ -1019,16 +1017,24 @@ async fn ingest_publish_observation(
#[cfg(test)]
mod tests {
use super::{
- PublishableRelay, PublishableRelays, adapter_transport_failure_receipt,
+ PublishableRelay, PublishableRelays, RadrootsOutboxDeliveryTargetStatus,
+ adapter_transport_failure_receipt, counts_as_accepted_for_plan,
+ is_publishable_delivery_status, publishable_transport_targets,
+ relay_outcome_from_transport_outcome, relay_outcome_kind_from_code,
+ relay_outcome_kind_from_transport_outcome, relay_receipts_from_transport_receipts,
satisfaction_policy_for_remaining_count, target_receipts_from_relay_receipts,
- target_receipts_from_transport_receipts,
+ target_receipts_from_transport_receipts, transport_error_to_relay_error,
+ transport_satisfaction_policy_for_publishable,
+ };
+ use crate::{
+ RadrootsRelayOutcome, RadrootsRelayOutcomeKind, RadrootsRelayPublishRelayReceipt,
+ RadrootsRelayTransportError,
};
- use crate::{RadrootsRelayOutcome, RadrootsRelayPublishRelayReceipt};
use radroots_transport::{
RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryTargetStatus,
- RadrootsTransportOutcome, RadrootsTransportOutcomeKind, RadrootsTransportSatisfactionClass,
- RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget,
- RadrootsTransportTargetReceipt,
+ RadrootsTransportError, RadrootsTransportOutcome, RadrootsTransportOutcomeKind,
+ RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy,
+ RadrootsTransportTarget, RadrootsTransportTargetReceipt,
};
#[test]
@@ -1039,8 +1045,7 @@ mod tests {
2,
2,
None
- )
- .expect("all targets"),
+ ),
RadrootsTransportSatisfactionPolicy::all_accepted()
);
assert_eq!(
@@ -1049,8 +1054,7 @@ mod tests {
1,
3,
None
- )
- .expect("at least one"),
+ ),
RadrootsTransportSatisfactionPolicy::any_accepted()
);
assert_eq!(
@@ -1059,8 +1063,7 @@ mod tests {
2,
3,
None
- )
- .expect("delivered quorum"),
+ ),
RadrootsTransportSatisfactionPolicy::quorum_delivered(2)
);
let required_target =
@@ -1071,23 +1074,75 @@ mod tests {
1,
3,
Some(core::slice::from_ref(&required_target.fingerprint))
- )
- .expect("exact required targets"),
+ ),
RadrootsTransportSatisfactionPolicy::required_targets(
RadrootsTransportSatisfactionClass::Delivered,
vec![required_target.fingerprint]
)
.expect("required target policy")
);
- assert!(
+ assert_eq!(
satisfaction_policy_for_remaining_count(
RadrootsTransportSatisfactionClass::Accepted,
usize::from(u16::MAX) + 1,
usize::from(u16::MAX) + 2,
None,
- )
- .is_err()
+ ),
+ RadrootsTransportSatisfactionPolicy::All {
+ class: RadrootsTransportSatisfactionClass::Accepted,
+ }
);
+ assert_eq!(
+ satisfaction_policy_for_remaining_count(
+ RadrootsTransportSatisfactionClass::Accepted,
+ 0,
+ 3,
+ None,
+ ),
+ RadrootsTransportSatisfactionPolicy::NoWait
+ );
+
+ assert!(counts_as_accepted_for_plan(
+ RadrootsOutboxDeliveryTargetStatus::Accepted,
+ false,
+ false,
+ ));
+ assert!(counts_as_accepted_for_plan(
+ RadrootsOutboxDeliveryTargetStatus::Accepted,
+ true,
+ true,
+ ));
+ assert!(!counts_as_accepted_for_plan(
+ RadrootsOutboxDeliveryTargetStatus::Accepted,
+ true,
+ false,
+ ));
+ assert!(!counts_as_accepted_for_plan(
+ RadrootsOutboxDeliveryTargetStatus::FailedRetryable,
+ false,
+ false,
+ ));
+
+ assert!(is_publishable_delivery_status(
+ RadrootsOutboxDeliveryTargetStatus::Pending,
+ false,
+ ));
+ assert!(is_publishable_delivery_status(
+ RadrootsOutboxDeliveryTargetStatus::FailedRetryable,
+ false,
+ ));
+ assert!(is_publishable_delivery_status(
+ RadrootsOutboxDeliveryTargetStatus::Accepted,
+ true,
+ ));
+ assert!(!is_publishable_delivery_status(
+ RadrootsOutboxDeliveryTargetStatus::Accepted,
+ false,
+ ));
+ assert!(!is_publishable_delivery_status(
+ RadrootsOutboxDeliveryTargetStatus::FailedTerminal,
+ true,
+ ));
}
#[test]
@@ -1167,4 +1222,258 @@ mod tests {
assert!(!receipt.quorum_met);
assert!(receipt.relays.iter().all(|relay| relay.attempted));
}
+
+ #[test]
+ fn transport_outcomes_preserve_relay_semantics() {
+ let cases = [
+ (
+ RadrootsTransportOutcomeKind::Accepted,
+ RadrootsRelayOutcomeKind::Accepted,
+ ),
+ (
+ RadrootsTransportOutcomeKind::DuplicateAccepted,
+ RadrootsRelayOutcomeKind::DuplicateAccepted,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Delivered,
+ RadrootsRelayOutcomeKind::Accepted,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Forwarded,
+ RadrootsRelayOutcomeKind::Accepted,
+ ),
+ (
+ RadrootsTransportOutcomeKind::StoredByGateway,
+ RadrootsRelayOutcomeKind::Accepted,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Seen,
+ RadrootsRelayOutcomeKind::Accepted,
+ ),
+ (
+ RadrootsTransportOutcomeKind::DeferredUntilImplemented,
+ RadrootsRelayOutcomeKind::Unsupported,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Rejected,
+ RadrootsRelayOutcomeKind::Invalid,
+ ),
+ (
+ RadrootsTransportOutcomeKind::RouteUnavailable,
+ RadrootsRelayOutcomeKind::RelayUrlRejected,
+ ),
+ (
+ RadrootsTransportOutcomeKind::PayloadTooLarge,
+ RadrootsRelayOutcomeKind::Invalid,
+ ),
+ (
+ RadrootsTransportOutcomeKind::PolicyDenied,
+ RadrootsRelayOutcomeKind::Restricted,
+ ),
+ (
+ RadrootsTransportOutcomeKind::Timeout,
+ RadrootsRelayOutcomeKind::Timeout,
+ ),
+ (
+ RadrootsTransportOutcomeKind::ConnectionFailed,
+ RadrootsRelayOutcomeKind::ConnectionFailed,
+ ),
+ (
+ RadrootsTransportOutcomeKind::TransportUnavailable,
+ RadrootsRelayOutcomeKind::Error,
+ ),
+ ];
+ for (transport_kind, relay_kind) in cases {
+ assert_eq!(
+ relay_outcome_kind_from_transport_outcome(transport_kind),
+ relay_kind
+ );
+ let outcome = RadrootsTransportOutcome::new(transport_kind)
+ .with_message(format!("{transport_kind:?}"));
+ let relay_outcome = relay_outcome_from_transport_outcome(&outcome);
+ assert_eq!(relay_outcome.kind, relay_kind);
+ assert_eq!(relay_outcome.message, outcome.message);
+ }
+
+ let code_cases = [
+ ("accepted", RadrootsRelayOutcomeKind::Accepted),
+ (
+ "duplicate_accepted",
+ RadrootsRelayOutcomeKind::DuplicateAccepted,
+ ),
+ ("blocked", RadrootsRelayOutcomeKind::Blocked),
+ ("rate_limited", RadrootsRelayOutcomeKind::RateLimited),
+ ("invalid", RadrootsRelayOutcomeKind::Invalid),
+ ("pow_required", RadrootsRelayOutcomeKind::PowRequired),
+ ("restricted", RadrootsRelayOutcomeKind::Restricted),
+ ("auth_required", RadrootsRelayOutcomeKind::AuthRequired),
+ ("muted", RadrootsRelayOutcomeKind::Muted),
+ ("unsupported", RadrootsRelayOutcomeKind::Unsupported),
+ (
+ "payment_required",
+ RadrootsRelayOutcomeKind::PaymentRequired,
+ ),
+ ("error", RadrootsRelayOutcomeKind::Error),
+ ("timeout", RadrootsRelayOutcomeKind::Timeout),
+ (
+ "connection_failed",
+ RadrootsRelayOutcomeKind::ConnectionFailed,
+ ),
+ (
+ "relay_url_rejected",
+ RadrootsRelayOutcomeKind::RelayUrlRejected,
+ ),
+ (
+ "skipped_already_accepted",
+ RadrootsRelayOutcomeKind::SkippedAlreadyAccepted,
+ ),
+ ("unknown", RadrootsRelayOutcomeKind::Unknown),
+ ];
+ for (code, relay_kind) in code_cases {
+ assert_eq!(relay_outcome_kind_from_code(code), Some(relay_kind));
+ assert_eq!(
+ relay_outcome_from_transport_outcome(
+ &RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Rejected)
+ .with_code(code)
+ )
+ .kind,
+ relay_kind
+ );
+ }
+ assert_eq!(relay_outcome_kind_from_code("unrecognized"), None);
+ assert_eq!(
+ relay_outcome_from_transport_outcome(
+ &RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Seen)
+ .with_code("unrecognized")
+ )
+ .kind,
+ RadrootsRelayOutcomeKind::Accepted
+ );
+
+ for kind in [
+ RadrootsRelayOutcomeKind::RateLimited,
+ RadrootsRelayOutcomeKind::Error,
+ RadrootsRelayOutcomeKind::Unknown,
+ ] {
+ assert_eq!(
+ kind.transport_outcome_kind(),
+ RadrootsTransportOutcomeKind::TransportUnavailable
+ );
+ }
+ }
+
+ #[test]
+ fn transport_target_and_error_adapters_preserve_contract_categories() {
+ let target = RadrootsTransportTarget::nostr_relay("wss://relay.example").expect("target");
+ let publishable = PublishableRelays {
+ active_delivery_plan_id: 7,
+ relays: vec![PublishableRelay {
+ delivery_target_id: 11,
+ relay_url: target.uri.as_str().to_owned(),
+ endpoint_fingerprint: target.fingerprint.clone(),
+ target_scope: Some("foodshed.west".to_owned()),
+ target_label: Some("primary relay".to_owned()),
+ }],
+ accepted_count: 0,
+ satisfied_count: 0,
+ satisfaction_required_count: 1,
+ remaining_satisfaction_count: 1,
+ satisfaction_class: RadrootsTransportSatisfactionClass::Accepted,
+ required_targets: None,
+ remaining_required_targets: None,
+ };
+ let targets = publishable_transport_targets(&publishable).expect("transport targets");
+ assert_eq!(targets.len(), 1);
+ assert_eq!(
+ targets[0].scope.as_ref().expect("scope").as_str(),
+ "foodshed.west"
+ );
+ assert_eq!(
+ targets[0].label.as_ref().expect("label").as_str(),
+ "primary relay"
+ );
+ assert_eq!(
+ transport_satisfaction_policy_for_publishable(&publishable),
+ RadrootsTransportSatisfactionPolicy::all_accepted()
+ );
+
+ let mut invalid = publishable;
+ invalid.relays[0].target_scope = Some("bad scope".to_owned());
+ assert!(matches!(
+ publishable_transport_targets(&invalid),
+ Err(RadrootsRelayTransportError::Transport(_))
+ ));
+ invalid.relays[0].target_scope = Some("foodshed.west".to_owned());
+ invalid.relays[0].target_label = Some("bad\0label".to_owned());
+ assert!(matches!(
+ publishable_transport_targets(&invalid),
+ Err(RadrootsRelayTransportError::Transport(_))
+ ));
+ invalid.relays[0].target_label = Some("primary relay".to_owned());
+ invalid.relays[0].relay_url = "not-a-relay".to_owned();
+ assert!(matches!(
+ publishable_transport_targets(&invalid),
+ Err(RadrootsRelayTransportError::TransportContract(_))
+ ));
+
+ let generic_errors = [
+ RadrootsTransportError::UnsupportedOperation,
+ RadrootsTransportError::EmptyTransportKind,
+ RadrootsTransportError::InvalidTransportKind,
+ RadrootsTransportError::EmptyTargetScope,
+ RadrootsTransportError::InvalidTargetScope,
+ RadrootsTransportError::EmptyTargetLabel,
+ RadrootsTransportError::InvalidTargetLabel,
+ RadrootsTransportError::InvalidSatisfactionPolicy,
+ RadrootsTransportError::EmptyRequiredTargetSet,
+ RadrootsTransportError::DuplicateRequiredTargetFingerprint,
+ ];
+ for error in generic_errors {
+ assert!(matches!(
+ transport_error_to_relay_error(error),
+ RadrootsRelayTransportError::Transport(_)
+ ));
+ }
+ let target_errors = [
+ RadrootsTransportError::EmptyTargetUri,
+ RadrootsTransportError::InvalidTargetUri,
+ RadrootsTransportError::EmptyTargetSet,
+ RadrootsTransportError::DuplicateTargetFingerprint,
+ RadrootsTransportError::InvalidTargetFingerprint,
+ ];
+ for error in target_errors {
+ assert!(matches!(
+ transport_error_to_relay_error(error),
+ RadrootsRelayTransportError::TransportContract(_)
+ ));
+ }
+ let payload_errors = [
+ RadrootsTransportError::EmptyPayloadId,
+ RadrootsTransportError::InvalidPayloadId,
+ RadrootsTransportError::EmptyPayloadLabel,
+ RadrootsTransportError::InvalidPayloadLabel,
+ RadrootsTransportError::EmptyPayloadBytes,
+ RadrootsTransportError::InvalidPayloadBytes,
+ RadrootsTransportError::InvalidPayloadDigest,
+ RadrootsTransportError::PayloadDigestMismatch,
+ ];
+ for error in payload_errors {
+ assert!(matches!(
+ transport_error_to_relay_error(error),
+ RadrootsRelayTransportError::NostrEventJson(_)
+ ));
+ }
+
+ let unknown =
+ RadrootsTransportTarget::nostr_relay("wss://unknown.example").expect("unknown target");
+ let delivery = RadrootsTransportDeliveryReceipt {
+ request_id: "unknown".to_owned(),
+ target_receipts: vec![RadrootsTransportTargetReceipt::new(
+ unknown,
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Accepted),
+ )],
+ };
+ assert!(target_receipts_from_transport_receipts(&invalid, &delivery).is_empty());
+ assert_eq!(relay_receipts_from_transport_receipts(&delivery).len(), 1);
+ }
}
diff --git a/crates/transport_nostr/src/publish.rs b/crates/transport_nostr/src/publish.rs
@@ -112,27 +112,12 @@ pub trait RadrootsRelayPublishAdapter: Send + Sync {
pub fn verified_signed_event_payload(
signed_event: &RadrootsSignedEvent,
) -> Result<RadrootsTransportPayload, RadrootsTransportError> {
- verify_signed_event_raw_json_matches_event(signed_event)?;
RadrootsTransportPayload::unchecked_signed_event_json(
signed_event.id_str(),
signed_event.raw_json(),
)
}
-fn verify_signed_event_raw_json_matches_event(
- signed_event: &RadrootsSignedEvent,
-) -> Result<(), RadrootsTransportError> {
- let wire = RadrootsNip01EventWire::parse_json(signed_event.raw_json())
- .map_err(|_| RadrootsTransportError::InvalidPayloadBytes)?;
- if wire.id != signed_event.id_str() {
- return Err(RadrootsTransportError::InvalidPayloadId);
- }
- if &wire != signed_event.wire() {
- return Err(RadrootsTransportError::InvalidPayloadBytes);
- }
- Ok(())
-}
-
impl<A> RadrootsRelayPublishAdapter for &A
where
A: RadrootsRelayPublishAdapter + ?Sized,
@@ -266,6 +251,114 @@ fn nostr_error_to_transport_error(error: RadrootsRelayTransportError) -> Radroot
}
}
+#[cfg(test)]
+mod contract_tests {
+ use super::nostr_error_to_transport_error;
+ use crate::RadrootsRelayTransportError;
+ use radroots_transport::RadrootsTransportError;
+
+ #[test]
+ fn relay_errors_map_to_stable_transport_contract_categories() {
+ assert_eq!(
+ nostr_error_to_transport_error(RadrootsRelayTransportError::TransportContract(
+ "contract".to_owned(),
+ )),
+ RadrootsTransportError::InvalidPayloadBytes
+ );
+
+ let target_errors = [
+ RadrootsRelayTransportError::RelayUrlParse {
+ url: "bad".to_owned(),
+ reason: "parse".to_owned(),
+ },
+ RadrootsRelayTransportError::WsRequiresLocalhostPolicy {
+ url: "ws://relay.example".to_owned(),
+ },
+ RadrootsRelayTransportError::UnsupportedRelayScheme {
+ url: "https://relay.example".to_owned(),
+ scheme: "https".to_owned(),
+ },
+ RadrootsRelayTransportError::EmptyRelayHost {
+ url: "wss://".to_owned(),
+ },
+ RadrootsRelayTransportError::RelayUrlUserinfo {
+ url: "wss://user@relay.example".to_owned(),
+ },
+ RadrootsRelayTransportError::RelayUrlQueryOrFragment {
+ url: "wss://relay.example?x=1".to_owned(),
+ },
+ RadrootsRelayTransportError::RelayUrlForbiddenDestination {
+ url: "wss://127.0.0.1".to_owned(),
+ reason: "loopback".to_owned(),
+ },
+ RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination {
+ url: "wss://relay.example".to_owned(),
+ address: "127.0.0.1".to_owned(),
+ reason: "loopback".to_owned(),
+ },
+ RadrootsRelayTransportError::EmptyTargetSet,
+ ];
+ for error in target_errors {
+ assert_eq!(
+ nostr_error_to_transport_error(error),
+ RadrootsTransportError::InvalidTargetUri
+ );
+ }
+
+ assert_eq!(
+ nostr_error_to_transport_error(RadrootsRelayTransportError::NostrEventJson(
+ "event".to_owned(),
+ )),
+ RadrootsTransportError::InvalidPayloadBytes
+ );
+ let json_error = serde_json::from_str::<serde_json::Value>("{").expect_err("invalid json");
+ assert_eq!(
+ nostr_error_to_transport_error(RadrootsRelayTransportError::Json(json_error)),
+ RadrootsTransportError::InvalidPayloadBytes
+ );
+ assert_eq!(
+ nostr_error_to_transport_error(RadrootsRelayTransportError::Transport(
+ "offline".to_owned(),
+ )),
+ RadrootsTransportError::InvalidTransportKind
+ );
+ assert_eq!(
+ nostr_error_to_transport_error(RadrootsRelayTransportError::EmptyFetchFilters),
+ RadrootsTransportError::InvalidTransportKind
+ );
+ assert_eq!(
+ nostr_error_to_transport_error(RadrootsRelayTransportError::InvalidFetchLimit {
+ field: "max_events",
+ }),
+ RadrootsTransportError::InvalidTransportKind
+ );
+
+ #[cfg(feature = "storage")]
+ {
+ assert_eq!(
+ nostr_error_to_transport_error(RadrootsRelayTransportError::EventStore(
+ radroots_event_store::RadrootsEventStoreError::MissingEvent(
+ "missing".to_owned(),
+ ),
+ )),
+ RadrootsTransportError::InvalidTransportKind
+ );
+ assert_eq!(
+ nostr_error_to_transport_error(RadrootsRelayTransportError::Outbox(
+ radroots_outbox::RadrootsOutboxError::EventNotFound(1),
+ )),
+ RadrootsTransportError::InvalidTransportKind
+ );
+ assert_eq!(
+ nostr_error_to_transport_error(
+ RadrootsRelayTransportError::MissingSignedOutboxEvent(1),
+ ),
+ RadrootsTransportError::InvalidTransportKind
+ );
+ }
+ }
+}
+
fn signed_event_from_transport_payload(
payload: &RadrootsTransportPayload,
) -> Result<RadrootsSignedEvent, RadrootsTransportError> {
diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs
@@ -17,11 +17,13 @@ use radroots_outbox::{
RadrootsOutboxOperationStatus,
};
use radroots_transport::{
- RadrootsTransport, RadrootsTransportDeliveryRequest, RadrootsTransportError,
- RadrootsTransportFetchRequest, RadrootsTransportKind, RadrootsTransportMeshScopeId,
- RadrootsTransportPayload, RadrootsTransportSatisfactionClass,
- RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget, RadrootsTransportTargetLabel,
- RadrootsTransportTargetSet,
+ RadrootsTransport, RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest,
+ RadrootsTransportDeliveryTargetStatus, RadrootsTransportError, RadrootsTransportFetchReceipt,
+ RadrootsTransportFetchRequest, RadrootsTransportFuture, RadrootsTransportImplementationState,
+ RadrootsTransportKind, RadrootsTransportMeshScopeId, RadrootsTransportOutcome,
+ RadrootsTransportOutcomeKind, RadrootsTransportPayload, RadrootsTransportSatisfactionClass,
+ RadrootsTransportSatisfactionPolicy, RadrootsTransportStatus, RadrootsTransportTarget,
+ RadrootsTransportTargetLabel, RadrootsTransportTargetReceipt, RadrootsTransportTargetSet,
};
use radroots_transport_nostr::{
RadrootsMockRelayFetchAdapter, RadrootsMockRelayPublishAdapter, RadrootsNostrTransport,
@@ -31,7 +33,8 @@ use radroots_transport_nostr::{
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, verified_signed_event_payload,
+ publish_claimed_outbox_event, publish_claimed_outbox_event_with_transport,
+ publish_signed_event, verified_signed_event_payload,
};
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
@@ -139,6 +142,62 @@ impl RadrootsRelayPublishAdapter for UnknownRelayReceiptPublishAdapter {
}
}
+#[derive(Clone)]
+struct ScriptedTransport {
+ outcomes: Vec<RadrootsTransportOutcome>,
+}
+
+impl ScriptedTransport {
+ fn new(outcomes: Vec<RadrootsTransportOutcome>) -> Self {
+ Self { outcomes }
+ }
+}
+
+impl RadrootsTransport for ScriptedTransport {
+ fn transport_kind(&self) -> RadrootsTransportKind {
+ RadrootsTransportKind::Nostr
+ }
+
+ fn status<'a>(&'a self) -> RadrootsTransportFuture<'a, RadrootsTransportStatus> {
+ Box::pin(async {
+ Ok(RadrootsTransportStatus::new(
+ RadrootsTransportKind::Nostr,
+ true,
+ RadrootsTransportImplementationState::Real,
+ true,
+ "scripted",
+ ))
+ })
+ }
+
+ fn deliver<'a>(
+ &'a self,
+ request: RadrootsTransportDeliveryRequest,
+ ) -> RadrootsTransportFuture<'a, RadrootsTransportDeliveryReceipt> {
+ Box::pin(async move {
+ let target_receipts = request
+ .target_set
+ .targets()
+ .iter()
+ .cloned()
+ .zip(self.outcomes.iter().cloned())
+ .map(|(target, outcome)| RadrootsTransportTargetReceipt::new(target, outcome))
+ .collect();
+ Ok(RadrootsTransportDeliveryReceipt {
+ request_id: request.request_id,
+ target_receipts,
+ })
+ })
+ }
+
+ fn fetch<'a>(
+ &'a self,
+ _request: RadrootsTransportFetchRequest,
+ ) -> RadrootsTransportFuture<'a, RadrootsTransportFetchReceipt> {
+ Box::pin(async { Err(RadrootsTransportError::UnsupportedOperation) })
+ }
+}
+
fn fixture_keys() -> RadrootsNostrKeys {
let secret_key =
RadrootsNostrSecretKey::from_hex(FIXTURE_ALICE_SECRET_KEY_HEX).expect("secret key");
@@ -726,7 +785,16 @@ async fn mock_publish_preserves_exact_raw_json_and_counts_outcomes() {
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 expected_status = RadrootsTransportStatus::new(
+ RadrootsTransportKind::Nostr,
+ true,
+ RadrootsTransportImplementationState::Real,
+ true,
+ "fixture ready",
+ );
+ let transport = RadrootsNostrTransport::new(&adapter).with_status(expected_status.clone());
+ assert_eq!(transport.transport_kind(), RadrootsTransportKind::Nostr);
+ assert!(transport.adapter().captured_raw_events().is_empty());
let target = nostr_target(RELAY_PRIMARY_WSS);
let request = RadrootsTransportDeliveryRequest::new(
"facade-request-1",
@@ -743,8 +811,7 @@ async fn nostr_transport_facade_delivers_signed_event_payloads() {
adapter.captured_raw_events(),
vec![signed.raw_json().to_owned()]
);
- assert!(status.capabilities.deliver);
- assert!(!status.capabilities.fetch);
+ assert_eq!(status, expected_status);
assert_eq!(receipt.request_id, "facade-request-1");
assert_eq!(receipt.target_receipts.len(), 1);
assert_eq!(receipt.target_receipts[0].target, target);
@@ -825,6 +892,154 @@ async fn nostr_transport_facade_rejects_unsupported_payloads_and_targets() {
.await
.expect_err("target rejected");
assert_eq!(target_error, RadrootsTransportError::InvalidTargetUri);
+
+ let invalid_json_error = transport
+ .deliver(RadrootsTransportDeliveryRequest::new(
+ "facade-request-invalid-json",
+ RadrootsTransportPayload::SignedEventJson {
+ event_id: signed.id_str().to_owned(),
+ raw_json: "{}".to_owned(),
+ digest: "00".repeat(32),
+ },
+ RadrootsTransportTargetSet::new(vec![nostr_target(RELAY_PRIMARY_WSS)])
+ .expect("targets"),
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ ))
+ .await
+ .expect_err("invalid event json rejected");
+ assert_eq!(
+ invalid_json_error,
+ RadrootsTransportError::InvalidPayloadBytes
+ );
+
+ let mismatched_id_error = transport
+ .deliver(RadrootsTransportDeliveryRequest::new(
+ "facade-request-mismatched-id",
+ RadrootsTransportPayload::SignedEventJson {
+ event_id: "00".repeat(32),
+ raw_json: signed.raw_json().to_owned(),
+ digest: "00".repeat(32),
+ },
+ RadrootsTransportTargetSet::new(vec![nostr_target(RELAY_PRIMARY_WSS)])
+ .expect("targets"),
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ ))
+ .await
+ .expect_err("mismatched event id rejected");
+ assert_eq!(
+ mismatched_id_error,
+ RadrootsTransportError::InvalidPayloadId
+ );
+
+ let mut tampered =
+ serde_json::from_str::<serde_json::Value>(signed.raw_json()).expect("signed event json");
+ tampered["content"] = serde_json::Value::String("tampered".to_owned());
+ let tampered_raw = serde_json::to_string(&tampered).expect("tampered event json");
+ let tampered_error = transport
+ .deliver(RadrootsTransportDeliveryRequest::new(
+ "facade-request-tampered-event",
+ RadrootsTransportPayload::SignedEventJson {
+ event_id: signed.id_str().to_owned(),
+ raw_json: tampered_raw,
+ digest: "00".repeat(32),
+ },
+ RadrootsTransportTargetSet::new(vec![nostr_target(RELAY_PRIMARY_WSS)])
+ .expect("targets"),
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ ))
+ .await
+ .expect_err("tampered event rejected");
+ assert_eq!(tampered_error, RadrootsTransportError::InvalidPayloadBytes);
+
+ let forbidden_target_error = transport
+ .deliver(RadrootsTransportDeliveryRequest::new(
+ "facade-request-forbidden-target",
+ RadrootsTransportPayload::unchecked_signed_event_json(
+ signed.id_str(),
+ signed.raw_json(),
+ )
+ .expect("payload"),
+ RadrootsTransportTargetSet::new(vec![nostr_target("wss://127.0.0.1")])
+ .expect("targets"),
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ ))
+ .await
+ .expect_err("forbidden relay rejected");
+ assert_eq!(
+ forbidden_target_error,
+ RadrootsTransportError::InvalidTargetUri
+ );
+
+ let local_receipt = transport
+ .deliver(RadrootsTransportDeliveryRequest::new(
+ "facade-request-local-relay",
+ RadrootsTransportPayload::unchecked_signed_event_json(
+ signed.id_str(),
+ signed.raw_json(),
+ )
+ .expect("payload"),
+ RadrootsTransportTargetSet::new(vec![nostr_target("ws://127.0.0.1:21002")])
+ .expect("targets"),
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ ))
+ .await
+ .expect("localhost relay accepted");
+ assert_eq!(local_receipt.target_receipts.len(), 1);
+}
+
+#[tokio::test]
+async fn nostr_transport_facade_preserves_adapter_failure_and_omission_evidence() {
+ let signed = signed_post("facade failures");
+ let payload =
+ RadrootsTransportPayload::unchecked_signed_event_json(signed.id_str(), signed.raw_json())
+ .expect("payload");
+ let targets = RadrootsTransportTargetSet::new(vec![
+ nostr_target(RELAY_PRIMARY_WSS),
+ nostr_target(RELAY_SECONDARY_WSS),
+ ])
+ .expect("targets");
+
+ let transport = RadrootsNostrTransport::new(TransportFailurePublishAdapter);
+ let failed = transport
+ .deliver(RadrootsTransportDeliveryRequest::new(
+ "facade-transport-failure",
+ payload.clone(),
+ targets.clone(),
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ ))
+ .await
+ .expect("failure receipts");
+ assert_eq!(failed.target_receipts.len(), 2);
+ assert!(failed.target_receipts.iter().all(|receipt| {
+ receipt.outcome.kind == RadrootsTransportOutcomeKind::ConnectionFailed
+ && receipt.status == RadrootsTransportDeliveryTargetStatus::FailedRetryable
+ }));
+
+ let partial = RadrootsNostrTransport::new(PartialPublishAdapter)
+ .deliver(RadrootsTransportDeliveryRequest::new(
+ "facade-partial",
+ payload.clone(),
+ targets.clone(),
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ ))
+ .await
+ .expect("partial receipts");
+ assert_eq!(partial.target_receipts.len(), 2);
+ assert_eq!(
+ partial.target_receipts[1].outcome.kind,
+ RadrootsTransportOutcomeKind::RouteUnavailable
+ );
+
+ let error = RadrootsNostrTransport::new(NostrJsonFailurePublishAdapter)
+ .deliver(RadrootsTransportDeliveryRequest::new(
+ "facade-json-failure",
+ payload,
+ targets,
+ RadrootsTransportSatisfactionPolicy::all_accepted(),
+ ))
+ .await
+ .expect_err("adapter JSON error");
+ assert_eq!(error, RadrootsTransportError::InvalidPayloadBytes);
}
#[tokio::test]
@@ -1000,7 +1215,7 @@ async fn publish_all_policy_uses_requested_target_count() {
let receipt = publish_signed_event(
&PartialPublishAdapter,
- RadrootsRelayPublishRequest::new(signed, targets, 1_080)
+ RadrootsRelayPublishRequest::new(signed.clone(), targets.clone(), 1_080)
.with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::all_accepted()),
)
.await
@@ -1010,6 +1225,16 @@ async fn publish_all_policy_uses_requested_target_count() {
assert_eq!(receipt.accepted_count, 1);
assert_eq!(receipt.quorum, 2);
assert!(!receipt.quorum_met);
+
+ let no_wait = publish_signed_event(
+ &PartialPublishAdapter,
+ RadrootsRelayPublishRequest::new(signed, targets, 1_081)
+ .with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::NoWait),
+ )
+ .await
+ .expect("no-wait publish");
+ assert_eq!(no_wait.quorum, 0);
+ assert!(no_wait.quorum_met);
}
#[test]
@@ -1747,6 +1972,424 @@ async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() {
}
#[tokio::test]
+async fn outbox_transport_facade_persists_partial_success_and_retryable_failures() {
+ let signed = signed_post("transport facade outbox");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at(),
+ signed.tags_as_vec(),
+ signed.content().to_owned(),
+ signed.pubkey_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(all_accepted_outbox_operation_input(
+ draft,
+ [RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "transport-sign", 2_000, 1_000)
+ .await
+ .expect("sign claim")
+ .expect("sign claim");
+ let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ let publish_claim = outbox
+ .claim_next_ready_event("publisher", "transport-publish", 3_000, 1_100)
+ .await
+ .expect("publish claim")
+ .expect("publish claim");
+
+ let adapter = RadrootsMockRelayPublishAdapter::new()
+ .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted())
+ .with_outcome(
+ RELAY_SECONDARY_WSS,
+ RadrootsRelayOutcome::timeout("timeout: transport facade"),
+ );
+ let transport = RadrootsNostrTransport::new(adapter);
+ let published = publish_claimed_outbox_event_with_transport(
+ &outbox,
+ &store,
+ &transport,
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_200,
+ )
+ .await
+ .expect("transport publish");
+
+ assert_eq!(published.event_id, signed.id_str());
+ assert_eq!(published.attempted_count, 2);
+ assert_eq!(published.accepted_count, 1);
+ assert_eq!(published.retryable_count, 1);
+ assert_eq!(published.terminal_count, 0);
+ assert!(!published.quorum_met);
+ assert_eq!(published.relay_receipts.len(), 2);
+ let targets = outbox
+ .delivery_targets(receipt.outbox_event_id)
+ .await
+ .expect("targets");
+ assert_eq!(
+ targets
+ .iter()
+ .find(|target| target.endpoint_uri.as_str() == RELAY_PRIMARY_WSS)
+ .expect("primary")
+ .status,
+ RadrootsOutboxDeliveryTargetStatus::Accepted
+ );
+ assert_eq!(
+ targets
+ .iter()
+ .find(|target| target.endpoint_uri.as_str() == RELAY_SECONDARY_WSS)
+ .expect("secondary")
+ .status,
+ RadrootsOutboxDeliveryTargetStatus::FailedRetryable
+ );
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable);
+ let observations = store
+ .observations_for_event(signed.id_str())
+ .await
+ .expect("observations");
+ assert_outbox_publish_observations(&observations, 1);
+}
+
+#[tokio::test]
+async fn outbox_transport_facade_persists_every_delivery_status() {
+ let signed = signed_post("transport outcome matrix");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at(),
+ signed.tags_as_vec(),
+ signed.content().to_owned(),
+ signed.pubkey_str(),
+ )
+ .expect("draft");
+ let relays = (0..14)
+ .map(|index| format!("wss://relay-{index}.example.com"))
+ .collect::<Vec<_>>();
+ let receipt = outbox
+ .enqueue_operation(all_accepted_outbox_operation_input(draft, &relays))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "matrix-sign", 2_000, 1_000)
+ .await
+ .expect("sign claim")
+ .expect("sign claim");
+ let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ let publish_claim = outbox
+ .claim_next_ready_event("publisher", "matrix-publish", 3_000, 1_100)
+ .await
+ .expect("publish claim")
+ .expect("publish claim");
+ let outcomes = [
+ RadrootsTransportOutcomeKind::Accepted,
+ RadrootsTransportOutcomeKind::DuplicateAccepted,
+ RadrootsTransportOutcomeKind::Delivered,
+ RadrootsTransportOutcomeKind::Forwarded,
+ RadrootsTransportOutcomeKind::StoredByGateway,
+ RadrootsTransportOutcomeKind::Seen,
+ RadrootsTransportOutcomeKind::DeferredUntilImplemented,
+ RadrootsTransportOutcomeKind::Rejected,
+ RadrootsTransportOutcomeKind::RouteUnavailable,
+ RadrootsTransportOutcomeKind::PayloadTooLarge,
+ RadrootsTransportOutcomeKind::PolicyDenied,
+ RadrootsTransportOutcomeKind::Timeout,
+ RadrootsTransportOutcomeKind::ConnectionFailed,
+ RadrootsTransportOutcomeKind::TransportUnavailable,
+ ]
+ .into_iter()
+ .map(RadrootsTransportOutcome::new)
+ .collect();
+ let transport = ScriptedTransport::new(outcomes);
+
+ let published = publish_claimed_outbox_event_with_transport(
+ &outbox,
+ &store,
+ &transport,
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_200,
+ )
+ .await
+ .expect("transport publish");
+
+ assert_eq!(published.event_id, signed.id_str());
+ assert_eq!(published.attempted_count, 13);
+ assert_eq!(published.accepted_count, 6);
+ assert_eq!(published.retryable_count, 3);
+ assert_eq!(published.terminal_count, 5);
+ assert!(!published.quorum_met);
+ assert_eq!(published.target_receipts.len(), 14);
+ assert_eq!(published.relay_receipts.len(), 14);
+ let targets = outbox
+ .delivery_targets(receipt.outbox_event_id)
+ .await
+ .expect("targets");
+ let expected_statuses = [
+ RadrootsOutboxDeliveryTargetStatus::Accepted,
+ RadrootsOutboxDeliveryTargetStatus::Accepted,
+ RadrootsOutboxDeliveryTargetStatus::Delivered,
+ RadrootsOutboxDeliveryTargetStatus::Forwarded,
+ RadrootsOutboxDeliveryTargetStatus::StoredByGateway,
+ RadrootsOutboxDeliveryTargetStatus::Seen,
+ RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented,
+ RadrootsOutboxDeliveryTargetStatus::FailedTerminal,
+ RadrootsOutboxDeliveryTargetStatus::FailedTerminal,
+ RadrootsOutboxDeliveryTargetStatus::FailedTerminal,
+ RadrootsOutboxDeliveryTargetStatus::SkippedPolicyDenied,
+ RadrootsOutboxDeliveryTargetStatus::FailedRetryable,
+ RadrootsOutboxDeliveryTargetStatus::FailedRetryable,
+ RadrootsOutboxDeliveryTargetStatus::FailedRetryable,
+ ];
+ assert_eq!(targets.len(), expected_statuses.len());
+ for (target, expected_status) in targets.iter().zip(expected_statuses) {
+ assert_eq!(target.status, expected_status);
+ }
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable);
+ let observations = store
+ .observations_for_event(signed.id_str())
+ .await
+ .expect("observations");
+ assert_outbox_publish_observations(&observations, 6);
+}
+
+#[tokio::test]
+async fn outbox_transport_facade_rejects_pending_receipts() {
+ let signed = signed_post("pending transport receipt");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at(),
+ signed.tags_as_vec(),
+ signed.content().to_owned(),
+ signed.pubkey_str(),
+ )
+ .expect("draft");
+ outbox
+ .enqueue_operation(all_accepted_outbox_operation_input(
+ draft,
+ [RELAY_PRIMARY_WSS],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "pending-sign", 2_000, 1_000)
+ .await
+ .expect("sign claim")
+ .expect("sign claim");
+ complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ let publish_claim = outbox
+ .claim_next_ready_event("publisher", "pending-publish", 3_000, 1_100)
+ .await
+ .expect("publish claim")
+ .expect("publish claim");
+ let transport = ScriptedTransport::new(vec![
+ RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Accepted)
+ .with_target_status(RadrootsTransportDeliveryTargetStatus::Pending),
+ ]);
+
+ let error = publish_claimed_outbox_event_with_transport(
+ &outbox,
+ &store,
+ &transport,
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_200,
+ )
+ .await
+ .expect_err("pending receipt rejected");
+ assert!(matches!(
+ error,
+ RadrootsRelayTransportError::TransportContract(_)
+ ));
+}
+
+#[tokio::test]
+async fn outbox_transport_facade_handles_empty_and_invalid_claim_plans() {
+ let signed = signed_post("transport plan edges");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at(),
+ signed.tags_as_vec(),
+ signed.content().to_owned(),
+ signed.pubkey_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(all_accepted_outbox_operation_input(
+ draft,
+ [RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "plan-sign", 2_000, 1_000)
+ .await
+ .expect("sign claim")
+ .expect("sign claim");
+ let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await;
+ let publish_claim = outbox
+ .claim_next_ready_event("publisher", "plan-publish", 3_000, 1_100)
+ .await
+ .expect("publish claim")
+ .expect("publish claim");
+ for target in &publish_claim.delivery_targets {
+ outbox
+ .mark_delivery_target_accepted(
+ publish_claim.outbox_event_id,
+ publish_claim.claim_token.as_str(),
+ target.delivery_target_id,
+ 2_150,
+ )
+ .await
+ .expect("accepted target");
+ }
+ let published = publish_claimed_outbox_event_with_transport(
+ &outbox,
+ &store,
+ &ScriptedTransport::new(Vec::new()),
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_200,
+ )
+ .await
+ .expect("already satisfied publish");
+ assert_eq!(published.event_id, signed.id_str());
+ assert_eq!(published.attempted_count, 0);
+ assert_eq!(published.accepted_count, 2);
+ assert!(published.quorum_met);
+
+ let second_signed = signed_post("invalid claimed plan");
+ let second_draft = RadrootsEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ second_signed.created_at(),
+ second_signed.tags_as_vec(),
+ second_signed.content().to_owned(),
+ second_signed.pubkey_str(),
+ )
+ .expect("draft");
+ outbox
+ .enqueue_operation(all_accepted_outbox_operation_input(
+ second_draft,
+ [RELAY_PRIMARY_WSS],
+ ))
+ .await
+ .expect("second enqueue");
+ let second_claimed = outbox
+ .claim_next_ready_event("signer", "invalid-plan-sign", 3_000, 2_200)
+ .await
+ .expect("second sign claim")
+ .expect("second sign claim");
+ complete_claimed_signing(&outbox, &second_claimed, 2_300).await;
+ let second_publish_claim = outbox
+ .claim_next_ready_event("publisher", "invalid-plan-publish", 4_000, 2_300)
+ .await
+ .expect("second publish claim")
+ .expect("second publish claim");
+ let mut invalid_claim = second_publish_claim.clone();
+ invalid_claim.active_delivery_plan_id = None;
+ let error = publish_claimed_outbox_event_with_transport(
+ &outbox,
+ &store,
+ &ScriptedTransport::new(Vec::new()),
+ &invalid_claim,
+ RadrootsOutboxPublishPolicy::new(3_500),
+ 2_400,
+ )
+ .await
+ .expect_err("missing plan rejected");
+ assert!(matches!(error, RadrootsRelayTransportError::Transport(_)));
+
+ invalid_claim.active_delivery_plan_id = Some(i64::MAX);
+ let error = publish_claimed_outbox_event_with_transport(
+ &outbox,
+ &store,
+ &ScriptedTransport::new(Vec::new()),
+ &invalid_claim,
+ RadrootsOutboxPublishPolicy::new(3_500),
+ 2_401,
+ )
+ .await
+ .expect_err("unknown plan rejected");
+ assert!(matches!(error, RadrootsRelayTransportError::Transport(_)));
+
+ let event = outbox
+ .get_event(receipt.outbox_event_id)
+ .await
+ .expect("event")
+ .expect("event");
+ assert_eq!(event.state, RadrootsOutboxEventState::Published);
+}
+
+#[tokio::test]
+async fn outbox_transport_facade_requires_signed_claims() {
+ let signed = signed_post("missing transport signature");
+ let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
+ let store = RadrootsEventStore::open_memory().await.expect("store");
+ let draft = RadrootsEventDraft::new(
+ "radroots.social.post.v1",
+ KIND_POST,
+ signed.created_at(),
+ signed.tags_as_vec(),
+ signed.content().to_owned(),
+ signed.pubkey_str(),
+ )
+ .expect("draft");
+ let receipt = outbox
+ .enqueue_operation(all_accepted_outbox_operation_input(
+ draft,
+ [RELAY_PRIMARY_WSS],
+ ))
+ .await
+ .expect("enqueue");
+ let claimed = outbox
+ .claim_next_ready_event("signer", "unsigned-transport", 2_000, 1_000)
+ .await
+ .expect("claim")
+ .expect("claim");
+
+ let error = publish_claimed_outbox_event_with_transport(
+ &outbox,
+ &store,
+ &ScriptedTransport::new(Vec::new()),
+ &claimed,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 1_100,
+ )
+ .await
+ .expect_err("missing signature rejected");
+ assert!(matches!(
+ error,
+ RadrootsRelayTransportError::MissingSignedOutboxEvent(event_id)
+ if event_id == receipt.outbox_event_id
+ ));
+}
+
+#[tokio::test]
async fn outbox_publish_fans_out_endpoint_receipts_to_scoped_logical_targets() {
let signed = signed_post("scoped duplicate relay");
let outbox = RadrootsOutbox::open_memory().await.expect("outbox");
@@ -2085,6 +2728,7 @@ async fn outbox_publish_required_targets_fan_out_same_endpoint_scoped_receipts()
.expect("draft");
let required = scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.west", "West foodshed");
let optional = scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.east", "East foodshed");
+ let terminal = scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.closed", "Closed foodshed");
let receipt = outbox
.enqueue_operation(RadrootsOutboxOperationInput::new(
"publish_post",
@@ -2097,7 +2741,12 @@ async fn outbox_publish_required_targets_fan_out_same_endpoint_scoped_receipts()
vec![required.fingerprint.clone()],
)
.expect("required target policy"),
- vec![required.clone(), optional.clone()],
+ vec![
+ required.clone(),
+ optional.clone(),
+ terminal.clone(),
+ RadrootsTransportTarget::reticulum().expect("reticulum target"),
+ ],
),
1_000,
))
@@ -2114,6 +2763,35 @@ async fn outbox_publish_required_targets_fan_out_same_endpoint_scoped_receipts()
.await
.expect("claim")
.expect("publish claim");
+ let optional_record = publish_claim
+ .delivery_targets
+ .iter()
+ .find(|target| target.endpoint_fingerprint == optional.fingerprint)
+ .expect("optional target");
+ outbox
+ .mark_delivery_target_accepted(
+ publish_claim.outbox_event_id,
+ publish_claim.claim_token.as_str(),
+ optional_record.delivery_target_id,
+ 2_150,
+ )
+ .await
+ .expect("optional target accepted");
+ let terminal_record = publish_claim
+ .delivery_targets
+ .iter()
+ .find(|target| target.endpoint_fingerprint == terminal.fingerprint)
+ .expect("terminal target");
+ outbox
+ .mark_delivery_target_failed_terminal(
+ publish_claim.outbox_event_id,
+ publish_claim.claim_token.as_str(),
+ terminal_record.delivery_target_id,
+ "terminal target",
+ 2_151,
+ )
+ .await
+ .expect("terminal target completed");
let adapter = RadrootsMockRelayPublishAdapter::new()
.with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted());
@@ -2122,7 +2800,7 @@ async fn outbox_publish_required_targets_fan_out_same_endpoint_scoped_receipts()
&store,
&adapter,
&publish_claim,
- RadrootsOutboxPublishPolicy::new(2_500),
+ RadrootsOutboxPublishPolicy::new(2_500).republish_accepted_relays(true),
2_200,
)
.await
@@ -2152,10 +2830,24 @@ async fn outbox_publish_required_targets_fan_out_same_endpoint_scoped_receipts()
.delivery_targets(receipt.outbox_event_id)
.await
.expect("targets");
- assert_eq!(targets.len(), 2);
- assert!(targets.iter().all(|target| {
- target.endpoint_uri.as_str() == RELAY_PRIMARY_WSS
- && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted
+ assert_eq!(targets.len(), 4);
+ assert!(
+ targets
+ .iter()
+ .filter(|target| { target.transport_kind == RadrootsTransportKind::Nostr })
+ .filter(|target| target.endpoint_fingerprint != terminal.fingerprint)
+ .all(|target| {
+ target.endpoint_uri.as_str() == RELAY_PRIMARY_WSS
+ && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted
+ })
+ );
+ assert!(targets.iter().any(|target| {
+ target.endpoint_fingerprint == terminal.fingerprint
+ && target.status == RadrootsOutboxDeliveryTargetStatus::FailedTerminal
+ }));
+ assert!(targets.iter().any(|target| {
+ target.transport_kind == RadrootsTransportKind::Reticulum
+ && target.status == RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented
}));
}
@@ -2901,6 +3593,21 @@ async fn outbox_publish_rejects_invalid_relay_target_uri_before_adapter_publish(
RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. }
));
assert!(adapter.captured_raw_events().is_empty());
+
+ let transport_error = publish_claimed_outbox_event_with_transport(
+ &outbox,
+ &store,
+ &ScriptedTransport::new(Vec::new()),
+ &publish_claim,
+ RadrootsOutboxPublishPolicy::new(2_500),
+ 2_201,
+ )
+ .await
+ .expect_err("invalid transport relay target");
+ assert!(matches!(
+ transport_error,
+ RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. }
+ ));
let event = outbox
.get_event(receipt.outbox_event_id)
.await