lib

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

commit 8b8ee1758467a3d0cae14c547cd5b9d7cc9f45a6
parent a3aafcec59e95be41e25a1f585ecd6d30c60fc9b
Author: triesap <tyson@radroots.org>
Date:   Fri, 10 Jul 2026 03:16:22 +0000

transport: make Nostr outbox receipts target aware

- replace the direct outbox publish receipt with logical delivery-target counts and target receipts
- fan out endpoint relay outcomes to every matching scoped Nostr delivery target before lifecycle evaluation
- keep raw relay publish receipts endpoint-scoped as separate physical evidence
- add scoped duplicate-endpoint outbox coverage with persisted metadata and lifecycle assertions

Diffstat:
Mcrates/transport_nostr/src/lib.rs | 3++-
Mcrates/transport_nostr/src/outbox.rs | 194++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------------
Mcrates/transport_nostr/tests/transport.rs | 195+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
3 files changed, 301 insertions(+), 91 deletions(-)

diff --git a/crates/transport_nostr/src/lib.rs b/crates/transport_nostr/src/lib.rs @@ -24,7 +24,8 @@ pub use fetch::{ }; #[cfg(feature = "storage")] pub use outbox::{ - RadrootsOutboxPublishPolicy, RadrootsOutboxPublishReceipt, publish_claimed_outbox_event, + RadrootsOutboxPublishPolicy, RadrootsOutboxPublishReceipt, RadrootsOutboxPublishTargetReceipt, + publish_claimed_outbox_event, }; pub use outcome::{RadrootsRelayOutcome, RadrootsRelayOutcomeKind}; #[cfg(feature = "client")] diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs @@ -49,7 +49,25 @@ impl RadrootsOutboxPublishPolicy { #[derive(Clone, Debug, PartialEq, Eq)] pub struct RadrootsOutboxPublishReceipt { pub local_ingest: RadrootsOutboxEventStoreIngestReceipt, - pub publish: RadrootsRelayPublishReceipt, + pub event_id: String, + pub attempted_count: usize, + pub accepted_count: usize, + pub retryable_count: usize, + pub terminal_count: usize, + pub quorum: usize, + pub quorum_met: bool, + pub target_receipts: Vec<RadrootsOutboxPublishTargetReceipt>, + pub relay_receipts: Vec<RadrootsRelayPublishRelayReceipt>, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RadrootsOutboxPublishTargetReceipt { + pub delivery_target_id: i64, + pub endpoint_uri: String, + pub target_scope: Option<String>, + pub target_label: Option<String>, + pub attempted: bool, + pub outcome: RadrootsRelayOutcome, } pub async fn publish_claimed_outbox_event<A>( @@ -86,7 +104,8 @@ where now_ms, ) .await?; - let publish = RadrootsRelayPublishReceipt { + return Ok(RadrootsOutboxPublishReceipt { + local_ingest, event_id: signed_event.id, attempted_count: 0, accepted_count: publishable.accepted_count, @@ -94,11 +113,8 @@ where terminal_count: 0, quorum: publishable.satisfaction_required_count, quorum_met: publishable.accepted_count >= publishable.satisfaction_required_count, - relays: Vec::new(), - }; - return Ok(RadrootsOutboxPublishReceipt { - local_ingest, - publish, + target_receipts: Vec::new(), + relay_receipts: Vec::new(), }); } let targets = RadrootsRelayTargetSet::new( @@ -111,7 +127,7 @@ where let target_strings = targets.relay_strings(); let satisfaction_policy = satisfaction_policy_for_required_accept_count( publishable.required_accept_count, - publishable.relays.len(), + targets.len(), )?; let active_delivery_plan_id = publishable.active_delivery_plan_id; let request = RadrootsRelayPublishRequest::new(signed_event.clone(), targets, now_ms) @@ -132,55 +148,64 @@ where ), Err(error) => return Err(error), }; + let target_receipts = target_receipts_from_relay_receipts(&publishable, &publish.relays); - for relay in &publish.relays { - if let Some(target) = publishable.target_for_relay(relay.relay_url.as_str()) { - if relay.outcome.counts_toward_quorum() { - outbox - .mark_delivery_target_accepted( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - target.delivery_target_id, - now_ms, - ) - .await?; - ingest_publish_observation( - event_store, - &signed_event, - relay.relay_url.as_str(), - relay.outcome.message.as_deref(), + 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?; - } else if relay.outcome.is_retryable() { - outbox - .mark_delivery_target_failed_retryable( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - target.delivery_target_id, - relay - .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.delivery_target_id, - relay - .outcome - .message - .as_deref() - .unwrap_or("relay publish terminal"), - now_ms, - ) - .await?; - } + } + } + + for relay in &publish.relays { + if relay.outcome.counts_toward_quorum() + && publishable + .targets_for_relay(relay.relay_url.as_str()) + .next() + .is_some() + { + ingest_publish_observation( + event_store, + &signed_event, + relay.relay_url.as_str(), + relay.outcome.message.as_deref(), + now_ms, + ) + .await?; } } @@ -197,7 +222,31 @@ where Ok(RadrootsOutboxPublishReceipt { local_ingest, - publish, + event_id: publish.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: target_receipts + .iter() + .filter(|receipt| receipt.outcome.counts_toward_quorum()) + .count() + >= publishable.required_accept_count, + target_receipts, + relay_receipts: publish.relays, }) } @@ -237,16 +286,41 @@ struct PublishableRelays { } impl PublishableRelays { - fn target_for_relay(&self, relay_url: &str) -> Option<&PublishableRelay> { + fn targets_for_relay<'a>( + &'a self, + relay_url: &'a str, + ) -> impl Iterator<Item = &'a PublishableRelay> + 'a { self.relays .iter() - .find(|target| target.relay_url == relay_url) + .filter(move |target| target.relay_url == relay_url) } } struct PublishableRelay { delivery_target_id: i64, relay_url: String, + target_scope: Option<String>, + target_label: Option<String>, +} + +fn target_receipts_from_relay_receipts( + publishable: &PublishableRelays, + relay_receipts: &[RadrootsRelayPublishRelayReceipt], +) -> Vec<RadrootsOutboxPublishTargetReceipt> { + let mut target_receipts = Vec::new(); + for relay_receipt in relay_receipts { + for target in publishable.targets_for_relay(relay_receipt.relay_url.as_str()) { + target_receipts.push(RadrootsOutboxPublishTargetReceipt { + delivery_target_id: target.delivery_target_id, + endpoint_uri: target.relay_url.clone(), + target_scope: target.target_scope.clone(), + target_label: target.target_label.clone(), + attempted: relay_receipt.attempted, + outcome: relay_receipt.outcome.clone(), + }); + } + } + target_receipts } async fn publishable_relays( @@ -323,6 +397,14 @@ async fn publishable_relays( relays.push(PublishableRelay { delivery_target_id: target.delivery_target_id, relay_url: target.endpoint_uri.as_str().to_owned(), + target_scope: target + .target_scope + .as_ref() + .map(|scope| scope.as_str().to_owned()), + target_label: target + .target_label + .as_ref() + .map(|label| label.as_str().to_owned()), }); } } diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -17,8 +17,8 @@ use radroots_outbox::{ RadrootsOutboxOperationStatus, }; use radroots_transport::{ - RadrootsTransportKind, RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy, - RadrootsTransportTarget, + RadrootsTransportKind, RadrootsTransportMeshScopeId, RadrootsTransportSatisfactionClass, + RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget, RadrootsTransportTargetLabel, }; use radroots_transport_nostr::{ RadrootsMockRelayFetchAdapter, RadrootsMockRelayPublishAdapter, RadrootsOutboxPublishPolicy, @@ -206,6 +206,16 @@ fn nostr_target(relay_url: &str) -> RadrootsTransportTarget { RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, relay_url).expect("nostr target") } +fn scoped_nostr_target(relay_url: &str, scope: &str, label: &str) -> RadrootsTransportTarget { + RadrootsTransportTarget::new_with_metadata( + RadrootsTransportKind::Nostr, + relay_url, + Some(RadrootsTransportMeshScopeId::parse(scope).expect("target scope")), + Some(RadrootsTransportTargetLabel::parse(label).expect("target label")), + ) + .expect("scoped nostr target") +} + fn outbox_operation_input<I, S>( draft: RadrootsFrozenEventDraft, relays: I, @@ -1456,9 +1466,9 @@ async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() { .await .expect("publish"); - assert_eq!(first.publish.attempted_count, 3); - assert_eq!(first.publish.accepted_count, 2); - assert!(!first.publish.quorum_met); + assert_eq!(first.attempted_count, 3); + assert_eq!(first.accepted_count, 2); + assert!(!first.quorum_met); let event = outbox .get_event(receipt.outbox_event_id) .await @@ -1514,7 +1524,7 @@ async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() { .expect("retry publish"); assert_eq!(second.local_ingest.event_id, signed.id); - assert_eq!(second.publish.attempted_count, 1); + assert_eq!(second.attempted_count, 1); assert_eq!(retry_adapter.captured_raw_events().len(), 1); let event = outbox @@ -1538,6 +1548,120 @@ async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() { } #[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"); + let store = RadrootsEventStore::open_memory().await.expect("store"); + let draft = RadrootsFrozenEventDraft::new( + "radroots.social.post.v1", + KIND_POST, + signed.created_at, + signed.tags.clone(), + signed.content.clone(), + signed.pubkey.as_str(), + ) + .expect("draft"); + let receipt = outbox + .enqueue_operation(RadrootsOutboxOperationInput::new( + "publish_post", + draft, + RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 2, + RadrootsTransportSatisfactionPolicy::all_accepted(), + vec![ + scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.west", "West foodshed"), + scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.east", "East foodshed"), + ], + ), + 1_000, + )) + .await + .expect("enqueue"); + let claimed = outbox + .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) + .await + .expect("claim") + .expect("claim"); + let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; + let publish_claim = outbox + .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) + .await + .expect("claim") + .expect("publish claim"); + let adapter = RadrootsMockRelayPublishAdapter::new() + .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()); + + let published = publish_claimed_outbox_event( + &outbox, + &store, + &adapter, + &publish_claim, + RadrootsOutboxPublishPolicy::new(2_500), + 2_200, + ) + .await + .expect("publish"); + + assert_eq!(published.local_ingest.event_id, signed.id); + assert_eq!(published.event_id, signed.id); + assert_eq!(published.attempted_count, 2); + assert_eq!(published.accepted_count, 2); + assert_eq!(published.retryable_count, 0); + assert_eq!(published.terminal_count, 0); + assert_eq!(published.quorum, 2); + assert!(published.quorum_met); + assert_eq!(published.relay_receipts.len(), 1); + assert_eq!(published.relay_receipts[0].relay_url, RELAY_PRIMARY_WSS); + assert_eq!(published.target_receipts.len(), 2); + assert!( + published + .target_receipts + .iter() + .all(|target| target.endpoint_uri == RELAY_PRIMARY_WSS && target.attempted) + ); + assert!(published.target_receipts.iter().any(|target| { + target.target_scope.as_deref() == Some("foodshed.west") + && target.target_label.as_deref() == Some("West foodshed") + })); + assert!(published.target_receipts.iter().any(|target| { + target.target_scope.as_deref() == Some("foodshed.east") + && target.target_label.as_deref() == Some("East foodshed") + })); + assert_eq!(adapter.captured_raw_events().len(), 1); + + let event = outbox + .get_event(receipt.outbox_event_id) + .await + .expect("event") + .expect("event"); + assert_eq!(event.state, RadrootsOutboxEventState::Published); + let targets = outbox + .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 + && target.attempt_count == 1 + })); + assert!(targets.iter().any(|target| { + target.target_scope.as_ref().map(|scope| scope.as_str()) == Some("foodshed.west") + && target.target_label.as_ref().map(|label| label.as_str()) == Some("West foodshed") + })); + assert!(targets.iter().any(|target| { + target.target_scope.as_ref().map(|scope| scope.as_str()) == Some("foodshed.east") + && target.target_label.as_ref().map(|label| label.as_str()) == Some("East foodshed") + })); + let observations = store + .observations_for_event(signed.id.as_str()) + .await + .expect("observations"); + assert_outbox_publish_observations(&observations, 1); +} + +#[tokio::test] async fn outbox_transport_publish_failure_releases_retryable_claim() { let signed = signed_post("adapter transport failure"); let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); @@ -1581,15 +1705,14 @@ async fn outbox_transport_publish_failure_releases_retryable_claim() { .await .expect("publish"); - assert_eq!(published.publish.attempted_count, 2); - assert_eq!(published.publish.accepted_count, 0); - assert_eq!(published.publish.retryable_count, 2); - assert_eq!(published.publish.terminal_count, 0); - assert!(!published.publish.quorum_met); + assert_eq!(published.attempted_count, 2); + assert_eq!(published.accepted_count, 0); + assert_eq!(published.retryable_count, 2); + assert_eq!(published.terminal_count, 0); + assert!(!published.quorum_met); assert!( published - .publish - .relays + .relay_receipts .iter() .all(|relay| relay.outcome.kind == RadrootsRelayOutcomeKind::ConnectionFailed) ); @@ -1694,12 +1817,13 @@ async fn outbox_publish_marks_published_without_adapter_when_all_relays_already_ .expect("publish"); assert_eq!(published.local_ingest.event_id, signed.id); - assert_eq!(published.publish.event_id, signed.id); - assert_eq!(published.publish.attempted_count, 0); - assert_eq!(published.publish.accepted_count, 2); - assert_eq!(published.publish.quorum, 2); - assert!(published.publish.quorum_met); - assert!(published.publish.relays.is_empty()); + assert_eq!(published.event_id, signed.id); + assert_eq!(published.attempted_count, 0); + assert_eq!(published.accepted_count, 2); + assert_eq!(published.quorum, 2); + assert!(published.quorum_met); + assert!(published.target_receipts.is_empty()); + assert!(published.relay_receipts.is_empty()); assert!(adapter.captured_raw_events().is_empty()); let event = outbox @@ -1761,8 +1885,11 @@ async fn outbox_publish_ignores_unknown_adapter_receipts() { .await .expect("publish"); - assert_eq!(published.publish.attempted_count, 2); - assert!(published.publish.quorum_met); + assert_eq!(published.attempted_count, 1); + assert_eq!(published.accepted_count, 1); + assert_eq!(published.target_receipts.len(), 1); + assert_eq!(published.relay_receipts.len(), 2); + assert!(published.quorum_met); let event = outbox .get_event(receipt.outbox_event_id) .await @@ -1839,7 +1966,7 @@ async fn outbox_publish_skips_non_nostr_targets() { .await .expect("publish"); - assert_eq!(published.publish.attempted_count, 1); + assert_eq!(published.attempted_count, 1); assert_eq!(adapter.captured_raw_events().len(), 1); let event = outbox .get_event(receipt.outbox_event_id) @@ -1917,10 +2044,10 @@ async fn outbox_publish_marks_published_when_delivery_plan_satisfaction_is_met_w .await .expect("publish"); - assert_eq!(published.publish.quorum, 2); - assert_eq!(published.publish.accepted_count, 2); - assert_eq!(published.publish.terminal_count, 1); - assert!(published.publish.quorum_met); + assert_eq!(published.quorum, 2); + assert_eq!(published.accepted_count, 2); + assert_eq!(published.terminal_count, 1); + assert!(published.quorum_met); let event = outbox .get_event(receipt.outbox_event_id) @@ -2023,10 +2150,10 @@ async fn outbox_publish_republishes_accepted_relays_when_policy_requests_it() { .expect("publish"); assert_eq!(published.local_ingest.event_id, signed.id); - assert_eq!(published.publish.attempted_count, 2); - assert_eq!(published.publish.accepted_count, 2); - assert_eq!(published.publish.quorum, 1); - assert!(published.publish.quorum_met); + assert_eq!(published.attempted_count, 2); + assert_eq!(published.accepted_count, 2); + assert_eq!(published.quorum, 1); + assert!(published.quorum_met); assert_eq!(adapter.captured_raw_events().len(), 1); let event = outbox @@ -2113,10 +2240,10 @@ async fn outbox_publish_republish_policy_keeps_terminal_targets_excluded() { .await .expect("publish"); - assert_eq!(published.publish.attempted_count, 1); - assert_eq!(published.publish.accepted_count, 1); - assert_eq!(published.publish.quorum, 1); - assert!(published.publish.quorum_met); + assert_eq!(published.attempted_count, 1); + assert_eq!(published.accepted_count, 1); + assert_eq!(published.quorum, 1); + assert!(published.quorum_met); assert_eq!(adapter.captured_raw_events().len(), 1); let event = outbox .get_event(receipt.outbox_event_id)