lib

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

commit ecc79cd00cf47c0ccde8d9e5ed2da51645873926
parent 742246bc06d734a17c0a6d43551432c2c7eb4ba9
Author: triesap <tyson@radroots.org>
Date:   Mon, 27 Jul 2026 10:56:09 +0000

transport: preserve recoverable publication effects

- persist per-target dispatch intent, receipts, and uncertain recovery state
- repair accepted observations idempotently without republishing signed bytes
- retain typed transport failures and stable retry dispatch identities
- cover persistence boundaries, concurrent targets, and empty target states

Diffstat:
MCargo.lock | 2++
Mcrates/outbox/contracts/phase1_publication_v1.descriptor.json | 4++++
Mcrates/outbox/contracts/phase1_publication_v1.manifest.json | 14+++++++-------
Mcrates/outbox/contracts/phase1_publication_v1.manifest.schema.json | 2+-
Mcrates/outbox/contracts/phase1_publication_v1.manifest.sha256 | 2+-
Mcrates/outbox/src/phase1_publication.rs | 453++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/transport_nostr/Cargo.toml | 3+++
Mcrates/transport_nostr/src/error.rs | 13++++++++++++-
Mcrates/transport_nostr/src/lib.rs | 8+++++---
Mcrates/transport_nostr/src/outbox.rs | 679++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
Mcrates/transport_nostr/src/publish.rs | 12++++++++++++
Mcrates/transport_nostr/tests/phase1_outbox_publication.rs | 355+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
Mcrates/transport_nostr/tests/transport.rs | 24++++++++++++++----------
Mtools/xtask/src/contract/outbox_phase1_publication.rs | 6+++---
14 files changed, 1352 insertions(+), 225 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -5257,6 +5257,8 @@ dependencies = [ "serde", "serde_json", "sha2", + "sqlx", + "tempfile", "thiserror 1.0.69", "tokio", "url", diff --git a/crates/outbox/contracts/phase1_publication_v1.descriptor.json b/crates/outbox/contracts/phase1_publication_v1.descriptor.json @@ -77,6 +77,9 @@ { "id": "dispatch-published", "scope": "event", "from": "dispatching", "to": "published", "revision_cas": true, "lease_predicate": "target-matching-live-token", "durable_side_effect": "persist-target-receipt", "retry_class": "none", "repair_edge": false, "terminal_destination": true }, { "id": "dispatch-waiting", "scope": "event", "from": "dispatching", "to": "signed-ready", "revision_cas": true, "lease_predicate": "target-matching-live-token", "durable_side_effect": "persist-target-result", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, { "id": "dispatch-exhausted", "scope": "event", "from": "dispatching", "to": "failed-terminal", "revision_cas": true, "lease_predicate": "target-matching-live-token", "durable_side_effect": "persist-target-result", "retry_class": "terminal", "repair_edge": false, "terminal_destination": true }, + { "id": "recover-expired-dispatch-waiting", "scope": "event", "from": "dispatching", "to": "signed-ready", "revision_cas": true, "lease_predicate": "expired-target-lease", "durable_side_effect": "persist-uncertain-intent", "retry_class": "repair", "repair_edge": true, "terminal_destination": false }, + { "id": "recover-expired-dispatch-published", "scope": "event", "from": "published", "to": "published", "revision_cas": true, "lease_predicate": "expired-target-lease", "durable_side_effect": "persist-uncertain-intent", "retry_class": "repair", "repair_edge": true, "terminal_destination": true }, + { "id": "complete-published-target-result", "scope": "event", "from": "published", "to": "published", "revision_cas": true, "lease_predicate": "target-matching-live-token", "durable_side_effect": "persist-target-result", "retry_class": "none", "repair_edge": false, "terminal_destination": true }, { "id": "claim-target", "scope": "target", "from": "pending", "to": "in-flight", "revision_cas": true, "lease_predicate": "absent-or-expired", "durable_side_effect": "persist-dispatch-intent", "retry_class": "none", "repair_edge": false, "terminal_destination": false }, { "id": "retry-target", "scope": "target", "from": "failed-retryable", "to": "in-flight", "revision_cas": true, "lease_predicate": "absent-or-expired", "durable_side_effect": "reuse-dispatch-intent", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, { "id": "repair-uncertain-target", "scope": "target", "from": "uncertain", "to": "in-flight", "revision_cas": true, "lease_predicate": "absent-or-expired", "durable_side_effect": "reuse-dispatch-intent", "retry_class": "repair", "repair_edge": true, "terminal_destination": false }, @@ -85,6 +88,7 @@ { "id": "target-retryable", "scope": "target", "from": "in-flight", "to": "failed-retryable", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-bounded-error", "retry_class": "retryable", "repair_edge": false, "terminal_destination": false }, { "id": "target-terminal", "scope": "target", "from": "in-flight", "to": "failed-terminal", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-bounded-error", "retry_class": "terminal", "repair_edge": false, "terminal_destination": true }, { "id": "target-uncertain", "scope": "target", "from": "in-flight", "to": "uncertain", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "persist-bounded-error", "retry_class": "repair", "repair_edge": true, "terminal_destination": false }, + { "id": "expire-target-dispatch", "scope": "target", "from": "in-flight", "to": "uncertain", "revision_cas": true, "lease_predicate": "expired", "durable_side_effect": "persist-uncertain-intent", "retry_class": "repair", "repair_edge": true, "terminal_destination": false }, { "id": "target-cancelled", "scope": "target", "from": "in-flight", "to": "cancelled", "revision_cas": true, "lease_predicate": "matching-live-token", "durable_side_effect": "clear-claim", "retry_class": "terminal", "repair_edge": false, "terminal_destination": true }, { "id": "observation-repaired", "scope": "target", "from": "accepted-observation-pending", "to": "accepted-observed", "revision_cas": true, "lease_predicate": "repair-revision-cas", "durable_side_effect": "complete-repair-and-receipt", "retry_class": "repair", "repair_edge": true, "terminal_destination": true } ], diff --git a/crates/outbox/contracts/phase1_publication_v1.manifest.json b/crates/outbox/contracts/phase1_publication_v1.manifest.json @@ -1,14 +1,14 @@ { "contract_id": "radroots_outbox.phase1_publication.v1", "descriptor": { - "byte_length": 11023, + "byte_length": 12212, "path": "crates/outbox/contracts/phase1_publication_v1.descriptor.json", - "sha256": "7c8f1a99e4c7001af5238d164c04acfd70afc6e7133a447602ae5b5d2a1f6622" + "sha256": "5e74ebeb6bfeb33286a5cdfe4e89413d0554fa476761faeca39f7398dd122bbf" }, "manifest_schema": { "byte_length": 3806, "path": "crates/outbox/contracts/phase1_publication_v1.manifest.schema.json", - "sha256": "05a810c798676197608cdfd2bc081fb003e4b878c285a338eed6abe636c0817d" + "sha256": "f1182af3c11f8f8f3fed48d91b48e4ffc99507a380a7cae3f7e88ce740043aa2" }, "migration": { "down": { @@ -64,9 +64,9 @@ }, { "file": { - "byte_length": 118318, + "byte_length": 133841, "path": "crates/outbox/src/phase1_publication.rs", - "sha256": "5ee6745c8eb36fdc18e8578880ca246cc1579e24ae7e7b89559ceaf1ffcbbc73" + "sha256": "41489c91ca667b8c1b7306925d254c43a0e7a498bd0b40c20b542c26a6e46417" }, "role": "phase1_publication_runtime" }, @@ -106,7 +106,7 @@ "file": { "byte_length": 24248, "path": "tools/xtask/src/contract/outbox_phase1_publication.rs", - "sha256": "8a2a09944d78bd67b76a7e1857e6815af2c4ebe712bb14ed48eb66e4cb156dd3" + "sha256": "8088984cf3b478852125cd865f1689daccd768feb8e66f637cf385fb1187cc46" }, "role": "contract_governance" }, @@ -147,6 +147,6 @@ "event_state_count": 9, "stable_error_count": 25, "target_state_count": 8, - "transition_count": 28 + "transition_count": 32 } } diff --git a/crates/outbox/contracts/phase1_publication_v1.manifest.schema.json b/crates/outbox/contracts/phase1_publication_v1.manifest.schema.json @@ -150,7 +150,7 @@ "const": 8 }, "transition_count": { - "const": 28 + "const": 32 } }, "required": [ diff --git a/crates/outbox/contracts/phase1_publication_v1.manifest.sha256 b/crates/outbox/contracts/phase1_publication_v1.manifest.sha256 @@ -1 +1 @@ -ecb9260074e3196e028a62ed683a9126b0d5a70a44eaf4ccd45be9fb97d5a06b +b9e4b58a8745555d0bcd66128de509e788b90dcb1d49ec73b9b6fa81a0ed0307 diff --git a/crates/outbox/src/phase1_publication.rs b/crates/outbox/src/phase1_publication.rs @@ -421,6 +421,39 @@ pub const RADROOTS_PHASE1_PUBLICATION_TRANSITIONS: &[RadrootsPhase1PublicationTr true ), transition!( + "recover-expired-dispatch-waiting", + Event, + "dispatching", + "signed-ready", + "expired-target-lease", + "persist-uncertain-intent", + Repair, + true, + false + ), + transition!( + "recover-expired-dispatch-published", + Event, + "published", + "published", + "expired-target-lease", + "persist-uncertain-intent", + Repair, + true, + true + ), + transition!( + "complete-published-target-result", + Event, + "published", + "published", + "target-matching-live-token", + "persist-target-result", + None, + false, + true + ), + transition!( "claim-target", Target, "pending", @@ -509,6 +542,17 @@ pub const RADROOTS_PHASE1_PUBLICATION_TRANSITIONS: &[RadrootsPhase1PublicationTr false ), transition!( + "expire-target-dispatch", + Target, + "in-flight", + "uncertain", + "expired", + "persist-uncertain-intent", + Repair, + true, + false + ), + transition!( "target-cancelled", Target, "in-flight", @@ -1287,6 +1331,8 @@ impl RadrootsOutbox { publication_id: i64, now_ms: i64, ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { + self.recover_expired_phase1_dispatches(publication_id, now_ms) + .await?; match self.load_phase1_publication(publication_id).await { Ok(record) => Ok(record), Err(error) if is_persisted_authority_error(&error) => { @@ -1322,6 +1368,99 @@ impl RadrootsOutbox { } } + async fn recover_expired_phase1_dispatches( + &self, + publication_id: i64, + now_ms: i64, + ) -> Result<(), RadrootsPhase1PublicationError> { + validate_time(now_ms)?; + let mut transaction = self.pool.begin().await?; + let expired = sqlx::query( + "SELECT target_id, dispatch_digest + FROM outbox_phase1_delivery_target + WHERE publication_id = ? AND state = 'in-flight' + AND claim_expires_at_ms <= ? + ORDER BY target_id", + ) + .bind(publication_id) + .bind(now_ms) + .fetch_all(&mut *transaction) + .await?; + if expired.is_empty() { + transaction.commit().await?; + return Ok(()); + } + + for row in expired { + let target_id: i64 = row.try_get("target_id")?; + let dispatch_digest = blob32(&row, "dispatch_digest")?; + let intent_affected = sqlx::query( + "UPDATE outbox_phase1_dispatch_intent + SET state = 'uncertain', state_revision = state_revision + 1, + updated_at_ms = ? + WHERE intent_digest = ? AND target_id = ? AND state = 'in-flight'", + ) + .bind(now_ms) + .bind(dispatch_digest.as_slice()) + .bind(target_id) + .execute(&mut *transaction) + .await? + .rows_affected(); + let target_affected = sqlx::query( + "UPDATE outbox_phase1_delivery_target + SET state = 'uncertain', state_revision = state_revision + 1, + claim_token = NULL, claim_expires_at_ms = NULL, + last_error = 'dispatch lease expired before durable completion', + next_attempt_after_ms = ?, updated_at_ms = ? + WHERE target_id = ? AND publication_id = ? AND state = 'in-flight' + AND claim_expires_at_ms <= ? AND dispatch_digest = ?", + ) + .bind(now_ms) + .bind(now_ms) + .bind(target_id) + .bind(publication_id) + .bind(now_ms) + .bind(dispatch_digest.as_slice()) + .execute(&mut *transaction) + .await? + .rows_affected(); + if intent_affected != 1 || target_affected != 1 { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + } + + let required_target_count = usize_from_i64( + sqlx::query_scalar::<_, i64>( + "SELECT required_target_count FROM outbox_phase1_publication + WHERE publication_id = ? AND state IN ('dispatching', 'published')", + ) + .bind(publication_id) + .fetch_optional(&mut *transaction) + .await? + .ok_or(RadrootsPhase1PublicationError::StoredAuthorityInvalid)?, + "required_target_count", + )?; + let destination = + aggregate_publication_state(&mut transaction, publication_id, required_target_count) + .await?; + let publication_affected = sqlx::query( + "UPDATE outbox_phase1_publication + SET state = ?, state_revision = state_revision + 1, updated_at_ms = ? + WHERE publication_id = ? AND state IN ('dispatching', 'published')", + ) + .bind(destination.as_str()) + .bind(now_ms) + .bind(publication_id) + .execute(&mut *transaction) + .await? + .rows_affected(); + if publication_affected != 1 { + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } + transaction.commit().await?; + Ok(()) + } + pub async fn claim_phase1_publication_for_signing( &self, publication_id: i64, @@ -1811,6 +1950,30 @@ impl RadrootsOutbox { .iter() .find(|target| target.target_id == target_id) .ok_or(RadrootsPhase1PublicationError::TargetNotFound { target_id })?; + if target.state == RadrootsPhase1PublicationTargetState::AcceptedObserved + && target.revision == observed_target_revision.saturating_add(1) + { + let repair_digest = repair_digest(&target.dispatch_digest); + let receipt_digest = receipt_digest(&target.dispatch_digest, "accepted-observed"); + let authority: Option<(String, i64)> = sqlx::query_as( + "SELECT repair.state, + (SELECT COUNT(*) FROM outbox_phase1_target_receipt + WHERE receipt_digest = ? AND target_id = ? + AND observation_kind = 'accepted-observed') + FROM outbox_phase1_observation_repair AS repair + WHERE repair.repair_digest = ? AND repair.target_id = ?", + ) + .bind(receipt_digest.as_slice()) + .bind(target_id) + .bind(repair_digest.as_slice()) + .bind(target_id) + .fetch_optional(&self.pool) + .await?; + if authority == Some(("complete".to_owned(), 1)) { + return Ok(record); + } + return Err(RadrootsPhase1PublicationError::StoredAuthorityInvalid); + } if target.revision != observed_target_revision || target.state != RadrootsPhase1PublicationTargetState::AcceptedObservationPending { @@ -1901,7 +2064,13 @@ async fn complete_target( ) -> Result<RadrootsPhase1PublicationRecord, RadrootsPhase1PublicationError> { validate_time(now_ms)?; let record = outbox.load_phase1_publication(claim.publication_id).await?; - if record.revision != claim.publication_revision { + if record.revision < claim.publication_revision + || !matches!( + record.state, + RadrootsPhase1PublicationEventState::Dispatching + | RadrootsPhase1PublicationEventState::Published + ) + { return Err(RadrootsPhase1PublicationError::RevisionConflict); } let target = record @@ -2019,13 +2188,13 @@ async fn complete_target( ) .await?; let publication_affected = sqlx::query( - "UPDATE outbox_phase1_publication SET state = ?, state_revision = state_revision + 1, updated_at_ms = ? WHERE publication_id = ? AND state_revision = ? AND state = 'dispatching'", + "UPDATE outbox_phase1_publication SET state = ?, state_revision = state_revision + 1, updated_at_ms = ? WHERE publication_id = ? AND state_revision = ? AND state IN ('dispatching', 'published')", ) .bind(destination_event.as_str()) .bind(now_ms) .bind(claim.publication_id) .bind(i64_from_u64( - claim.publication_revision, + record.revision, "publication_revision", )?) .execute(&mut *transaction) @@ -2553,7 +2722,7 @@ mod tests { assert!(!transition.lease_predicate.is_empty()); assert!(!transition.durable_side_effect.is_empty()); } - assert_eq!(ids.len(), 28); + assert_eq!(ids.len(), 32); } #[test] @@ -3232,6 +3401,282 @@ mod tests { } #[tokio::test] + async fn phase1_publication_partial_effect_boundaries_roll_back_and_recover() { + let temp = tempfile::tempdir().unwrap(); + let path = temp.path().join("partial-effect.sqlite"); + let outbox = RadrootsOutbox::open_file(&path).await.unwrap(); + let ready = ready_update(); + let receipt = outbox + .enqueue_phase1_publication( + &ready, + &RadrootsPhase1PublicationTargetPolicy::new(["wss://relay.example"], 1).unwrap(), + 100, + ) + .await + .unwrap(); + let signing_claim = outbox + .claim_phase1_publication_for_signing( + receipt.record().publication_id(), + receipt.record().revision(), + 101, + 100, + ) + .await + .unwrap(); + let preflight = outbox + .preflight_phase1_publication_signing(&signing_claim, 102) + .await + .unwrap(); + let signed = outbox + .complete_phase1_publication_signing(&preflight, &signed_update(&ready), 103) + .await + .unwrap(); + let target = &signed.targets()[0]; + let expired_claim = outbox + .claim_phase1_publication_target( + signed.publication_id(), + signed.revision(), + target.target_id(), + target.revision(), + 104, + 10, + ) + .await + .unwrap(); + outbox.close().await; + let outbox = RadrootsOutbox::open_file(&path).await.unwrap(); + assert_eq!( + outbox + .claim_phase1_publication_target( + signed.publication_id(), + expired_claim.publication_revision(), + target.target_id(), + expired_claim.target_revision(), + 114, + 10, + ) + .await + .unwrap_err() + .code(), + "phase1_publication_revision_conflict" + ); + let recovered = outbox + .load_phase1_publication(signed.publication_id()) + .await + .unwrap(); + assert_eq!( + recovered.targets()[0].state(), + RadrootsPhase1PublicationTargetState::Uncertain + ); + assert_eq!( + sqlx::query_scalar::<_, String>( + "SELECT state FROM outbox_phase1_dispatch_intent WHERE target_id = ?", + ) + .bind(target.target_id()) + .fetch_one(outbox.pool()) + .await + .unwrap(), + "uncertain" + ); + let retry_claim = outbox + .claim_phase1_publication_target( + recovered.publication_id(), + recovered.revision(), + recovered.targets()[0].target_id(), + recovered.targets()[0].revision(), + 115, + 100, + ) + .await + .unwrap(); + assert_eq!( + retry_claim.dispatch_digest(), + expired_claim.dispatch_digest() + ); + + sqlx::query( + "CREATE TEMP TRIGGER fail_phase1_pending_receipt + BEFORE INSERT ON outbox_phase1_target_receipt + BEGIN SELECT RAISE(ABORT, 'injected pending receipt failure'); END", + ) + .execute(outbox.pool()) + .await + .unwrap(); + assert!(matches!( + outbox + .complete_phase1_target_accepted_pending(&retry_claim, 116) + .await, + Err(RadrootsPhase1PublicationError::Sqlite(_)) + )); + let still_in_flight = outbox + .load_phase1_publication(recovered.publication_id()) + .await + .unwrap(); + assert_eq!( + still_in_flight.targets()[0].state(), + RadrootsPhase1PublicationTargetState::InFlight + ); + sqlx::query("DROP TRIGGER fail_phase1_pending_receipt") + .execute(outbox.pool()) + .await + .unwrap(); + + let pending = outbox + .complete_phase1_target_accepted_pending(&retry_claim, 117) + .await + .unwrap(); + let pending_target = &pending.targets()[0]; + sqlx::query( + "CREATE TEMP TRIGGER fail_phase1_observed_receipt + BEFORE INSERT ON outbox_phase1_target_receipt + WHEN NEW.observation_kind = 'accepted-observed' + BEGIN SELECT RAISE(ABORT, 'injected observed receipt failure'); END", + ) + .execute(outbox.pool()) + .await + .unwrap(); + assert!(matches!( + outbox + .complete_phase1_observation_repair( + pending.publication_id(), + pending_target.target_id(), + pending_target.revision(), + 118, + ) + .await, + Err(RadrootsPhase1PublicationError::Sqlite(_)) + )); + let repair_still_pending = outbox + .load_phase1_publication(pending.publication_id()) + .await + .unwrap(); + assert_eq!( + repair_still_pending.targets()[0].state(), + RadrootsPhase1PublicationTargetState::AcceptedObservationPending + ); + sqlx::query("DROP TRIGGER fail_phase1_observed_receipt") + .execute(outbox.pool()) + .await + .unwrap(); + + let completed = outbox + .complete_phase1_observation_repair( + pending.publication_id(), + pending_target.target_id(), + pending_target.revision(), + 119, + ) + .await + .unwrap(); + let replayed = outbox + .complete_phase1_observation_repair( + pending.publication_id(), + pending_target.target_id(), + pending_target.revision(), + 120, + ) + .await + .unwrap(); + assert_eq!(replayed.targets(), completed.targets()); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM outbox_phase1_target_receipt WHERE target_id = ?", + ) + .bind(pending_target.target_id()) + .fetch_one(outbox.pool()) + .await + .unwrap(), + 2 + ); + } + + #[tokio::test] + async fn phase1_publication_partial_effect_concurrent_target_claims_complete_independently() { + let outbox = RadrootsOutbox::open_memory().await.unwrap(); + let ready = ready_update(); + let receipt = outbox + .enqueue_phase1_publication( + &ready, + &RadrootsPhase1PublicationTargetPolicy::new( + ["wss://relay-a.example", "wss://relay-b.example"], + 1, + ) + .unwrap(), + 100, + ) + .await + .unwrap(); + let signing_claim = outbox + .claim_phase1_publication_for_signing( + receipt.record().publication_id(), + receipt.record().revision(), + 101, + 100, + ) + .await + .unwrap(); + let preflight = outbox + .preflight_phase1_publication_signing(&signing_claim, 102) + .await + .unwrap(); + let signed = outbox + .complete_phase1_publication_signing(&preflight, &signed_update(&ready), 103) + .await + .unwrap(); + let first = outbox + .claim_phase1_publication_target( + signed.publication_id(), + signed.revision(), + signed.targets()[0].target_id(), + signed.targets()[0].revision(), + 104, + 100, + ) + .await + .unwrap(); + let after_first_claim = outbox + .load_phase1_publication(signed.publication_id()) + .await + .unwrap(); + let second = outbox + .claim_phase1_publication_target( + after_first_claim.publication_id(), + after_first_claim.revision(), + after_first_claim.targets()[1].target_id(), + after_first_claim.targets()[1].revision(), + 105, + 100, + ) + .await + .unwrap(); + + let published = outbox + .complete_phase1_target_accepted_pending(&first, 106) + .await + .unwrap(); + assert_eq!( + published.state(), + RadrootsPhase1PublicationEventState::Published + ); + let completed = outbox + .fail_phase1_target_retryable(&second, 107, 108, "relay unavailable") + .await + .unwrap(); + assert_eq!( + completed.state(), + RadrootsPhase1PublicationEventState::Published + ); + assert_eq!( + completed.targets()[0].state(), + RadrootsPhase1PublicationTargetState::AcceptedObservationPending + ); + assert_eq!( + completed.targets()[1].state(), + RadrootsPhase1PublicationTargetState::FailedRetryable + ); + } + + #[tokio::test] async fn phase1_publication_reopen_upgrade_and_exact_rollback_are_transactional() { let temp = tempfile::tempdir().unwrap(); let path = temp.path().join("phase1.sqlite"); diff --git a/crates/transport_nostr/Cargo.toml b/crates/transport_nostr/Cargo.toml @@ -23,6 +23,7 @@ client = [ storage = [ "dep:radroots_event_store", "dep:radroots_outbox", + "dep:sqlx", "radroots_outbox/event-store-adapter", "client", ] @@ -58,6 +59,7 @@ hex = { workspace = true } nostr = { workspace = true } serde = { workspace = true, features = ["derive", "std"] } serde_json = { workspace = true, features = ["std"] } +sqlx = { workspace = true, optional = true } thiserror = { workspace = true } tokio = { workspace = true, optional = true, features = ["rt"] } url = { workspace = true } @@ -67,6 +69,7 @@ radroots_authority = { workspace = true, features = ["local_signer"] } radroots_blossom = { workspace = true, features = ["serde", "std"] } radroots_event_codec = { workspace = true, features = ["serde_json"] } sha2 = { workspace = true } +tempfile = { workspace = true } tokio = { workspace = true, features = ["macros", "rt"] } [[test]] diff --git a/crates/transport_nostr/src/error.rs b/crates/transport_nostr/src/error.rs @@ -41,6 +41,10 @@ pub enum RadrootsRelayTransportError { #[error("Relay target set contains duplicate URL `{url}`")] DuplicateRelayUrl { url: String }, + #[cfg(feature = "storage")] + #[error("Outbox delivery plan {delivery_plan_id} has no valid publishable Nostr targets")] + InvalidEmptyPublishableTargetSet { delivery_plan_id: i64 }, + #[error("Relay fetch item contains invalid relay URL `{url}`: {reason}")] InvalidFetchItemRelayUrl { url: String, reason: String }, @@ -124,6 +128,9 @@ pub enum RadrootsRelayTransportError { #[error("Transport contract error: {0}")] TransportContract(String), + #[error("Transport contract error: {0}")] + TransportContractError(radroots_transport::RadrootsTransportError), + #[error("JSON error: {0}")] Json(#[from] serde_json::Error), @@ -151,6 +158,10 @@ pub enum RadrootsRelayTransportError { Outbox(#[from] radroots_outbox::RadrootsOutboxError), #[cfg(feature = "storage")] + #[error("Phase 1 publication error: {0}")] + Phase1Publication(#[from] radroots_outbox::RadrootsPhase1PublicationError), + + #[cfg(feature = "storage")] #[error("Outbox claim {0} does not contain a signed event")] MissingSignedOutboxEvent(i64), @@ -170,6 +181,6 @@ pub(crate) fn ensure_nonnegative_timestamp( impl From<radroots_transport::RadrootsTransportError> for RadrootsRelayTransportError { fn from(value: radroots_transport::RadrootsTransportError) -> Self { - Self::TransportContract(value.to_string()) + Self::TransportContractError(value) } } diff --git a/crates/transport_nostr/src/lib.rs b/crates/transport_nostr/src/lib.rs @@ -29,10 +29,12 @@ pub use fetch::{ }; #[cfg(feature = "storage")] pub use outbox::{ - RadrootsOutboxPublishPolicy, RadrootsOutboxPublishReceipt, RadrootsOutboxPublishTargetReceipt, - phase1_publication_delivery_request, publish_claimed_outbox_event, - publish_claimed_outbox_event_with_transport, + RadrootsOutboxPublishEmptyTargetState, RadrootsOutboxPublishPolicy, + RadrootsOutboxPublishReceipt, RadrootsOutboxPublishTargetReceipt, + execute_claimed_phase1_publication_target_with_transport, phase1_publication_delivery_request, + publish_claimed_outbox_event, publish_claimed_outbox_event_with_transport, publish_claimed_phase1_publication_target_with_transport, + repair_phase1_publication_observation, }; pub use outcome::{RadrootsRelayOutcome, RadrootsRelayOutcomeKind}; #[cfg(feature = "client")] diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs @@ -9,13 +9,14 @@ use crate::{ }; use radroots_event::draft::RadrootsVerifiedSignedEvent; use radroots_event_store::{ - RadrootsEventIngest, RadrootsEventStore, RadrootsTransportObservation, + RadrootsEventIngest, RadrootsEventStore, RadrootsEventStoreError, RadrootsTransportObservation, RadrootsTransportObservationType, }; use radroots_outbox::{ - RadrootsOutbox, RadrootsOutboxClaimedEvent, RadrootsOutboxDeliveryTargetRecord, - RadrootsOutboxDeliveryTargetStatus, RadrootsOutboxEventStoreIngestReceipt, - RadrootsPhase1PublicationTargetClaim, + RadrootsOutbox, RadrootsOutboxClaimedEvent, RadrootsOutboxDeliveryPlanStatus, + RadrootsOutboxDeliveryTargetRecord, RadrootsOutboxDeliveryTargetStatus, + RadrootsOutboxEventStoreIngestReceipt, RadrootsPhase1PublicationRecord, + RadrootsPhase1PublicationTargetClaim, RadrootsPhase1PublicationTargetState, }; use radroots_transport::{ RadrootsTransport, RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest, @@ -75,10 +76,18 @@ pub struct RadrootsOutboxPublishReceipt { terminal_count: usize, quorum: usize, quorum_met: bool, + empty_target_state: Option<RadrootsOutboxPublishEmptyTargetState>, target_receipts: Vec<RadrootsOutboxPublishTargetReceipt>, relay_receipts: Vec<RadrootsRelayPublishRelayReceipt>, } +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum RadrootsOutboxPublishEmptyTargetState { + AlreadySatisfied, + Terminal, + Cancelled, +} + #[derive(Clone, Debug, PartialEq, Eq)] pub struct RadrootsOutboxPublishTargetReceipt { delivery_target_id: i64, @@ -124,6 +133,10 @@ impl RadrootsOutboxPublishReceipt { self.quorum_met } + pub const fn empty_target_state(&self) -> Option<RadrootsOutboxPublishEmptyTargetState> { + self.empty_target_state + } + pub fn target_receipts(&self) -> &[RadrootsOutboxPublishTargetReceipt] { self.target_receipts.as_slice() } @@ -214,6 +227,154 @@ where Ok(receipt) } +pub async fn execute_claimed_phase1_publication_target_with_transport<T>( + outbox: &RadrootsOutbox, + event_store: &RadrootsEventStore, + transport: &T, + claim: &RadrootsPhase1PublicationTargetClaim, + next_attempt_after_ms: i64, + now_ms: i64, +) -> Result<RadrootsPhase1PublicationRecord, RadrootsRelayTransportError> +where + T: RadrootsTransport + ?Sized, +{ + ensure_nonnegative_timestamp("now_ms", now_ms)?; + ensure_nonnegative_timestamp("next_attempt_after_ms", next_attempt_after_ms)?; + let transport_kind = transport.transport_kind(); + if transport_kind != RadrootsTransportKind::Nostr { + let error = RadrootsRelayTransportError::UnexpectedTransportKind { + expected: "nostr", + actual: transport_kind.canonical_label(), + }; + outbox + .fail_phase1_target_retryable( + claim, + now_ms, + next_attempt_after_ms, + error.to_string().as_str(), + ) + .await?; + return Err(error); + } + let request = match phase1_publication_delivery_request(claim, now_ms) { + Ok(request) => request, + Err(error) => { + outbox + .fail_phase1_target_terminal(claim, now_ms, error.to_string().as_str()) + .await?; + return Err(error); + } + }; + let receipt = match transport.deliver(request.clone()).await { + Ok(receipt) => receipt, + Err(error) => { + let relay_error = transport_error_to_relay_error(error); + outbox + .mark_phase1_target_uncertain(claim, now_ms, relay_error.to_string().as_str()) + .await?; + return Err(relay_error); + } + }; + if let Err(error) = receipt.validate_for_request(&request) { + let relay_error = transport_error_to_relay_error(error); + outbox + .mark_phase1_target_uncertain(claim, now_ms, relay_error.to_string().as_str()) + .await?; + return Err(relay_error); + } + let target_receipt = receipt + .target_receipts() + .first() + .ok_or(RadrootsTransportError::MissingDeliveryTargetReceipt)?; + let outcome = target_receipt.outcome(); + let diagnostic = outcome.message().unwrap_or(outcome.code()); + if target_receipt.status().counts_as_accepted_satisfaction() { + let pending = match outbox + .complete_phase1_target_accepted_pending(claim, now_ms) + .await + { + Ok(record) => record, + Err(primary) => { + outbox + .mark_phase1_target_uncertain( + claim, + now_ms, + "remote acceptance could not be persisted", + ) + .await?; + return Err(primary.into()); + } + }; + let pending_target = pending + .targets() + .iter() + .find(|target| target.target_id() == claim.target_id()) + .ok_or( + radroots_outbox::RadrootsPhase1PublicationError::TargetNotFound { + target_id: claim.target_id(), + }, + )?; + ingest_publish_observation( + event_store, + claim.signed_event(), + claim.endpoint_uri(), + now_ms, + ) + .await?; + return outbox + .complete_phase1_observation_repair( + pending.publication_id(), + pending_target.target_id(), + pending_target.revision(), + now_ms, + ) + .await + .map_err(Into::into); + } + if target_receipt.status().is_retryable_failure() + || target_receipt.status().is_deferred_until_implemented() + { + return outbox + .fail_phase1_target_retryable(claim, now_ms, next_attempt_after_ms, diagnostic) + .await + .map_err(Into::into); + } + outbox + .fail_phase1_target_terminal(claim, now_ms, diagnostic) + .await + .map_err(Into::into) +} + +pub async fn repair_phase1_publication_observation( + outbox: &RadrootsOutbox, + event_store: &RadrootsEventStore, + publication_id: i64, + target_id: i64, + now_ms: i64, +) -> Result<RadrootsPhase1PublicationRecord, RadrootsRelayTransportError> { + ensure_nonnegative_timestamp("now_ms", now_ms)?; + let record = outbox.load_phase1_publication(publication_id).await?; + let target = record + .targets() + .iter() + .find(|target| target.target_id() == target_id) + .ok_or(radroots_outbox::RadrootsPhase1PublicationError::TargetNotFound { target_id })?; + if target.state() == RadrootsPhase1PublicationTargetState::AcceptedObserved { + return Ok(record); + } + if target.state() != RadrootsPhase1PublicationTargetState::AcceptedObservationPending { + return Err(radroots_outbox::RadrootsPhase1PublicationError::StateConflict.into()); + } + let signed_event = record + .signed_event() + .ok_or(radroots_outbox::RadrootsPhase1PublicationError::StoredAuthorityInvalid)?; + ingest_publish_observation(event_store, signed_event, target.endpoint_uri(), now_ms).await?; + outbox + .complete_phase1_observation_repair(publication_id, target_id, target.revision(), now_ms) + .await + .map_err(Into::into) +} + pub async fn publish_claimed_outbox_event<A>( outbox: &RadrootsOutbox, event_store: &RadrootsEventStore, @@ -240,28 +401,16 @@ where .await?; let publishable = publishable_relays(outbox, claimed, policy.republish_accepted_relays).await?; if publishable.relays.is_empty() { - outbox - .complete_publish_attempt( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - "relay publish incomplete", - "relay publish terminal", - policy.next_attempt_after_ms, - now_ms, - ) - .await?; - return Ok(RadrootsOutboxPublishReceipt { + return complete_empty_publishable_attempt(EmptyPublishEffects { + outbox, + claimed, local_ingest, event_id: signed_event.signed_event().id_str().to_owned(), - attempted_count: 0, - accepted_count: publishable.accepted_count, - retryable_count: 0, - terminal_count: 0, - quorum: publishable.remaining_satisfaction_count, - quorum_met: publishable.satisfied_count >= publishable.satisfaction_required_count, - target_receipts: Vec::new(), - relay_receipts: Vec::new(), - }); + publishable: &publishable, + next_attempt_after_ms: policy.next_attempt_after_ms, + now_ms, + }) + .await; } let targets = RadrootsRelayTargetSet::new( unique_publishable_relay_urls(&publishable), @@ -273,7 +422,6 @@ where .with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::no_wait()) .try_with_idempotency_key(outbox_publish_idempotency_key( claimed.outbox_event_id, - claimed.attempt_count, signed_event.signed_event().id_str(), active_delivery_plan_id, ))?; @@ -288,38 +436,18 @@ where Err(error) => return Err(error), }; let target_receipts = target_receipts_from_relay_receipts(&publishable, publish.relays()); - - for target_receipt in &target_receipts { - complete_outbox_delivery_target(outbox, claimed, target_receipt, now_ms).await?; - } - - for relay in publish.relays() { - if relay - .outcome() - .kind() - .transport_outcome_kind() - .target_status() - .counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted) - && publishable - .targets_for_relay(relay.relay_url()) - .next() - .is_some() - { - ingest_publish_observation(event_store, &signed_event, relay.relay_url(), now_ms) - .await?; - } - } - - outbox - .complete_publish_attempt( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - "relay publish incomplete", - "relay publish terminal", - policy.next_attempt_after_ms, - now_ms, - ) - .await?; + persist_legacy_publish_effects(LegacyPublishEffects { + outbox, + event_store, + claimed, + signed_event: &signed_event, + publishable: &publishable, + target_receipts: &target_receipts, + relay_receipts: publish.relays(), + next_attempt_after_ms: policy.next_attempt_after_ms, + now_ms, + }) + .await?; let (event_id, relay_receipts) = publish.into_event_id_and_relays(); Ok(RadrootsOutboxPublishReceipt { @@ -344,6 +472,7 @@ where quorum: publishable.remaining_satisfaction_count, quorum_met: publishable.satisfied_count_after_receipts(&target_receipts) >= publishable.satisfaction_required_count, + empty_target_state: None, target_receipts, relay_receipts, }) @@ -382,28 +511,16 @@ where .await?; let publishable = publishable_relays(outbox, claimed, policy.republish_accepted_relays).await?; if publishable.relays.is_empty() { - outbox - .complete_publish_attempt( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - "relay publish incomplete", - "relay publish terminal", - policy.next_attempt_after_ms, - now_ms, - ) - .await?; - return Ok(RadrootsOutboxPublishReceipt { + return complete_empty_publishable_attempt(EmptyPublishEffects { + outbox, + claimed, local_ingest, event_id: signed_event.signed_event().id_str().to_owned(), - attempted_count: 0, - accepted_count: publishable.accepted_count, - retryable_count: 0, - terminal_count: 0, - quorum: publishable.remaining_satisfaction_count, - quorum_met: publishable.satisfied_count >= publishable.satisfaction_required_count, - target_receipts: Vec::new(), - relay_receipts: Vec::new(), - }); + publishable: &publishable, + next_attempt_after_ms: policy.next_attempt_after_ms, + now_ms, + }) + .await; } RadrootsRelayTargetSet::new( unique_publishable_relay_urls(&publishable), @@ -414,7 +531,6 @@ where let satisfaction_policy = transport_satisfaction_policy_for_publishable(&publishable); let request_id = outbox_publish_idempotency_key( claimed.outbox_event_id, - claimed.attempt_count, signed_event.signed_event().id_str(), publishable.active_delivery_plan_id, ); @@ -433,38 +549,18 @@ where .map_err(transport_error_to_relay_error)?; let relay_receipts = relay_receipts_from_transport_receipts(&delivery)?; let target_receipts = target_receipts_from_transport_receipts(&publishable, &delivery)?; - - for target_receipt in &target_receipts { - complete_outbox_delivery_target(outbox, claimed, target_receipt, now_ms).await?; - } - - for relay in &relay_receipts { - if relay - .outcome() - .kind() - .transport_outcome_kind() - .target_status() - .counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted) - && publishable - .targets_for_relay(relay.relay_url()) - .next() - .is_some() - { - ingest_publish_observation(event_store, &signed_event, relay.relay_url(), now_ms) - .await?; - } - } - - outbox - .complete_publish_attempt( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - "relay publish incomplete", - "relay publish terminal", - policy.next_attempt_after_ms, - now_ms, - ) - .await?; + persist_legacy_publish_effects(LegacyPublishEffects { + outbox, + event_store, + claimed, + signed_event: &signed_event, + publishable: &publishable, + target_receipts: &target_receipts, + relay_receipts: &relay_receipts, + next_attempt_after_ms: policy.next_attempt_after_ms, + now_ms, + }) + .await?; Ok(RadrootsOutboxPublishReceipt { local_ingest, @@ -488,6 +584,7 @@ where quorum: publishable.remaining_satisfaction_count, quorum_met: publishable.satisfied_count_after_receipts(&target_receipts) >= publishable.satisfaction_required_count, + empty_target_state: None, target_receipts, relay_receipts, }) @@ -511,6 +608,138 @@ fn adapter_transport_failure_receipt( RadrootsRelayPublishReceipt::new(event_id, quorum, false, relays) } +struct LegacyPublishEffects<'a> { + outbox: &'a RadrootsOutbox, + event_store: &'a RadrootsEventStore, + claimed: &'a RadrootsOutboxClaimedEvent, + signed_event: &'a RadrootsVerifiedSignedEvent, + publishable: &'a PublishableRelays, + target_receipts: &'a [RadrootsOutboxPublishTargetReceipt], + relay_receipts: &'a [RadrootsRelayPublishRelayReceipt], + next_attempt_after_ms: i64, + now_ms: i64, +} + +async fn persist_legacy_publish_effects( + effects: LegacyPublishEffects<'_>, +) -> Result<(), RadrootsRelayTransportError> { + for target_receipt in effects.target_receipts { + complete_outbox_delivery_target( + effects.outbox, + effects.claimed, + target_receipt, + effects.now_ms, + ) + .await?; + } + for relay in effects.relay_receipts { + if relay + .outcome() + .kind() + .transport_outcome_kind() + .target_status() + .counts_as_satisfied(RadrootsTransportSatisfactionClass::Accepted) + && effects + .publishable + .targets_for_relay(relay.relay_url()) + .next() + .is_some() + { + ingest_publish_observation( + effects.event_store, + effects.signed_event, + relay.relay_url(), + effects.now_ms, + ) + .await?; + } + } + effects + .outbox + .complete_publish_attempt( + effects.claimed.outbox_event_id, + effects.claimed.claim_token.as_str(), + "relay publish incomplete", + "relay publish terminal", + effects.next_attempt_after_ms, + effects.now_ms, + ) + .await?; + Ok(()) +} + +struct EmptyPublishEffects<'a> { + outbox: &'a RadrootsOutbox, + claimed: &'a RadrootsOutboxClaimedEvent, + local_ingest: RadrootsOutboxEventStoreIngestReceipt, + event_id: String, + publishable: &'a PublishableRelays, + next_attempt_after_ms: i64, + now_ms: i64, +} + +async fn complete_empty_publishable_attempt( + effects: EmptyPublishEffects<'_>, +) -> Result<RadrootsOutboxPublishReceipt, RadrootsRelayTransportError> { + let empty_target_state = effects.publishable.empty_target_state.ok_or( + RadrootsRelayTransportError::InvalidEmptyPublishableTargetSet { + delivery_plan_id: effects.publishable.active_delivery_plan_id, + }, + )?; + match empty_target_state { + RadrootsOutboxPublishEmptyTargetState::AlreadySatisfied => { + effects + .outbox + .complete_publish_attempt( + effects.claimed.outbox_event_id, + effects.claimed.claim_token.as_str(), + "delivery plan already satisfied", + "delivery plan already satisfied", + effects.next_attempt_after_ms, + effects.now_ms, + ) + .await?; + } + RadrootsOutboxPublishEmptyTargetState::Terminal => { + effects + .outbox + .complete_publish_attempt( + effects.claimed.outbox_event_id, + effects.claimed.claim_token.as_str(), + "delivery plan has no recoverable Nostr targets", + "delivery plan has no recoverable Nostr targets", + effects.next_attempt_after_ms, + effects.now_ms, + ) + .await?; + } + RadrootsOutboxPublishEmptyTargetState::Cancelled => { + effects + .outbox + .cancel_claimed_event( + effects.claimed.outbox_event_id, + effects.claimed.claim_token.as_str(), + effects.now_ms, + ) + .await?; + } + } + Ok(RadrootsOutboxPublishReceipt { + local_ingest: effects.local_ingest, + event_id: effects.event_id, + attempted_count: 0, + accepted_count: effects.publishable.accepted_count, + retryable_count: 0, + terminal_count: 0, + quorum: effects.publishable.remaining_satisfaction_count, + quorum_met: effects.publishable.satisfied_count + >= effects.publishable.satisfaction_required_count, + empty_target_state: Some(empty_target_state), + target_receipts: Vec::new(), + relay_receipts: Vec::new(), + }) +} + struct PublishableRelays { active_delivery_plan_id: i64, relays: Vec<PublishableRelay>, @@ -521,6 +750,7 @@ struct PublishableRelays { satisfaction_class: RadrootsTransportSatisfactionClass, required_targets: Option<Vec<RadrootsTransportTargetFingerprint>>, remaining_required_targets: Option<Vec<RadrootsTransportTargetFingerprint>>, + empty_target_state: Option<RadrootsOutboxPublishEmptyTargetState>, } impl PublishableRelays { @@ -862,60 +1092,7 @@ fn transport_satisfaction_policy_for_publishable( } fn transport_error_to_relay_error(error: RadrootsTransportError) -> RadrootsRelayTransportError { - match error { - RadrootsTransportError::UnsupportedOperation - | RadrootsTransportError::EmptyTransportKind - | RadrootsTransportError::InvalidTransportKind - | RadrootsTransportError::EmptyTargetScope - | RadrootsTransportError::InvalidTargetScope - | RadrootsTransportError::EmptyTargetLabel - | RadrootsTransportError::InvalidTargetLabel - | RadrootsTransportError::InvalidSatisfactionPolicy - | RadrootsTransportError::EmptyRequiredTargetSet - | RadrootsTransportError::DuplicateRequiredTargetFingerprint - | RadrootsTransportError::RequiredTargetNotRequested - | RadrootsTransportError::EmptyDeliveryRequestId - | RadrootsTransportError::InvalidDeliveryRequestId - | RadrootsTransportError::EmptyFetchRequestId - | RadrootsTransportError::InvalidFetchRequestId - | RadrootsTransportError::InvalidPayloadSignature - | RadrootsTransportError::InvalidDeliveryTimestamp => { - RadrootsRelayTransportError::Transport(error.to_string()) - } - RadrootsTransportError::EmptyTargetUri - | RadrootsTransportError::InvalidTargetUri - | RadrootsTransportError::ResourceLimitExceeded { .. } - | RadrootsTransportError::EmptyTargetSet - | RadrootsTransportError::DuplicateTargetFingerprint - | RadrootsTransportError::InvalidTargetFingerprint - | RadrootsTransportError::UnexpectedDeliveryTargetReceipt - | RadrootsTransportError::DuplicateDeliveryTargetReceipt - | RadrootsTransportError::MissingDeliveryTargetReceipt - | RadrootsTransportError::DeliveryTargetReceiptStatusMismatch - | RadrootsTransportError::DeliveryTargetReceiptAttemptMismatch - | RadrootsTransportError::TransportOutcomeStatusMismatch - | RadrootsTransportError::TransportOutcomeCodeMismatch - | RadrootsTransportError::TransportOutcomeRetryClassMismatch - | RadrootsTransportError::DeliveryReceiptRequestIdMismatch - | RadrootsTransportError::DeliveryReceiptTargetSetMismatch - | RadrootsTransportError::UnexpectedFetchTargetReceipt - | RadrootsTransportError::DuplicateFetchTargetReceipt - | RadrootsTransportError::MissingFetchTargetReceipt - | RadrootsTransportError::FetchReceiptRequestIdMismatch - | RadrootsTransportError::FetchReceiptTargetSetMismatch => { - RadrootsRelayTransportError::TransportContract(error.to_string()) - } - RadrootsTransportError::EmptyPayloadId - | RadrootsTransportError::InvalidPayloadId - | RadrootsTransportError::EmptyPayloadLabel - | RadrootsTransportError::InvalidPayloadLabel - | RadrootsTransportError::EmptyPayloadBytes - | RadrootsTransportError::InvalidPayloadBytes - | RadrootsTransportError::InvalidPayloadDigest - | RadrootsTransportError::PayloadDigestMismatch => { - RadrootsRelayTransportError::NostrEventJson(error.to_string()) - } - } + error.into() } async fn publishable_relays( @@ -1049,6 +1226,16 @@ async fn publishable_relays( }); } } + let empty_target_state = if relays.is_empty() { + Some(classify_empty_publishable_target_set( + plan.status, + remaining_satisfaction_count, + active_targets.iter().map(|target| target.status), + active_delivery_plan_id, + )?) + } else { + None + }; Ok(PublishableRelays { active_delivery_plan_id, relays, @@ -1059,18 +1246,58 @@ async fn publishable_relays( satisfaction_class, required_targets, remaining_required_targets, + empty_target_state, }) } +fn classify_empty_publishable_target_set( + plan_status: RadrootsOutboxDeliveryPlanStatus, + remaining_satisfaction_count: usize, + target_statuses: impl IntoIterator<Item = RadrootsOutboxDeliveryTargetStatus>, + delivery_plan_id: i64, +) -> Result<RadrootsOutboxPublishEmptyTargetState, RadrootsRelayTransportError> { + let mut target_statuses = target_statuses.into_iter(); + let terminal_targets_only = target_statuses + .next() + .is_some_and(is_terminal_empty_target_status) + && target_statuses.all(is_terminal_empty_target_status); + match plan_status { + RadrootsOutboxDeliveryPlanStatus::Complete => { + Ok(RadrootsOutboxPublishEmptyTargetState::AlreadySatisfied) + } + RadrootsOutboxDeliveryPlanStatus::FailedTerminal => { + Ok(RadrootsOutboxPublishEmptyTargetState::Terminal) + } + RadrootsOutboxDeliveryPlanStatus::Cancelled => { + Ok(RadrootsOutboxPublishEmptyTargetState::Cancelled) + } + RadrootsOutboxDeliveryPlanStatus::Queued if remaining_satisfaction_count == 0 => { + Ok(RadrootsOutboxPublishEmptyTargetState::AlreadySatisfied) + } + RadrootsOutboxDeliveryPlanStatus::Queued if terminal_targets_only => { + Ok(RadrootsOutboxPublishEmptyTargetState::Terminal) + } + RadrootsOutboxDeliveryPlanStatus::Queued => { + Err(RadrootsRelayTransportError::InvalidEmptyPublishableTargetSet { delivery_plan_id }) + } + } +} + +fn is_terminal_empty_target_status(status: RadrootsOutboxDeliveryTargetStatus) -> bool { + matches!( + status, + RadrootsOutboxDeliveryTargetStatus::FailedTerminal + | RadrootsOutboxDeliveryTargetStatus::SkippedPolicyDenied + | RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented + ) +} + fn outbox_publish_idempotency_key( outbox_event_id: i64, - attempt_count: i64, event_id: &str, active_delivery_plan_id: i64, ) -> String { - format!( - "radroots-nostr-outbox-{outbox_event_id}-{attempt_count}-{event_id}-{active_delivery_plan_id}" - ) + format!("radroots-nostr-outbox-{outbox_event_id}-{event_id}-{active_delivery_plan_id}") } fn counts_as_accepted_for_plan( @@ -1135,21 +1362,52 @@ async fn ingest_publish_observation( RadrootsTransportObservationType::PublishAck, observed_at_ms, )?; + let event_id = signed_event.signed_event().id_str().to_owned(); + let transport_kind = observation.transport_kind().canonical_label(); + let endpoint_fingerprint = observation.endpoint_fingerprint().as_str().to_owned(); + let observation_type = observation.observation_type().as_str(); let ingest = RadrootsEventIngest::from_signed_event( signed_event.signed_event().clone(), observed_at_ms, )? .with_observation(observation); - event_store.ingest_event(ingest).await?; + let mut transaction = event_store.begin_write_transaction().await?; + let already_present = sqlx::query_scalar::<_, i64>( + "SELECT 1 FROM event_transport_observation + WHERE event_id = ? AND transport_kind = ? + AND endpoint_fingerprint = ? AND observation_type = ?", + ) + .bind(event_id) + .bind(transport_kind) + .bind(endpoint_fingerprint) + .bind(observation_type) + .fetch_optional(&mut *transaction) + .await + .map_err(RadrootsEventStoreError::from)? + .is_some(); + if !already_present + && let Err(error) = event_store + .ingest_event_in_transaction(&mut transaction, ingest) + .await + { + let _ = transaction.rollback().await; + return Err(error.into()); + } + transaction + .commit() + .await + .map_err(RadrootsEventStoreError::from)?; Ok(()) } #[cfg(test)] mod tests { use super::{ - PublishableRelay, PublishableRelays, RadrootsOutboxDeliveryTargetStatus, - adapter_transport_failure_receipt, counts_as_accepted_for_plan, - is_publishable_delivery_status, publishable_transport_targets, + PublishableRelay, PublishableRelays, RadrootsOutboxDeliveryPlanStatus, + RadrootsOutboxDeliveryTargetStatus, RadrootsOutboxPublishEmptyTargetState, + adapter_transport_failure_receipt, classify_empty_publishable_target_set, + counts_as_accepted_for_plan, is_publishable_delivery_status, + outbox_publish_idempotency_key, publishable_transport_targets, relay_outcome_from_transport_outcome, relay_outcome_kind_from_transport_outcome, relay_receipts_from_transport_receipts, satisfaction_policy_for_remaining_count, target_receipts_from_relay_receipts, target_receipts_from_transport_receipts, @@ -1170,6 +1428,10 @@ mod tests { #[test] fn internal_outbox_publish_helpers_cover_policy_edges() { assert_eq!( + outbox_publish_idempotency_key(7, "event-id", 11), + "radroots-nostr-outbox-7-event-id-11" + ); + assert_eq!( satisfaction_policy_for_remaining_count( RadrootsTransportSatisfactionClass::Accepted, 2, @@ -1230,6 +1492,63 @@ mod tests { RadrootsTransportSatisfactionPolicy::no_wait() ); + let empty_target_cases = [ + ( + RadrootsOutboxDeliveryPlanStatus::Complete, + 1, + Vec::new(), + RadrootsOutboxPublishEmptyTargetState::AlreadySatisfied, + ), + ( + RadrootsOutboxDeliveryPlanStatus::Queued, + 0, + vec![RadrootsOutboxDeliveryTargetStatus::Accepted], + RadrootsOutboxPublishEmptyTargetState::AlreadySatisfied, + ), + ( + RadrootsOutboxDeliveryPlanStatus::FailedTerminal, + 1, + vec![RadrootsOutboxDeliveryTargetStatus::Pending], + RadrootsOutboxPublishEmptyTargetState::Terminal, + ), + ( + RadrootsOutboxDeliveryPlanStatus::Queued, + 1, + vec![ + RadrootsOutboxDeliveryTargetStatus::FailedTerminal, + RadrootsOutboxDeliveryTargetStatus::SkippedPolicyDenied, + RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented, + ], + RadrootsOutboxPublishEmptyTargetState::Terminal, + ), + ( + RadrootsOutboxDeliveryPlanStatus::Cancelled, + 1, + vec![RadrootsOutboxDeliveryTargetStatus::Pending], + RadrootsOutboxPublishEmptyTargetState::Cancelled, + ), + ]; + for (status, remaining, target_statuses, expected) in empty_target_cases { + assert_eq!( + classify_empty_publishable_target_set(status, remaining, target_statuses, 7) + .expect("classified empty target state"), + expected + ); + } + assert!(matches!( + classify_empty_publishable_target_set( + RadrootsOutboxDeliveryPlanStatus::Queued, + 1, + [RadrootsOutboxDeliveryTargetStatus::Accepted], + 7, + ), + Err( + RadrootsRelayTransportError::InvalidEmptyPublishableTargetSet { + delivery_plan_id: 7 + } + ) + )); + assert!(counts_as_accepted_for_plan( RadrootsOutboxDeliveryTargetStatus::Accepted, false, @@ -1292,6 +1611,7 @@ mod tests { satisfaction_class: RadrootsTransportSatisfactionClass::Delivered, required_targets: None, remaining_required_targets: None, + empty_target_state: None, }; let accepted_relay_receipts = target_receipts_from_relay_receipts( @@ -1465,6 +1785,7 @@ mod tests { satisfaction_class: RadrootsTransportSatisfactionClass::Accepted, required_targets: None, remaining_required_targets: None, + empty_target_state: None, }; let targets = publishable_transport_targets(&publishable).expect("transport targets"); assert_eq!(targets.len(), 1); @@ -1479,19 +1800,19 @@ mod tests { invalid.relays[0].target_scope = Some("bad scope".to_owned()); assert!(matches!( publishable_transport_targets(&invalid), - Err(RadrootsRelayTransportError::Transport(_)) + Err(RadrootsRelayTransportError::TransportContractError(_)) )); invalid.relays[0].target_scope = Some("foodshed.west".to_owned()); invalid.relays[0].target_label = Some("bad\0label".to_owned()); assert!(matches!( publishable_transport_targets(&invalid), - Err(RadrootsRelayTransportError::Transport(_)) + Err(RadrootsRelayTransportError::TransportContractError(_)) )); invalid.relays[0].target_label = Some("primary relay".to_owned()); invalid.relays[0].relay_url = "not-a-relay".to_owned(); assert!(matches!( publishable_transport_targets(&invalid), - Err(RadrootsRelayTransportError::TransportContract(_)) + Err(RadrootsRelayTransportError::TransportContractError(_)) )); let generic_errors = [ @@ -1509,7 +1830,7 @@ mod tests { for error in generic_errors { assert!(matches!( transport_error_to_relay_error(error), - RadrootsRelayTransportError::Transport(_) + RadrootsRelayTransportError::TransportContractError(_) )); } let target_errors = [ @@ -1522,7 +1843,7 @@ mod tests { for error in target_errors { assert!(matches!( transport_error_to_relay_error(error), - RadrootsRelayTransportError::TransportContract(_) + RadrootsRelayTransportError::TransportContractError(_) )); } assert!(matches!( @@ -1531,7 +1852,7 @@ mod tests { max: 1, actual: 2, }), - RadrootsRelayTransportError::TransportContract(_) + RadrootsRelayTransportError::TransportContractError(_) )); let payload_errors = [ RadrootsTransportError::EmptyPayloadId, @@ -1546,7 +1867,7 @@ mod tests { for error in payload_errors { assert!(matches!( transport_error_to_relay_error(error), - RadrootsRelayTransportError::NostrEventJson(_) + RadrootsRelayTransportError::TransportContractError(_) )); } diff --git a/crates/transport_nostr/src/publish.rs b/crates/transport_nostr/src/publish.rs @@ -637,6 +637,7 @@ where fn nostr_error_to_transport_error(error: RadrootsRelayTransportError) -> RadrootsTransportError { match error { + RadrootsRelayTransportError::TransportContractError(error) => error, RadrootsRelayTransportError::TransportContract(_) => { RadrootsTransportError::InvalidPayloadBytes } @@ -648,6 +649,10 @@ fn nostr_error_to_transport_error(error: RadrootsRelayTransportError) -> Radroot | RadrootsRelayTransportError::ConflictingFetchTerminalRelayUrl { .. } => { RadrootsTransportError::InvalidTransportKind } + #[cfg(feature = "storage")] + RadrootsRelayTransportError::InvalidEmptyPublishableTargetSet { .. } => { + RadrootsTransportError::InvalidTransportKind + } RadrootsRelayTransportError::RelayUrlParse { .. } | RadrootsRelayTransportError::WsRequiresLocalhostPolicy { .. } | RadrootsRelayTransportError::UnsupportedRelayScheme { .. } @@ -687,6 +692,7 @@ fn nostr_error_to_transport_error(error: RadrootsRelayTransportError) -> Radroot #[cfg(feature = "storage")] RadrootsRelayTransportError::EventStore(_) | RadrootsRelayTransportError::Outbox(_) + | RadrootsRelayTransportError::Phase1Publication(_) | RadrootsRelayTransportError::MissingSignedOutboxEvent(_) | RadrootsRelayTransportError::MissingPersistedFetchReceiptEventId | RadrootsRelayTransportError::MissingStoredEventVisibility { .. } @@ -705,6 +711,12 @@ mod contract_tests { #[test] fn relay_errors_map_to_stable_transport_contract_categories() { assert_eq!( + nostr_error_to_transport_error( + RadrootsTransportError::DeliveryReceiptRequestIdMismatch.into(), + ), + RadrootsTransportError::DeliveryReceiptRequestIdMismatch + ); + assert_eq!( nostr_error_to_transport_error(RadrootsRelayTransportError::TransportContract( "contract".to_owned(), )), diff --git a/crates/transport_nostr/tests/phase1_outbox_publication.rs b/crates/transport_nostr/tests/phase1_outbox_publication.rs @@ -36,15 +36,27 @@ use radroots_event_codec::wire::publication::{ RadrootsPhase1PublicationDraft, RadrootsPhase1PublicationMediaReference, allowlist::allow_phase1_publication_artifact, bind_phase1_publication_media_readiness, }; +use radroots_event_store::RadrootsEventStore; use radroots_nostr::prelude::{RadrootsNostrKeys, RadrootsNostrSecretKey}; -use radroots_outbox::{RadrootsOutbox, RadrootsPhase1PublicationTargetPolicy}; -use radroots_transport::{RadrootsTransportDeliveryTargetStatus, RadrootsTransportPayload}; +use radroots_outbox::{ + RadrootsOutbox, RadrootsPhase1PublicationTargetPolicy, RadrootsPhase1PublicationTargetState, +}; +use radroots_transport::{ + RadrootsTransport, RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest, + RadrootsTransportError, RadrootsTransportFetchReceipt, RadrootsTransportFetchRequest, + RadrootsTransportFuture, RadrootsTransportImplementationState, RadrootsTransportKind, + RadrootsTransportOutcome, RadrootsTransportOutcomeKind, RadrootsTransportPayload, + RadrootsTransportStatus, RadrootsTransportTargetReceipt, +}; use radroots_transport_nostr::{ RadrootsMockRelayPublishAdapter, RadrootsNostrTransport, RadrootsRelayOutcome, - phase1_publication_delivery_request, publish_claimed_phase1_publication_target_with_transport, + RadrootsRelayTransportError, execute_claimed_phase1_publication_target_with_transport, + phase1_publication_delivery_request, repair_phase1_publication_observation, }; use serde::Serialize; use sha2::{Digest, Sha256}; +use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; +use std::sync::atomic::{AtomicUsize, Ordering}; const SECRET_KEY: &str = "10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5"; const PUBLIC_KEY: &str = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; @@ -97,9 +109,105 @@ impl RadrootsPhase1PublicationSigner for CountingPhase1Signer { } } +struct ForgedPhase1ReceiptTransport; + +impl RadrootsTransport for ForgedPhase1ReceiptTransport { + fn transport_kind(&self) -> RadrootsTransportKind { + RadrootsTransportKind::Nostr + } + + fn status<'a>(&'a self) -> RadrootsTransportFuture<'a, RadrootsTransportStatus> { + Box::pin(async { + RadrootsTransportStatus::new( + RadrootsTransportKind::Nostr, + true, + RadrootsTransportImplementationState::Real, + true, + "forged Phase 1 receipt fixture", + ) + }) + } + + fn deliver<'a>( + &'a self, + request: RadrootsTransportDeliveryRequest, + ) -> RadrootsTransportFuture<'a, RadrootsTransportDeliveryReceipt> { + Box::pin(async move { + let target = request.target_set().targets()[0].clone(); + RadrootsTransportDeliveryReceipt::new( + "forged-phase1-request", + request.target_set().clone(), + vec![RadrootsTransportTargetReceipt::new( + target, + RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Accepted), + )], + ) + }) + } + + fn fetch<'a>( + &'a self, + _request: RadrootsTransportFetchRequest, + ) -> RadrootsTransportFuture<'a, RadrootsTransportFetchReceipt> { + Box::pin(async { Err(RadrootsTransportError::UnsupportedOperation) }) + } +} + +#[derive(Default)] +struct DuplicateAcceptedPhase1Transport { + invocations: AtomicUsize, +} + +impl DuplicateAcceptedPhase1Transport { + fn invocations(&self) -> usize { + self.invocations.load(Ordering::SeqCst) + } +} + +impl RadrootsTransport for DuplicateAcceptedPhase1Transport { + fn transport_kind(&self) -> RadrootsTransportKind { + RadrootsTransportKind::Nostr + } + + fn status<'a>(&'a self) -> RadrootsTransportFuture<'a, RadrootsTransportStatus> { + Box::pin(async { + RadrootsTransportStatus::new( + RadrootsTransportKind::Nostr, + true, + RadrootsTransportImplementationState::Real, + true, + "duplicate accepted Phase 1 fixture", + ) + }) + } + + fn deliver<'a>( + &'a self, + request: RadrootsTransportDeliveryRequest, + ) -> RadrootsTransportFuture<'a, RadrootsTransportDeliveryReceipt> { + Box::pin(async move { + self.invocations.fetch_add(1, Ordering::SeqCst); + let target = request.target_set().targets()[0].clone(); + let receipt = RadrootsTransportTargetReceipt::skipped( + target, + RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::DuplicateAccepted), + )?; + RadrootsTransportDeliveryReceipt::for_request(&request, vec![receipt]) + }) + } + + fn fetch<'a>( + &'a self, + _request: RadrootsTransportFetchRequest, + ) -> RadrootsTransportFuture<'a, RadrootsTransportFetchReceipt> { + Box::pin(async { Err(RadrootsTransportError::UnsupportedOperation) }) + } +} + #[tokio::test] async fn outbox_publication_all_seven_leaves_reuse_exact_bytes_and_dispatch_identity() { let outbox = RadrootsOutbox::open_memory().await.unwrap(); + let event_store = RadrootsEventStore::open_memory().await.unwrap(); let signer = CountingPhase1Signer::fixture(); let policy = RadrootsPhase1PublicationTargetPolicy::new([RELAY], 1).unwrap(); let ready_artifacts = all_ready_artifacts(); @@ -179,26 +287,24 @@ async fn outbox_publication_all_seven_leaves_reuse_exact_bytes_and_dispatch_iden .expect("bounded relay outcome"), ); let first_transport = RadrootsNostrTransport::new(first_adapter.clone()); - let first_receipt = publish_claimed_phase1_publication_target_with_transport( + let retryable = execute_claimed_phase1_publication_target_with_transport( + &outbox, + &event_store, &first_transport, &first_claim, + base + 7, base + 5, ) .await .unwrap(); assert_eq!( - first_receipt.target_receipts()[0].status(), - RadrootsTransportDeliveryTargetStatus::FailedRetryable + retryable.targets()[0].state(), + RadrootsPhase1PublicationTargetState::FailedRetryable ); assert_eq!( first_adapter.captured_raw_events().as_slice(), core::slice::from_ref(&exact_raw) ); - - let retryable = outbox - .fail_phase1_target_retryable(&first_claim, base + 6, base + 7, "relay unavailable") - .await - .unwrap(); let retry_target = retryable .targets() .iter() @@ -222,29 +328,246 @@ async fn outbox_publication_all_seven_leaves_reuse_exact_bytes_and_dispatch_iden let retry_adapter = RadrootsMockRelayPublishAdapter::new(); let retry_transport = RadrootsNostrTransport::new(retry_adapter.clone()); - let retry_receipt = publish_claimed_phase1_publication_target_with_transport( + let accepted = execute_claimed_phase1_publication_target_with_transport( + &outbox, + &event_store, &retry_transport, &retry_claim, + base + 10, base + 8, ) .await .unwrap(); assert_eq!( - retry_receipt.target_receipts()[0].status(), - RadrootsTransportDeliveryTargetStatus::Accepted + accepted.targets()[0].state(), + RadrootsPhase1PublicationTargetState::AcceptedObserved ); assert_eq!( retry_adapter.captured_raw_events().as_slice(), core::slice::from_ref(&exact_raw) ); - outbox - .complete_phase1_target_accepted_observed(&retry_claim, base + 9) + let observations = event_store + .observations_for_event(verified.signed_event().id_str()) .await .unwrap(); + assert_eq!(observations.len(), 1); + assert_eq!(observations[0].observation_count, 1); assert_eq!(signer.invocations(), index + 1, "retry must not re-sign"); } } +#[tokio::test] +async fn outbox_publication_partial_effect_repairs_observation_without_republish() { + let temp = tempfile::tempdir().unwrap(); + let outbox_path = temp.path().join("phase1-partial-effect.sqlite"); + let outbox = RadrootsOutbox::open_file(&outbox_path).await.unwrap(); + let injection_pool = SqlitePoolOptions::new() + .max_connections(1) + .connect_with(SqliteConnectOptions::new().filename(&outbox_path)) + .await + .unwrap(); + let event_store = RadrootsEventStore::open_memory().await.unwrap(); + let signer = CountingPhase1Signer::fixture(); + let ready = all_ready_artifacts().remove(0); + let enqueue = outbox + .enqueue_phase1_publication( + &ready, + &RadrootsPhase1PublicationTargetPolicy::new([RELAY], 1).unwrap(), + 1_000, + ) + .await + .unwrap(); + let signing_claim = outbox + .claim_phase1_publication_for_signing( + enqueue.record().publication_id(), + enqueue.record().revision(), + 1_001, + 100, + ) + .await + .unwrap(); + let preflight = outbox + .preflight_phase1_publication_signing(&signing_claim, 1_002) + .await + .unwrap(); + let contract = event_contract(ready.artifact().event_contract_id()).unwrap(); + let actor = RadrootsActorContext::test(PUBLIC_KEY, [contract.author_role]).unwrap(); + let verified = sign_authorized_phase1_publication(&actor, &signer, &ready).unwrap(); + let signed = outbox + .complete_phase1_publication_signing(&preflight, &verified, 1_003) + .await + .unwrap(); + let target = &signed.targets()[0]; + let claim = outbox + .claim_phase1_publication_target( + signed.publication_id(), + signed.revision(), + target.target_id(), + target.revision(), + 1_004, + 100, + ) + .await + .unwrap(); + let forged_error = execute_claimed_phase1_publication_target_with_transport( + &outbox, + &event_store, + &ForgedPhase1ReceiptTransport, + &claim, + 1_010, + 1_005, + ) + .await + .unwrap_err(); + assert!(matches!( + forged_error, + RadrootsRelayTransportError::TransportContractError( + RadrootsTransportError::DeliveryReceiptRequestIdMismatch + ) + )); + let uncertain = outbox + .load_phase1_publication(signed.publication_id()) + .await + .unwrap(); + assert_eq!( + uncertain.targets()[0].state(), + RadrootsPhase1PublicationTargetState::Uncertain + ); + let retry_claim = outbox + .claim_phase1_publication_target( + uncertain.publication_id(), + uncertain.revision(), + uncertain.targets()[0].target_id(), + uncertain.targets()[0].revision(), + 1_006, + 100, + ) + .await + .unwrap(); + assert_eq!(retry_claim.dispatch_digest(), claim.dispatch_digest()); + sqlx::query( + "CREATE TRIGGER fail_phase1_accepted_receipt + BEFORE INSERT ON outbox_phase1_target_receipt + BEGIN SELECT RAISE(ABORT, 'injected accepted receipt failure'); END", + ) + .execute(&injection_pool) + .await + .unwrap(); + let accepted_adapter = RadrootsMockRelayPublishAdapter::new(); + let accepted_transport = RadrootsNostrTransport::new(accepted_adapter.clone()); + let persistence_failure = execute_claimed_phase1_publication_target_with_transport( + &outbox, + &event_store, + &accepted_transport, + &retry_claim, + 1_010, + 1_007, + ) + .await; + assert!( + matches!( + persistence_failure, + Err(RadrootsRelayTransportError::Phase1Publication(_)) + ), + "unexpected persistence result: {persistence_failure:?}" + ); + assert_eq!(accepted_adapter.captured_raw_events().len(), 1); + sqlx::query("DROP TRIGGER fail_phase1_accepted_receipt") + .execute(&injection_pool) + .await + .unwrap(); + let uncertain_after_acceptance = outbox + .load_phase1_publication(signed.publication_id()) + .await + .unwrap(); + assert_eq!( + uncertain_after_acceptance.targets()[0].state(), + RadrootsPhase1PublicationTargetState::Uncertain + ); + let duplicate_claim = outbox + .claim_phase1_publication_target( + uncertain_after_acceptance.publication_id(), + uncertain_after_acceptance.revision(), + uncertain_after_acceptance.targets()[0].target_id(), + uncertain_after_acceptance.targets()[0].revision(), + 1_008, + 100, + ) + .await + .unwrap(); + assert_eq!(duplicate_claim.dispatch_digest(), claim.dispatch_digest()); + sqlx::query( + "CREATE TEMP TRIGGER fail_phase1_publish_observation + BEFORE INSERT ON event_transport_observation + BEGIN SELECT RAISE(ABORT, 'injected observation failure'); END", + ) + .execute(event_store.pool()) + .await + .unwrap(); + let duplicate_transport = DuplicateAcceptedPhase1Transport::default(); + let observation_failure = execute_claimed_phase1_publication_target_with_transport( + &outbox, + &event_store, + &duplicate_transport, + &duplicate_claim, + 1_010, + 1_009, + ) + .await; + assert!( + matches!( + observation_failure, + Err(RadrootsRelayTransportError::EventStore(_)) + ), + "unexpected observation result: {observation_failure:?}" + ); + let pending = outbox + .load_phase1_publication(signed.publication_id()) + .await + .unwrap(); + assert_eq!( + pending.targets()[0].state(), + RadrootsPhase1PublicationTargetState::AcceptedObservationPending + ); + assert_eq!(duplicate_transport.invocations(), 1); + + sqlx::query("DROP TRIGGER fail_phase1_publish_observation") + .execute(event_store.pool()) + .await + .unwrap(); + let repaired = repair_phase1_publication_observation( + &outbox, + &event_store, + pending.publication_id(), + pending.targets()[0].target_id(), + 1_010, + ) + .await + .unwrap(); + assert_eq!( + repaired.targets()[0].state(), + RadrootsPhase1PublicationTargetState::AcceptedObserved + ); + let replayed = repair_phase1_publication_observation( + &outbox, + &event_store, + repaired.publication_id(), + repaired.targets()[0].target_id(), + 1_011, + ) + .await + .unwrap(); + assert_eq!(replayed.targets(), repaired.targets()); + assert_eq!(accepted_adapter.captured_raw_events().len(), 1); + assert_eq!(duplicate_transport.invocations(), 1); + let observations = event_store + .observations_for_event(verified.signed_event().id_str()) + .await + .unwrap(); + assert_eq!(observations.len(), 1); + assert_eq!(observations[0].observation_count, 1); +} + fn assert_payload_exact(payload: &RadrootsTransportPayload, expected_raw: &str) { let (_, raw_json) = payload .signed_event_json_parts() diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -34,15 +34,15 @@ use radroots_transport_nostr::{ RADROOTS_RELAY_FETCH_FILTER_LIMIT_MAX, RADROOTS_RELAY_FETCH_FILTER_SET_JSON_BYTE_LIMIT_MAX, RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX, RADROOTS_RELAY_FETCH_RAW_JSON_BYTE_LIMIT_MAX, RADROOTS_RELAY_FETCH_TIMEOUT_MS_MAX, RadrootsMockRelayFetchAdapter, - RadrootsMockRelayPublishAdapter, RadrootsNostrTransport, RadrootsOutboxPublishPolicy, - RadrootsRelayFetchEventAdmission, RadrootsRelayFetchEventValidStream, - RadrootsRelayFetchEventVerification, RadrootsRelayFetchEventVisibility, - RadrootsRelayFetchFailure, RadrootsRelayFetchFilters, RadrootsRelayFetchItem, - RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchReceipt, - RadrootsRelayFetchRelayOutcome, RadrootsRelayFetchRequest, RadrootsRelayFetchedEvent, - RadrootsRelayOutcome, RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, - RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, RadrootsRelayTargetSet, - RadrootsRelayTransportError, RadrootsRelayUrl, RadrootsRelayUrlPolicy, + RadrootsMockRelayPublishAdapter, RadrootsNostrTransport, RadrootsOutboxPublishEmptyTargetState, + RadrootsOutboxPublishPolicy, RadrootsRelayFetchEventAdmission, + RadrootsRelayFetchEventValidStream, RadrootsRelayFetchEventVerification, + RadrootsRelayFetchEventVisibility, RadrootsRelayFetchFailure, RadrootsRelayFetchFilters, + RadrootsRelayFetchItem, RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, + RadrootsRelayFetchReceipt, RadrootsRelayFetchRelayOutcome, RadrootsRelayFetchRequest, + RadrootsRelayFetchedEvent, RadrootsRelayOutcome, RadrootsRelayOutcomeKind, + RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, + RadrootsRelayTargetSet, RadrootsRelayTransportError, RadrootsRelayUrl, RadrootsRelayUrlPolicy, fetch_and_ingest_relay_events, fetch_relay_events, fetch_relay_events_blocking, publish_claimed_outbox_event, publish_claimed_outbox_event_with_transport, publish_signed_event, verified_signed_event_payload, @@ -4479,7 +4479,7 @@ async fn outbox_transport_facade_rejects_receipts_forged_for_another_request() { .expect_err("forged receipt rejected"); assert!(matches!( error, - RadrootsRelayTransportError::TransportContract(_) + RadrootsRelayTransportError::TransportContractError(_) )); let event = outbox .get_event(receipt.outbox_event_id) @@ -4543,6 +4543,10 @@ async fn outbox_transport_facade_handles_empty_and_invalid_claim_plans() { assert_eq!(published.attempted_count(), 0); assert_eq!(published.accepted_count(), 2); assert!(published.quorum_met()); + assert_eq!( + published.empty_target_state(), + Some(RadrootsOutboxPublishEmptyTargetState::AlreadySatisfied) + ); let second_draft = generic_draft("invalid claimed plan"); outbox diff --git a/tools/xtask/src/contract/outbox_phase1_publication.rs b/tools/xtask/src/contract/outbox_phase1_publication.rs @@ -391,8 +391,8 @@ fn validate_identity(identity: &Identity, domain: &str, preimage: &[&str]) -> Re } fn validate_transitions(descriptor: &Descriptor) -> Result<(), String> { - if descriptor.transitions.len() != 28 { - return Err("Phase 1 publication transition inventory must contain 28 entries".to_owned()); + if descriptor.transitions.len() != 32 { + return Err("Phase 1 publication transition inventory must contain 32 entries".to_owned()); } let event_states = descriptor.event_states.iter().collect::<BTreeSet<_>>(); let target_states = descriptor.target_states.iter().collect::<BTreeSet<_>>(); @@ -565,7 +565,7 @@ fn manifest_schema() -> Value { "properties": { "event_state_count": { "const": 9 }, "target_state_count": { "const": 8 }, - "transition_count": { "const": 28 }, + "transition_count": { "const": 32 }, "stable_error_count": { "const": 25 } } },