lib

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

commit fa2fcb6804bb0c433c093399cb7dc089fae18ced
parent 6c56aed578ec0ba7cd005a3909f3047a439a7e44
Author: triesap <tyson@radroots.org>
Date:   Mon, 13 Jul 2026 06:49:34 +0000

transport: align downstream satisfaction semantics

- rename transport publish outcome satisfaction to accepted delivery semantics
- preserve transport target statuses through Nostr outbox completion
- add strict success outbox status completion and monotonic upgrades
- cover class-aware outbox publish and strict status completion tests

Diffstat:
Mcrates/outbox/src/store.rs | 332+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcrates/transport_nostr/src/outbox.rs | 498+++++++++++++++++++++++++++++++++++++++++++++++++++++++------------------------
Mcrates/transport_nostr/tests/transport.rs | 2+-
Mcrates/transport_publish_protocol/src/lib.rs | 22++++++++++++----------
4 files changed, 691 insertions(+), 163 deletions(-)

diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs @@ -774,6 +774,78 @@ impl RadrootsOutbox { .await } + pub async fn mark_delivery_target_delivered( + &self, + outbox_event_id: i64, + claim_token: &str, + delivery_target_id: i64, + attempted_at_ms: i64, + ) -> Result<(), RadrootsOutboxError> { + self.mark_delivery_target_status( + outbox_event_id, + claim_token, + delivery_target_id, + RadrootsOutboxDeliveryTargetStatus::Delivered, + None, + attempted_at_ms, + ) + .await + } + + pub async fn mark_delivery_target_forwarded( + &self, + outbox_event_id: i64, + claim_token: &str, + delivery_target_id: i64, + attempted_at_ms: i64, + ) -> Result<(), RadrootsOutboxError> { + self.mark_delivery_target_status( + outbox_event_id, + claim_token, + delivery_target_id, + RadrootsOutboxDeliveryTargetStatus::Forwarded, + None, + attempted_at_ms, + ) + .await + } + + pub async fn mark_delivery_target_stored_by_gateway( + &self, + outbox_event_id: i64, + claim_token: &str, + delivery_target_id: i64, + attempted_at_ms: i64, + ) -> Result<(), RadrootsOutboxError> { + self.mark_delivery_target_status( + outbox_event_id, + claim_token, + delivery_target_id, + RadrootsOutboxDeliveryTargetStatus::StoredByGateway, + None, + attempted_at_ms, + ) + .await + } + + pub async fn mark_delivery_target_seen( + &self, + outbox_event_id: i64, + claim_token: &str, + delivery_target_id: i64, + attempted_at_ms: i64, + ) -> Result<(), RadrootsOutboxError> { + self.mark_delivery_target_status( + outbox_event_id, + claim_token, + delivery_target_id, + RadrootsOutboxDeliveryTargetStatus::Seen, + None, + attempted_at_ms, + ) + .await + } + pub async fn mark_delivery_target_failed_retryable( &self, outbox_event_id: i64, @@ -1130,11 +1202,13 @@ impl RadrootsOutbox { tx.commit().await?; return Ok(()); } - return Err(RadrootsOutboxError::DeliveryTargetStatusConflict { - delivery_target_id, - current_status: current_status.as_str(), - requested_status: status.as_str(), - }); + if !delivery_target_status_can_advance(current_status, status) { + return Err(RadrootsOutboxError::DeliveryTargetStatusConflict { + delivery_target_id, + current_status: current_status.as_str(), + requested_status: status.as_str(), + }); + } } let completed_at_ms = status.is_completed().then_some(attempted_at_ms); let outcome_kind = outcome_kind_for_status(status); @@ -2306,6 +2380,27 @@ fn outcome_kind_for_status( } } +fn delivery_target_status_can_advance( + current: RadrootsOutboxDeliveryTargetStatus, + requested: RadrootsOutboxDeliveryTargetStatus, +) -> bool { + matches!( + (current, requested), + ( + RadrootsOutboxDeliveryTargetStatus::Accepted, + RadrootsOutboxDeliveryTargetStatus::Forwarded + | RadrootsOutboxDeliveryTargetStatus::StoredByGateway + | RadrootsOutboxDeliveryTargetStatus::Seen + | RadrootsOutboxDeliveryTargetStatus::Delivered + ) | ( + RadrootsOutboxDeliveryTargetStatus::Forwarded + | RadrootsOutboxDeliveryTargetStatus::StoredByGateway + | RadrootsOutboxDeliveryTargetStatus::Seen, + RadrootsOutboxDeliveryTargetStatus::Delivered + ) + ) +} + fn parse_transport_outcome_kind( value: &str, field: &'static str, @@ -4814,6 +4909,233 @@ mod tests { } #[tokio::test] + async fn accepted_delivery_target_can_advance_to_strict_success_status() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "accepted then delivered"); + let signed_event = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("signed event"); + let receipt = outbox + .enqueue_signed_operation(RadrootsOutboxSignedOperationInput::new( + "publish_post", + draft, + signed_event, + RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 1, + RadrootsTransportSatisfactionPolicy::all_delivered(), + vec![nostr_target(NOSTR_PRIMARY_WSS)], + ), + 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 target_id = claimed.delivery_targets[0].delivery_target_id; + + outbox + .mark_delivery_target_accepted(receipt.outbox_event_id, "claim-a", target_id, 1_100) + .await + .expect("accepted"); + outbox + .mark_delivery_target_delivered(receipt.outbox_event_id, "claim-a", target_id, 1_110) + .await + .expect("delivered"); + + 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_eq!( + targets[0].status, + RadrootsOutboxDeliveryTargetStatus::Delivered + ); + assert_eq!(targets[0].attempt_count, 2); + assert_eq!(targets[0].completed_at_ms, Some(1_110)); + assert_eq!( + targets[0].last_outcome_kind, + Some(RadrootsTransportOutcomeKind::Delivered) + ); + let attempts = outbox.delivery_attempts(target_id).await.expect("attempts"); + assert_eq!(attempts.len(), 2); + assert_eq!( + attempts[0].outcome_kind, + RadrootsTransportOutcomeKind::Accepted + ); + assert_eq!( + attempts[1].outcome_kind, + RadrootsTransportOutcomeKind::Delivered + ); + let plans = outbox + .delivery_plans(receipt.outbox_event_id) + .await + .expect("plans"); + assert_eq!(plans[0].status, RadrootsOutboxDeliveryPlanStatus::Complete); + } + + #[tokio::test] + async fn claimed_delivery_targets_can_complete_as_strict_transport_successes() { + for (content, mark, expected_status, satisfaction_policy) in [ + ( + "transport delivered", + "delivered", + RadrootsOutboxDeliveryTargetStatus::Delivered, + RadrootsTransportSatisfactionPolicy::all_delivered(), + ), + ( + "transport forwarded", + "forwarded", + RadrootsOutboxDeliveryTargetStatus::Forwarded, + RadrootsTransportSatisfactionPolicy::all_forwarded(), + ), + ( + "transport stored", + "stored", + RadrootsOutboxDeliveryTargetStatus::StoredByGateway, + RadrootsTransportSatisfactionPolicy::all_stored(), + ), + ( + "transport seen", + "seen", + RadrootsOutboxDeliveryTargetStatus::Seen, + RadrootsTransportSatisfactionPolicy::all_seen(), + ), + ] { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, content); + let signed_event = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("signed event"); + let receipt = outbox + .enqueue_signed_operation(RadrootsOutboxSignedOperationInput::new( + "publish_post", + draft, + signed_event, + RadrootsOutboxDeliveryPlanInput::new( + "transport.nostr.local", + 1, + satisfaction_policy, + vec![nostr_target(NOSTR_PRIMARY_WSS)], + ), + 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 target_id = claimed.delivery_targets[0].delivery_target_id; + + match mark { + "delivered" => { + outbox + .mark_delivery_target_delivered( + receipt.outbox_event_id, + "claim-a", + target_id, + 1_100, + ) + .await + .expect("delivered"); + } + "forwarded" => { + outbox + .mark_delivery_target_forwarded( + receipt.outbox_event_id, + "claim-a", + target_id, + 1_100, + ) + .await + .expect("forwarded"); + } + "stored" => { + outbox + .mark_delivery_target_stored_by_gateway( + receipt.outbox_event_id, + "claim-a", + target_id, + 1_100, + ) + .await + .expect("stored"); + } + "seen" => { + outbox + .mark_delivery_target_seen( + receipt.outbox_event_id, + "claim-a", + target_id, + 1_100, + ) + .await + .expect("seen"); + } + _ => unreachable!(), + } + + 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 operation = outbox + .get_operation(receipt.operation_id) + .await + .expect("operation") + .expect("operation"); + assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete); + let targets = outbox + .delivery_targets(receipt.outbox_event_id) + .await + .expect("targets"); + assert_eq!(targets[0].status, expected_status); + assert_eq!( + targets[0].last_outcome_kind, + Some(outcome_kind_for_status(expected_status)) + ); + let attempts = outbox.delivery_attempts(target_id).await.expect("attempts"); + assert_eq!(attempts.len(), 1); + assert_eq!( + attempts[0].outcome_kind, + outcome_kind_for_status(expected_status) + ); + let plans = outbox + .delivery_plans(receipt.outbox_event_id) + .await + .expect("plans"); + assert_eq!(plans[0].status, RadrootsOutboxDeliveryPlanStatus::Complete); + } + } + + #[tokio::test] async fn claimed_delivery_targets_can_complete_as_preview_outcomes() { for ( content, diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs @@ -72,6 +72,7 @@ pub struct RadrootsOutboxPublishTargetReceipt { pub target_scope: Option<String>, pub target_label: Option<String>, pub attempted: bool, + pub transport_status: RadrootsTransportDeliveryTargetStatus, pub outcome: RadrootsRelayOutcome, } @@ -116,8 +117,8 @@ where accepted_count: publishable.accepted_count, retryable_count: 0, terminal_count: 0, - quorum: publishable.satisfaction_required_count, - quorum_met: publishable.accepted_count >= publishable.satisfaction_required_count, + quorum: publishable.remaining_satisfaction_count, + quorum_met: publishable.satisfied_count >= publishable.satisfaction_required_count, target_receipts: Vec::new(), relay_receipts: Vec::new(), }); @@ -130,10 +131,11 @@ where policy.relay_url_policy, )?; let target_strings = targets.relay_strings(); - let satisfaction_policy = satisfaction_policy_for_required_accept_count( - publishable.required_accept_count, + let satisfaction_policy = satisfaction_policy_for_remaining_count( + publishable.satisfaction_class, + publishable.remaining_satisfaction_count, targets.len(), - publishable.required_targets.is_some(), + 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) @@ -149,7 +151,7 @@ where Err(RadrootsRelayTransportError::Transport(message)) => adapter_transport_failure_receipt( signed_event.id.clone(), target_strings, - publishable.required_accept_count, + publishable.remaining_satisfaction_count, message, ), Err(error) => return Err(error), @@ -157,48 +159,15 @@ where let target_receipts = target_receipts_from_relay_receipts(&publishable, &publish.relays); 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?; - } + complete_outbox_delivery_target(outbox, claimed, target_receipt, now_ms).await?; } for relay in &publish.relays { - if relay.outcome.counts_toward_quorum() + if relay + .outcome + .to_transport_outcome() + .status + .counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted) && publishable .targets_for_relay(relay.relay_url.as_str()) .next() @@ -245,7 +214,7 @@ where .iter() .filter(|receipt| receipt.outcome.is_terminal_failure()) .count(), - quorum: publishable.required_accept_count, + quorum: publishable.remaining_satisfaction_count, quorum_met: publishable.satisfied_count_after_receipts(&target_receipts) >= publishable.satisfaction_required_count, target_receipts, @@ -294,8 +263,8 @@ where accepted_count: publishable.accepted_count, retryable_count: 0, terminal_count: 0, - quorum: publishable.satisfaction_required_count, - quorum_met: publishable.accepted_count >= publishable.satisfaction_required_count, + quorum: publishable.remaining_satisfaction_count, + quorum_met: publishable.satisfied_count >= publishable.satisfaction_required_count, target_receipts: Vec::new(), relay_receipts: Vec::new(), }); @@ -336,48 +305,14 @@ where let target_receipts = target_receipts_from_transport_receipts(&publishable, &delivery); for target_receipt in &target_receipts { - if target_receipt.outcome.counts_toward_quorum() { - outbox - .mark_delivery_target_accepted( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - target_receipt.delivery_target_id, - now_ms, - ) - .await?; - } else if target_receipt.outcome.is_retryable() { - outbox - .mark_delivery_target_failed_retryable( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - target_receipt.delivery_target_id, - target_receipt - .outcome - .message - .as_deref() - .unwrap_or("relay publish retryable"), - now_ms, - ) - .await?; - } else { - outbox - .mark_delivery_target_failed_terminal( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - target_receipt.delivery_target_id, - target_receipt - .outcome - .message - .as_deref() - .unwrap_or("relay publish terminal"), - now_ms, - ) - .await?; - } + complete_outbox_delivery_target(outbox, claimed, target_receipt, now_ms).await?; } for target_receipt in &target_receipts { - if target_receipt.outcome.counts_toward_quorum() { + if target_receipt + .transport_status + .counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted) + { ingest_publish_observation( event_store, &signed_event, @@ -421,7 +356,7 @@ where .iter() .filter(|receipt| receipt.outcome.is_terminal_failure()) .count(), - quorum: publishable.required_accept_count, + quorum: publishable.remaining_satisfaction_count, quorum_met: publishable.satisfied_count_after_receipts(&target_receipts) >= publishable.satisfaction_required_count, target_receipts, @@ -460,9 +395,12 @@ struct PublishableRelays { active_delivery_plan_id: i64, relays: Vec<PublishableRelay>, accepted_count: usize, + satisfied_count: usize, satisfaction_required_count: usize, - required_accept_count: usize, + remaining_satisfaction_count: usize, + satisfaction_class: RadrootsTransportSatisfactionClass, required_targets: Option<Vec<RadrootsTransportTargetFingerprint>>, + remaining_required_targets: Option<Vec<RadrootsTransportTargetFingerprint>>, } impl PublishableRelays { @@ -479,11 +417,13 @@ impl PublishableRelays { &self, target_receipts: &[RadrootsOutboxPublishTargetReceipt], ) -> usize { - self.accepted_count + self.satisfied_count + target_receipts .iter() .filter(|receipt| { - receipt.outcome.counts_toward_quorum() + receipt + .transport_status + .counts_as_satisfied(self.satisfaction_class) && self.required_targets.as_ref().is_none_or(|required| { required .iter() @@ -516,6 +456,7 @@ fn target_receipts_from_relay_receipts( target_scope: target.target_scope.clone(), target_label: target.target_label.clone(), attempted: relay_receipt.attempted, + transport_status: relay_receipt.outcome.to_transport_outcome().status, outcome: relay_receipt.outcome.clone(), }); } @@ -543,12 +484,155 @@ fn target_receipts_from_transport_receipts( target_label: target.target_label.clone(), attempted: receipt.status != RadrootsTransportDeliveryTargetStatus::SkippedPolicyDenied, + transport_status: receipt.status, outcome: relay_outcome_from_transport_outcome(&receipt.outcome), }) }) .collect() } +async fn complete_outbox_delivery_target( + outbox: &RadrootsOutbox, + claimed: &RadrootsOutboxClaimedEvent, + receipt: &RadrootsOutboxPublishTargetReceipt, + now_ms: i64, +) -> Result<(), RadrootsRelayTransportError> { + match receipt.transport_status { + RadrootsTransportDeliveryTargetStatus::Pending => { + return Err(RadrootsRelayTransportError::TransportContract( + "outbox publish receipt cannot complete a target with pending transport status" + .to_owned(), + )); + } + RadrootsTransportDeliveryTargetStatus::Accepted => { + outbox + .mark_delivery_target_accepted( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + now_ms, + ) + .await?; + } + RadrootsTransportDeliveryTargetStatus::Delivered => { + outbox + .mark_delivery_target_delivered( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + now_ms, + ) + .await?; + } + RadrootsTransportDeliveryTargetStatus::Forwarded => { + outbox + .mark_delivery_target_forwarded( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + now_ms, + ) + .await?; + } + RadrootsTransportDeliveryTargetStatus::StoredByGateway => { + outbox + .mark_delivery_target_stored_by_gateway( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + now_ms, + ) + .await?; + } + RadrootsTransportDeliveryTargetStatus::Seen => { + outbox + .mark_delivery_target_seen( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + now_ms, + ) + .await?; + } + RadrootsTransportDeliveryTargetStatus::DeferredUntilImplemented => { + outbox + .mark_delivery_target_deferred_until_implemented( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + receipt + .outcome + .message + .as_deref() + .unwrap_or("relay publish deferred until implemented"), + now_ms, + ) + .await?; + } + RadrootsTransportDeliveryTargetStatus::PreviewUnavailable => { + outbox + .mark_delivery_target_preview_unavailable( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + receipt + .outcome + .message + .as_deref() + .unwrap_or("relay publish preview unavailable"), + now_ms, + ) + .await?; + } + RadrootsTransportDeliveryTargetStatus::SkippedPolicyDenied => { + outbox + .mark_delivery_target_skipped_policy_denied( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + receipt + .outcome + .message + .as_deref() + .unwrap_or("relay publish skipped by policy"), + now_ms, + ) + .await?; + } + RadrootsTransportDeliveryTargetStatus::FailedRetryable => { + outbox + .mark_delivery_target_failed_retryable( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + receipt + .outcome + .message + .as_deref() + .unwrap_or("relay publish retryable"), + now_ms, + ) + .await?; + } + RadrootsTransportDeliveryTargetStatus::FailedTerminal => { + outbox + .mark_delivery_target_failed_terminal( + claimed.outbox_event_id, + claimed.claim_token.as_str(), + receipt.delivery_target_id, + receipt + .outcome + .message + .as_deref() + .unwrap_or("relay publish terminal"), + now_ms, + ) + .await?; + } + } + Ok(()) +} + fn relay_receipts_from_transport_receipts( delivery: &RadrootsTransportDeliveryReceipt, ) -> Vec<RadrootsRelayPublishRelayReceipt> { @@ -663,22 +747,11 @@ fn publishable_transport_targets( fn transport_satisfaction_policy_for_publishable( publishable: &PublishableRelays, ) -> Result<RadrootsTransportSatisfactionPolicy, RadrootsRelayTransportError> { - if publishable.required_targets.is_some() { - let required_targets = publishable - .relays - .iter() - .map(|relay| relay.endpoint_fingerprint.clone()) - .collect::<Vec<_>>(); - return RadrootsTransportSatisfactionPolicy::required_targets( - RadrootsTransportSatisfactionClass::Accepted, - required_targets, - ) - .map_err(transport_error_to_relay_error); - } - satisfaction_policy_for_required_accept_count( - publishable.required_accept_count, + satisfaction_policy_for_remaining_count( + publishable.satisfaction_class, + publishable.remaining_satisfaction_count, publishable.relays.len(), - false, + publishable.remaining_required_targets.as_deref(), ) } @@ -763,26 +836,38 @@ async fn publishable_relays( active_delivery_plan_id ))); } - let satisfied_count = plan + let satisfaction_class = plan .satisfaction_policy .target_satisfaction_class() - .map(|satisfaction_class| { - active_targets - .iter() - .filter(|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) - }) - .count() + .unwrap_or(RadrootsTransportSatisfactionClass::Accepted); + let satisfied_count = active_targets + .iter() + .filter(|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) }) - .unwrap_or(0); - let required_accept_count = + .count(); + let remaining_satisfaction_count = (plan.required_success_count as usize).saturating_sub(satisfied_count); + let remaining_required_targets = required_targets.as_ref().map(|required_targets| { + required_targets + .iter() + .filter(|required| { + !active_targets.iter().any(|target| { + target.endpoint_fingerprint == **required + && target + .status + .counts_as_transport_satisfaction(satisfaction_class) + }) + }) + .cloned() + .collect::<Vec<_>>() + }); let mut relays = Vec::new(); let mut accepted_count = 0usize; for target in &active_targets { @@ -803,7 +888,7 @@ async fn publishable_relays( } let can_contribute_to_satisfaction = required_targets.is_none() || required_for_satisfaction; - if required_accept_count > 0 + if remaining_satisfaction_count > 0 && can_contribute_to_satisfaction && (target.status.is_ready_for_attempt() || (republish_accepted_relays @@ -862,9 +947,12 @@ async fn publishable_relays( active_delivery_plan_id, relays, accepted_count, + satisfied_count, satisfaction_required_count, - required_accept_count, + remaining_satisfaction_count, + satisfaction_class, required_targets, + remaining_required_targets, }) } @@ -883,32 +971,44 @@ fn is_nostr_target(target: &RadrootsOutboxDeliveryTargetRecord) -> bool { target.transport_kind == RadrootsTransportKind::Nostr } -fn satisfaction_policy_for_required_accept_count( - required_accept_count: usize, +fn satisfaction_policy_for_remaining_count( + satisfaction_class: RadrootsTransportSatisfactionClass, + remaining_satisfaction_count: usize, target_count: usize, - exact_required_targets: bool, + exact_required_targets: Option<&[RadrootsTransportTargetFingerprint]>, ) -> Result<RadrootsTransportSatisfactionPolicy, RadrootsRelayTransportError> { - if exact_required_targets { - return Ok(RadrootsTransportSatisfactionPolicy::no_wait()); + if let Some(targets) = exact_required_targets { + return RadrootsTransportSatisfactionPolicy::required_targets( + satisfaction_class, + targets.to_vec(), + ) + .map_err(transport_error_to_relay_error); } - if required_accept_count >= target_count { - return Ok(RadrootsTransportSatisfactionPolicy::all_accepted()); + if remaining_satisfaction_count >= target_count { + return Ok(RadrootsTransportSatisfactionPolicy::All { + class: satisfaction_class, + }); } - if required_accept_count == 0 { + if remaining_satisfaction_count == 0 { return Err(RadrootsRelayTransportError::Transport( - "required Nostr relay acceptance count must be greater than zero".to_owned(), + "required Nostr relay satisfaction count must be greater than zero".to_owned(), )); } - if required_accept_count == 1 { - return Ok(RadrootsTransportSatisfactionPolicy::any_accepted()); + if remaining_satisfaction_count == 1 { + return Ok(RadrootsTransportSatisfactionPolicy::Any { + class: satisfaction_class, + }); } - let count = u16::try_from(required_accept_count).map_err(|_| { + let count = u16::try_from(remaining_satisfaction_count).map_err(|_| { RadrootsRelayTransportError::Transport( - "required Nostr relay acceptance count exceeds supported transport policy range" + "required Nostr relay satisfaction count exceeds supported transport policy range" .to_owned(), ) })?; - Ok(RadrootsTransportSatisfactionPolicy::quorum_accepted(count)) + Ok(RadrootsTransportSatisfactionPolicy::Quorum { + class: satisfaction_class, + threshold: count, + }) } async fn ingest_publish_observation( @@ -949,35 +1049,139 @@ fn event_from_signed(signed_event: &RadrootsSignedEvent) -> RadrootsEventEnvelop #[cfg(test)] mod tests { - use super::{adapter_transport_failure_receipt, satisfaction_policy_for_required_accept_count}; - use radroots_transport::RadrootsTransportSatisfactionPolicy; + use super::{ + PublishableRelay, PublishableRelays, adapter_transport_failure_receipt, + satisfaction_policy_for_remaining_count, target_receipts_from_relay_receipts, + target_receipts_from_transport_receipts, + }; + use crate::{RadrootsRelayOutcome, RadrootsRelayPublishRelayReceipt}; + use radroots_transport::{ + RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryTargetStatus, + RadrootsTransportKind, RadrootsTransportOutcome, RadrootsTransportOutcomeKind, + RadrootsTransportSatisfactionClass, RadrootsTransportSatisfactionPolicy, + RadrootsTransportTarget, RadrootsTransportTargetReceipt, + }; #[test] fn internal_outbox_publish_helpers_cover_policy_edges() { assert_eq!( - satisfaction_policy_for_required_accept_count(2, 2, false).expect("all targets"), + satisfaction_policy_for_remaining_count( + RadrootsTransportSatisfactionClass::Accepted, + 2, + 2, + None + ) + .expect("all targets"), RadrootsTransportSatisfactionPolicy::all_accepted() ); assert_eq!( - satisfaction_policy_for_required_accept_count(1, 3, false).expect("at least one"), + satisfaction_policy_for_remaining_count( + RadrootsTransportSatisfactionClass::Accepted, + 1, + 3, + None + ) + .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() + satisfaction_policy_for_remaining_count( + RadrootsTransportSatisfactionClass::Delivered, + 2, + 3, + None + ) + .expect("delivered quorum"), + RadrootsTransportSatisfactionPolicy::quorum_delivered(2) + ); + let required_target = + RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, "wss://relay.example") + .expect("required target"); + assert_eq!( + satisfaction_policy_for_remaining_count( + RadrootsTransportSatisfactionClass::Delivered, + 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!( - satisfaction_policy_for_required_accept_count( + satisfaction_policy_for_remaining_count( + RadrootsTransportSatisfactionClass::Accepted, usize::from(u16::MAX) + 1, usize::from(u16::MAX) + 2, - false, + None, ) .is_err() ); } #[test] + fn outbox_publish_satisfaction_counts_use_active_transport_class() { + let target = + RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, "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: None, + target_label: None, + }], + accepted_count: 0, + satisfied_count: 0, + satisfaction_required_count: 1, + remaining_satisfaction_count: 1, + satisfaction_class: RadrootsTransportSatisfactionClass::Delivered, + required_targets: None, + remaining_required_targets: None, + }; + + let accepted_relay_receipts = target_receipts_from_relay_receipts( + &publishable, + &[RadrootsRelayPublishRelayReceipt::attempted( + target.uri.as_str(), + RadrootsRelayOutcome::accepted(), + )], + ); + assert_eq!( + accepted_relay_receipts[0].transport_status, + RadrootsTransportDeliveryTargetStatus::Accepted + ); + assert_eq!( + publishable.satisfied_count_after_receipts(&accepted_relay_receipts), + 0 + ); + + let delivered_transport_receipts = target_receipts_from_transport_receipts( + &publishable, + &RadrootsTransportDeliveryReceipt { + request_id: "request-1".to_owned(), + target_receipts: vec![RadrootsTransportTargetReceipt::new( + target, + RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Delivered), + )], + }, + ); + assert_eq!( + delivered_transport_receipts[0].transport_status, + RadrootsTransportDeliveryTargetStatus::Delivered + ); + assert_eq!( + publishable.satisfied_count_after_receipts(&delivered_transport_receipts), + 1 + ); + } + + #[test] fn adapter_transport_failure_receipts_preserve_each_target() { let receipt = adapter_transport_failure_receipt( "event-1".to_owned(), diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -2237,7 +2237,7 @@ async fn outbox_publish_marks_published_without_adapter_when_all_relays_already_ 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_eq!(published.quorum, 0); assert!(published.quorum_met); assert!(published.target_receipts.is_empty()); assert!(published.relay_receipts.is_empty()); diff --git a/crates/transport_publish_protocol/src/lib.rs b/crates/transport_publish_protocol/src/lib.rs @@ -730,7 +730,7 @@ pub enum TransportPublishOutcomeKind { } impl TransportPublishOutcomeKind { - pub fn counts_toward_satisfaction(self) -> bool { + pub fn counts_toward_accepted_delivery(self) -> bool { matches!( self, Self::Accepted | Self::DuplicateAccepted | Self::SkippedAlreadyAccepted @@ -902,7 +902,7 @@ impl TransportPublishJobView { let acknowledged_count = self .targets .iter() - .filter(|target| target.outcome_kind.counts_toward_satisfaction()) + .filter(|target| target.outcome_kind.counts_toward_accepted_delivery()) .count(); let retryable_count = self .targets @@ -1302,7 +1302,7 @@ fn validate_job_status_state( required_outcomes.iter().any(|outcome| { target_outcome_fingerprint(outcome, 0).is_ok_and(|fingerprint| { fingerprint == *required - && outcome.outcome_kind.counts_toward_satisfaction() + && outcome.outcome_kind.counts_toward_accepted_delivery() }) }) }); @@ -1452,7 +1452,7 @@ mod tests { target_count: targets.len(), acknowledged_count: targets .iter() - .filter(|target| target.outcome_kind.counts_toward_satisfaction()) + .filter(|target| target.outcome_kind.counts_toward_accepted_delivery()) .count(), retryable_count: targets .iter() @@ -1774,8 +1774,10 @@ mod tests { #[test] fn outcome_kinds_classify_satisfaction_retry_and_terminal() { - assert!(TransportPublishOutcomeKind::Accepted.counts_toward_satisfaction()); - assert!(TransportPublishOutcomeKind::SkippedAlreadyAccepted.counts_toward_satisfaction()); + assert!(TransportPublishOutcomeKind::Accepted.counts_toward_accepted_delivery()); + assert!( + TransportPublishOutcomeKind::SkippedAlreadyAccepted.counts_toward_accepted_delivery() + ); assert!(TransportPublishOutcomeKind::Timeout.is_retryable()); assert!(TransportPublishOutcomeKind::PreviewUnavailable.is_deferred_preview()); assert!(TransportPublishOutcomeKind::DeferredUntilImplemented.is_deferred_preview()); @@ -2482,22 +2484,22 @@ mod tests { ]; for kind in satisfied { - assert!(kind.counts_toward_satisfaction()); + assert!(kind.counts_toward_accepted_delivery()); assert!(!kind.is_retryable()); assert!(!kind.is_terminal_failure()); } for kind in retryable { - assert!(!kind.counts_toward_satisfaction()); + assert!(!kind.counts_toward_accepted_delivery()); assert!(kind.is_retryable()); assert!(!kind.is_terminal_failure()); } for kind in terminal { - assert!(!kind.counts_toward_satisfaction()); + assert!(!kind.counts_toward_accepted_delivery()); assert!(!kind.is_retryable()); assert!(kind.is_terminal_failure()); } for kind in deferred_preview { - assert!(!kind.counts_toward_satisfaction()); + assert!(!kind.counts_toward_accepted_delivery()); assert!(!kind.is_retryable()); assert!(!kind.is_terminal_failure()); assert!(kind.is_deferred_preview());