lib

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

commit fafbf7d57688812ec9903737752cbd254edc865b
parent e5f6aac45c3357d6457eb8c1822bda90ad4123c3
Author: triesap <tyson@radroots.org>
Date:   Fri, 10 Jul 2026 08:36:48 +0000

outbox: preserve required target idempotency

- canonicalize required target policies before outbox plan hashing
- evaluate required target delivery state from exact persisted fingerprints
- keep direct Nostr outbox publish satisfaction scoped to required targets
- cover optional target success, retry, and scoped same-relay fan-out

Diffstat:
Mcrates/outbox/src/store.rs | 283++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
Mcrates/transport_nostr/src/outbox.rs | 107++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
Mcrates/transport_nostr/tests/transport.rs | 299+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
3 files changed, 657 insertions(+), 32 deletions(-)

diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs @@ -1346,17 +1346,14 @@ fn prepare_delivery_plan( }) }) .collect::<Result<Vec<_>, RadrootsOutboxError>>()?; - let required_success_count = - input - .satisfaction_policy - .required_target_count(satisfaction_target_count( - &prepared_targets, - total_target_count, - ))? as i64; - let initial_status = - initial_delivery_plan_status(&input.satisfaction_policy, &prepared_targets); + let satisfaction_policy = canonical_satisfaction_policy(&input.satisfaction_policy)?; + validate_required_targets_belong_to_plan(&satisfaction_policy, &prepared_targets)?; + let required_success_count = satisfaction_policy.required_target_count( + satisfaction_target_count(&prepared_targets, total_target_count), + )? as i64; + let initial_status = initial_delivery_plan_status(&satisfaction_policy, &prepared_targets); let target_policy_fingerprint = target_policy_fingerprint( - &input.satisfaction_policy, + &satisfaction_policy, input.reticulum_preview_behavior, &prepared_targets, ); @@ -1370,7 +1367,7 @@ fn prepare_delivery_plan( transport_profile_id: input.transport_profile_id.trim().to_owned(), target_policy_fingerprint, target_policy_version: input.target_policy_version, - satisfaction_policy: input.satisfaction_policy.clone(), + satisfaction_policy, required_success_count, delivery_plan_idempotency_digest, initial_status, @@ -1378,6 +1375,40 @@ fn prepare_delivery_plan( }) } +fn canonical_satisfaction_policy( + policy: &RadrootsTransportSatisfactionPolicy, +) -> Result<RadrootsTransportSatisfactionPolicy, RadrootsOutboxError> { + match policy { + RadrootsTransportSatisfactionPolicy::RequiredTargets { class, targets } => Ok( + RadrootsTransportSatisfactionPolicy::required_targets(*class, targets.clone())?, + ), + RadrootsTransportSatisfactionPolicy::NoWait + | RadrootsTransportSatisfactionPolicy::Any { .. } + | RadrootsTransportSatisfactionPolicy::All { .. } + | RadrootsTransportSatisfactionPolicy::Quorum { .. } => Ok(policy.clone()), + } +} + +fn validate_required_targets_belong_to_plan( + policy: &RadrootsTransportSatisfactionPolicy, + prepared_targets: &[PreparedDeliveryTarget], +) -> Result<(), RadrootsOutboxError> { + let RadrootsTransportSatisfactionPolicy::RequiredTargets { targets, .. } = policy else { + return Ok(()); + }; + if targets.iter().all(|required| { + prepared_targets + .iter() + .any(|target| target.target.fingerprint == *required) + }) { + Ok(()) + } else { + Err(RadrootsOutboxError::Transport( + RadrootsTransportError::InvalidSatisfactionPolicy, + )) + } +} + fn satisfaction_target_count( prepared_targets: &[PreparedDeliveryTarget], total_target_count: usize, @@ -1984,21 +2015,22 @@ async fn evaluate_delivery_plans( for plan in plans { let targets = delivery_targets_for_plan_tx(tx, plan.delivery_plan_id).await?; let satisfied_count = outbox_satisfied_target_count(&plan.satisfaction_policy, &targets); - let ready_count = targets + let status_targets = delivery_plan_status_targets(&plan.satisfaction_policy, &targets); + let ready_count = status_targets .iter() .filter(|target| target.status.is_ready_for_attempt()) .count(); - let deferred_count = targets + let deferred_count = status_targets .iter() .filter(|target| target.status.is_deferred_preview()) .count(); - let preview_unavailable_count = targets + let preview_unavailable_count = status_targets .iter() .filter(|target| { target.status == RadrootsOutboxDeliveryTargetStatus::PreviewUnavailable }) .count(); - let terminal_failure_count = targets + let terminal_failure_count = status_targets .iter() .filter(|target| target.status.is_terminal_failure()) .count(); @@ -2050,6 +2082,29 @@ async fn evaluate_delivery_plans( }) } +fn delivery_plan_status_targets<'a>( + policy: &RadrootsTransportSatisfactionPolicy, + targets: &'a [RadrootsOutboxDeliveryTargetRecord], +) -> Vec<&'a RadrootsOutboxDeliveryTargetRecord> { + match policy { + RadrootsTransportSatisfactionPolicy::RequiredTargets { + targets: required_targets, + .. + } => targets + .iter() + .filter(|target| { + required_targets + .iter() + .any(|required| target.endpoint_fingerprint == *required) + }) + .collect(), + RadrootsTransportSatisfactionPolicy::NoWait + | RadrootsTransportSatisfactionPolicy::Any { .. } + | RadrootsTransportSatisfactionPolicy::All { .. } + | RadrootsTransportSatisfactionPolicy::Quorum { .. } => targets.iter().collect(), + } +} + fn outbox_satisfied_target_count( policy: &RadrootsTransportSatisfactionPolicy, targets: &[RadrootsOutboxDeliveryTargetRecord], @@ -2392,11 +2447,12 @@ fn satisfaction_policy_storage_value(policy: &RadrootsTransportSatisfactionPolic ) } RadrootsTransportSatisfactionPolicy::RequiredTargets { class, targets } => { - let fingerprints = targets + let mut fingerprints = targets .iter() .map(RadrootsTransportTargetFingerprint::as_str) - .collect::<Vec<_>>() - .join(","); + .collect::<Vec<_>>(); + fingerprints.sort(); + let fingerprints = fingerprints.join(","); format!( "required_{}:{fingerprints}", satisfaction_class_storage_value(*class) @@ -2571,13 +2627,11 @@ mod tests { .expect("required targets policy"); let stored = satisfaction_policy_storage_value(&policy); + let mut fingerprints = vec![first.fingerprint.as_str(), second.fingerprint.as_str()]; + fingerprints.sort(); assert_eq!( stored, - format!( - "required_delivered:{},{}", - first.fingerprint.as_str(), - second.fingerprint.as_str() - ) + format!("required_delivered:{},{}", fingerprints[0], fingerprints[1]) ); assert_eq!( parse_satisfaction_policy(stored.as_str(), 2).expect("parse required targets"), @@ -2589,6 +2643,60 @@ mod tests { )); } + #[test] + fn required_target_policy_idempotency_is_order_independent() { + let first = nostr_target("wss://required-one.example"); + let second = nostr_target("wss://required-two.example"); + let first_policy = RadrootsTransportSatisfactionPolicy::RequiredTargets { + class: RadrootsTransportSatisfactionClass::Accepted, + targets: vec![second.fingerprint.clone(), first.fingerprint.clone()], + }; + let second_policy = RadrootsTransportSatisfactionPolicy::RequiredTargets { + class: RadrootsTransportSatisfactionClass::Accepted, + targets: vec![first.fingerprint.clone(), second.fingerprint.clone()], + }; + let targets = vec![first, second]; + + let first_prepared = prepare_delivery_plan( + "event-required-target-order", + &RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 1, + first_policy, + targets.clone(), + ), + ) + .expect("first delivery plan"); + let second_prepared = prepare_delivery_plan( + "event-required-target-order", + &RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 1, + second_policy, + targets, + ), + ) + .expect("second delivery plan"); + + assert_eq!( + first_prepared.satisfaction_policy, + second_prepared.satisfaction_policy + ); + assert_eq!(first_prepared.required_success_count, 2); + assert_eq!( + first_prepared.target_policy_fingerprint, + second_prepared.target_policy_fingerprint + ); + assert_eq!( + first_prepared.delivery_plan_idempotency_digest, + second_prepared.delivery_plan_idempotency_digest + ); + assert_eq!( + satisfaction_policy_storage_value(&first_prepared.satisfaction_policy), + satisfaction_policy_storage_value(&second_prepared.satisfaction_policy) + ); + } + fn malformed_reticulum_target(uri: &str) -> RadrootsTransportTarget { let endpoint_uri = RadrootsTransportTargetUri::parse(uri).expect("target uri"); let endpoint_fingerprint = RadrootsTransportTargetFingerprint::from_target( @@ -3362,6 +3470,135 @@ mod tests { } #[tokio::test] + async fn required_target_plan_evaluation_ignores_optional_retryable_failure() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let draft = post_draft( + FIXTURE_ALICE_PUBLIC_KEY_HEX, + "required target optional failure", + ); + let signed_event = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("signed event"); + let optional = nostr_target(NOSTR_PRIMARY_WSS); + let required = nostr_target(NOSTR_SECONDARY_WSS); + let receipt = outbox + .enqueue_signed_operation(RadrootsOutboxSignedOperationInput::new( + "publish_post", + draft, + signed_event, + RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 1, + RadrootsTransportSatisfactionPolicy::required_targets( + RadrootsTransportSatisfactionClass::Accepted, + vec![required.fingerprint.clone()], + ) + .expect("required target policy"), + vec![optional.clone(), required.clone()], + ), + true, + 1_007, + 1_000, + )) + .await + .expect("enqueue"); + let claimed = outbox + .claim_next_ready_signed_event("publisher", "claim-a", 2_000, 1_000) + .await + .expect("claim") + .expect("claimed"); + let optional_claim = claimed + .delivery_targets + .iter() + .find(|target| target.endpoint_fingerprint == optional.fingerprint) + .expect("optional target"); + let required_claim = claimed + .delivery_targets + .iter() + .find(|target| target.endpoint_fingerprint == required.fingerprint) + .expect("required target"); + outbox + .mark_delivery_target_failed_retryable( + receipt.outbox_event_id, + "claim-a", + optional_claim.delivery_target_id, + "optional relay timeout", + 1_100, + ) + .await + .expect("optional retryable"); + outbox + .mark_delivery_target_accepted( + receipt.outbox_event_id, + "claim-a", + required_claim.delivery_target_id, + 1_110, + ) + .await + .expect("required accepted"); + + let state = outbox + .complete_publish_attempt( + receipt.outbox_event_id, + "claim-a", + "retryable", + "terminal", + 2_500, + 1_200, + ) + .await + .expect("complete"); + + assert_eq!(state, RadrootsOutboxEventState::Published); + let targets = outbox + .delivery_targets(receipt.outbox_event_id) + .await + .expect("targets"); + assert!(targets.iter().any(|target| { + target.endpoint_fingerprint == optional.fingerprint + && target.status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable + })); + assert!(targets.iter().any(|target| { + target.endpoint_fingerprint == required.fingerprint + && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted + })); + } + + #[tokio::test] + async fn enqueue_rejects_required_targets_outside_delivery_plan_before_persistence() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "stale required target"); + let missing = nostr_target(NOSTR_SECONDARY_WSS); + + let err = outbox + .enqueue_operation(RadrootsOutboxOperationInput::new( + "publish_post", + draft, + RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 1, + RadrootsTransportSatisfactionPolicy::required_targets( + RadrootsTransportSatisfactionClass::Accepted, + vec![missing.fingerprint], + ) + .expect("required target policy"), + vec![nostr_target(NOSTR_PRIMARY_WSS)], + ), + 1_000, + )) + .await + .expect_err("missing required target"); + + assert!(matches!( + err, + RadrootsOutboxError::Transport(RadrootsTransportError::InvalidSatisfactionPolicy) + )); + assert_eq!(table_count(&outbox, "outbox_operations").await, 0); + assert_eq!(table_count(&outbox, "outbox_event").await, 0); + assert_eq!(table_count(&outbox, "outbox_delivery_plan").await, 0); + assert_eq!(table_count(&outbox, "outbox_delivery_target").await, 0); + } + + #[tokio::test] async fn terminal_delivery_target_updates_are_idempotent_and_conflicting_rewrites_fail() { for terminal_status in [ RadrootsOutboxDeliveryTargetStatus::Accepted, diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs @@ -17,6 +17,7 @@ use radroots_outbox::{ }; use radroots_transport::{ RadrootsTransportKind, RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy, + RadrootsTransportTargetFingerprint, }; #[derive(Clone, Debug, PartialEq, Eq)] @@ -64,6 +65,7 @@ pub struct RadrootsOutboxPublishReceipt { pub struct RadrootsOutboxPublishTargetReceipt { pub delivery_target_id: i64, pub endpoint_uri: String, + pub endpoint_fingerprint: RadrootsTransportTargetFingerprint, pub target_scope: Option<String>, pub target_label: Option<String>, pub attempted: bool, @@ -128,6 +130,7 @@ where let satisfaction_policy = satisfaction_policy_for_required_accept_count( publishable.required_accept_count, targets.len(), + publishable.required_targets.is_some(), )?; let active_delivery_plan_id = publishable.active_delivery_plan_id; let request = RadrootsRelayPublishRequest::new(signed_event.clone(), targets, now_ms) @@ -240,11 +243,8 @@ where .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, + quorum_met: publishable.satisfied_count_after_receipts(&target_receipts) + >= publishable.satisfaction_required_count, target_receipts, relay_receipts: publish.relays, }) @@ -283,6 +283,7 @@ struct PublishableRelays { accepted_count: usize, satisfaction_required_count: usize, required_accept_count: usize, + required_targets: Option<Vec<RadrootsTransportTargetFingerprint>>, } impl PublishableRelays { @@ -294,11 +295,30 @@ impl PublishableRelays { .iter() .filter(move |target| target.relay_url == relay_url) } + + fn satisfied_count_after_receipts( + &self, + target_receipts: &[RadrootsOutboxPublishTargetReceipt], + ) -> usize { + self.accepted_count + + target_receipts + .iter() + .filter(|receipt| { + receipt.outcome.counts_toward_quorum() + && self.required_targets.as_ref().is_none_or(|required| { + required + .iter() + .any(|fingerprint| receipt.endpoint_fingerprint == *fingerprint) + }) + }) + .count() + } } struct PublishableRelay { delivery_target_id: i64, relay_url: String, + endpoint_fingerprint: RadrootsTransportTargetFingerprint, target_scope: Option<String>, target_label: Option<String>, } @@ -313,6 +333,7 @@ fn target_receipts_from_relay_receipts( target_receipts.push(RadrootsOutboxPublishTargetReceipt { delivery_target_id: target.delivery_target_id, endpoint_uri: target.relay_url.clone(), + endpoint_fingerprint: target.endpoint_fingerprint.clone(), target_scope: target.target_scope.clone(), target_label: target.target_label.clone(), attempted: relay_receipt.attempted, @@ -346,6 +367,15 @@ async fn publishable_relays( )) })?; let satisfaction_required_count = plan.required_success_count as usize; + let required_targets = match &plan.satisfaction_policy { + RadrootsTransportSatisfactionPolicy::RequiredTargets { targets, .. } => { + Some(targets.clone()) + } + RadrootsTransportSatisfactionPolicy::NoWait + | RadrootsTransportSatisfactionPolicy::Any { .. } + | RadrootsTransportSatisfactionPolicy::All { .. } + | RadrootsTransportSatisfactionPolicy::Quorum { .. } => None, + }; let active_targets = targets .iter() .filter(|target| target.delivery_plan_id == active_delivery_plan_id) @@ -368,7 +398,11 @@ async fn publishable_relays( active_targets .iter() .filter(|target| { - target + required_targets.as_ref().is_none_or(|required| { + required + .iter() + .any(|fingerprint| target.endpoint_fingerprint == *fingerprint) + }) && target .status .counts_as_transport_satisfaction(satisfaction_class) }) @@ -379,17 +413,26 @@ async fn publishable_relays( (plan.required_success_count as usize).saturating_sub(satisfied_count); let mut relays = Vec::new(); let mut accepted_count = 0usize; - for target in active_targets { + for target in &active_targets { if !is_nostr_target(target) { continue; } + let required_for_satisfaction = required_targets.as_ref().is_some_and(|required| { + required + .iter() + .any(|fingerprint| target.endpoint_fingerprint == *fingerprint) + }); if target .status .counts_as_transport_satisfaction(RadrootsTransportSatisfactionClass::Accepted) + && (required_targets.is_none() || required_for_satisfaction) { accepted_count += 1; } + let can_contribute_to_satisfaction = + required_targets.is_none() || required_for_satisfaction; if required_accept_count > 0 + && can_contribute_to_satisfaction && (target.status.is_ready_for_attempt() || (republish_accepted_relays && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted)) @@ -397,6 +440,41 @@ async fn publishable_relays( relays.push(PublishableRelay { delivery_target_id: target.delivery_target_id, relay_url: target.endpoint_uri.as_str().to_owned(), + endpoint_fingerprint: target.endpoint_fingerprint.clone(), + 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()), + }); + } + } + if required_targets.is_some() { + let selected_relay_urls = relays + .iter() + .map(|relay| relay.relay_url.clone()) + .collect::<Vec<_>>(); + for target in &active_targets { + if !is_nostr_target(target) + || !selected_relay_urls + .iter() + .any(|relay_url| relay_url == target.endpoint_uri.as_str()) + || 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)) + { + continue; + } + relays.push(PublishableRelay { + delivery_target_id: target.delivery_target_id, + relay_url: target.endpoint_uri.as_str().to_owned(), + endpoint_fingerprint: target.endpoint_fingerprint.clone(), target_scope: target .target_scope .as_ref() @@ -414,6 +492,7 @@ async fn publishable_relays( accepted_count, satisfaction_required_count, required_accept_count, + required_targets, }) } @@ -435,7 +514,11 @@ fn is_nostr_target(target: &RadrootsOutboxDeliveryTargetRecord) -> bool { fn satisfaction_policy_for_required_accept_count( required_accept_count: usize, target_count: usize, + exact_required_targets: bool, ) -> Result<RadrootsTransportSatisfactionPolicy, RadrootsRelayTransportError> { + if exact_required_targets { + return Ok(RadrootsTransportSatisfactionPolicy::no_wait()); + } if required_accept_count >= target_count { return Ok(RadrootsTransportSatisfactionPolicy::all_accepted()); } @@ -500,17 +583,23 @@ mod tests { #[test] fn internal_outbox_publish_helpers_cover_policy_edges() { assert_eq!( - satisfaction_policy_for_required_accept_count(2, 2).expect("all targets"), + satisfaction_policy_for_required_accept_count(2, 2, false).expect("all targets"), RadrootsTransportSatisfactionPolicy::all_accepted() ); assert_eq!( - satisfaction_policy_for_required_accept_count(1, 3).expect("at least one"), + satisfaction_policy_for_required_accept_count(1, 3, false).expect("at least one"), RadrootsTransportSatisfactionPolicy::any_accepted() ); + assert_eq!( + satisfaction_policy_for_required_accept_count(1, 3, true) + .expect("exact required targets"), + RadrootsTransportSatisfactionPolicy::no_wait() + ); assert!( satisfaction_policy_for_required_accept_count( usize::from(u16::MAX) + 1, usize::from(u16::MAX) + 2, + false, ) .is_err() ); diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -1662,6 +1662,305 @@ async fn outbox_publish_fans_out_endpoint_receipts_to_scoped_logical_targets() { } #[tokio::test] +async fn outbox_publish_required_target_failure_is_not_satisfied_by_optional_success() { + let signed = signed_post("required target optional success"); + 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 optional = nostr_target(RELAY_PRIMARY_WSS); + let required = nostr_target(RELAY_SECONDARY_WSS); + let receipt = outbox + .enqueue_operation(RadrootsOutboxOperationInput::new( + "publish_post", + draft, + RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 1, + RadrootsTransportSatisfactionPolicy::required_targets( + RadrootsTransportSatisfactionClass::Accepted, + vec![required.fingerprint.clone()], + ) + .expect("required target policy"), + vec![optional.clone(), required.clone()], + ), + 1_000, + )) + .await + .expect("enqueue"); + let claimed = outbox + .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) + .await + .expect("claim") + .expect("claim"); + 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 optional_target_id = publish_claim + .delivery_targets + .iter() + .find(|target| target.endpoint_fingerprint == optional.fingerprint) + .expect("optional target") + .delivery_target_id; + outbox + .mark_delivery_target_accepted( + publish_claim.outbox_event_id, + publish_claim.claim_token.as_str(), + optional_target_id, + 2_000, + ) + .await + .expect("optional accepted"); + + let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome( + RELAY_SECONDARY_WSS, + RadrootsRelayOutcome::timeout("required relay timeout"), + ); + let published = publish_claimed_outbox_event( + &outbox, + &store, + &adapter, + &publish_claim, + RadrootsOutboxPublishPolicy::new(2_500), + 2_200, + ) + .await + .expect("publish"); + + assert_eq!(published.attempted_count, 1); + assert_eq!(published.accepted_count, 0); + assert_eq!(published.retryable_count, 1); + assert_eq!(published.quorum, 1); + assert!(!published.quorum_met); + assert_eq!(published.relay_receipts.len(), 1); + assert_eq!(published.relay_receipts[0].relay_url, RELAY_SECONDARY_WSS); + let event = outbox + .get_event(receipt.outbox_event_id) + .await + .expect("event") + .expect("event"); + assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable); + let targets = outbox + .delivery_targets(receipt.outbox_event_id) + .await + .expect("targets"); + assert!(targets.iter().any(|target| { + target.endpoint_fingerprint == optional.fingerprint + && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted + })); + assert!(targets.iter().any(|target| { + target.endpoint_fingerprint == required.fingerprint + && target.status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable + })); +} + +#[tokio::test] +async fn outbox_publish_required_target_success_is_not_blocked_by_optional_retryable_failure() { + let signed = signed_post("required target optional failure"); + 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 optional = nostr_target(RELAY_PRIMARY_WSS); + let required = nostr_target(RELAY_SECONDARY_WSS); + let receipt = outbox + .enqueue_operation(RadrootsOutboxOperationInput::new( + "publish_post", + draft, + RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 1, + RadrootsTransportSatisfactionPolicy::required_targets( + RadrootsTransportSatisfactionClass::Accepted, + vec![required.fingerprint.clone()], + ) + .expect("required target policy"), + vec![optional.clone(), required.clone()], + ), + 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 optional_target_id = publish_claim + .delivery_targets + .iter() + .find(|target| target.endpoint_fingerprint == optional.fingerprint) + .expect("optional target") + .delivery_target_id; + outbox + .mark_delivery_target_failed_retryable( + publish_claim.outbox_event_id, + publish_claim.claim_token.as_str(), + optional_target_id, + "optional relay timeout", + 2_000, + ) + .await + .expect("optional retryable"); + + let adapter = RadrootsMockRelayPublishAdapter::new() + .with_outcome(RELAY_SECONDARY_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.attempted_count, 1); + assert_eq!(published.accepted_count, 1); + assert_eq!(published.retryable_count, 0); + assert_eq!(published.quorum, 1); + assert!(published.quorum_met); + 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!(targets.iter().any(|target| { + target.endpoint_fingerprint == optional.fingerprint + && target.status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable + })); + assert!(targets.iter().any(|target| { + target.endpoint_fingerprint == required.fingerprint + && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted + })); + let observations = store + .observations_for_event(signed.id.as_str()) + .await + .expect("observations"); + assert_outbox_publish_observations(&observations, 1); +} + +#[tokio::test] +async fn outbox_publish_required_targets_fan_out_same_endpoint_scoped_receipts() { + let signed = signed_post("required target 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 required = scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.west", "West foodshed"); + let optional = scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.east", "East foodshed"); + let receipt = outbox + .enqueue_operation(RadrootsOutboxOperationInput::new( + "publish_post", + draft, + RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 1, + RadrootsTransportSatisfactionPolicy::required_targets( + RadrootsTransportSatisfactionClass::Accepted, + vec![required.fingerprint.clone()], + ) + .expect("required target policy"), + vec![required.clone(), optional.clone()], + ), + 1_000, + )) + .await + .expect("enqueue"); + let claimed = outbox + .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) + .await + .expect("claim") + .expect("claim"); + 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.attempted_count, 2); + assert_eq!(published.accepted_count, 2); + assert_eq!(published.quorum, 1); + assert!(published.quorum_met); + assert_eq!(published.relay_receipts.len(), 1); + assert_eq!(published.target_receipts.len(), 2); + assert!(published.target_receipts.iter().any(|target| { + target.endpoint_fingerprint == required.fingerprint + && target.target_scope.as_deref() == Some("foodshed.west") + })); + assert!(published.target_receipts.iter().any(|target| { + target.endpoint_fingerprint == optional.fingerprint + && target.target_scope.as_deref() == Some("foodshed.east") + })); + 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 + })); +} + +#[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");