lib

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

commit 06b5d28cd6a4c4388117db0757e4544c5fef95a0
parent ed0710cc7011a29b6f8cc6dcfa5b9f5f2b3887c0
Author: triesap <tyson@radroots.org>
Date:   Tue, 30 Jun 2026 09:58:29 +0000

sdk: preflight workflow idempotency

Diffstat:
Mcrates/sdk/src/workflow_runtime.rs | 32+++++++++++++++++++++++---------
Mcrates/sdk/tests/farms_runtime.rs | 63++++++++++++++++++++++++++++++++++++++++++++++++++-------------
Mcrates/sdk/tests/listings_runtime.rs | 63++++++++++++++++++++++++++++++++++++++++++++++++++-------------
Mcrates/sdk/tests/orders_runtime.rs | 61+++++++++++++++++++++++++++++++++++++++++++++++++------------
Mcrates/sdk/tests/unit/orders_runtime_tests.rs | 40++++++++++++++++++++++++++++------------
Mcrates/sdk/tests/unit/workflow_runtime_tests.rs | 85+++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------------
6 files changed, 257 insertions(+), 87 deletions(-)

diff --git a/crates/sdk/src/workflow_runtime.rs b/crates/sdk/src/workflow_runtime.rs @@ -12,6 +12,7 @@ use radroots_events::{ ids::RadrootsEventId, }; use radroots_outbox::{RadrootsOutboxEnqueueStatus, RadrootsOutboxSignedOperationInput}; +#[cfg(test)] use sha2::{Digest, Sha256}; pub(crate) struct SdkWorkflowEnqueueRequest<'a> { @@ -76,24 +77,35 @@ async fn enqueue_signed_workflow_event( let observed_at_ms = sdk_now_ms(sdk)?; let signed_event_id = RadrootsEventId::parse(request.frozen_draft.expected_event_id.as_str()) .expect("frozen workflow draft has a valid expected event id"); - let event = event_from_signed(&signed_event); - let ingest = RadrootsEventIngest::new(event, observed_at_ms) - .with_raw_json(signed_event.raw_json.clone()); - let ingest_receipt = sdk._event_store.ingest_event(ingest).await?; - let canonical_target_relays = target_relays.canonical_relays.clone(); + let allow_empty_target_relays = target_relays.allow_empty_target_relays; let target_relay_values = target_relays.relays; - let partial_failure_digest_prefix = outbox_idempotency_digest_prefix( + let idempotency_key_for_enqueue = idempotency_key.clone(); + let preflight_input = signed_outbox_input( request.operation_kind, request.frozen_draft, - canonical_target_relays.as_slice(), + signed_event.clone(), + target_relay_values.clone(), + idempotency_key, + allow_empty_target_relays, + false, + observed_at_ms, ); + let preflight = sdk + ._outbox + .preflight_signed_operation_idempotency(&preflight_input) + .await?; + let partial_failure_digest_prefix = digest_prefix(preflight.idempotency_digest.as_str()); + let event = event_from_signed(&signed_event); + let ingest = RadrootsEventIngest::new(event, observed_at_ms) + .with_raw_json(signed_event.raw_json.clone()); + let ingest_receipt = sdk._event_store.ingest_event(ingest).await?; let outbox_input = signed_outbox_input( request.operation_kind, request.frozen_draft, signed_event, target_relay_values, - idempotency_key, - target_relays.allow_empty_target_relays, + idempotency_key_for_enqueue, + allow_empty_target_relays, ingest_receipt.inserted, observed_at_ms, ); @@ -175,6 +187,7 @@ fn resolved_target_relays( } #[derive(serde::Serialize)] +#[cfg(test)] struct SdkWorkflowOutboxDigestInput<'a> { operation_kind: &'static str, expected_pubkey: &'a str, @@ -182,6 +195,7 @@ struct SdkWorkflowOutboxDigestInput<'a> { target_relays: &'a [String], } +#[cfg(test)] fn outbox_idempotency_digest_prefix( operation_kind: &'static str, frozen_draft: &RadrootsFrozenEventDraft, diff --git a/crates/sdk/tests/farms_runtime.rs b/crates/sdk/tests/farms_runtime.rs @@ -20,9 +20,9 @@ use radroots_sdk::{ FarmPrivateLocationSetResult, FarmPrivateLocationUpsertRequest, Geocoder, GeocoderLocalityQuery, PushOutboxEventState, PushOutboxRelayOutcomeKind, PushOutboxRequest, RadrootsClient, RadrootsSdkError, RadrootsSdkErrorClass, RadrootsSdkGeoNamesErrorKind, - RadrootsSdkPartialLocalMutationFailure, RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, - SdkExactLocation, SdkIdempotencyKey, SdkMutationState, SdkPublicLocality, SdkRelayTargetPolicy, - SdkRelayTargetSet, SdkRelayUrlPolicy, StorageStatusRequest, + RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, SdkExactLocation, SdkIdempotencyKey, + SdkMutationState, SdkPublicLocality, SdkRelayTargetPolicy, SdkRelayTargetSet, + SdkRelayUrlPolicy, StorageStatusRequest, }; use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; @@ -824,7 +824,7 @@ async fn farm_enqueue_publish_pushes_queued_event_with_mock_relay_sync() { } #[tokio::test] -async fn farm_enqueue_publish_reports_partial_local_mutation_after_outbox_conflict() { +async fn farm_enqueue_publish_reports_preflight_idempotency_conflict_without_mutation() { let (_tempdir, sdk) = directory_sdk().await; let first = FarmEnqueuePublishRequest::new( farmer_actor(), @@ -837,6 +837,29 @@ async fn farm_enqueue_publish_reports_partial_local_mutation_after_outbox_confli .enqueue_publish_with_explicit_signer(first, &FixtureSigner::new(FARMER)) .await .expect("first enqueue"); + let paths = sdk.storage_paths().expect("paths"); + let event_store = RadrootsEventStore::open_file(&paths.event_store_path) + .await + .expect("event store"); + let outbox = RadrootsOutbox::open_file(&paths.outbox_path) + .await + .expect("outbox"); + assert_eq!( + event_store + .status_summary() + .await + .expect("event store status") + .total_events, + 1 + ); + assert_eq!( + outbox + .status_summary(0) + .await + .expect("outbox status") + .total_events, + 1 + ); let second = FarmEnqueuePublishRequest::new( farmer_actor(), @@ -849,20 +872,34 @@ async fn farm_enqueue_publish_reports_partial_local_mutation_after_outbox_confli .farms() .enqueue_publish_with_explicit_signer(second, &FixtureSigner::new(FARMER)) .await - .expect_err("partial"); + .expect_err("conflict"); assert!(matches!( error, - RadrootsSdkError::PartialLocalMutation(ref partial) - if partial.stored - && !partial.queued - && partial.event_id.is_some() - && partial.operation_kind == FARM_PUBLISH_OPERATION_KIND - && partial.idempotency_digest_prefix.is_some() - && partial.failure == RadrootsSdkPartialLocalMutationFailure::OutboxIdempotencyConflict - && partial.recovery == RadrootsSdkRecoveryAction::RetryOperationWithSameIdempotencyKey + RadrootsSdkError::IdempotencyConflict { ref operation_kind, .. } + if operation_kind == FARM_PUBLISH_OPERATION_KIND )); + assert_eq!( + error.recovery_actions(), + vec![RadrootsSdkRecoveryAction::RetryOperationWithSameIdempotencyKey] + ); assert!(!error.to_string().contains("farm-idem-e")); + assert_eq!( + event_store + .status_summary() + .await + .expect("event store status after conflict") + .total_events, + 1 + ); + assert_eq!( + outbox + .status_summary(0) + .await + .expect("outbox status after conflict") + .total_events, + 1 + ); } #[tokio::test] diff --git a/crates/sdk/tests/listings_runtime.rs b/crates/sdk/tests/listings_runtime.rs @@ -19,9 +19,9 @@ use radroots_events::{ use radroots_outbox::{RadrootsOutbox, RadrootsOutboxEventState}; use radroots_sdk::{ LISTING_PUBLISH_OPERATION_KIND, ListingEnqueuePublishRequest, ListingPreparePublishRequest, - RadrootsClient, RadrootsSdkError, RadrootsSdkPartialLocalMutationFailure, - RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, SdkIdempotencyKey, SdkMutationState, - SdkRelayTargetPolicy, SdkRelayTargetSet, SdkRelayUrlPolicy, + RadrootsClient, RadrootsSdkError, RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, + SdkIdempotencyKey, SdkMutationState, SdkRelayTargetPolicy, SdkRelayTargetSet, + SdkRelayUrlPolicy, }; use radroots_trade::listing::RadrootsListingDraftDocumentV1; @@ -640,7 +640,7 @@ async fn enqueue_publish_returns_sanitized_signer_errors() { } #[tokio::test] -async fn enqueue_publish_reports_partial_local_mutation_after_outbox_conflict() { +async fn enqueue_publish_reports_preflight_idempotency_conflict_without_mutation() { let (_tempdir, sdk) = directory_sdk().await; let first = ListingEnqueuePublishRequest::new( actor(), @@ -653,6 +653,29 @@ async fn enqueue_publish_reports_partial_local_mutation_after_outbox_conflict() .enqueue_publish_with_explicit_signer(first, &FixtureSigner::new(SELLER)) .await .expect("first enqueue"); + let paths = sdk.storage_paths().expect("paths"); + let event_store = RadrootsEventStore::open_file(&paths.event_store_path) + .await + .expect("event store"); + let outbox = RadrootsOutbox::open_file(&paths.outbox_path) + .await + .expect("outbox"); + assert_eq!( + event_store + .status_summary() + .await + .expect("event store status") + .total_events, + 1 + ); + assert_eq!( + outbox + .status_summary(0) + .await + .expect("outbox status") + .total_events, + 1 + ); let second = ListingEnqueuePublishRequest::new( actor(), @@ -665,20 +688,34 @@ async fn enqueue_publish_reports_partial_local_mutation_after_outbox_conflict() .listings() .enqueue_publish_with_explicit_signer(second, &FixtureSigner::new(SELLER)) .await - .expect_err("partial"); + .expect_err("conflict"); assert!(matches!( error, - RadrootsSdkError::PartialLocalMutation(ref partial) - if partial.stored - && !partial.queued - && partial.event_id.is_some() - && partial.operation_kind == LISTING_PUBLISH_OPERATION_KIND - && partial.idempotency_digest_prefix.is_some() - && partial.failure == RadrootsSdkPartialLocalMutationFailure::OutboxIdempotencyConflict - && partial.recovery == RadrootsSdkRecoveryAction::RetryOperationWithSameIdempotencyKey + RadrootsSdkError::IdempotencyConflict { ref operation_kind, .. } + if operation_kind == LISTING_PUBLISH_OPERATION_KIND )); + assert_eq!( + error.recovery_actions(), + vec![RadrootsSdkRecoveryAction::RetryOperationWithSameIdempotencyKey] + ); assert!(!error.to_string().contains("idem-d")); + assert_eq!( + event_store + .status_summary() + .await + .expect("event store status after conflict") + .total_events, + 1 + ); + assert_eq!( + outbox + .status_summary(0) + .await + .expect("outbox status after conflict") + .total_events, + 1 + ); } #[tokio::test] diff --git a/crates/sdk/tests/orders_runtime.rs b/crates/sdk/tests/orders_runtime.rs @@ -32,11 +32,11 @@ use radroots_nostr::prelude::{ use radroots_outbox::RadrootsOutbox; use radroots_sdk::{ AckPolicy, DvmValidationReceiptIngestRequest, PublishMode, RadrootsClient, RadrootsSdkError, - RadrootsSdkPartialLocalMutationFailure, RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, - RelayResolutionPolicy, SdkMutationState, SdkRelayTargetSet, SdkRelayUrlPolicy, - SdkTradeStatusIssue, SdkTradeStatusIssueKind, SdkTradeStatusSource, TRADE_STATUS_DEFAULT_LIMIT, - TRADE_STATUS_MAX_LIMIT, TRADE_SUBMIT_OPERATION_KIND, TradeAcceptRequest, TradeCancelRequest, - TradeDeclineRequest, TradeEvidenceIngestRequest, TradeMutationOutcome, TradeProposeRequest, + RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, RelayResolutionPolicy, SdkMutationState, + SdkRelayTargetSet, SdkRelayUrlPolicy, SdkTradeStatusIssue, SdkTradeStatusIssueKind, + SdkTradeStatusSource, TRADE_STATUS_DEFAULT_LIMIT, TRADE_STATUS_MAX_LIMIT, + TRADE_SUBMIT_OPERATION_KIND, TradeAcceptRequest, TradeCancelRequest, TradeDeclineRequest, + TradeEvidenceIngestRequest, TradeMutationOutcome, TradeProposeRequest, TradeRequestEvidenceIngestRequest, TradeResyncRequest, TradeRevisionDecisionRequest, TradeRevisionProposalRequest, TradeSellerInboxRequest, TradeStatusKind, TradeStatusNextActionKind, TradeStatusRequest, @@ -1148,6 +1148,13 @@ async fn trade_product_propose_idempotency_replays_same_payload_and_conflicts_di let tempdir = tempfile::tempdir().expect("tempdir"); let storage_root = tempdir.path().join("sdk"); let buyer_sdk = directory_sdk_with_signer(storage_root.as_path(), BUYER_SECRET_KEY_HEX).await; + let storage_paths = buyer_sdk.storage_paths().expect("storage paths"); + let store = RadrootsEventStore::open_file(&storage_paths.event_store_path) + .await + .expect("event store"); + let outbox = RadrootsOutbox::open_file(&storage_paths.outbox_path) + .await + .expect("outbox"); let request = TradeProposeRequest::new( buyer_actor(), listing_event_ptr(), @@ -1185,6 +1192,22 @@ async fn trade_product_propose_idempotency_replays_same_payload_and_conflicts_di .idempotency .safe_to_retry_with_same_idempotency_key ); + assert_eq!( + store + .status_summary() + .await + .expect("event store summary") + .total_events, + 1 + ); + assert_eq!( + outbox + .status_summary(i64::MAX) + .await + .expect("outbox summary") + .total_events, + 1 + ); let conflict = buyer_sdk .trades() @@ -1206,14 +1229,28 @@ async fn trade_product_propose_idempotency_replays_same_payload_and_conflicts_di assert!(matches!( conflict, - RadrootsSdkError::PartialLocalMutation(ref partial) - if partial.stored - && !partial.queued - && partial.operation_kind == TRADE_SUBMIT_OPERATION_KIND - && partial.failure == RadrootsSdkPartialLocalMutationFailure::OutboxIdempotencyConflict - && partial.recovery == RadrootsSdkRecoveryAction::RetryOperationWithSameIdempotencyKey + RadrootsSdkError::IdempotencyConflict { + ref operation_kind, + .. + } if operation_kind == TRADE_SUBMIT_OPERATION_KIND )); - assert_eq!(conflict.code(), "partial_local_mutation"); + assert_eq!(conflict.code(), "idempotency_conflict"); + assert_eq!( + store + .status_summary() + .await + .expect("event store summary") + .total_events, + 1 + ); + assert_eq!( + outbox + .status_summary(i64::MAX) + .await + .expect("outbox summary") + .total_events, + 1 + ); } #[cfg(all(feature = "signer-adapters", feature = "local-signer"))] diff --git a/crates/sdk/tests/unit/orders_runtime_tests.rs b/crates/sdk/tests/unit/orders_runtime_tests.rs @@ -653,15 +653,16 @@ fn assert_error_display<T: core::fmt::Debug>(result: Result<T, RadrootsSdkError> assert!(result.unwrap_err().to_string().contains(expected)); } -fn assert_partial_outbox_enqueue(error: RadrootsSdkError, operation_kind: &str) { - assert!(matches!( - error, - RadrootsSdkError::PartialLocalMutation(partial) - if partial.operation_kind == operation_kind - && partial.stored - && !partial.queued - && partial.failure == crate::RadrootsSdkPartialLocalMutationFailure::OutboxEnqueue - )); +fn assert_outbox_preflight_error(error: RadrootsSdkError) { + assert!(matches!(error, RadrootsSdkError::Outbox { .. })); +} + +async fn local_event_count(sdk: &RadrootsClient) -> i64 { + sdk._event_store + .status_summary() + .await + .expect("event store summary") + .total_events } #[test] @@ -3389,6 +3390,7 @@ async fn prepared_lifecycle_enqueues_report_closed_outbox_after_preflight() { proposal_payload, )) .expect("proposal plan"); + let proposal_events_before = local_event_count(&proposal_sdk).await; proposal_sdk._outbox.pool().close().await; let proposal_error = proposal_sdk .trades() @@ -3403,7 +3405,11 @@ async fn prepared_lifecycle_enqueues_report_closed_outbox_after_preflight() { ) .await .expect_err("closed outbox proposal"); - assert_partial_outbox_enqueue(proposal_error, TRADE_REVISION_PROPOSAL_OPERATION_KIND); + assert_outbox_preflight_error(proposal_error); + assert_eq!( + local_event_count(&proposal_sdk).await, + proposal_events_before + ); let revision_sdk = prepared_order_sdk().await; let revision_submit = @@ -3445,6 +3451,7 @@ async fn prepared_lifecycle_enqueues_report_closed_outbox_after_preflight() { fixture_revision_decision(&proposal_payload, &proposal.signed_event_id), )) .expect("revision plan"); + let revision_events_before = local_event_count(&revision_sdk).await; revision_sdk._outbox.pool().close().await; let revision_error = revision_sdk .trades() @@ -3459,7 +3466,11 @@ async fn prepared_lifecycle_enqueues_report_closed_outbox_after_preflight() { ) .await .expect_err("closed outbox revision"); - assert_partial_outbox_enqueue(revision_error, TRADE_REVISION_DECISION_OPERATION_KIND); + assert_outbox_preflight_error(revision_error); + assert_eq!( + local_event_count(&revision_sdk).await, + revision_events_before + ); let cancellation_sdk = prepared_order_sdk().await; let cancellation_submit = @@ -3473,6 +3484,7 @@ async fn prepared_lifecycle_enqueues_report_closed_outbox_after_preflight() { fixture_cancellation("order-closed-outbox-cancellation"), )) .expect("cancellation plan"); + let cancellation_events_before = local_event_count(&cancellation_sdk).await; cancellation_sdk._outbox.pool().close().await; let cancellation_error = cancellation_sdk .trades() @@ -3487,7 +3499,11 @@ async fn prepared_lifecycle_enqueues_report_closed_outbox_after_preflight() { ) .await .expect_err("closed outbox cancellation"); - assert_partial_outbox_enqueue(cancellation_error, TRADE_CANCELLATION_OPERATION_KIND); + assert_outbox_preflight_error(cancellation_error); + assert_eq!( + local_event_count(&cancellation_sdk).await, + cancellation_events_before + ); } #[tokio::test] diff --git a/crates/sdk/tests/unit/workflow_runtime_tests.rs b/crates/sdk/tests/unit/workflow_runtime_tests.rs @@ -166,6 +166,22 @@ async fn enqueue_signed_workflow_stores_signed_event_and_reports_idempotency_con assert!(receipt.outbox_operation_id > 0); assert!(receipt.outbox_event_id > 0); assert_eq!(receipt.idempotency_digest_prefix.len(), 12); + assert_eq!( + sdk._event_store + .status_summary() + .await + .expect("event store summary") + .total_events, + 1 + ); + assert_eq!( + sdk._outbox + .status_summary(i64::MAX) + .await + .expect("outbox summary") + .total_events, + 1 + ); let second_draft = frozen_draft_for_d_tag(FARMER_PUBLIC_KEY_HEX, "workflow-conflict"); let error = match enqueue_signed_workflow( @@ -185,21 +201,29 @@ async fn enqueue_signed_workflow_stores_signed_event_and_reports_idempotency_con Ok(_) => panic!("expected idempotency conflict"), }; - match error { - RadrootsSdkError::PartialLocalMutation(partial) => { - assert!(partial.stored); - assert!(!partial.queued); - assert_eq!( - partial.failure, - crate::RadrootsSdkPartialLocalMutationFailure::OutboxIdempotencyConflict - ); - assert_eq!( - partial.idempotency_digest_prefix.as_deref().map(str::len), - Some(12) - ); - } - other => panic!("unexpected workflow error: {other:?}"), - } + assert!(matches!( + error, + RadrootsSdkError::IdempotencyConflict { + operation_kind, + .. + } if operation_kind == "workflow.test.v1" + )); + assert_eq!( + sdk._event_store + .status_summary() + .await + .expect("event store summary") + .total_events, + 1 + ); + assert_eq!( + sdk._outbox + .status_summary(i64::MAX) + .await + .expect("outbox summary") + .total_events, + 1 + ); } #[cfg(feature = "signer-adapters")] @@ -240,12 +264,20 @@ async fn enqueue_configured_signed_workflow_uses_sdk_signer_provider() { } #[tokio::test] -async fn enqueue_signed_workflow_reports_partial_mutation_when_outbox_fails() { +async fn enqueue_signed_workflow_reports_outbox_preflight_failure_without_mutation() { let sdk = crate::RadrootsClient::builder() .relay_url("wss://relay.example.com") .build() .await .expect("sdk"); + assert_eq!( + sdk._event_store + .status_summary() + .await + .expect("event store summary") + .total_events, + 0 + ); sdk._outbox.pool().close().await; let actor = RadrootsActorContext::test(FARMER_PUBLIC_KEY_HEX, [RadrootsActorRole::Farmer]) .expect("actor"); @@ -263,18 +295,15 @@ async fn enqueue_signed_workflow_reports_partial_mutation_when_outbox_fails() { Ok(_) => panic!("expected closed outbox error"), }; - match error { - RadrootsSdkError::PartialLocalMutation(partial) => { - assert!(partial.stored); - assert!(!partial.queued); - assert_eq!(partial.operation_kind, "workflow.test.v1"); - assert_eq!( - partial.failure, - crate::RadrootsSdkPartialLocalMutationFailure::OutboxEnqueue - ); - } - other => panic!("unexpected workflow error: {other:?}"), - } + assert!(matches!(error, RadrootsSdkError::Outbox { .. })); + assert_eq!( + sdk._event_store + .status_summary() + .await + .expect("event store summary") + .total_events, + 0 + ); } #[tokio::test]