lib

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

commit f54abd7f296006006d855ab37a20eba56e2c78d6
parent 296fbbfad980580370d2f851b3f07544e50b5a75
Author: triesap <tyson@radroots.org>
Date:   Tue, 30 Jun 2026 03:34:29 +0000

sdk: support trade product publish outcomes

Diffstat:
Mcrates/sdk/src/lib.rs | 8++++----
Mcrates/sdk/src/orders_runtime.rs | 426+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------
Mcrates/sdk/src/sync_runtime.rs | 28++++++++++++++++++++++++++--
Mcrates/sdk/tests/orders_runtime.rs | 139++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------------
Mcrates/sdk/tests/source_boundary.rs | 1+
Mcrates/sdk/tests/sync_runtime.rs | 2++
Acrates/sdk/tests/trade_product_publish_runtime.rs | 269+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/sdk/tests/unit/orders_runtime_tests.rs | 8++++++--
Mcrates/sdk/tests/unit/sync_runtime_tests.rs | 20++++++++++++++++----
9 files changed, 742 insertions(+), 159 deletions(-)

diff --git a/crates/sdk/src/lib.rs b/crates/sdk/src/lib.rs @@ -111,10 +111,10 @@ pub use crate::orders_runtime::{ TradeAcceptRequest, TradeCancelRequest, TradeCancellationEnqueueRequest, TradeCancellationPlan, TradeCancellationPrepareRequest, TradeCancellationReceipt, TradeDecisionEnqueueRequest, TradeDecisionPlan, TradeDecisionPrepareRequest, TradeDecisionReceipt, TradeDeclineRequest, - TradeEvidenceIngestReceipt, TradeEvidenceIngestRequest, TradeProposeRequest, - TradeRequestEvidenceIngestReceipt, TradeRequestEvidenceIngestRequest, TradeResyncReceipt, - TradeResyncRequest, TradeRevisionDecisionEnqueueRequest, TradeRevisionDecisionPlan, - TradeRevisionDecisionPrepareRequest, TradeRevisionDecisionReceipt, + TradeEvidenceIngestReceipt, TradeEvidenceIngestRequest, TradeMutationOutcome, + TradeProposeRequest, TradeRequestEvidenceIngestReceipt, TradeRequestEvidenceIngestRequest, + TradeResyncReceipt, TradeResyncRequest, TradeRevisionDecisionEnqueueRequest, + TradeRevisionDecisionPlan, TradeRevisionDecisionPrepareRequest, TradeRevisionDecisionReceipt, TradeRevisionDecisionRequest, TradeRevisionProposalEnqueueRequest, TradeRevisionProposalPlan, TradeRevisionProposalPrepareRequest, TradeRevisionProposalReceipt, TradeRevisionProposalRequest, TradeSellerInboxReceipt, TradeSellerInboxRequest, diff --git a/crates/sdk/src/orders_runtime.rs b/crates/sdk/src/orders_runtime.rs @@ -2,9 +2,10 @@ use crate::workflow_runtime::enqueue_configured_signed_workflow; #[cfg(feature = "runtime")] use crate::{ - AckPolicy, PublishMode, RadrootsSdkError, RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, - RelayResolutionPolicy, SdkIdempotencyKey, SdkMutationState, SdkRelayUrlPolicy, - TradeBuyerClient, TradeResyncClient, TradeSellerClient, TradeStatusClient, TradesClient, order, + AckPolicy, PublishMode, PushOutboxReceipt, PushOutboxRequest, RadrootsSdkError, + RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, RelayResolutionPolicy, SdkIdempotencyKey, + SdkMutationState, SdkRelayUrlPolicy, TradeBuyerClient, TradeResyncClient, TradeSellerClient, + TradeStatusClient, TradesClient, order, workflow_runtime::{SdkWorkflowEnqueueRequest, enqueue_signed_workflow}, }; #[cfg(feature = "runtime")] @@ -921,6 +922,22 @@ pub struct TradeCancellationReceipt { } #[cfg(feature = "runtime")] +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] +#[serde(rename_all = "snake_case", tag = "mode")] +pub enum TradeMutationOutcome<Plan, Receipt> { + DryRun { + plan: Plan, + }, + Enqueued { + receipt: Receipt, + }, + Published { + receipt: Receipt, + publish: PushOutboxReceipt, + }, +} + +#[cfg(feature = "runtime")] #[derive(Clone, Debug, serde::Serialize)] #[non_exhaustive] pub struct TradeProposeRequest { @@ -2446,52 +2463,98 @@ impl<'sdk> TradeBuyerClient<'sdk> { pub async fn propose_trade( &self, request: TradeProposeRequest, - ) -> Result<TradeSubmitReceipt, RadrootsSdkError> { - trades_client(self.sdk) - .enqueue_submit(TradeSubmitEnqueueRequest { - actor: request.actor, - listing_event: request.listing_event, - order: request.order, - target_relays: request.target_relays, - publish_mode: request.publish_mode, - ack_policy: request.ack_policy, - idempotency_key: request.idempotency_key, - created_at: request.created_at, - }) - .await + ) -> Result<TradeMutationOutcome<TradeSubmitPlan, TradeSubmitReceipt>, RadrootsSdkError> { + validate_trade_product_publish_policy(request.publish_mode, request.ack_policy)?; + let TradeProposeRequest { + actor, + listing_event, + order, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + created_at, + } = request; + let client = trades_client(self.sdk); + let plan = client.prepare_submit(TradeSubmitPrepareRequest { + actor: actor.clone(), + listing_event, + order, + created_at, + })?; + if publish_mode == PublishMode::DryRun { + return Ok(TradeMutationOutcome::DryRun { plan }); + } + let receipt = client + .enqueue_prepared_submit( + &actor, + plan, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + ) + .await?; + trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await } pub async fn cancel_trade( &self, request: TradeCancelRequest, - ) -> Result<TradeCancellationReceipt, RadrootsSdkError> { - let context = trade_mutation_context(self.sdk, request.locator, "trade.cancel").await?; + ) -> Result< + TradeMutationOutcome<TradeCancellationPlan, TradeCancellationReceipt>, + RadrootsSdkError, + > { + validate_trade_product_publish_policy(request.publish_mode, request.ack_policy)?; + let TradeCancelRequest { + actor, + locator, + reason, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + created_at, + } = request; + let context = trade_mutation_context(self.sdk, locator, "trade.cancel").await?; let cancellation = RadrootsOrderCancellation { order_id: context.order_id.clone(), listing_addr: context.listing_addr.clone(), buyer_pubkey: context.buyer_pubkey.clone(), seller_pubkey: context.seller_pubkey.clone(), - reason: request.reason, + reason, }; - trades_client(self.sdk) - .enqueue_cancellation(TradeCancellationEnqueueRequest { - actor: request.actor, - root_event: event_ptr(&context.root_event_id), - previous_event: event_ptr(&context.previous_event_id), - cancellation, - target_relays: request.target_relays, - publish_mode: request.publish_mode, - ack_policy: request.ack_policy, - idempotency_key: request.idempotency_key, - created_at: request.created_at, - }) - .await + let client = trades_client(self.sdk); + let plan = client.prepare_cancellation(TradeCancellationPrepareRequest { + actor: actor.clone(), + root_event: event_ptr(&context.root_event_id), + previous_event: event_ptr(&context.previous_event_id), + cancellation, + created_at, + })?; + if publish_mode == PublishMode::DryRun { + return Ok(TradeMutationOutcome::DryRun { plan }); + } + let receipt = client + .enqueue_prepared_cancellation( + &actor, + plan, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + ) + .await?; + trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await } pub async fn accept_revision( &self, mut request: TradeRevisionDecisionRequest, - ) -> Result<TradeRevisionDecisionReceipt, RadrootsSdkError> { + ) -> Result< + TradeMutationOutcome<TradeRevisionDecisionPlan, TradeRevisionDecisionReceipt>, + RadrootsSdkError, + > { request.decision = RadrootsOrderRevisionOutcome::Accepted; self.decide_revision(request).await } @@ -2499,44 +2562,70 @@ impl<'sdk> TradeBuyerClient<'sdk> { pub async fn decline_revision( &self, request: TradeRevisionDecisionRequest, - ) -> Result<TradeRevisionDecisionReceipt, RadrootsSdkError> { + ) -> Result< + TradeMutationOutcome<TradeRevisionDecisionPlan, TradeRevisionDecisionReceipt>, + RadrootsSdkError, + > { self.decide_revision(request).await } async fn decide_revision( &self, request: TradeRevisionDecisionRequest, - ) -> Result<TradeRevisionDecisionReceipt, RadrootsSdkError> { - let context = - trade_mutation_context(self.sdk, request.locator, "trade.revision_decision").await?; + ) -> Result< + TradeMutationOutcome<TradeRevisionDecisionPlan, TradeRevisionDecisionReceipt>, + RadrootsSdkError, + > { + validate_trade_product_publish_policy(request.publish_mode, request.ack_policy)?; + let TradeRevisionDecisionRequest { + actor, + locator, + revision_id, + decision, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + created_at, + } = request; + let context = trade_mutation_context(self.sdk, locator, "trade.revision_decision").await?; let previous_event_id = context.pending_revision_event_id.clone().ok_or_else(|| { RadrootsSdkError::InvalidRequest { message: "trade revision decision requires a pending revision".to_owned(), } })?; let decision = RadrootsOrderRevisionDecision { - revision_id: request.revision_id, + revision_id, order_id: context.order_id.clone(), listing_addr: context.listing_addr.clone(), buyer_pubkey: context.buyer_pubkey.clone(), seller_pubkey: context.seller_pubkey.clone(), root_event_id: context.root_event_id.clone(), prev_event_id: previous_event_id.clone(), - decision: request.decision, + decision, }; - trades_client(self.sdk) - .enqueue_revision_decision(TradeRevisionDecisionEnqueueRequest { - actor: request.actor, - root_event: event_ptr(&context.root_event_id), - previous_event: event_ptr(&previous_event_id), - decision, - target_relays: request.target_relays, - publish_mode: request.publish_mode, - ack_policy: request.ack_policy, - idempotency_key: request.idempotency_key, - created_at: request.created_at, - }) - .await + let client = trades_client(self.sdk); + let plan = client.prepare_revision_decision(TradeRevisionDecisionPrepareRequest { + actor: actor.clone(), + root_event: event_ptr(&context.root_event_id), + previous_event: event_ptr(&previous_event_id), + decision, + created_at, + })?; + if publish_mode == PublishMode::DryRun { + return Ok(TradeMutationOutcome::DryRun { plan }); + } + let receipt = client + .enqueue_prepared_revision_decision( + &actor, + plan, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + ) + .await?; + trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await } } @@ -2617,90 +2706,155 @@ impl<'sdk> TradeSellerClient<'sdk> { pub async fn accept_trade( &self, request: TradeAcceptRequest, - ) -> Result<TradeDecisionReceipt, RadrootsSdkError> { - let context = trade_mutation_context(self.sdk, request.locator, "trade.accept").await?; + ) -> Result<TradeMutationOutcome<TradeDecisionPlan, TradeDecisionReceipt>, RadrootsSdkError> + { + validate_trade_product_publish_policy(request.publish_mode, request.ack_policy)?; + let TradeAcceptRequest { + actor, + locator, + inventory_commitments, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + created_at, + } = request; + let context = trade_mutation_context(self.sdk, locator, "trade.accept").await?; let decision = RadrootsOrderDecision { order_id: context.order_id.clone(), listing_addr: context.listing_addr.clone(), buyer_pubkey: context.buyer_pubkey.clone(), seller_pubkey: context.seller_pubkey.clone(), decision: RadrootsOrderDecisionOutcome::Accepted { - inventory_commitments: request.inventory_commitments, + inventory_commitments, }, }; - trades_client(self.sdk) - .enqueue_decision(TradeDecisionEnqueueRequest { - actor: request.actor, - request_event: event_ptr(&context.root_event_id), - decision, - target_relays: request.target_relays, - publish_mode: request.publish_mode, - ack_policy: request.ack_policy, - idempotency_key: request.idempotency_key, - created_at: request.created_at, - }) - .await + let client = trades_client(self.sdk); + let plan = client.prepare_decision(TradeDecisionPrepareRequest { + actor: actor.clone(), + request_event: event_ptr(&context.root_event_id), + decision, + created_at, + })?; + if publish_mode == PublishMode::DryRun { + return Ok(TradeMutationOutcome::DryRun { plan }); + } + let receipt = client + .enqueue_prepared_decision( + &actor, + plan, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + ) + .await?; + trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await } pub async fn decline_trade( &self, request: TradeDeclineRequest, - ) -> Result<TradeDecisionReceipt, RadrootsSdkError> { - let context = trade_mutation_context(self.sdk, request.locator, "trade.decline").await?; + ) -> Result<TradeMutationOutcome<TradeDecisionPlan, TradeDecisionReceipt>, RadrootsSdkError> + { + validate_trade_product_publish_policy(request.publish_mode, request.ack_policy)?; + let TradeDeclineRequest { + actor, + locator, + reason, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + created_at, + } = request; + let context = trade_mutation_context(self.sdk, locator, "trade.decline").await?; let decision = RadrootsOrderDecision { order_id: context.order_id.clone(), listing_addr: context.listing_addr.clone(), buyer_pubkey: context.buyer_pubkey.clone(), seller_pubkey: context.seller_pubkey.clone(), - decision: RadrootsOrderDecisionOutcome::Declined { - reason: request.reason, - }, + decision: RadrootsOrderDecisionOutcome::Declined { reason }, }; - trades_client(self.sdk) - .enqueue_decision(TradeDecisionEnqueueRequest { - actor: request.actor, - request_event: event_ptr(&context.root_event_id), - decision, - target_relays: request.target_relays, - publish_mode: request.publish_mode, - ack_policy: request.ack_policy, - idempotency_key: request.idempotency_key, - created_at: request.created_at, - }) - .await + let client = trades_client(self.sdk); + let plan = client.prepare_decision(TradeDecisionPrepareRequest { + actor: actor.clone(), + request_event: event_ptr(&context.root_event_id), + decision, + created_at, + })?; + if publish_mode == PublishMode::DryRun { + return Ok(TradeMutationOutcome::DryRun { plan }); + } + let receipt = client + .enqueue_prepared_decision( + &actor, + plan, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + ) + .await?; + trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await } pub async fn propose_revision( &self, request: TradeRevisionProposalRequest, - ) -> Result<TradeRevisionProposalReceipt, RadrootsSdkError> { - let context = - trade_mutation_context(self.sdk, request.locator, "trade.propose_revision").await?; + ) -> Result< + TradeMutationOutcome<TradeRevisionProposalPlan, TradeRevisionProposalReceipt>, + RadrootsSdkError, + > { + validate_trade_product_publish_policy(request.publish_mode, request.ack_policy)?; + let TradeRevisionProposalRequest { + actor, + locator, + revision_id, + items, + economics, + reason, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + created_at, + } = request; + let context = trade_mutation_context(self.sdk, locator, "trade.propose_revision").await?; let proposal = RadrootsOrderRevisionProposal { - revision_id: request.revision_id, + revision_id, order_id: context.order_id.clone(), listing_addr: context.listing_addr.clone(), buyer_pubkey: context.buyer_pubkey.clone(), seller_pubkey: context.seller_pubkey.clone(), root_event_id: context.root_event_id.clone(), prev_event_id: context.previous_event_id.clone(), - items: request.items, - economics: request.economics, - reason: request.reason, + items, + economics, + reason, }; - trades_client(self.sdk) - .enqueue_revision_proposal(TradeRevisionProposalEnqueueRequest { - actor: request.actor, - root_event: event_ptr(&context.root_event_id), - previous_event: event_ptr(&context.previous_event_id), - proposal, - target_relays: request.target_relays, - publish_mode: request.publish_mode, - ack_policy: request.ack_policy, - idempotency_key: request.idempotency_key, - created_at: request.created_at, - }) - .await + let client = trades_client(self.sdk); + let plan = client.prepare_revision_proposal(TradeRevisionProposalPrepareRequest { + actor: actor.clone(), + root_event: event_ptr(&context.root_event_id), + previous_event: event_ptr(&context.previous_event_id), + proposal, + created_at, + })?; + if publish_mode == PublishMode::DryRun { + return Ok(TradeMutationOutcome::DryRun { plan }); + } + let receipt = client + .enqueue_prepared_revision_proposal( + &actor, + plan, + target_relays, + publish_mode, + ack_policy, + idempotency_key, + ) + .await?; + trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await } } @@ -2801,6 +2955,65 @@ fn event_ptr(event_id: &RadrootsEventId) -> RadrootsNostrEventPtr { } #[cfg(feature = "runtime")] +async fn trade_product_post_enqueue_outcome<Plan, Receipt>( + sdk: &crate::RadrootsClient, + publish_mode: PublishMode, + ack_policy: AckPolicy, + receipt: Receipt, +) -> Result<TradeMutationOutcome<Plan, Receipt>, RadrootsSdkError> { + match publish_mode { + PublishMode::DryRun => Err(RadrootsSdkError::InvalidRequest { + message: "trade product dry-run must return before enqueue".to_owned(), + }), + PublishMode::EnqueueOnly => Ok(TradeMutationOutcome::Enqueued { receipt }), + PublishMode::EnqueueAndPublish => { + let publish = sdk + .sync() + .push_outbox(push_request_for_ack_policy(ack_policy)?) + .await?; + Ok(TradeMutationOutcome::Published { receipt, publish }) + } + } +} + +#[cfg(feature = "runtime")] +fn push_request_for_ack_policy( + ack_policy: AckPolicy, +) -> Result<PushOutboxRequest, RadrootsSdkError> { + let request = PushOutboxRequest::new().with_limit(1); + match ack_policy { + AckPolicy::NoWait => Err(RadrootsSdkError::InvalidRequest { + message: "trade enqueue-and-publish requires a relay acknowledgement policy".to_owned(), + }), + AckPolicy::AtLeastOneRelay => Ok(request.with_accepted_quorum(1)), + AckPolicy::AllRelays => Ok(request), + AckPolicy::Quorum { required } => Ok(request.with_accepted_quorum(usize::from(required))), + } +} + +#[cfg(feature = "runtime")] +fn validate_trade_product_publish_policy( + publish_mode: PublishMode, + ack_policy: AckPolicy, +) -> Result<(), RadrootsSdkError> { + match publish_mode { + PublishMode::DryRun | PublishMode::EnqueueOnly if ack_policy != AckPolicy::NoWait => { + Err(RadrootsSdkError::InvalidRequest { + message: "trade dry-run and enqueue-only modes require no-wait acknowledgement" + .to_owned(), + }) + } + PublishMode::EnqueueAndPublish if ack_policy == AckPolicy::NoWait => { + Err(RadrootsSdkError::InvalidRequest { + message: "trade enqueue-and-publish requires a relay acknowledgement policy" + .to_owned(), + }) + } + _ => Ok(()), + } +} + +#[cfg(feature = "runtime")] fn validate_trade_enqueue_policy( publish_mode: PublishMode, ack_policy: AckPolicy, @@ -2816,10 +3029,9 @@ fn validate_trade_enqueue_policy( .to_owned(), }); } - if publish_mode == PublishMode::EnqueueAndPublish { + if publish_mode == PublishMode::EnqueueAndPublish && ack_policy == AckPolicy::NoWait { return Err(RadrootsSdkError::InvalidRequest { - message: "trade enqueue-and-publish mode requires publish receipt orchestration" - .to_owned(), + message: "trade enqueue-and-publish requires a relay acknowledgement policy".to_owned(), }); } Ok(()) diff --git a/crates/sdk/src/sync_runtime.rs b/crates/sdk/src/sync_runtime.rs @@ -220,6 +220,7 @@ impl Default for SdkRelayAuthPolicy { pub struct PushOutboxRequest { pub limit: usize, pub republish_accepted_relays: bool, + pub accepted_quorum: Option<usize>, pub relay_url_policy: SdkRelayUrlPolicy, pub auth_policy: SdkRelayAuthPolicy, pub claim_ttl_ms: i64, @@ -232,6 +233,7 @@ impl Default for PushOutboxRequest { Self { limit: PUSH_OUTBOX_DEFAULT_LIMIT, republish_accepted_relays: false, + accepted_quorum: None, relay_url_policy: SdkRelayUrlPolicy::Public, auth_policy: SdkRelayAuthPolicy::DetectOnly, claim_ttl_ms: PUSH_OUTBOX_DEFAULT_CLAIM_TTL_MS, @@ -256,6 +258,11 @@ impl PushOutboxRequest { self } + pub fn with_accepted_quorum(mut self, accepted_quorum: usize) -> Self { + self.accepted_quorum = Some(accepted_quorum); + self + } + pub fn with_relay_url_policy(mut self, policy: SdkRelayUrlPolicy) -> Self { self.relay_url_policy = policy; self @@ -297,6 +304,11 @@ impl PushOutboxRequest { message: "push_outbox next attempt delay must be positive".to_owned(), }); } + if self.accepted_quorum == Some(0) { + return Err(RadrootsSdkError::InvalidRequest { + message: "push_outbox accepted quorum must be positive".to_owned(), + }); + } Ok(()) } } @@ -522,6 +534,10 @@ impl<'sdk> SyncClient<'sdk> { ) .republish_accepted_relays(request.republish_accepted_relays) .relay_url_policy(request.relay_url_policy.relay_transport_policy()); + let policy = match request.accepted_quorum { + Some(accepted_quorum) => policy.with_accepted_quorum(accepted_quorum), + None => policy, + }; let publish = publish_claimed_outbox_event( &self.sdk._outbox, &self.sdk._event_store, @@ -569,6 +585,7 @@ impl<'sdk> SyncClient<'sdk> { self, adapter, &claimed, + request.accepted_quorum, request.next_attempt_delay_ms, publish_now_ms, ) @@ -606,6 +623,7 @@ async fn push_proxy_claimed_outbox_event( sync: &SyncClient<'_>, adapter: &RadrootsdProxyPublishAdapter, claimed: &RadrootsOutboxClaimedEvent, + accepted_quorum: Option<usize>, next_attempt_delay_ms: i64, now_ms: i64, ) -> Result<RadrootsRelayPublishReceipt, RadrootsSdkError> { @@ -626,7 +644,7 @@ async fn push_proxy_claimed_outbox_event( let request = RadrootsdProxyPublishRequest { signed_event: signed_event.clone(), relays: claimed.target_relays.clone(), - delivery_policy: proxy_delivery_policy(claimed.target_relays.len()), + delivery_policy: proxy_delivery_policy(claimed.target_relays.len(), accepted_quorum), idempotency_key: Some(proxy_outbox_idempotency_key( claimed.outbox_event_id, claimed.attempt_count, @@ -656,7 +674,13 @@ async fn push_proxy_claimed_outbox_event( } #[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] -fn proxy_delivery_policy(target_count: usize) -> PublishDeliveryPolicy { +fn proxy_delivery_policy( + target_count: usize, + accepted_quorum: Option<usize>, +) -> PublishDeliveryPolicy { + if let Some(quorum) = accepted_quorum { + return PublishDeliveryPolicy::Quorum { quorum }; + } if target_count == 0 { PublishDeliveryPolicy::Any } else { diff --git a/crates/sdk/tests/orders_runtime.rs b/crates/sdk/tests/orders_runtime.rs @@ -45,11 +45,12 @@ use radroots_sdk::{ TRADE_REVISION_PROPOSAL_OPERATION_KIND, TRADE_STATUS_DEFAULT_LIMIT, TRADE_STATUS_MAX_LIMIT, TRADE_SUBMIT_OPERATION_KIND, TradeAcceptRequest, TradeCancellationEnqueueRequest, TradeCancellationPrepareRequest, TradeDecisionEnqueueRequest, TradeDecisionPrepareRequest, - TradeEvidenceIngestRequest, TradeProposeRequest, TradeRequestEvidenceIngestRequest, - TradeResyncRequest, TradeRevisionDecisionEnqueueRequest, TradeRevisionDecisionPrepareRequest, - TradeRevisionProposalEnqueueRequest, TradeRevisionProposalPrepareRequest, - TradeSellerInboxRequest, TradeStatusKind, TradeStatusNextActionKind, TradeStatusRequest, - TradeSubmitEnqueueRequest, TradeSubmitPrepareRequest, TradeWorkflowKind, + TradeEvidenceIngestRequest, TradeMutationOutcome, TradeProposeRequest, + TradeRequestEvidenceIngestRequest, TradeResyncRequest, TradeRevisionDecisionEnqueueRequest, + TradeRevisionDecisionPrepareRequest, TradeRevisionProposalEnqueueRequest, + TradeRevisionProposalPrepareRequest, TradeSellerInboxRequest, TradeStatusKind, + TradeStatusNextActionKind, TradeStatusRequest, TradeSubmitEnqueueRequest, + TradeSubmitPrepareRequest, TradeWorkflowKind, }; #[cfg(all(feature = "signer-adapters", feature = "local-signer"))] use radroots_sdk::{RadrootsSdkLocalKeySigner, RadrootsSdkSignerProvider}; @@ -459,6 +460,14 @@ fn explicit_trade_relays() -> RelayResolutionPolicy { ) } +fn expect_enqueued<Plan, Receipt>(outcome: TradeMutationOutcome<Plan, Receipt>) -> Receipt { + match outcome { + TradeMutationOutcome::Enqueued { receipt } => receipt, + TradeMutationOutcome::DryRun { .. } => panic!("expected enqueue outcome"), + TradeMutationOutcome::Published { .. } => panic!("expected enqueue outcome"), + } +} + fn deterministic_event_id(raw: &str) -> RadrootsEventId { let mut bytes = [0u8; 32]; for (index, byte) in raw.bytes().enumerate() { @@ -821,23 +830,25 @@ async fn trade_product_clients_propose_inbox_accept_status_and_resync() { 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 propose_receipt = buyer_sdk - .trades() - .buyer() - .propose_trade( - TradeProposeRequest::new( - buyer_actor(), - listing_event_ptr(), - order_request("trade-product-facade-flow"), - explicit_trade_relays(), - PublishMode::EnqueueOnly, - AckPolicy::NoWait, + let propose_receipt = expect_enqueued( + buyer_sdk + .trades() + .buyer() + .propose_trade( + TradeProposeRequest::new( + buyer_actor(), + listing_event_ptr(), + order_request("trade-product-facade-flow"), + explicit_trade_relays(), + PublishMode::EnqueueOnly, + AckPolicy::NoWait, + ) + .try_with_idempotency_key("trade-product-facade-propose") + .expect("propose idempotency"), ) - .try_with_idempotency_key("trade-product-facade-propose") - .expect("propose idempotency"), - ) - .await - .expect("propose trade"); + .await + .expect("propose trade"), + ); assert_eq!( propose_receipt.order_id.as_str(), @@ -882,26 +893,28 @@ async fn trade_product_clients_propose_inbox_accept_status_and_resync() { Some(propose_receipt.signed_event_id.as_str()) ); - let accept_receipt = seller_sdk - .trades() - .seller() - .accept_trade( - TradeAcceptRequest::new( - seller_actor(), - propose_receipt.locator.clone(), - vec![RadrootsOrderInventoryCommitment { - bin_id: "bin-1".parse().expect("bin id"), - bin_count: 2, - }], - explicit_trade_relays(), - PublishMode::EnqueueOnly, - AckPolicy::NoWait, + let accept_receipt = expect_enqueued( + seller_sdk + .trades() + .seller() + .accept_trade( + TradeAcceptRequest::new( + seller_actor(), + propose_receipt.locator.clone(), + vec![RadrootsOrderInventoryCommitment { + bin_id: "bin-1".parse().expect("bin id"), + bin_count: 2, + }], + explicit_trade_relays(), + PublishMode::EnqueueOnly, + AckPolicy::NoWait, + ) + .try_with_idempotency_key("trade-product-facade-accept") + .expect("accept idempotency"), ) - .try_with_idempotency_key("trade-product-facade-accept") - .expect("accept idempotency"), - ) - .await - .expect("accept trade"); + .await + .expect("accept trade"), + ); assert_eq!(accept_receipt.order_id, propose_receipt.order_id); assert_eq!(accept_receipt.locator, propose_receipt.locator); @@ -937,6 +950,52 @@ async fn trade_product_clients_propose_inbox_accept_status_and_resync() { ); } +#[cfg(feature = "signer-adapters")] +#[tokio::test] +async fn trade_product_propose_dry_run_returns_plan_without_local_side_effects() { + let (_tempdir, sdk, store) = directory_sdk_and_store().await; + let outcome = sdk + .trades() + .buyer() + .propose_trade(TradeProposeRequest::new( + buyer_actor(), + listing_event_ptr(), + order_request("trade-product-dry-run"), + explicit_trade_relays(), + PublishMode::DryRun, + AckPolicy::NoWait, + )) + .await + .expect("dry-run proposal"); + let plan = match outcome { + TradeMutationOutcome::DryRun { plan } => plan, + TradeMutationOutcome::Enqueued { .. } => panic!("expected dry-run outcome"), + TradeMutationOutcome::Published { .. } => panic!("expected dry-run outcome"), + }; + + assert_eq!(plan.order_id.as_str(), "trade-product-dry-run"); + assert_eq!(plan.frozen_draft.kind, KIND_ORDER_REQUEST); + assert_eq!(plan.expected_event_id, plan.workflow.expected_event_id); + assert_eq!( + store + .status_summary() + .await + .expect("event store status") + .total_events, + 0 + ); + let outbox = RadrootsOutbox::open_file(&sdk.storage_paths().expect("paths").outbox_path) + .await + .expect("outbox"); + assert!( + outbox + .claim_next_ready_event("worker", "claim", 2_000, 1_700_000_000_000) + .await + .expect("claim") + .is_none() + ); +} + #[tokio::test] async fn order_submit_enqueue_returns_sanitized_signer_errors_before_mutation() { let (_tempdir, sdk, store) = directory_sdk_and_store().await; diff --git a/crates/sdk/tests/source_boundary.rs b/crates/sdk/tests/source_boundary.rs @@ -60,6 +60,7 @@ const REQUIRED_TRADE_RUNTIME_EXPORTS: &[&str] = &[ "TradeDeclineRequest", "TradeEvidenceIngestReceipt", "TradeEvidenceIngestRequest", + "TradeMutationOutcome", "TradeProposeRequest", "TradeRequestEvidenceIngestReceipt", "TradeRequestEvidenceIngestRequest", diff --git a/crates/sdk/tests/sync_runtime.rs b/crates/sdk/tests/sync_runtime.rs @@ -1700,6 +1700,7 @@ fn push_outbox_contract_dtos_serialize_deterministically() { let request = PushOutboxRequest::new() .with_limit(2) .republish_accepted_relays(true) + .with_accepted_quorum(1) .with_relay_url_policy(SdkRelayUrlPolicy::Localhost) .with_auth_policy(SdkRelayAuthPolicy::DetectOnly) .with_claim_ttl_ms(1_000) @@ -1709,6 +1710,7 @@ fn push_outbox_contract_dtos_serialize_deterministically() { serde_json::json!({ "limit": 2, "republish_accepted_relays": true, + "accepted_quorum": 1, "relay_url_policy": "localhost", "auth_policy": "detect_only", "claim_ttl_ms": 1000, diff --git a/crates/sdk/tests/trade_product_publish_runtime.rs b/crates/sdk/tests/trade_product_publish_runtime.rs @@ -0,0 +1,269 @@ +#![cfg(all( + feature = "runtime", + feature = "signer-adapters", + feature = "local-signer", + feature = "radrootsd-proxy" +))] + +use radroots_authority::RadrootsActorContext; +use radroots_core::{ + RadrootsCoreCurrency, RadrootsCoreDecimal, RadrootsCoreMoney, RadrootsCoreUnit, +}; +use radroots_events::{contract::RadrootsActorRole, kinds::KIND_LISTING}; +use radroots_nostr::prelude::{RadrootsNostrKeys, RadrootsNostrSecretKey}; +use radroots_sdk::protocol::events::RadrootsNostrEventPtr; +use radroots_sdk::protocol::order::{ + RadrootsListingAddress, RadrootsOrderEconomicItem, RadrootsOrderEconomicLine, + RadrootsOrderEconomics, RadrootsOrderItem, RadrootsOrderPricingBasis, RadrootsOrderRequest, +}; +use radroots_sdk::{ + AckPolicy, PublishMode, PushOutboxRelayOutcomeKind, RadrootsClient, RadrootsSdkLocalKeySigner, + RadrootsSdkSignerProvider, RadrootsSdkTimestamp, RelayResolutionPolicy, SdkPublishTransport, + SdkRelayTargetSet, SdkRelayUrlPolicy, TradeMutationOutcome, TradeProposeRequest, + adapters::radrootsd::RadrootsdProxyConfig, +}; +use std::{ + io::{Read, Write}, + net::{TcpListener, TcpStream}, + thread::JoinHandle, +}; + +const BUYER_SECRET_KEY_HEX: &str = + "10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5"; +const BUYER_PUBLIC_KEY_HEX: &str = + "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; +const SELLER_PUBLIC_KEY_HEX: &str = + "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af"; +const RELAY: &str = "wss://relay.radroots.test"; + +struct RecordedProxyRequest { + body: String, +} + +fn spawn_trade_publish_proxy_server() -> (String, JoinHandle<RecordedProxyRequest>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind proxy server"); + let endpoint = format!("http://{}/rpc", listener.local_addr().expect("addr")); + let handle = std::thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("accept"); + let body = read_proxy_request_body(&mut stream); + write_proxy_accept_response(&mut stream, body.as_str()); + RecordedProxyRequest { body } + }); + (endpoint, handle) +} + +fn read_proxy_request_body(stream: &mut TcpStream) -> String { + let mut request = Vec::new(); + let mut buffer = [0u8; 1024]; + loop { + let read = stream.read(&mut buffer).expect("read request"); + if read == 0 { + break; + } + request.extend_from_slice(&buffer[..read]); + if request.windows(4).any(|window| window == b"\r\n\r\n") { + let headers_end = request + .windows(4) + .position(|window| window == b"\r\n\r\n") + .expect("headers end") + + 4; + let header_text = String::from_utf8_lossy(&request[..headers_end]); + let content_length = header_text + .lines() + .find_map(|line| { + let (name, value) = line.split_once(':')?; + name.eq_ignore_ascii_case("content-length") + .then(|| value.trim().parse::<usize>().expect("content length")) + }) + .unwrap_or(0); + while request.len() < headers_end + content_length { + let read = stream.read(&mut buffer).expect("read body"); + if read == 0 { + break; + } + request.extend_from_slice(&buffer[..read]); + } + break; + } + } + let request_text = String::from_utf8_lossy(&request); + let (_, body) = request_text.split_once("\r\n\r\n").expect("request body"); + body.to_owned() +} + +fn write_proxy_accept_response(stream: &mut TcpStream, body: &str) { + let body_json: serde_json::Value = serde_json::from_str(body).expect("body json"); + let event = &body_json["params"]["event"]; + let response_body = serde_json::json!({ + "jsonrpc": "2.0", + "id": body_json["id"], + "result": { + "deduplicated": false, + "job": { + "job_id": "trade-product-publish-job", + "status": "delivery_satisfied", + "terminal": true, + "delivery_satisfied": true, + "event_id": event["id"], + "pubkey": event["pubkey"], + "event_kind": event["kind"], + "relay_policy": body_json["params"]["relay_policy"], + "delivery_policy": body_json["params"]["delivery_policy"], + "relay_count": 1, + "acknowledged_count": 1, + "retryable_count": 0, + "terminal_count": 0, + "requested_at_ms": 1700000000000i64, + "completed_at_ms": 1700000000100i64, + "relays": [{ + "relay_url": RELAY, + "source": "request", + "attempted": true, + "outcome_kind": "accepted", + "message": "accepted" + }] + } + } + }) + .to_string(); + let response = format!( + "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}", + response_body.len(), + response_body + ); + stream + .write_all(response.as_bytes()) + .expect("write response"); +} + +fn buyer_actor() -> RadrootsActorContext { + RadrootsActorContext::test(BUYER_PUBLIC_KEY_HEX, [RadrootsActorRole::Buyer]).expect("actor") +} + +fn listing_address() -> RadrootsListingAddress { + RadrootsListingAddress::parse(format!( + "{KIND_LISTING}:{SELLER_PUBLIC_KEY_HEX}:AAAAAAAAAAAAAAAAAAAAAg" + )) + .expect("listing address") +} + +fn listing_event_ptr() -> RadrootsNostrEventPtr { + RadrootsNostrEventPtr { + id: "6ccf12d1e56c21065d239bc3d46c0000cd000095d20000d9000073cd00009600".to_owned(), + relays: Some(RELAY.to_owned()), + } +} + +fn decimal(raw: &str) -> RadrootsCoreDecimal { + raw.parse().expect("decimal") +} + +fn usd(raw: &str) -> RadrootsCoreMoney { + RadrootsCoreMoney::new(decimal(raw), RadrootsCoreCurrency::USD) +} + +fn economics() -> RadrootsOrderEconomics { + RadrootsOrderEconomics { + quote_id: "quote-1".parse().expect("quote id"), + quote_version: 1, + pricing_basis: RadrootsOrderPricingBasis::ListingEvent, + currency: RadrootsCoreCurrency::USD, + items: vec![RadrootsOrderEconomicItem { + bin_id: "bin-1".parse().expect("bin id"), + bin_count: 2, + quantity_amount: decimal("1"), + quantity_unit: RadrootsCoreUnit::Each, + unit_price_amount: decimal("5"), + unit_price_currency: RadrootsCoreCurrency::USD, + line_subtotal: usd("10"), + }], + discounts: Vec::<RadrootsOrderEconomicLine>::new(), + adjustments: Vec::<RadrootsOrderEconomicLine>::new(), + subtotal: usd("10"), + discount_total: usd("0"), + adjustment_total: usd("0"), + total: usd("10"), + } +} + +fn order_request(raw_order_id: &str) -> RadrootsOrderRequest { + RadrootsOrderRequest { + order_id: raw_order_id.parse().expect("order id"), + listing_addr: listing_address(), + buyer_pubkey: BUYER_PUBLIC_KEY_HEX.parse().expect("buyer pubkey"), + seller_pubkey: SELLER_PUBLIC_KEY_HEX.parse().expect("seller pubkey"), + items: vec![RadrootsOrderItem { + bin_id: "bin-1".parse().expect("bin id"), + bin_count: 2, + }], + economics: economics(), + } +} + +fn explicit_trade_relays() -> RelayResolutionPolicy { + RelayResolutionPolicy::explicit( + SdkRelayTargetSet::new([RELAY], SdkRelayUrlPolicy::Public).expect("target relays"), + ) +} + +#[tokio::test] +async fn trade_product_propose_enqueue_and_publish_uses_ack_policy() { + let (endpoint, handle) = spawn_trade_publish_proxy_server(); + let tempdir = tempfile::tempdir().expect("tempdir"); + let secret_key = RadrootsNostrSecretKey::from_hex(BUYER_SECRET_KEY_HEX).expect("secret key"); + let signer_keys = RadrootsNostrKeys::new(secret_key); + let sdk = RadrootsClient::builder() + .directory_storage(tempdir.path().join("sdk")) + .fixed_clock(RadrootsSdkTimestamp::from_unix_seconds(1_700_000_000)) + .signer_provider(RadrootsSdkSignerProvider::LocalKey( + RadrootsSdkLocalKeySigner::new(signer_keys).expect("local signer"), + )) + .publish_transport(SdkPublishTransport::RadrootsdProxy( + RadrootsdProxyConfig::new(endpoint), + )) + .build() + .await + .expect("sdk"); + + let outcome = sdk + .trades() + .buyer() + .propose_trade( + TradeProposeRequest::new( + buyer_actor(), + listing_event_ptr(), + order_request("trade-product-publish"), + explicit_trade_relays(), + PublishMode::EnqueueAndPublish, + AckPolicy::AtLeastOneRelay, + ) + .try_with_idempotency_key("trade-product-publish") + .expect("idempotency"), + ) + .await + .expect("publish proposal"); + let (receipt, publish) = match outcome { + TradeMutationOutcome::Published { receipt, publish } => (receipt, publish), + TradeMutationOutcome::DryRun { .. } => panic!("expected published outcome"), + TradeMutationOutcome::Enqueued { .. } => panic!("expected published outcome"), + }; + + assert_eq!(receipt.order_id.as_str(), "trade-product-publish"); + assert_eq!(publish.attempted_events, 1); + assert_eq!(publish.published_events, 1); + assert_eq!(publish.events.len(), 1); + assert_eq!(publish.events[0].outbox_event_id, receipt.outbox_event_id); + assert_eq!(publish.events[0].quorum, 1); + assert!(publish.events[0].quorum_met); + assert_eq!( + publish.events[0].relays[0].outcome_kind, + PushOutboxRelayOutcomeKind::Accepted + ); + + let recorded = handle.join().expect("proxy request"); + let body: serde_json::Value = serde_json::from_str(recorded.body.as_str()).expect("body"); + assert_eq!(body["method"], "publish.event"); + assert_eq!(body["params"]["delivery_policy"]["mode"], "quorum"); + assert_eq!(body["params"]["delivery_policy"]["quorum"], 1); + assert_eq!(body["params"]["relays"], serde_json::json!([RELAY])); +} diff --git a/crates/sdk/tests/unit/orders_runtime_tests.rs b/crates/sdk/tests/unit/orders_runtime_tests.rs @@ -2220,11 +2220,15 @@ fn trade_enqueue_policy_rejects_publish_modes_without_matching_side_effects() { if message == "trade enqueue-only publish mode only supports no-wait acknowledgement" )); assert!(matches!( - validate_trade_enqueue_policy(PublishMode::EnqueueAndPublish, AckPolicy::AtLeastOneRelay), + validate_trade_enqueue_policy(PublishMode::EnqueueAndPublish, AckPolicy::NoWait), Err(RadrootsSdkError::InvalidRequest { ref message }) - if message == "trade enqueue-and-publish mode requires publish receipt orchestration" + if message == "trade enqueue-and-publish requires a relay acknowledgement policy" )); assert!(validate_trade_enqueue_policy(PublishMode::EnqueueOnly, AckPolicy::NoWait).is_ok()); + assert!( + validate_trade_enqueue_policy(PublishMode::EnqueueAndPublish, AckPolicy::AtLeastOneRelay) + .is_ok() + ); } #[tokio::test] diff --git a/crates/sdk/tests/unit/sync_runtime_tests.rs b/crates/sdk/tests/unit/sync_runtime_tests.rs @@ -350,12 +350,14 @@ fn push_outbox_request_builders_validate_all_bounds() { let request = super::PushOutboxRequest::new() .with_limit(2) .republish_accepted_relays(true) + .with_accepted_quorum(2) .with_relay_url_policy(crate::SdkRelayUrlPolicy::Localhost) .with_auth_policy(SdkRelayAuthPolicy::DetectOnly) .with_claim_ttl_ms(7) .with_next_attempt_delay_ms(11); assert_eq!(request.limit, 2); assert!(request.republish_accepted_relays); + assert_eq!(request.accepted_quorum, Some(2)); request.validate().expect("valid request"); assert!(matches!( @@ -370,6 +372,12 @@ fn push_outbox_request_builders_validate_all_bounds() { )); assert!(matches!( super::PushOutboxRequest::new() + .with_accepted_quorum(0) + .validate(), + Err(RadrootsSdkError::InvalidRequest { message }) if message.contains("accepted quorum") + )); + assert!(matches!( + super::PushOutboxRequest::new() .with_claim_ttl_ms(0) .validate(), Err(RadrootsSdkError::InvalidRequest { message }) if message.contains("TTL") @@ -490,8 +498,12 @@ async fn proxy_push_empty_queue_and_private_helpers_are_deterministic() { .expect("empty proxy push"); assert_eq!(receipt.attempted_events, 0); - assert_eq!(proxy_delivery_policy(0), PublishDeliveryPolicy::Any); - assert_eq!(proxy_delivery_policy(2), PublishDeliveryPolicy::All); + assert_eq!(proxy_delivery_policy(0, None), PublishDeliveryPolicy::Any); + assert_eq!(proxy_delivery_policy(2, None), PublishDeliveryPolicy::All); + assert_eq!( + proxy_delivery_policy(3, Some(2)), + PublishDeliveryPolicy::Quorum { quorum: 2 } + ); assert_eq!( proxy_outbox_idempotency_key(7, 3, "event-id"), "radroots-sdk-outbox-7-3-event-id" @@ -578,7 +590,7 @@ async fn proxy_push_reports_missing_signed_claim_before_daemon_publish() { }; assert!(matches!( - push_proxy_claimed_outbox_event(&sync, &adapter, &claimed, 60_000, 1_700_000_000_000) + push_proxy_claimed_outbox_event(&sync, &adapter, &claimed, None, 60_000, 1_700_000_000_000) .await, Err(RadrootsSdkError::RelayTransport { message }) if message.contains("Outbox claim 41 does not contain a signed event") @@ -593,7 +605,7 @@ async fn proxy_claim_publish_marks_retryable_transport_errors() { let adapter = RadrootsdProxyPublishAdapter::new(RadrootsdProxyConfig::new("http://127.0.0.1:9/rpc")); let receipt = - push_proxy_claimed_outbox_event(&sync, &adapter, &claimed, 60_000, 1_700_000_000_000) + push_proxy_claimed_outbox_event(&sync, &adapter, &claimed, None, 60_000, 1_700_000_000_000) .await .expect("transport error receipt");