lib

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

commit b28d22eaf1be1bedbc417a305c16497162c6ea8c
parent f54abd7f296006006d855ab37a20eba56e2c78d6
Author: triesap <tyson@radroots.org>
Date:   Tue, 30 Jun 2026 03:41:06 +0000

sdk: publish targeted trade outbox events

Diffstat:
Mcrates/sdk/src/orders_runtime.rs | 60++++++++++++++++++++++++++++++++++++++++++++++++++++--------
Mcrates/sdk/src/sync_runtime.rs | 80+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------
Mcrates/sdk/tests/sync_runtime.rs | 44+++++++++++++++++++++++++++++++++++++++++++-
Mcrates/sdk/tests/unit/sync_runtime_tests.rs | 10+++++++++-
4 files changed, 164 insertions(+), 30 deletions(-)

diff --git a/crates/sdk/src/orders_runtime.rs b/crates/sdk/src/orders_runtime.rs @@ -2495,7 +2495,14 @@ impl<'sdk> TradeBuyerClient<'sdk> { idempotency_key, ) .await?; - trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await + trade_product_post_enqueue_outcome( + self.sdk, + publish_mode, + ack_policy, + receipt.outbox_event_id, + receipt, + ) + .await } pub async fn cancel_trade( @@ -2545,7 +2552,14 @@ impl<'sdk> TradeBuyerClient<'sdk> { idempotency_key, ) .await?; - trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await + trade_product_post_enqueue_outcome( + self.sdk, + publish_mode, + ack_policy, + receipt.outbox_event_id, + receipt, + ) + .await } pub async fn accept_revision( @@ -2625,7 +2639,14 @@ impl<'sdk> TradeBuyerClient<'sdk> { idempotency_key, ) .await?; - trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await + trade_product_post_enqueue_outcome( + self.sdk, + publish_mode, + ack_policy, + receipt.outbox_event_id, + receipt, + ) + .await } } @@ -2749,7 +2770,14 @@ impl<'sdk> TradeSellerClient<'sdk> { idempotency_key, ) .await?; - trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await + trade_product_post_enqueue_outcome( + self.sdk, + publish_mode, + ack_policy, + receipt.outbox_event_id, + receipt, + ) + .await } pub async fn decline_trade( @@ -2796,7 +2824,14 @@ impl<'sdk> TradeSellerClient<'sdk> { idempotency_key, ) .await?; - trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await + trade_product_post_enqueue_outcome( + self.sdk, + publish_mode, + ack_policy, + receipt.outbox_event_id, + receipt, + ) + .await } pub async fn propose_revision( @@ -2854,7 +2889,14 @@ impl<'sdk> TradeSellerClient<'sdk> { idempotency_key, ) .await?; - trade_product_post_enqueue_outcome(self.sdk, publish_mode, ack_policy, receipt).await + trade_product_post_enqueue_outcome( + self.sdk, + publish_mode, + ack_policy, + receipt.outbox_event_id, + receipt, + ) + .await } } @@ -2959,6 +3001,7 @@ async fn trade_product_post_enqueue_outcome<Plan, Receipt>( sdk: &crate::RadrootsClient, publish_mode: PublishMode, ack_policy: AckPolicy, + outbox_event_id: i64, receipt: Receipt, ) -> Result<TradeMutationOutcome<Plan, Receipt>, RadrootsSdkError> { match publish_mode { @@ -2969,7 +3012,7 @@ async fn trade_product_post_enqueue_outcome<Plan, Receipt>( PublishMode::EnqueueAndPublish => { let publish = sdk .sync() - .push_outbox(push_request_for_ack_policy(ack_policy)?) + .push_outbox(push_request_for_ack_policy(ack_policy, outbox_event_id)?) .await?; Ok(TradeMutationOutcome::Published { receipt, publish }) } @@ -2979,8 +3022,9 @@ async fn trade_product_post_enqueue_outcome<Plan, Receipt>( #[cfg(feature = "runtime")] fn push_request_for_ack_policy( ack_policy: AckPolicy, + outbox_event_id: i64, ) -> Result<PushOutboxRequest, RadrootsSdkError> { - let request = PushOutboxRequest::new().with_limit(1); + let request = PushOutboxRequest::new().with_outbox_event_id(outbox_event_id); match ack_policy { AckPolicy::NoWait => Err(RadrootsSdkError::InvalidRequest { message: "trade enqueue-and-publish requires a relay acknowledgement policy".to_owned(), diff --git a/crates/sdk/src/sync_runtime.rs b/crates/sdk/src/sync_runtime.rs @@ -219,6 +219,7 @@ impl Default for SdkRelayAuthPolicy { #[non_exhaustive] pub struct PushOutboxRequest { pub limit: usize, + pub outbox_event_id: Option<i64>, pub republish_accepted_relays: bool, pub accepted_quorum: Option<usize>, pub relay_url_policy: SdkRelayUrlPolicy, @@ -232,6 +233,7 @@ impl Default for PushOutboxRequest { fn default() -> Self { Self { limit: PUSH_OUTBOX_DEFAULT_LIMIT, + outbox_event_id: None, republish_accepted_relays: false, accepted_quorum: None, relay_url_policy: SdkRelayUrlPolicy::Public, @@ -253,6 +255,12 @@ impl PushOutboxRequest { self } + pub fn with_outbox_event_id(mut self, outbox_event_id: i64) -> Self { + self.outbox_event_id = Some(outbox_event_id); + self.limit = 1; + self + } + pub fn republish_accepted_relays(mut self, enabled: bool) -> Self { self.republish_accepted_relays = enabled; self @@ -294,6 +302,13 @@ impl PushOutboxRequest { message: format!("push_outbox limit must be between 1 and {PUSH_OUTBOX_MAX_LIMIT}"), }); } + if let Some(outbox_event_id) = self.outbox_event_id + && outbox_event_id <= 0 + { + return Err(RadrootsSdkError::InvalidRequest { + message: "push_outbox outbox event id must be positive".to_owned(), + }); + } if self.claim_ttl_ms <= 0 { return Err(RadrootsSdkError::InvalidRequest { message: "push_outbox claim TTL must be positive".to_owned(), @@ -515,16 +530,13 @@ impl<'sdk> SyncClient<'sdk> { for _ in 0..request.limit { let claim_now_ms = sdk_now_ms(self.sdk)?; let claim_token = push_outbox_claim_token(); - let Some(claimed) = self - .sdk - ._outbox - .claim_next_ready_signed_event( - CLAIM_OWNER, - claim_token.as_str(), - claim_now_ms.saturating_add(request.claim_ttl_ms), - claim_now_ms, - ) - .await? + let Some(claimed) = claim_ready_signed_event_for_push( + self.sdk, + &request, + claim_token.as_str(), + claim_now_ms, + ) + .await? else { break; }; @@ -567,16 +579,13 @@ impl<'sdk> SyncClient<'sdk> { for _ in 0..request.limit { let claim_now_ms = sdk_now_ms(self.sdk)?; let claim_token = push_outbox_claim_token(); - let Some(claimed) = self - .sdk - ._outbox - .claim_next_ready_signed_event( - CLAIM_OWNER, - claim_token.as_str(), - claim_now_ms.saturating_add(request.claim_ttl_ms), - claim_now_ms, - ) - .await? + let Some(claimed) = claim_ready_signed_event_for_push( + self.sdk, + &request, + claim_token.as_str(), + claim_now_ms, + ) + .await? else { break; }; @@ -601,6 +610,37 @@ impl<'sdk> SyncClient<'sdk> { } #[cfg(feature = "runtime")] +async fn claim_ready_signed_event_for_push( + sdk: &RadrootsClient, + request: &PushOutboxRequest, + claim_token: &str, + claim_now_ms: i64, +) -> Result<Option<radroots_outbox::RadrootsOutboxClaimedEvent>, RadrootsSdkError> { + let claim_expires_at_ms = claim_now_ms.saturating_add(request.claim_ttl_ms); + match request.outbox_event_id { + Some(outbox_event_id) => Ok(sdk + ._outbox + .claim_ready_signed_event( + outbox_event_id, + CLAIM_OWNER, + claim_token, + claim_expires_at_ms, + claim_now_ms, + ) + .await?), + None => Ok(sdk + ._outbox + .claim_next_ready_signed_event( + CLAIM_OWNER, + claim_token, + claim_expires_at_ms, + claim_now_ms, + ) + .await?), + } +} + +#[cfg(feature = "runtime")] pub(crate) async fn refresh_product_projections_for_sdk( sdk: &RadrootsClient, request: SyncProjectionRefreshRequest, diff --git a/crates/sdk/tests/sync_runtime.rs b/crates/sdk/tests/sync_runtime.rs @@ -1699,6 +1699,7 @@ async fn product_push_outbox_radrootsd_proxy_error_and_terminal_paths_update_out fn push_outbox_contract_dtos_serialize_deterministically() { let request = PushOutboxRequest::new() .with_limit(2) + .with_outbox_event_id(7) .republish_accepted_relays(true) .with_accepted_quorum(1) .with_relay_url_policy(SdkRelayUrlPolicy::Localhost) @@ -1708,7 +1709,8 @@ fn push_outbox_contract_dtos_serialize_deterministically() { assert_eq!( serde_json::to_value(&request).expect("request json"), serde_json::json!({ - "limit": 2, + "limit": 1, + "outbox_event_id": 7, "republish_accepted_relays": true, "accepted_quorum": 1, "relay_url_policy": "localhost", @@ -1890,6 +1892,46 @@ async fn push_outbox_with_adapter_uses_queued_targets_without_builder_relays() { } #[tokio::test] +async fn push_outbox_with_adapter_can_publish_targeted_ready_event() { + let (_tempdir, sdk) = directory_sdk(&[]).await; + let older_outbox_event_id = + enqueue_listing(&sdk, LISTING_A_D_TAG, "Earlier Coffee", &[RELAY_A]).await; + let targeted_outbox_event_id = + enqueue_listing(&sdk, LISTING_B_D_TAG, "Target Coffee", &[RELAY_B]).await; + let adapter = RadrootsMockRelayPublishAdapter::new(); + + let receipt = sdk + .sync() + .push_outbox_with_adapter( + &adapter, + PushOutboxRequest::new().with_outbox_event_id(targeted_outbox_event_id), + ) + .await + .expect("targeted push"); + + assert_eq!(receipt.attempted_events, 1); + assert_eq!(receipt.published_events, 1); + assert_eq!(receipt.events[0].outbox_event_id, targeted_outbox_event_id); + assert_eq!(adapter.captured_raw_events().len(), 1); + + let outbox = RadrootsOutbox::open_file(&sdk.storage_paths().expect("paths").outbox_path) + .await + .expect("outbox"); + let older = outbox + .get_event(older_outbox_event_id) + .await + .expect("older event") + .expect("older event"); + let targeted = outbox + .get_event(targeted_outbox_event_id) + .await + .expect("targeted event") + .expect("targeted event"); + assert_eq!(older.state, RadrootsOutboxEventState::Signed); + assert_eq!(targeted.state, RadrootsOutboxEventState::Published); +} + +#[tokio::test] async fn push_outbox_default_public_policy_rejects_queued_localhost_ws_targets() { let (_tempdir, sdk) = directory_sdk(&[]).await; enqueue_listing_with_policy( diff --git a/crates/sdk/tests/unit/sync_runtime_tests.rs b/crates/sdk/tests/unit/sync_runtime_tests.rs @@ -349,13 +349,15 @@ fn sync_status_summary_conversions_preserve_all_fields() { fn push_outbox_request_builders_validate_all_bounds() { let request = super::PushOutboxRequest::new() .with_limit(2) + .with_outbox_event_id(9) .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_eq!(request.limit, 1); + assert_eq!(request.outbox_event_id, Some(9)); assert!(request.republish_accepted_relays); assert_eq!(request.accepted_quorum, Some(2)); request.validate().expect("valid request"); @@ -378,6 +380,12 @@ fn push_outbox_request_builders_validate_all_bounds() { )); assert!(matches!( super::PushOutboxRequest::new() + .with_outbox_event_id(0) + .validate(), + Err(RadrootsSdkError::InvalidRequest { message }) if message.contains("outbox event id") + )); + assert!(matches!( + super::PushOutboxRequest::new() .with_claim_ttl_ms(0) .validate(), Err(RadrootsSdkError::InvalidRequest { message }) if message.contains("TTL")