lib

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

commit e17c55cbed11a0742cbb880781441952d56ed0fa
parent 240d6b2a8dcee40383d00a755f25b53cf297f335
Author: triesap <tyson@radroots.org>
Date:   Sat, 18 Jul 2026 09:52:56 +0000

outbox: close semantic storage coverage

- Cover semantic trade enqueue, preflight, and transactional idempotency paths.
- Exercise claim recovery, stale-worker guards, and publish lifecycle boundaries.
- Reject malformed persisted values, policies, and trade metadata.
- Remove redundant terminal-state branches and verify strict Clippy plus full coverage.

Diffstat:
Mcrates/outbox/src/model.rs | 86+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/outbox/src/store.rs | 1645++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------
2 files changed, 1590 insertions(+), 141 deletions(-)

diff --git a/crates/outbox/src/model.rs b/crates/outbox/src/model.rs @@ -401,6 +401,7 @@ pub struct RadrootsOutboxSignedTradeMutationInput { } impl RadrootsOutboxSignedTradeMutationInput { + #[allow(clippy::too_many_arguments)] pub fn new( operation_kind: impl Into<String>, trade_id: RadrootsTradeId, @@ -584,6 +585,7 @@ pub struct RadrootsOutboxStatusSummary { } #[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] mod tests { use super::*; @@ -782,6 +784,90 @@ mod tests { deferred_until_implemented ); assert_eq!(status.is_terminal_failure(), terminal_failure); + assert_eq!( + status.is_retryable_failure(), + status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable + ); } + + for (class, satisfying_statuses) in [ + ( + radroots_transport::RadrootsTransportSatisfactionClass::Forwarded, + vec![ + RadrootsOutboxDeliveryTargetStatus::Forwarded, + RadrootsOutboxDeliveryTargetStatus::Delivered, + ], + ), + ( + radroots_transport::RadrootsTransportSatisfactionClass::Stored, + vec![RadrootsOutboxDeliveryTargetStatus::StoredByGateway], + ), + ( + radroots_transport::RadrootsTransportSatisfactionClass::Seen, + vec![ + RadrootsOutboxDeliveryTargetStatus::Seen, + RadrootsOutboxDeliveryTargetStatus::Delivered, + ], + ), + ( + radroots_transport::RadrootsTransportSatisfactionClass::DurableOrObserved, + vec![ + RadrootsOutboxDeliveryTargetStatus::StoredByGateway, + RadrootsOutboxDeliveryTargetStatus::Seen, + RadrootsOutboxDeliveryTargetStatus::Delivered, + ], + ), + ] { + for status in satisfying_statuses { + assert!(status.counts_as_transport_satisfaction(class)); + } + assert!( + !RadrootsOutboxDeliveryTargetStatus::Pending + .counts_as_transport_satisfaction(class) + ); + } + + for invalid in ["", "unknown"] { + assert!(RadrootsOutboxOperationStatus::parse(invalid).is_err()); + assert!(RadrootsOutboxEventState::parse(invalid).is_err()); + assert!(RadrootsOutboxDeliveryPlanStatus::parse(invalid).is_err()); + assert!(RadrootsOutboxDeliveryTargetStatus::parse(invalid).is_err()); + } + + for state in [ + RadrootsOutboxEventState::DraftQueued, + RadrootsOutboxEventState::Signing, + RadrootsOutboxEventState::Signed, + RadrootsOutboxEventState::Publishing, + RadrootsOutboxEventState::SignRetryable, + RadrootsOutboxEventState::PublishRetryable, + ] { + assert!(!state.is_terminal()); + } + for state in [ + RadrootsOutboxEventState::Published, + RadrootsOutboxEventState::FailedTerminal, + RadrootsOutboxEventState::Cancelled, + ] { + assert!(state.is_terminal()); + } + + assert!( + RadrootsOutboxDeliveryTargetStatus::Forwarded + .counts_as_transport_satisfaction(RadrootsTransportSatisfactionClass::Forwarded,) + ); + assert!( + RadrootsOutboxDeliveryTargetStatus::StoredByGateway + .counts_as_transport_satisfaction(RadrootsTransportSatisfactionClass::Stored) + ); + assert!( + RadrootsOutboxDeliveryTargetStatus::Seen + .counts_as_transport_satisfaction(RadrootsTransportSatisfactionClass::Seen,) + ); + assert!( + RadrootsOutboxDeliveryTargetStatus::StoredByGateway.counts_as_transport_satisfaction( + RadrootsTransportSatisfactionClass::DurableOrObserved, + ) + ); } } diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs @@ -926,13 +926,9 @@ impl RadrootsOutbox { let outbox_event_ids = reticulum_event_ids_pool(&self.pool, outbox_event_id, limit).await?; let mut records = Vec::with_capacity(outbox_event_ids.len()); for outbox_event_id in outbox_event_ids { - let Some(event) = self.get_event(outbox_event_id).await? else { - continue; - }; + let event = self.get_event(outbox_event_id).await?; let targets = reticulum_targets_for_event_pool(&self.pool, outbox_event_id).await?; - if !targets.is_empty() { - records.push(RadrootsOutboxReticulumEventRecord { event, targets }); - } + records.extend(reticulum_event_record(event, targets)); } Ok(records) } @@ -1233,7 +1229,8 @@ impl RadrootsOutbox { "local:outbox", RadrootsTransportObservationType::LocalImport, observed_at_ms, - )?; + ) + .expect("the static local outbox transport URI must remain valid"); let ingest = RadrootsEventIngest::new(signed_event.clone(), observed_at_ms) .with_observation(observation); let receipt = event_store.ingest_event(ingest).await?; @@ -1763,7 +1760,6 @@ struct PlanInsertResult { struct PlanEvaluation { all_complete: bool, - any_failed_terminal: bool, any_ready: bool, } @@ -1793,13 +1789,6 @@ fn publish_lifecycle_from_plan_evaluation<'a>( Some(retryable_error), next_attempt_after_ms, ) - } else if evaluation.any_failed_terminal { - ( - RadrootsOutboxEventState::FailedTerminal, - Some(RadrootsOutboxOperationStatus::FailedTerminal), - Some(terminal_error), - now_ms, - ) } else { ( RadrootsOutboxEventState::FailedTerminal, @@ -1979,19 +1968,12 @@ fn validate_unique_targets(targets: &[RadrootsTransportTarget]) -> Result<(), Ra fn initial_delivery_target_status( target: &RadrootsTransportTarget, - reticulum_behavior: RadrootsOutboxReticulumBehavior, + _reticulum_behavior: RadrootsOutboxReticulumBehavior, ) -> RadrootsOutboxDeliveryTargetStatus { if target.kind != RadrootsTransportKind::Reticulum { return RadrootsOutboxDeliveryTargetStatus::Pending; } - match reticulum_behavior { - RadrootsOutboxReticulumBehavior::RejectDeliveryAttempts => { - RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented - } - RadrootsOutboxReticulumBehavior::DeferDeliveryPlans => { - RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented - } - } + RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented } fn initial_delivery_plan_status( @@ -2007,12 +1989,7 @@ fn initial_delivery_plan_status( { return RadrootsOutboxDeliveryPlanStatus::Queued; } - if prepared_targets.iter().all(|target| { - target.initial_status == RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented - }) { - return RadrootsOutboxDeliveryPlanStatus::FailedTerminal; - } - RadrootsOutboxDeliveryPlanStatus::Queued + RadrootsOutboxDeliveryPlanStatus::FailedTerminal } async fn existing_idempotent_operation( @@ -2564,7 +2541,6 @@ async fn evaluate_delivery_plans( ) -> Result<PlanEvaluation, RadrootsOutboxError> { let plans = delivery_plans_for_tx(tx, outbox_event_id).await?; let mut all_complete = !plans.is_empty(); - let mut any_failed_terminal = false; let mut any_ready = false; for plan in plans { let targets = delivery_targets_for_plan_tx(tx, plan.delivery_plan_id).await?; @@ -2574,29 +2550,16 @@ async fn evaluate_delivery_plans( .iter() .filter(|target| target.status.is_ready_for_attempt()) .count(); - let deferred_count = status_targets - .iter() - .filter(|target| target.status.is_deferred_until_implemented()) - .count(); - let terminal_failure_count = status_targets - .iter() - .filter(|target| target.status.is_terminal_failure()) - .count(); let plan_status = if satisfied_count >= plan.required_success_count { RadrootsOutboxDeliveryPlanStatus::Complete } else if ready_count > 0 { RadrootsOutboxDeliveryPlanStatus::Queued - } else if terminal_failure_count > 0 || deferred_count > 0 { - RadrootsOutboxDeliveryPlanStatus::FailedTerminal } else { RadrootsOutboxDeliveryPlanStatus::FailedTerminal }; if plan_status != RadrootsOutboxDeliveryPlanStatus::Complete { all_complete = false; } - if plan_status == RadrootsOutboxDeliveryPlanStatus::FailedTerminal { - any_failed_terminal = true; - } if ready_count > 0 { any_ready = true; } @@ -2613,7 +2576,6 @@ async fn evaluate_delivery_plans( } Ok(PlanEvaluation { all_complete, - any_failed_terminal, any_ready, }) } @@ -2664,6 +2626,16 @@ fn outbox_satisfied_target_count( } } +fn reticulum_event_record( + event: Option<RadrootsOutboxEventRecord>, + targets: Vec<RadrootsOutboxDeliveryTargetRecord>, +) -> Option<RadrootsOutboxReticulumEventRecord> { + if targets.is_empty() { + return None; + } + event.map(|event| RadrootsOutboxReticulumEventRecord { event, targets }) +} + fn operation_from_row( row: sqlx::sqlite::SqliteRow, ) -> Result<RadrootsOutboxOperationRecord, RadrootsOutboxError> { @@ -3008,10 +2980,16 @@ fn target_policy_fingerprint( }) .collect::<Vec<_>>(); target_inputs.sort_by(|left, right| { - left.endpoint_fingerprint - .cmp(right.endpoint_fingerprint) - .then_with(|| left.transport_kind.cmp(&right.transport_kind)) - .then_with(|| left.endpoint_uri.cmp(right.endpoint_uri)) + ( + left.endpoint_fingerprint, + left.transport_kind.as_str(), + left.endpoint_uri, + ) + .cmp(&( + right.endpoint_fingerprint, + right.transport_kind.as_str(), + right.endpoint_uri, + )) }); sha256_json(&TargetPolicyDigestInput { satisfaction_policy: satisfaction_policy_storage_value(satisfaction_policy), @@ -3067,14 +3045,6 @@ fn validate_trade_mutation_input( } let parsed = trade_mutation_from_canonical_content(draft.content()) .map_err(|_| RadrootsOutboxError::TradeMutationMetadataMismatch { field: "content" })?; - if parsed.contract_id != draft.contract_id() { - return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { - field: "contract_id", - }); - } - if parsed.mutation_kind().nostr_kind() != draft.kind_u32() { - return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "kind" }); - } if parsed.author_pubkey.as_str() != draft.expected_pubkey_str() { return Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "author_pubkey", @@ -3249,12 +3219,8 @@ fn parse_required_target_policy( .split(',') .map(RadrootsTransportTargetFingerprint::parse) .collect::<Result<Vec<_>, _>>()?; - if required_success_count - != i64::try_from(targets.len()).map_err(|_| RadrootsOutboxError::IntegerRange { - field: "required_success_count", - value: required_success_count, - })? - { + let required_success_count = usize::from(required_count_u16(required_success_count)?); + if required_success_count != targets.len() { return Err(RadrootsOutboxError::InvalidStoredEnum { field: "outbox_delivery_plan.satisfaction_policy", value: stored.to_owned(), @@ -3281,6 +3247,7 @@ fn u32_from_i64(field: &'static str, value: i64) -> Result<u32, RadrootsOutboxEr } #[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] mod tests { use super::*; use radroots_event::ids::{ @@ -3294,7 +3261,7 @@ mod tests { RadrootsTradeCandidateLineV1, RadrootsTradeCandidateTermsV1, RadrootsTradeCanonicalMutationV1, RadrootsTradeEconomicAdjustmentV1, RadrootsTradeEconomicsProfileV1, RadrootsTradeMutationBodyV1, - RadrootsTradeMutationEnvelopeV1, canonical_trade_mutation_content, + RadrootsTradeMutationEnvelopeV1, canonical_jcs_value, canonical_trade_mutation_content, }; use radroots_nostr::prelude::{ RadrootsNostrKeys, RadrootsNostrSecretKey, radroots_nostr_sign_frozen_draft, @@ -3447,6 +3414,22 @@ mod tests { .expect("trade mutation draft") } + fn proposal_draft_with_content( + content: impl Into<String>, + expected_pubkey: &str, + ) -> RadrootsEventDraft { + let valid = trade_mutation_draft(&canonical_trade_proposal()); + RadrootsEventDraft::new( + valid.contract_id(), + valid.kind_u32(), + valid.created_at_u64(), + valid.tags_as_vec(), + content, + expected_pubkey, + ) + .expect("proposal draft with custom content") + } + fn signed_trade_mutation( canonical: &RadrootsTradeCanonicalMutationV1, ) -> (RadrootsEventDraft, RadrootsSignedEvent) { @@ -3601,6 +3584,254 @@ mod tests { ); } + #[test] + fn validation_and_persisted_value_parsers_reject_malformed_contract_data() { + let empty_profile = RadrootsOutboxDeliveryPlanInput::new( + " ", + 1, + RadrootsTransportSatisfactionPolicy::no_wait(), + Vec::new(), + ); + assert!(matches!( + prepare_delivery_plan("event-empty-profile", &empty_profile), + Err(RadrootsOutboxError::EmptyTransportProfileId) + )); + + for (stored, required_count) in [ + ("no_wait", 1), + ("all_unknown", 1), + ("any_unknown", 1), + ("quorum_accepted", 1), + ("quorum_accepted:2", 1), + ("quorum_unknown:1", 1), + ("required_unknown:value", 1), + ("unknown", 0), + ] { + assert!(matches!( + parse_satisfaction_policy(stored, required_count), + Err(RadrootsOutboxError::InvalidStoredEnum { .. }) + )); + } + assert!(parse_satisfaction_class_storage_value("unknown").is_none()); + assert!( + parse_required_target_policy( + "required_accepted:invalid", + "invalid", + RadrootsTransportSatisfactionClass::Accepted, + 1, + ) + .is_err() + ); + let target = nostr_target(NOSTR_PRIMARY_WSS); + assert_eq!( + parse_required_target_policy( + "required_accepted:valid", + target.fingerprint.as_str(), + RadrootsTransportSatisfactionClass::Accepted, + 1, + ) + .expect("valid required-target policy"), + RadrootsTransportSatisfactionPolicy::required_targets( + RadrootsTransportSatisfactionClass::Accepted, + vec![target.fingerprint.clone()], + ) + .expect("required-target policy"), + ); + assert!(matches!( + parse_required_target_policy( + "required_accepted:count-mismatch", + target.fingerprint.as_str(), + RadrootsTransportSatisfactionClass::Accepted, + 0, + ), + Err(RadrootsOutboxError::InvalidStoredEnum { .. }) + )); + let duplicate_fingerprints = format!( + "{},{}", + target.fingerprint.as_str(), + target.fingerprint.as_str() + ); + assert!( + parse_required_target_policy( + "required_accepted:duplicate", + duplicate_fingerprints.as_str(), + RadrootsTransportSatisfactionClass::Accepted, + 2, + ) + .is_err() + ); + assert!(required_count_u16(-1).is_err()); + assert!(required_count_u16(i64::from(u16::MAX) + 1).is_err()); + assert!(u32_from_i64("negative", -1).is_err()); + assert_eq!(bool_i64(false), 0); + assert_eq!(bool_i64(true), 1); + assert!(parse_optional_stored_trade_id("trade_id", None).is_ok()); + assert!(parse_optional_stored_trade_id("trade_id", Some("invalid".to_owned())).is_err()); + assert!( + parse_optional_stored_mutation_id("mutation_id", Some("invalid".to_owned())).is_err() + ); + assert!(matches!( + outcome_kind_for_status(RadrootsOutboxDeliveryTargetStatus::Pending), + RadrootsTransportOutcomeKind::TransportUnavailable + )); + assert!(parse_transport_outcome_kind("invalid", "outcome").is_err()); + for (stored, expected) in [ + ("accepted", RadrootsTransportOutcomeKind::Accepted), + ( + "duplicate_accepted", + RadrootsTransportOutcomeKind::DuplicateAccepted, + ), + ("delivered", RadrootsTransportOutcomeKind::Delivered), + ("forwarded", RadrootsTransportOutcomeKind::Forwarded), + ( + "stored_by_gateway", + RadrootsTransportOutcomeKind::StoredByGateway, + ), + ("seen", RadrootsTransportOutcomeKind::Seen), + ( + "deferred_until_implemented", + RadrootsTransportOutcomeKind::DeferredUntilImplemented, + ), + ("rejected", RadrootsTransportOutcomeKind::Rejected), + ( + "route_unavailable", + RadrootsTransportOutcomeKind::RouteUnavailable, + ), + ( + "payload_too_large", + RadrootsTransportOutcomeKind::PayloadTooLarge, + ), + ("policy_denied", RadrootsTransportOutcomeKind::PolicyDenied), + ("timeout", RadrootsTransportOutcomeKind::Timeout), + ( + "connection_failed", + RadrootsTransportOutcomeKind::ConnectionFailed, + ), + ( + "transport_unavailable", + RadrootsTransportOutcomeKind::TransportUnavailable, + ), + ] { + assert_eq!( + parse_transport_outcome_kind(stored, "outcome").expect("stored outcome"), + expected + ); + } + + let canonical = canonical_trade_proposal(); + let valid_draft = trade_mutation_draft(&canonical); + let valid_hash = sha256_hex(canonical.content.as_bytes()); + assert!(matches!( + validate_trade_mutation_input( + canonical.envelope.trade_id.as_str(), + canonical.mutation_id.as_str(), + valid_hash.as_str(), + &post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "not a trade mutation"), + ), + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "kind" }) + )); + for invalid_hash in ["short".to_owned(), "g".repeat(64)] { + assert!(matches!( + validate_trade_mutation_input( + canonical.envelope.trade_id.as_str(), + canonical.mutation_id.as_str(), + invalid_hash.as_str(), + &valid_draft, + ), + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "canonical_payload_sha256" + }) + )); + } + let invalid_content = proposal_draft_with_content( + format!(" {}", canonical.content), + FIXTURE_ALICE_PUBLIC_KEY_HEX, + ); + assert!(matches!( + validate_trade_mutation_input( + canonical.envelope.trade_id.as_str(), + canonical.mutation_id.as_str(), + valid_hash.as_str(), + &invalid_content, + ), + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "content" }) + )); + + let author_mismatch = + proposal_draft_with_content(canonical.content.clone(), hex_64('b').as_str()); + assert!(matches!( + validate_trade_mutation_input( + canonical.envelope.trade_id.as_str(), + canonical.mutation_id.as_str(), + valid_hash.as_str(), + &author_mismatch, + ), + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "author_pubkey" + }) + )); + assert!(matches!( + validate_trade_mutation_input( + hex_32('2').as_str(), + canonical.mutation_id.as_str(), + valid_hash.as_str(), + &valid_draft, + ), + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "trade_id" }) + )); + assert!(matches!( + validate_trade_mutation_input( + canonical.envelope.trade_id.as_str(), + hex_64('3').as_str(), + valid_hash.as_str(), + &valid_draft, + ), + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "mutation_id" + }) + )); + assert!(matches!( + validate_trade_mutation_input( + canonical.envelope.trade_id.as_str(), + canonical.mutation_id.as_str(), + "0".repeat(64).as_str(), + &valid_draft, + ), + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "canonical_payload_sha256" + }) + )); + + let envelope_without_mutation_id = proposal_envelope(); + let content_without_mutation_id = canonical_jcs_value( + &serde_json::to_value(&envelope_without_mutation_id).expect("proposal value"), + ) + .expect("canonical proposal without mutation id"); + let draft_without_mutation_id = proposal_draft_with_content( + content_without_mutation_id.clone(), + FIXTURE_ALICE_PUBLIC_KEY_HEX, + ); + assert!(matches!( + validate_trade_mutation_input( + envelope_without_mutation_id.trade_id.as_str(), + canonical.mutation_id.as_str(), + sha256_hex(content_without_mutation_id.as_bytes()).as_str(), + &draft_without_mutation_id, + ), + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { + field: "mutation_id" + }) + )); + + let empty_no_wait = RadrootsOutboxDeliveryPlanInput::new( + "transport.none", + 1, + RadrootsTransportSatisfactionPolicy::no_wait(), + Vec::new(), + ); + assert!(prepare_delivery_plan("event-no-wait", &empty_no_wait).is_ok()); + } + fn malformed_reticulum_target(uri: &str) -> RadrootsTransportTarget { let endpoint_uri = RadrootsTransportTargetUri::parse(uri).expect("target uri"); let endpoint_fingerprint = RadrootsTransportTargetFingerprint::from_target( @@ -3681,6 +3912,23 @@ mod tests { ) } + fn trade_mutation_input( + canonical: &RadrootsTradeCanonicalMutationV1, + draft: RadrootsEventDraft, + targets: Vec<RadrootsTransportTarget>, + created_at_ms: i64, + ) -> RadrootsOutboxTradeMutationInput { + RadrootsOutboxTradeMutationInput::new( + "publish_trade_mutation", + canonical.envelope.trade_id.clone(), + canonical.mutation_id.clone(), + sha256_hex(canonical.content.as_bytes()), + draft, + delivery_plan(targets), + created_at_ms, + ) + } + fn fixture_keys() -> RadrootsNostrKeys { let secret_key = RadrootsNostrSecretKey::from_hex(FIXTURE_ALICE_SECRET_KEY_HEX).expect("secret key"); @@ -3857,103 +4105,1055 @@ mod tests { } #[tokio::test] - async fn operation_and_delivery_plan_idempotency_are_split() { - let outbox = RadrootsOutbox::open_memory().await.expect("open"); - let draft = post_draft(hex_64('a').as_str(), "hello"); - let first = outbox - .enqueue_operation(operation_input(draft.clone(), 1_000).with_idempotency_key("idem-a")) - .await - .expect("first"); - let same_plan = outbox - .enqueue_operation(operation_input(draft.clone(), 1_100).with_idempotency_key("idem-a")) - .await - .expect("same"); - let new_plan = outbox - .enqueue_operation( - RadrootsOutboxOperationInput::new( - "publish_post", - draft.clone(), - delivery_plan(vec![nostr_target("wss://relay-3.example.com")]), - 1_200, - ) - .with_idempotency_key("idem-a"), - ) + async fn constructors_preflight_and_transactional_enqueue_cover_public_storage_surfaces() { + let directory = tempfile::tempdir().expect("tempdir"); + let file_outbox = RadrootsOutbox::open_file(directory.path().join("outbox.sqlite")) .await - .expect("new plan"); - - assert_eq!(first.status, RadrootsOutboxEnqueueStatus::Inserted); - assert_eq!(same_plan.status, RadrootsOutboxEnqueueStatus::Existing); - assert_eq!(new_plan.status, RadrootsOutboxEnqueueStatus::Inserted); - assert_eq!(first.operation_id, same_plan.operation_id); - assert_eq!(first.operation_id, new_plan.operation_id); - assert_eq!(first.outbox_event_id, new_plan.outbox_event_id); + .expect("file outbox"); assert_eq!( - first.operation_idempotency_digest, - new_plan.operation_idempotency_digest - ); - assert_ne!( - first.delivery_plan_idempotency_digest, - new_plan.delivery_plan_idempotency_digest + file_outbox + .pragma_journal_mode() + .await + .expect("file journal mode"), + "wal" ); - assert_eq!(table_count(&outbox, "outbox_operations").await, 1); - assert_eq!(table_count(&outbox, "outbox_event").await, 1); - assert_eq!(table_count(&outbox, "outbox_delivery_plan").await, 2); - let conflict = outbox - .enqueue_operation( - operation_input(post_draft(hex_64('a').as_str(), "changed"), 1_300) - .with_idempotency_key("idem-a"), - ) + let options = SqliteConnectOptions::from_str("sqlite::memory:").expect("memory options"); + let pool = SqlitePoolOptions::new() + .max_connections(1) + .connect_with(options) .await - .expect_err("conflict"); + .expect("memory pool"); + let outbox = RadrootsOutbox::open_pool(pool, false) + .await + .expect("pool outbox"); + + let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "preflight"); + let signed_event = radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft) + .expect("signed preflight event"); + let input = signed_operation_input(draft.clone(), signed_event.clone(), 1_000) + .with_idempotency_key("preflight-idem"); + outbox + .preflight_signed_operation_idempotency(&signed_operation_input( + draft.clone(), + signed_event.clone(), + 900, + )) + .await + .expect("preflight without idempotency key"); + outbox + .preflight_signed_operation_idempotency(&input) + .await + .expect("preflight before insert"); + outbox + .enqueue_signed_operation(input.clone()) + .await + .expect("signed insert"); + outbox + .preflight_signed_operation_idempotency(&input) + .await + .expect("matching existing preflight"); + + let changed_draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "preflight changed"); + let changed_signed = radroots_nostr_sign_frozen_draft(&fixture_keys(), &changed_draft) + .expect("changed signed event"); + let conflict = signed_operation_input(changed_draft, changed_signed, 1_100) + .with_idempotency_key("preflight-idem"); assert!(matches!( - conflict, - RadrootsOutboxError::IdempotencyConflict { .. } + outbox + .preflight_signed_operation_idempotency(&conflict) + .await, + Err(RadrootsOutboxError::IdempotencyConflict { .. }) )); + + let transaction_draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "transaction"); + let transaction_signed = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &transaction_draft) + .expect("transaction signed event"); + let transaction_input = + signed_operation_input(transaction_draft.clone(), transaction_signed, 2_000) + .with_idempotency_key("transaction-idem"); + let mut transaction = outbox.pool().begin().await.expect("begin transaction"); + let inserted = outbox + .enqueue_signed_operation_in_transaction(&mut transaction, transaction_input.clone()) + .await + .expect("transaction insert"); + let existing = outbox + .enqueue_signed_operation_in_transaction(&mut transaction, transaction_input.clone()) + .await + .expect("transaction existing"); + assert_eq!(inserted.operation_id, existing.operation_id); + + let transaction_changed = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "transaction changed"); + let transaction_changed_signed = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &transaction_changed) + .expect("transaction changed signed event"); + let transaction_conflict = + signed_operation_input(transaction_changed, transaction_changed_signed, 2_100) + .with_idempotency_key("transaction-idem"); + assert!(matches!( + outbox + .enqueue_signed_operation_in_transaction(&mut transaction, transaction_conflict,) + .await, + Err(RadrootsOutboxError::IdempotencyConflict { .. }) + )); + transaction.commit().await.expect("commit transaction"); } #[tokio::test] - async fn generic_enqueue_rejects_trade_mutation_drafts_before_persistence() { + async fn defensive_storage_decoding_and_remaining_idempotency_edges_are_explicit() { let outbox = RadrootsOutbox::open_memory().await.expect("open"); - let canonical = canonical_trade_proposal(); - let (draft, signed_event) = signed_trade_mutation(&canonical); - - let unsigned_err = outbox - .enqueue_operation(RadrootsOutboxOperationInput::new( - "publish_trade_mutation", + let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "defensive storage"); + let signed_event = radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft) + .expect("signed defensive event"); + let signed_receipt = outbox + .enqueue_signed_operation(signed_operation_input( draft.clone(), - delivery_plan(vec![nostr_target(NOSTR_PRIMARY_WSS)]), + signed_event.clone(), 1_000, )) .await - .expect_err("generic unsigned trade rejection"); - assert!(matches!( - unsigned_err, - RadrootsOutboxError::TradeMutationRequiresSemanticOutbox - )); + .expect("signed insert without idempotency key"); - let signed_err = outbox - .enqueue_signed_operation(RadrootsOutboxSignedOperationInput::new( - "publish_trade_mutation", - draft, - signed_event, - delivery_plan(vec![nostr_target(NOSTR_PRIMARY_WSS)]), - true, - 1_007, - 1_000, - )) + let keyed_draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "keyed public insert"); + let keyed_signed = radroots_nostr_sign_frozen_draft(&fixture_keys(), &keyed_draft) + .expect("keyed signed event"); + outbox + .enqueue_signed_operation( + signed_operation_input(keyed_draft, keyed_signed, 1_100) + .with_idempotency_key("public-signed-conflict"), + ) .await - .expect_err("generic signed trade rejection"); + .expect("keyed public signed insert"); + let changed_draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "keyed public conflict"); + let changed_signed = radroots_nostr_sign_frozen_draft(&fixture_keys(), &changed_draft) + .expect("changed signed event"); assert!(matches!( - signed_err, - RadrootsOutboxError::TradeMutationRequiresSemanticOutbox + outbox + .enqueue_signed_operation( + signed_operation_input(changed_draft, changed_signed, 1_200) + .with_idempotency_key("public-signed-conflict"), + ) + .await, + Err(RadrootsOutboxError::IdempotencyConflict { .. }) )); - 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); - } - #[tokio::test] + let canonical = canonical_trade_proposal(); + let invalid_trade_draft = post_draft( + FIXTURE_ALICE_PUBLIC_KEY_HEX, + "not a semantic trade mutation", + ); + let invalid_trade_signed = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &invalid_trade_draft) + .expect("invalid trade draft still forms a valid signed event"); + assert!(matches!( + outbox + .preflight_signed_trade_mutation_idempotency(&signed_trade_mutation_input( + &canonical, + invalid_trade_draft.clone(), + invalid_trade_signed, + vec![nostr_target("wss://invalid-trade-preflight.example")], + 1_300, + )) + .await, + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "kind" }) + )); + assert!(matches!( + outbox + .enqueue_trade_mutation_operation(trade_mutation_input( + &canonical, + invalid_trade_draft, + vec![nostr_target("wss://invalid-trade-enqueue.example")], + 1_400, + )) + .await, + Err(RadrootsOutboxError::TradeMutationMetadataMismatch { field: "kind" }) + )); + + let (trade_draft, trade_signed_event) = signed_trade_mutation(&canonical); + outbox + .preflight_signed_trade_mutation_idempotency(&signed_trade_mutation_input( + &canonical, + trade_draft.clone(), + trade_signed_event, + vec![nostr_target("wss://trade-preflight-no-key.example")], + 1_500, + )) + .await + .expect("trade preflight without idempotency key"); + let signed_trade_without_key = RadrootsOutbox::open_memory() + .await + .expect("open isolated signed trade outbox"); + let (isolated_trade_draft, isolated_trade_event) = signed_trade_mutation(&canonical); + signed_trade_without_key + .enqueue_signed_trade_mutation_operation(signed_trade_mutation_input( + &canonical, + isolated_trade_draft, + isolated_trade_event, + vec![nostr_target("wss://signed-trade-no-key.example")], + 1_550, + )) + .await + .expect("signed trade insert without idempotency key"); + let trade_receipt = outbox + .enqueue_trade_mutation_operation(trade_mutation_input( + &canonical, + trade_draft, + vec![nostr_target("wss://trade-storage.example")], + 1_600, + )) + .await + .expect("trade insert without idempotency key"); + + let transaction_draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "transaction no key"); + let transaction_signed = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &transaction_draft) + .expect("transaction signed event"); + let mut transaction = outbox.pool().begin().await.expect("begin transaction"); + outbox + .enqueue_signed_operation_in_transaction( + &mut transaction, + signed_operation_input(transaction_draft, transaction_signed, 1_700), + ) + .await + .expect("transaction insert without idempotency key"); + assert!(matches!( + claimed_event_identity_tx(&mut transaction, i64::MAX, "missing").await, + Err(RadrootsOutboxError::EventNotFound(_)) + )); + assert!(matches!( + claimed_event_identity_tx( + &mut transaction, + signed_receipt.outbox_event_id, + "wrong-token", + ) + .await, + Err(RadrootsOutboxError::ClaimTokenMismatch { .. }) + )); + assert!(matches!( + claimed_event_from_tx( + &mut transaction, + signed_receipt.outbox_event_id, + RadrootsOutboxEventState::Publishing, + None, + "unused-token", + ) + .await, + Err(RadrootsOutboxError::MissingActiveDeliveryPlan { .. }) + )); + transaction.commit().await.expect("commit transaction"); + + let event = outbox + .get_event(signed_receipt.outbox_event_id) + .await + .expect("event query") + .expect("event"); + let plans = outbox + .delivery_plans(signed_receipt.outbox_event_id) + .await + .expect("plans"); + let targets = outbox + .delivery_targets(signed_receipt.outbox_event_id) + .await + .expect("targets"); + assert!(reticulum_event_record(None, targets.clone()).is_none()); + assert!(reticulum_event_record(Some(event.clone()), Vec::new()).is_none()); + assert!(reticulum_event_record(Some(event.clone()), targets.clone()).is_some()); + assert_eq!( + outbox_satisfied_target_count( + &RadrootsTransportSatisfactionPolicy::no_wait(), + &targets, + ), + 0 + ); + assert_eq!( + outbox_satisfied_target_count( + &RadrootsTransportSatisfactionPolicy::any_accepted(), + &targets, + ), + 0 + ); + let mut accepted_target = targets[0].clone(); + accepted_target.status = RadrootsOutboxDeliveryTargetStatus::Accepted; + assert_eq!( + outbox_satisfied_target_count( + &RadrootsTransportSatisfactionPolicy::required_targets( + RadrootsTransportSatisfactionClass::Accepted, + vec![accepted_target.endpoint_fingerprint.clone()], + ) + .expect("required-target policy"), + &[accepted_target], + ), + 1 + ); + + sqlx::query("PRAGMA ignore_check_constraints = ON") + .execute(outbox.pool()) + .await + .expect("disable check constraints for defensive decoding"); + + sqlx::query("UPDATE outbox_operations SET trade_id = 'invalid' WHERE operation_id = ?") + .bind(trade_receipt.operation_id) + .execute(outbox.pool()) + .await + .expect("corrupt trade id"); + assert!( + outbox + .get_operation(trade_receipt.operation_id) + .await + .is_err() + ); + sqlx::query("UPDATE outbox_operations SET trade_id = ?, mutation_id = 'invalid' WHERE operation_id = ?") + .bind(canonical.envelope.trade_id.as_str()) + .bind(trade_receipt.operation_id) + .execute(outbox.pool()) + .await + .expect("restore trade id and corrupt mutation id"); + assert!( + outbox + .get_operation(trade_receipt.operation_id) + .await + .is_err() + ); + sqlx::query("UPDATE outbox_operations SET mutation_id = ? WHERE operation_id = ?") + .bind(canonical.mutation_id.as_str()) + .bind(trade_receipt.operation_id) + .execute(outbox.pool()) + .await + .expect("restore mutation id"); + + sqlx::query("UPDATE outbox_event SET raw_event_json = NULL WHERE outbox_event_id = ?") + .bind(signed_receipt.outbox_event_id) + .execute(outbox.pool()) + .await + .expect("remove raw signed event JSON"); + assert!(matches!( + outbox.get_event(signed_receipt.outbox_event_id).await, + Err(RadrootsOutboxError::StoredSignedEventMissingRawJson(_)) + )); + sqlx::query("UPDATE outbox_event SET raw_event_json = ?, signed_event_json = NULL WHERE outbox_event_id = ?") + .bind(signed_event.raw_json()) + .bind(signed_receipt.outbox_event_id) + .execute(outbox.pool()) + .await + .expect("restore raw JSON and remove signed wire JSON"); + assert!(matches!( + outbox.get_event(signed_receipt.outbox_event_id).await, + Err(RadrootsOutboxError::StoredRawEventMissingSignedEvent(_)) + )); + sqlx::query( + "UPDATE outbox_event SET signed_event_json = ?, event_id = ? WHERE outbox_event_id = ?", + ) + .bind(signed_event_wire_json(&signed_event).expect("signed wire JSON")) + .bind(hex_64('0')) + .bind(signed_receipt.outbox_event_id) + .execute(outbox.pool()) + .await + .expect("restore wire JSON and corrupt event id"); + assert!(matches!( + outbox.get_event(signed_receipt.outbox_event_id).await, + Err(RadrootsOutboxError::SignedEventIdMismatch { .. }) + )); + sqlx::query("UPDATE outbox_event SET event_id = ? WHERE outbox_event_id = ?") + .bind(event.event_id.as_str()) + .bind(signed_receipt.outbox_event_id) + .execute(outbox.pool()) + .await + .expect("restore event id"); + + sqlx::query("UPDATE outbox_delivery_plan SET satisfaction_policy = 'invalid' WHERE delivery_plan_id = ?") + .bind(plans[0].delivery_plan_id) + .execute(outbox.pool()) + .await + .expect("corrupt satisfaction policy"); + assert!( + outbox + .delivery_plans(signed_receipt.outbox_event_id) + .await + .is_err() + ); + sqlx::query("UPDATE outbox_delivery_plan SET satisfaction_policy = ?, target_policy_version = -1 WHERE delivery_plan_id = ?") + .bind(satisfaction_policy_storage_value(&plans[0].satisfaction_policy)) + .bind(plans[0].delivery_plan_id) + .execute(outbox.pool()) + .await + .expect("restore policy and corrupt version"); + assert!( + outbox + .delivery_plans(signed_receipt.outbox_event_id) + .await + .is_err() + ); + sqlx::query( + "UPDATE outbox_delivery_plan SET target_policy_version = ? WHERE delivery_plan_id = ?", + ) + .bind(i64::from(plans[0].target_policy_version)) + .bind(plans[0].delivery_plan_id) + .execute(outbox.pool()) + .await + .expect("restore target policy version"); + + sqlx::query("UPDATE outbox_delivery_target SET endpoint_fingerprint = 'invalid' WHERE delivery_target_id = ?") + .bind(targets[0].delivery_target_id) + .execute(outbox.pool()) + .await + .expect("corrupt target fingerprint"); + assert!( + outbox + .delivery_targets(signed_receipt.outbox_event_id) + .await + .is_err() + ); + sqlx::query("UPDATE outbox_delivery_target SET endpoint_fingerprint = ? WHERE delivery_target_id = ?") + .bind(targets[0].endpoint_fingerprint.as_str()) + .bind(targets[0].delivery_target_id) + .execute(outbox.pool()) + .await + .expect("restore target fingerprint"); + + let attempt = sqlx::query( + "INSERT INTO outbox_delivery_attempt(delivery_plan_id, delivery_target_id, status, outcome_kind, attempted_at_ms) VALUES (?, ?, 'accepted', 'accepted', ?)", + ) + .bind(plans[0].delivery_plan_id) + .bind(targets[0].delivery_target_id) + .bind(2_000_i64) + .execute(outbox.pool()) + .await + .expect("insert delivery attempt"); + sqlx::query("UPDATE outbox_delivery_attempt SET outcome_kind = 'invalid' WHERE delivery_attempt_id = ?") + .bind(attempt.last_insert_rowid()) + .execute(outbox.pool()) + .await + .expect("corrupt attempt outcome"); + assert!( + outbox + .delivery_attempts(targets[0].delivery_target_id) + .await + .is_err() + ); + + sqlx::query("PRAGMA ignore_check_constraints = OFF") + .execute(outbox.pool()) + .await + .expect("restore check constraints"); + } + + #[tokio::test] + async fn signed_plan_lifecycle_and_reticulum_limit_cover_boundary_states() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + assert!(matches!( + outbox.reticulum_events(None, usize::MAX).await, + Err(RadrootsOutboxError::IntegerRange { + field: "reticulum_events.limit", + .. + }) + )); + + let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "plan lifecycle"); + let signed_event = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("sign plan lifecycle"); + let first = outbox + .enqueue_signed_operation( + RadrootsOutboxSignedOperationInput::new( + "publish_post", + draft.clone(), + signed_event.clone(), + delivery_plan(vec![nostr_target("wss://plan-one.example")]), + true, + 1_007, + 1_000, + ) + .with_idempotency_key("plan-lifecycle"), + ) + .await + .expect("first plan"); + outbox + .enqueue_signed_operation( + RadrootsOutboxSignedOperationInput::new( + "publish_post", + draft, + signed_event, + delivery_plan(vec![nostr_target("wss://plan-two.example")]), + true, + 1_107, + 1_100, + ) + .with_idempotency_key("plan-lifecycle"), + ) + .await + .expect("second plan"); + let plans = outbox + .delivery_plans(first.outbox_event_id) + .await + .expect("plans"); + assert_eq!(plans.len(), 2); + + let mut transaction = outbox.pool().begin().await.expect("begin lifecycle"); + sqlx::query( + "UPDATE outbox_delivery_plan SET status = 'complete' WHERE delivery_plan_id = ?", + ) + .bind(plans[0].delivery_plan_id) + .execute(&mut *transaction) + .await + .expect("complete first plan"); + sqlx::query( + "UPDATE outbox_delivery_plan SET status = 'cancelled' WHERE delivery_plan_id = ?", + ) + .bind(plans[1].delivery_plan_id) + .execute(&mut *transaction) + .await + .expect("cancel second plan"); + assert_eq!( + signed_event_lifecycle_for_plans(&mut transaction, first.outbox_event_id) + .await + .expect("mixed lifecycle"), + ( + RadrootsOutboxEventState::Signed, + RadrootsOutboxOperationStatus::Queued, + ) + ); + sqlx::query( + "UPDATE outbox_delivery_plan SET status = 'cancelled' WHERE outbox_event_id = ?", + ) + .bind(first.outbox_event_id) + .execute(&mut *transaction) + .await + .expect("cancel all plans"); + assert_eq!( + signed_event_lifecycle_for_plans(&mut transaction, first.outbox_event_id) + .await + .expect("cancelled lifecycle"), + ( + RadrootsOutboxEventState::Cancelled, + RadrootsOutboxOperationStatus::Cancelled, + ) + ); + sqlx::query("DELETE FROM outbox_delivery_plan WHERE outbox_event_id = ?") + .bind(first.outbox_event_id) + .execute(&mut *transaction) + .await + .expect("delete plans"); + assert_eq!( + signed_event_lifecycle_for_plans(&mut transaction, first.outbox_event_id) + .await + .expect("empty lifecycle"), + ( + RadrootsOutboxEventState::Signed, + RadrootsOutboxOperationStatus::Queued, + ) + ); + transaction.rollback().await.expect("rollback lifecycle"); + } + + #[tokio::test] + async fn operation_and_delivery_plan_idempotency_are_split() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let draft = post_draft(hex_64('a').as_str(), "hello"); + let first = outbox + .enqueue_operation(operation_input(draft.clone(), 1_000).with_idempotency_key("idem-a")) + .await + .expect("first"); + let same_plan = outbox + .enqueue_operation(operation_input(draft.clone(), 1_100).with_idempotency_key("idem-a")) + .await + .expect("same"); + let new_plan = outbox + .enqueue_operation( + RadrootsOutboxOperationInput::new( + "publish_post", + draft.clone(), + delivery_plan(vec![nostr_target("wss://relay-3.example.com")]), + 1_200, + ) + .with_idempotency_key("idem-a"), + ) + .await + .expect("new plan"); + + assert_eq!(first.status, RadrootsOutboxEnqueueStatus::Inserted); + assert_eq!(same_plan.status, RadrootsOutboxEnqueueStatus::Existing); + assert_eq!(new_plan.status, RadrootsOutboxEnqueueStatus::Inserted); + assert_eq!(first.operation_id, same_plan.operation_id); + assert_eq!(first.operation_id, new_plan.operation_id); + assert_eq!(first.outbox_event_id, new_plan.outbox_event_id); + assert_eq!( + first.operation_idempotency_digest, + new_plan.operation_idempotency_digest + ); + assert_ne!( + first.delivery_plan_idempotency_digest, + new_plan.delivery_plan_idempotency_digest + ); + assert_eq!(table_count(&outbox, "outbox_operations").await, 1); + assert_eq!(table_count(&outbox, "outbox_event").await, 1); + assert_eq!(table_count(&outbox, "outbox_delivery_plan").await, 2); + + let conflict = outbox + .enqueue_operation( + operation_input(post_draft(hex_64('a').as_str(), "changed"), 1_300) + .with_idempotency_key("idem-a"), + ) + .await + .expect_err("conflict"); + assert!(matches!( + conflict, + RadrootsOutboxError::IdempotencyConflict { .. } + )); + } + + #[tokio::test] + async fn claim_retry_recovery_and_stale_update_guards_cover_worker_lifecycle() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + assert!( + outbox + .claim_next_ready_event("worker", "empty", 100, 0) + .await + .expect("empty next claim") + .is_none() + ); + assert!( + outbox + .claim_ready_signed_event(99_999, "worker", "missing", 100, 0) + .await + .expect("missing direct claim") + .is_none() + ); + + let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "retry signing"); + let receipt = outbox + .enqueue_operation(operation_input(draft.clone(), 1_000)) + .await + .expect("enqueue unsigned"); + let signed = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("sign queued event"); + let claimed = outbox + .claim_next_ready_event("signer", "sign-claim", 2_000, 1_000) + .await + .expect("claim unsigned") + .expect("unsigned claim"); + assert_eq!(claimed.outbox_event_id, receipt.outbox_event_id); + assert!(matches!( + outbox + .complete_signing( + receipt.outbox_event_id, + "wrong-claim", + signed.clone(), + 1_010, + ) + .await, + Err(RadrootsOutboxError::ClaimTokenMismatch { .. }) + )); + + let other_draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "different event"); + let other_signed = radroots_nostr_sign_frozen_draft(&fixture_keys(), &other_draft) + .expect("sign different event"); + assert!(matches!( + outbox + .complete_signing(receipt.outbox_event_id, "sign-claim", other_signed, 1_020,) + .await, + Err(RadrootsOutboxError::SignedEventIdMismatch { .. }) + )); + + sqlx::query( + "CREATE TRIGGER ignore_complete_signing BEFORE UPDATE OF signed_event_json ON outbox_event WHEN OLD.outbox_event_id = 1 AND NEW.signed_event_json IS NOT NULL BEGIN SELECT RAISE(IGNORE); END", + ) + .execute(outbox.pool()) + .await + .expect("complete signing trigger"); + assert!(matches!( + outbox + .complete_signing(receipt.outbox_event_id, "sign-claim", signed.clone(), 1_030,) + .await, + Err(RadrootsOutboxError::ClaimTokenMismatch { .. }) + )); + sqlx::query("DROP TRIGGER ignore_complete_signing") + .execute(outbox.pool()) + .await + .expect("drop complete signing trigger"); + + outbox + .mark_sign_retryable( + receipt.outbox_event_id, + "sign-claim", + "try signing again", + 1_040, + 1_035, + ) + .await + .expect("mark sign retryable"); + let reclaimed = outbox + .claim_next_ready_event("signer", "sign-claim-two", 2_100, 1_040) + .await + .expect("reclaim unsigned") + .expect("reclaimed unsigned"); + assert_eq!(reclaimed.state, RadrootsOutboxEventState::Signing); + outbox + .complete_signing(receipt.outbox_event_id, "sign-claim-two", signed, 1_050) + .await + .expect("complete signing"); + + let direct = outbox + .claim_ready_signed_event( + receipt.outbox_event_id, + "publisher", + "publish-claim", + 2_200, + 1_050, + ) + .await + .expect("claim signed directly") + .expect("direct signed claim"); + assert_eq!(direct.state, RadrootsOutboxEventState::Publishing); + assert!( + outbox + .claim_ready_signed_event( + receipt.outbox_event_id, + "publisher", + "publish-race", + 2_300, + 1_060, + ) + .await + .expect("stale direct claim") + .is_none() + ); + outbox + .mark_publish_retryable( + receipt.outbox_event_id, + "publish-claim", + "try publishing again", + 1_070, + 1_065, + ) + .await + .expect("mark publish retryable"); + + let expiring = outbox + .claim_next_ready_signed_event("publisher", "expiring", 1_080, 1_070) + .await + .expect("claim expiring") + .expect("expiring claim"); + assert_eq!(expiring.outbox_event_id, receipt.outbox_event_id); + assert_eq!( + outbox + .recover_expired_claims(1_080) + .await + .expect("recover expired claim"), + 1 + ); + assert_eq!( + outbox + .get_event(receipt.outbox_event_id) + .await + .expect("event after recovery") + .expect("event") + .state, + RadrootsOutboxEventState::PublishRetryable + ); + + let guarded = outbox + .claim_next_ready_signed_event("publisher", "guarded", 2_400, 1_080) + .await + .expect("claim for guarded retry") + .expect("guarded claim"); + assert_eq!(guarded.outbox_event_id, receipt.outbox_event_id); + sqlx::query( + "CREATE TRIGGER ignore_publish_retry BEFORE UPDATE OF state ON outbox_event WHEN NEW.state = 'publish_retryable' BEGIN SELECT RAISE(IGNORE); END", + ) + .execute(outbox.pool()) + .await + .expect("publish retry trigger"); + assert!(matches!( + outbox + .mark_publish_retryable( + receipt.outbox_event_id, + "guarded", + "ignored update", + 2_500, + 1_090, + ) + .await, + Err(RadrootsOutboxError::ClaimTokenMismatch { .. }) + )); + sqlx::query("DROP TRIGGER ignore_publish_retry") + .execute(outbox.pool()) + .await + .expect("drop publish retry trigger"); + } + + #[tokio::test] + async fn claim_compare_and_swap_guards_report_lost_worker_races() { + let unsigned_outbox = RadrootsOutbox::open_memory().await.expect("open unsigned"); + let unsigned_draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "unsigned race"); + unsigned_outbox + .enqueue_operation(operation_input(unsigned_draft, 1_000)) + .await + .expect("enqueue unsigned race"); + sqlx::query( + "CREATE TRIGGER ignore_unsigned_claim BEFORE UPDATE OF claim_token ON outbox_event WHEN NEW.claim_token = 'race-next' BEGIN SELECT RAISE(IGNORE); END", + ) + .execute(unsigned_outbox.pool()) + .await + .expect("unsigned claim trigger"); + assert!( + unsigned_outbox + .claim_next_ready_event("worker", "race-next", 2_000, 1_000) + .await + .expect("lost unsigned claim race") + .is_none() + ); + + let signed_outbox = RadrootsOutbox::open_memory().await.expect("open signed"); + let signed_draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "signed race"); + let signed_event = radroots_nostr_sign_frozen_draft(&fixture_keys(), &signed_draft) + .expect("sign race event"); + signed_outbox + .enqueue_signed_operation(signed_operation_input(signed_draft, signed_event, 1_000)) + .await + .expect("enqueue signed race"); + sqlx::query( + "CREATE TRIGGER ignore_signed_claim BEFORE UPDATE OF claim_token ON outbox_event WHEN NEW.claim_token = 'race-signed' BEGIN SELECT RAISE(IGNORE); END", + ) + .execute(signed_outbox.pool()) + .await + .expect("signed claim trigger"); + assert!( + signed_outbox + .claim_next_ready_signed_event("worker", "race-signed", 2_000, 1_000) + .await + .expect("lost signed claim race") + .is_none() + ); + } + + #[tokio::test] + async fn claimed_publish_mutations_enforce_plan_identity_and_atomic_updates() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let draft = post_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "publish guards"); + let signed_event = + radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft).expect("sign publish guards"); + let receipt = outbox + .enqueue_signed_operation(signed_operation_input(draft, signed_event, 1_000)) + .await + .expect("enqueue publish guards"); + let claimed = outbox + .claim_next_ready_signed_event("publisher", "publish-guards", 2_000, 1_000) + .await + .expect("claim publish guards") + .expect("claimed publish guards"); + let target_id = claimed.delivery_targets[0].delivery_target_id; + + sqlx::query( + "CREATE TRIGGER ignore_target_update BEFORE UPDATE OF status ON outbox_delivery_target WHEN NEW.status = 'accepted' BEGIN SELECT RAISE(IGNORE); END", + ) + .execute(outbox.pool()) + .await + .expect("target update trigger"); + assert!(matches!( + outbox + .mark_delivery_target_accepted( + receipt.outbox_event_id, + "publish-guards", + target_id, + 1_010, + ) + .await, + Err(RadrootsOutboxError::DeliveryTargetNotFound(_)) + )); + sqlx::query("DROP TRIGGER ignore_target_update") + .execute(outbox.pool()) + .await + .expect("drop target update trigger"); + + for (trigger_name, operation) in [ + ("ignore_complete_publish", "complete"), + ("ignore_terminal_publish", "terminal"), + ("ignore_cancel_publish", "cancel"), + ] { + let trigger = format!( + "CREATE TRIGGER {trigger_name} BEFORE UPDATE OF state ON outbox_event WHEN OLD.claim_token = 'publish-guards' AND NEW.claim_token IS NULL BEGIN SELECT RAISE(IGNORE); END" + ); + sqlx::query(sqlx::AssertSqlSafe(trigger)) + .execute(outbox.pool()) + .await + .expect("event update trigger"); + let error = match operation { + "complete" => outbox + .complete_publish_attempt( + receipt.outbox_event_id, + "publish-guards", + "retry", + "terminal", + 1_100, + 1_020, + ) + .await + .expect_err("ignored complete update"), + "terminal" => outbox + .mark_active_delivery_plan_failed_terminal( + receipt.outbox_event_id, + "publish-guards", + "terminal", + 1_020, + ) + .await + .expect_err("ignored terminal update"), + "cancel" => outbox + .cancel_claimed_event(receipt.outbox_event_id, "publish-guards", 1_020) + .await + .expect_err("ignored cancel update"), + _ => unreachable!(), + }; + assert!(matches!( + error, + RadrootsOutboxError::ClaimTokenMismatch { .. } + )); + let drop_trigger = format!("DROP TRIGGER {trigger_name}"); + sqlx::query(sqlx::AssertSqlSafe(drop_trigger)) + .execute(outbox.pool()) + .await + .expect("drop event update trigger"); + } + + sqlx::query("PRAGMA ignore_check_constraints = ON") + .execute(outbox.pool()) + .await + .expect("ignore target check constraint"); + sqlx::query( + "UPDATE outbox_delivery_target SET status = 'invalid' WHERE delivery_target_id = ?", + ) + .bind(target_id) + .execute(outbox.pool()) + .await + .expect("corrupt target status"); + assert!(matches!( + outbox + .mark_delivery_target_accepted( + receipt.outbox_event_id, + "publish-guards", + target_id, + 1_030, + ) + .await, + Err(RadrootsOutboxError::InvalidStoredEnum { .. }) + )); + sqlx::query( + "UPDATE outbox_delivery_target SET status = 'pending' WHERE delivery_target_id = ?", + ) + .bind(target_id) + .execute(outbox.pool()) + .await + .expect("restore target status"); + sqlx::query("PRAGMA ignore_check_constraints = OFF") + .execute(outbox.pool()) + .await + .expect("restore target check constraints"); + + sqlx::query( + "UPDATE outbox_event SET active_delivery_plan_id = NULL WHERE outbox_event_id = ?", + ) + .bind(receipt.outbox_event_id) + .execute(outbox.pool()) + .await + .expect("remove active plan"); + assert!(matches!( + outbox + .complete_publish_attempt( + receipt.outbox_event_id, + "publish-guards", + "retry", + "terminal", + 1_100, + 1_040, + ) + .await, + Err(RadrootsOutboxError::MissingActiveDeliveryPlan { .. }) + )); + assert!(matches!( + outbox + .mark_active_delivery_plan_failed_terminal( + receipt.outbox_event_id, + "publish-guards", + "terminal", + 1_040, + ) + .await, + Err(RadrootsOutboxError::MissingActiveDeliveryPlan { .. }) + )); + assert!(matches!( + outbox + .mark_delivery_target_accepted( + receipt.outbox_event_id, + "publish-guards", + target_id, + 1_040, + ) + .await, + Err(RadrootsOutboxError::MissingActiveDeliveryPlan { .. }) + )); + + assert!(matches!( + outbox + .mark_sign_retryable(99_999, "missing", "missing", 1_100, 1_050) + .await, + Err(RadrootsOutboxError::EventNotFound(_)) + )); + assert!(matches!( + outbox + .mark_sign_retryable( + receipt.outbox_event_id, + "wrong-token", + "wrong token", + 1_100, + 1_050, + ) + .await, + Err(RadrootsOutboxError::ClaimTokenMismatch { .. }) + )); + } + + #[tokio::test] + async fn generic_enqueue_rejects_trade_mutation_drafts_before_persistence() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let canonical = canonical_trade_proposal(); + let (draft, signed_event) = signed_trade_mutation(&canonical); + + let unsigned_err = outbox + .enqueue_operation(RadrootsOutboxOperationInput::new( + "publish_trade_mutation", + draft.clone(), + delivery_plan(vec![nostr_target(NOSTR_PRIMARY_WSS)]), + 1_000, + )) + .await + .expect_err("generic unsigned trade rejection"); + assert!(matches!( + unsigned_err, + RadrootsOutboxError::TradeMutationRequiresSemanticOutbox + )); + + let signed_err = outbox + .enqueue_signed_operation(RadrootsOutboxSignedOperationInput::new( + "publish_trade_mutation", + draft, + signed_event, + delivery_plan(vec![nostr_target(NOSTR_PRIMARY_WSS)]), + true, + 1_007, + 1_000, + )) + .await + .expect_err("generic signed trade rejection"); + assert!(matches!( + signed_err, + RadrootsOutboxError::TradeMutationRequiresSemanticOutbox + )); + 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); + } + + #[tokio::test] async fn semantic_trade_mutation_enqueue_persists_metadata_and_deduplicates_by_mutation() { let outbox = RadrootsOutbox::open_memory().await.expect("open"); let canonical = canonical_trade_proposal(); @@ -4033,6 +5233,169 @@ mod tests { } #[tokio::test] + async fn unsigned_trade_and_defensive_idempotency_paths_preserve_semantic_identity() { + let outbox = RadrootsOutbox::open_memory().await.expect("open"); + let canonical = canonical_trade_proposal(); + let (draft, signed_event) = signed_trade_mutation(&canonical); + let unsigned = |uri: &str, created_at_ms| { + trade_mutation_input( + &canonical, + draft.clone(), + vec![nostr_target(uri)], + created_at_ms, + ) + .with_idempotency_key("unsigned-trade-idem") + }; + let signed = |uri: &str, created_at_ms| { + signed_trade_mutation_input( + &canonical, + draft.clone(), + signed_event.clone(), + vec![nostr_target(uri)], + created_at_ms, + ) + .with_idempotency_key("unsigned-trade-idem") + }; + + let first = outbox + .enqueue_trade_mutation_operation(unsigned("wss://unsigned-one.example", 1_000)) + .await + .expect("unsigned insert"); + let existing = outbox + .enqueue_trade_mutation_operation(unsigned("wss://unsigned-one.example", 1_100)) + .await + .expect("unsigned existing plan"); + let inserted_plan = outbox + .enqueue_trade_mutation_operation(unsigned("wss://unsigned-two.example", 1_200)) + .await + .expect("unsigned new plan"); + assert_eq!(existing.status, RadrootsOutboxEnqueueStatus::Existing); + assert_eq!(inserted_plan.status, RadrootsOutboxEnqueueStatus::Inserted); + + outbox + .preflight_signed_trade_mutation_idempotency(&signed( + "wss://unsigned-three.example", + 1_300, + )) + .await + .expect("matching trade preflight"); + + sqlx::query("UPDATE outbox_operations SET mutation_id = ? WHERE operation_id = ?") + .bind("stored-other-mutation") + .bind(first.operation_id) + .execute(outbox.pool()) + .await + .expect("move stored mutation identity"); + let idempotent_existing = outbox + .enqueue_trade_mutation_operation(unsigned("wss://unsigned-one.example", 1_400)) + .await + .expect("unsigned idempotent existing plan"); + let idempotent_inserted = outbox + .enqueue_trade_mutation_operation(unsigned("wss://unsigned-four.example", 1_500)) + .await + .expect("unsigned idempotent new plan"); + assert_eq!( + idempotent_existing.status, + RadrootsOutboxEnqueueStatus::Existing + ); + assert_eq!( + idempotent_inserted.status, + RadrootsOutboxEnqueueStatus::Inserted + ); + + sqlx::query( + "UPDATE outbox_operations SET operation_idempotency_digest = 'corrupt' WHERE operation_id = ?", + ) + .bind(first.operation_id) + .execute(outbox.pool()) + .await + .expect("corrupt stored digest"); + assert!(matches!( + outbox + .enqueue_trade_mutation_operation(unsigned( + "wss://unsigned-conflict.example", + 1_600, + )) + .await, + Err(RadrootsOutboxError::IdempotencyConflict { .. }) + )); + assert!(matches!( + outbox + .preflight_signed_trade_mutation_idempotency(&signed( + "wss://signed-preflight-conflict.example", + 1_700, + )) + .await, + Err(RadrootsOutboxError::IdempotencyConflict { .. }) + )); + assert!(matches!( + outbox + .enqueue_signed_trade_mutation_operation(signed( + "wss://signed-idempotent-conflict.example", + 1_800, + )) + .await, + Err(RadrootsOutboxError::IdempotencyConflict { .. }) + )); + + sqlx::query( + "UPDATE outbox_operations SET operation_idempotency_digest = ? WHERE operation_id = ?", + ) + .bind(first.operation_idempotency_digest.as_str()) + .bind(first.operation_id) + .execute(outbox.pool()) + .await + .expect("restore stored digest"); + let signed_idempotent = outbox + .enqueue_signed_trade_mutation_operation(signed( + "wss://signed-idempotent.example", + 1_900, + )) + .await + .expect("signed idempotent branch"); + assert_eq!( + signed_idempotent.status, + RadrootsOutboxEnqueueStatus::Inserted + ); + + sqlx::query( + "UPDATE outbox_operations SET mutation_id = ?, operation_idempotency_digest = 'corrupt' WHERE operation_id = ?", + ) + .bind(canonical.mutation_id.as_str()) + .bind(first.operation_id) + .execute(outbox.pool()) + .await + .expect("restore mutation and corrupt digest"); + assert!(matches!( + outbox + .enqueue_trade_mutation_operation(unsigned( + "wss://unsigned-mutation-conflict.example", + 2_000, + )) + .await, + Err(RadrootsOutboxError::IdempotencyConflict { .. }) + )); + assert!(matches!( + outbox + .preflight_signed_trade_mutation_idempotency(&signed( + "wss://signed-mutation-preflight.example", + 2_100, + )) + .await, + Err(RadrootsOutboxError::IdempotencyConflict { .. }) + )); + assert!(matches!( + outbox + .enqueue_signed_trade_mutation_operation(signed( + "wss://signed-mutation-conflict.example", + 2_200, + )) + .await, + Err(RadrootsOutboxError::IdempotencyConflict { .. }) + )); + } + + #[tokio::test] async fn semantic_trade_mutation_enqueue_rejects_payload_hash_mismatch() { let outbox = RadrootsOutbox::open_memory().await.expect("open"); let canonical = canonical_trade_proposal();