commit f5038dbc581eb48f4b238b7a5e2732948281abfb
parent e0131fcb84123ce5939f4eec38d99f7b484023f5
Author: triesap <tyson@radroots.org>
Date: Tue, 30 Jun 2026 09:58:29 +0000
sdk: preflight workflow idempotency
Diffstat:
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]