commit 29a7d1f5b96eabc36dbab20eacc3b2947a5f22ee
parent 6060e2303856227a0f78207e95926f09b99b62f9
Author: triesap <tyson@radroots.org>
Date: Tue, 30 Jun 2026 03:41:06 +0000
sdk: publish targeted trade outbox events
Diffstat:
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")