sdk

Radroots SDK and bindings
git clone https://radroots.dev/git/sdk.git
Log | Files | Refs | README

commit 52be745797061ad23fe82f77470dcd58cbb5bcf0
parent bc9de9e176233b58a46261de2be9b66b1109fda7
Author: triesap <tyson@radroots.org>
Date:   Sun,  5 Jul 2026 04:33:09 +0000

status: add bounded trade status watches

- add typed trade status watch requests, updates, and cancellation receipts
- run watch producers through bounded Tokio channels with explicit cancellation
- expose watch constants and API from the SDK runtime surface
- cover finite refreshes, slow consumers, producer errors, and capacity rejection

Diffstat:
Mcrates/sdk/Cargo.toml | 3++-
Mcrates/sdk/src/lib.rs | 14+++++++++-----
Mcrates/sdk/src/orders_runtime.rs | 255++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/sdk/tests/orders_runtime.rs | 145+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
4 files changed, 403 insertions(+), 14 deletions(-)

diff --git a/crates/sdk/Cargo.toml b/crates/sdk/Cargo.toml @@ -55,6 +55,7 @@ signer-adapters = [ runtime = [ "std", "serde_json", + "dep:tokio", "dep:hex", "dep:radroots_authority", "dep:radroots_event_store", @@ -138,7 +139,7 @@ sqlx = { workspace = true, optional = true, default-features = false, features = "runtime-tokio", "sqlite", ] } -tokio = { workspace = true, optional = true, features = ["time"] } +tokio = { workspace = true, optional = true, features = ["macros", "rt", "sync", "time"] } uuid = { workspace = true, optional = true } [[example]] diff --git a/crates/sdk/src/lib.rs b/crates/sdk/src/lib.rs @@ -104,16 +104,20 @@ pub use crate::market_runtime::{ #[cfg(feature = "runtime")] pub use crate::orders_runtime::{ SdkTradeStatusIssue, SdkTradeStatusIssueKind, SdkTradeStatusSource, TRADE_STATUS_DEFAULT_LIMIT, - TRADE_STATUS_MAX_LIMIT, TRADE_STATUS_ROOT_SELECTOR_SEPARATOR, TradeEvidenceBranchReceipt, - TradeEvidenceIngestReceipt, TradeEvidenceIngestRequest, TradeEvidenceQueryBranch, - TradeEvidenceQueryBranchKind, TradeEvidenceQueryPlan, TradeEvidenceRelayFilter, - TradeEvidenceRelayTagFilter, TradeRequestEvidenceIngestReceipt, + TRADE_STATUS_MAX_LIMIT, TRADE_STATUS_ROOT_SELECTOR_SEPARATOR, + TRADE_STATUS_WATCH_DEFAULT_CAPACITY, TRADE_STATUS_WATCH_DEFAULT_REFRESH_INTERVAL_MS, + TRADE_STATUS_WATCH_MAX_CAPACITY, TRADE_STATUS_WATCH_MAX_REFRESH_INTERVAL_MS, + TradeEvidenceBranchReceipt, TradeEvidenceIngestReceipt, TradeEvidenceIngestRequest, + TradeEvidenceQueryBranch, TradeEvidenceQueryBranchKind, TradeEvidenceQueryPlan, + TradeEvidenceRelayFilter, TradeEvidenceRelayTagFilter, TradeRequestEvidenceIngestReceipt, TradeRequestEvidenceIngestRequest, TradeResyncEventImportReceipt, TradeResyncEvidenceReceipt, TradeResyncReceipt, TradeResyncRelayOutcomeKind, TradeResyncRelayOutcomeReceipt, TradeResyncRelayTransportOutcomeKind, TradeResyncRequest, TradeSellerInboxReceipt, TradeSellerInboxRequest, TradeStatusAmbiguityCandidate, TradeStatusEligibility, TradeStatusEvidenceSummary, TradeStatusKind, TradeStatusNextActionKind, TradeStatusReceipt, - TradeStatusRequest, TradeValidationReceiptEvent, TradeValidationReceiptInspectReceipt, + TradeStatusRequest, TradeStatusWatch, TradeStatusWatchCancelReceipt, + TradeStatusWatchCancelState, TradeStatusWatchRequest, TradeStatusWatchUpdate, + TradeValidationReceiptEvent, TradeValidationReceiptInspectReceipt, TradeValidationReceiptInspectRequest, TradeValidationReceiptInvalidCandidate, TradeValidationReceiptListReceipt, TradeValidationReceiptListRequest, TradeValidationReceiptRelayEvidenceReceipt, TradeValidationReceiptRelayOutcomeKind, diff --git a/crates/sdk/src/orders_runtime.rs b/crates/sdk/src/orders_runtime.rs @@ -109,13 +109,29 @@ use serde::Deserialize; #[cfg(feature = "runtime")] use serde::ser::SerializeStruct; #[cfg(feature = "runtime")] -use std::collections::{BTreeMap, BTreeSet}; +use std::{ + collections::{BTreeMap, BTreeSet}, + time::Duration, +}; +#[cfg(feature = "runtime")] +use tokio::{ + sync::{mpsc, oneshot}, + task::JoinHandle, +}; #[cfg(feature = "runtime")] pub const TRADE_STATUS_DEFAULT_LIMIT: u32 = 500; #[cfg(feature = "runtime")] pub const TRADE_STATUS_MAX_LIMIT: u32 = 1_000; #[cfg(feature = "runtime")] pub const TRADE_STATUS_ROOT_SELECTOR_SEPARATOR: char = '@'; +#[cfg(feature = "runtime")] +pub const TRADE_STATUS_WATCH_DEFAULT_CAPACITY: usize = 8; +#[cfg(feature = "runtime")] +pub const TRADE_STATUS_WATCH_MAX_CAPACITY: usize = 128; +#[cfg(feature = "runtime")] +pub const TRADE_STATUS_WATCH_DEFAULT_REFRESH_INTERVAL_MS: u64 = 1_000; +#[cfg(feature = "runtime")] +pub const TRADE_STATUS_WATCH_MAX_REFRESH_INTERVAL_MS: u64 = 60_000; #[cfg(any(feature = "signer-adapters", test))] pub const TRADE_SUBMIT_OPERATION_KIND: &str = "trade.submit.v1"; #[cfg(any(feature = "signer-adapters", test))] @@ -2041,6 +2057,177 @@ impl TradeStatusRequest { } #[cfg(feature = "runtime")] +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] +#[non_exhaustive] +pub struct TradeStatusWatchRequest { + pub status: TradeStatusRequest, + pub capacity: usize, + pub refresh_interval_ms: u64, + pub refresh_limit: Option<u32>, +} + +#[cfg(feature = "runtime")] +impl TradeStatusWatchRequest { + pub fn new(status: TradeStatusRequest) -> Self { + Self { + status, + capacity: TRADE_STATUS_WATCH_DEFAULT_CAPACITY, + refresh_interval_ms: TRADE_STATUS_WATCH_DEFAULT_REFRESH_INTERVAL_MS, + refresh_limit: None, + } + } + + pub fn parse(selector: &str) -> Result<Self, RadrootsSdkError> { + TradeStatusRequest::parse(selector).map(Self::new) + } + + pub fn with_capacity(mut self, capacity: usize) -> Self { + self.capacity = capacity; + self + } + + pub fn with_refresh_interval_ms(mut self, refresh_interval_ms: u64) -> Self { + self.refresh_interval_ms = refresh_interval_ms; + self + } + + pub fn with_refresh_limit(mut self, refresh_limit: u32) -> Self { + self.refresh_limit = Some(refresh_limit); + self + } + + pub fn without_refresh_limit(mut self) -> Self { + self.refresh_limit = None; + self + } + + fn validate(&self) -> Result<(), RadrootsSdkError> { + self.status.validate()?; + if self.capacity == 0 || self.capacity > TRADE_STATUS_WATCH_MAX_CAPACITY { + return Err(RadrootsSdkError::InvalidRequest { + message: format!( + "trade status watch capacity must be between 1 and {TRADE_STATUS_WATCH_MAX_CAPACITY}" + ), + }); + } + if self.refresh_interval_ms == 0 + || self.refresh_interval_ms > TRADE_STATUS_WATCH_MAX_REFRESH_INTERVAL_MS + { + return Err(RadrootsSdkError::InvalidRequest { + message: format!( + "trade status watch refresh interval must be between 1 and {TRADE_STATUS_WATCH_MAX_REFRESH_INTERVAL_MS} milliseconds" + ), + }); + } + if self.refresh_limit == Some(0) { + return Err(RadrootsSdkError::InvalidRequest { + message: "trade status watch refresh limit must be greater than zero".to_owned(), + }); + } + Ok(()) + } +} + +#[cfg(feature = "runtime")] +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] +pub struct TradeStatusWatchUpdate { + pub sequence: u64, + pub status: TradeStatusReceipt, +} + +#[cfg(feature = "runtime")] +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] +pub struct TradeStatusWatchCancelReceipt { + pub state: TradeStatusWatchCancelState, + pub buffered_updates_dropped: usize, +} + +#[cfg(feature = "runtime")] +#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)] +#[serde(rename_all = "snake_case")] +#[non_exhaustive] +pub enum TradeStatusWatchCancelState { + Cancelled, + AlreadyFinished, +} + +#[cfg(feature = "runtime")] +pub struct TradeStatusWatch { + receiver: mpsc::Receiver<Result<TradeStatusWatchUpdate, RadrootsSdkError>>, + cancel: Option<oneshot::Sender<()>>, + producer: Option<JoinHandle<()>>, + capacity: usize, +} + +#[cfg(feature = "runtime")] +impl TradeStatusWatch { + pub fn capacity(&self) -> usize { + self.capacity + } + + pub fn buffered_len(&self) -> usize { + self.receiver.len() + } + + pub async fn next(&mut self) -> Result<Option<TradeStatusWatchUpdate>, RadrootsSdkError> { + match self.receiver.recv().await { + Some(Ok(update)) => Ok(Some(update)), + Some(Err(error)) => Err(error), + None => Ok(None), + } + } + + pub async fn cancel(&mut self) -> TradeStatusWatchCancelReceipt { + let producer_active = self + .producer + .as_ref() + .map(|producer| !producer.is_finished()) + .unwrap_or(false); + let buffered_updates_dropped = self.drain_buffered_updates(); + let cancel_sent = self + .cancel + .take() + .map(|sender| sender.send(()).is_ok()) + .unwrap_or(false); + let producer_finished = match self.producer.take() { + Some(producer) => producer.await.is_ok(), + None => true, + }; + let state = if producer_active || cancel_sent || !producer_finished { + TradeStatusWatchCancelState::Cancelled + } else { + TradeStatusWatchCancelState::AlreadyFinished + }; + TradeStatusWatchCancelReceipt { + state, + buffered_updates_dropped, + } + } + + fn drain_buffered_updates(&mut self) -> usize { + self.receiver.close(); + let mut dropped = 0; + while self.receiver.try_recv().is_ok() { + dropped += 1; + } + dropped + } +} + +#[cfg(feature = "runtime")] +impl Drop for TradeStatusWatch { + fn drop(&mut self) { + self.drain_buffered_updates(); + if let Some(cancel) = self.cancel.take() { + let _ = cancel.send(()); + } + if let Some(producer) = self.producer.take() { + producer.abort(); + } + } +} + +#[cfg(feature = "runtime")] fn trade_status_selector_parts(selector: &str) -> Result<(&str, Option<&str>), RadrootsSdkError> { let selector = selector.trim(); if selector.is_empty() { @@ -3124,6 +3311,34 @@ impl<'sdk> TradesClient<'sdk> { } } + pub async fn watch( + &self, + request: TradeStatusWatchRequest, + ) -> Result<TradeStatusWatch, RadrootsSdkError> { + request.validate()?; + let handle = tokio::runtime::Handle::try_current().map_err(|_| { + RadrootsSdkError::InvalidRequest { + message: "trade status watch requires an active Tokio runtime".to_owned(), + } + })?; + let (updates, receiver) = mpsc::channel(request.capacity); + let (cancel, cancel_receiver) = oneshot::channel(); + let sdk = self.sdk.clone(); + let capacity = request.capacity; + let producer = handle.spawn(run_trade_status_watch( + sdk, + request, + updates, + cancel_receiver, + )); + Ok(TradeStatusWatch { + receiver, + cancel: Some(cancel), + producer: Some(producer), + capacity, + }) + } + #[cfg(all(feature = "signer-adapters", feature = "relay-runtime"))] pub async fn status_with_fetch_adapter<A>( &self, @@ -3260,6 +3475,44 @@ impl<'sdk> TradesClient<'sdk> { } #[cfg(feature = "runtime")] +async fn run_trade_status_watch( + sdk: crate::RadrootsClient, + request: TradeStatusWatchRequest, + updates: mpsc::Sender<Result<TradeStatusWatchUpdate, RadrootsSdkError>>, + mut cancel: oneshot::Receiver<()>, +) { + let refresh_interval = Duration::from_millis(request.refresh_interval_ms); + let mut sequence = 0_u64; + loop { + let status_result = sdk.trades().status(request.status.clone()).await; + let stop_after_update = status_result + .as_ref() + .map(|status| status.lifecycle_terminal) + .unwrap_or(true); + sequence += 1; + let update_result = status_result.map(|status| TradeStatusWatchUpdate { sequence, status }); + let send_result = tokio::select! { + _ = &mut cancel => break, + send_result = updates.send(update_result) => send_result, + }; + if send_result.is_err() || stop_after_update { + break; + } + if request + .refresh_limit + .map(|refresh_limit| sequence >= u64::from(refresh_limit)) + .unwrap_or(false) + { + break; + } + tokio::select! { + _ = &mut cancel => break, + _ = tokio::time::sleep(refresh_interval) => {} + } + } +} + +#[cfg(feature = "runtime")] impl<'sdk> TradeResyncClient<'sdk> { pub async fn resync( &self, diff --git a/crates/sdk/tests/orders_runtime.rs b/crates/sdk/tests/orders_runtime.rs @@ -48,13 +48,14 @@ use radroots_sdk::{ RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, RadrootsTradeValidationTrustPolicy, RadrootsTradeValidationTrustState, RelayResolutionPolicy, SdkMutationState, SdkRelayTargetSet, SdkRelayUrlPolicy, SdkTradeStatusIssue, SdkTradeStatusIssueKind, SdkTradeStatusSource, - TRADE_STATUS_DEFAULT_LIMIT, TRADE_STATUS_MAX_LIMIT, TRADE_SUBMIT_OPERATION_KIND, - TradeAcceptRequest, TradeCancelRequest, TradeDeclineRequest, TradeEvidenceIngestRequest, - TradeEvidenceMode, TradeEvidenceQueryBranchKind, TradeMutationOutcome, TradeProposeRequest, - TradeRequestEvidenceIngestRequest, TradeResyncRelayOutcomeKind, - TradeResyncRelayTransportOutcomeKind, TradeResyncRequest, TradeRevisionDecisionRequest, - TradeRevisionProposalRequest, TradeSellerInboxRequest, TradeStatusKind, - TradeStatusNextActionKind, TradeStatusRequest, TradeValidationReceiptInspectRequest, + TRADE_STATUS_DEFAULT_LIMIT, TRADE_STATUS_MAX_LIMIT, TRADE_STATUS_WATCH_MAX_CAPACITY, + TRADE_SUBMIT_OPERATION_KIND, TradeAcceptRequest, TradeCancelRequest, TradeDeclineRequest, + TradeEvidenceIngestRequest, TradeEvidenceMode, TradeEvidenceQueryBranchKind, + TradeMutationOutcome, TradeProposeRequest, TradeRequestEvidenceIngestRequest, + TradeResyncRelayOutcomeKind, TradeResyncRelayTransportOutcomeKind, TradeResyncRequest, + TradeRevisionDecisionRequest, TradeRevisionProposalRequest, TradeSellerInboxRequest, + TradeStatusKind, TradeStatusNextActionKind, TradeStatusRequest, TradeStatusWatchCancelState, + TradeStatusWatchRequest, TradeValidationReceiptInspectRequest, TradeValidationReceiptListRequest, TradeValidationReceiptVerifyRequest, }; use radroots_sdk::{PrivacyPreflightConfirmation, PrivacyPreflightStatus, ProductSensitivityField}; @@ -6833,6 +6834,136 @@ async fn order_status_maps_malformed_local_data_to_sanitized_error() { } #[tokio::test] +async fn trade_status_watch_emits_finite_refresh_window() { + let (_tempdir, sdk, store) = directory_sdk_and_store().await; + let order_id = "watch-finite-refresh-window"; + let request_event = signed_order_request_event(order_id, 820); + store + .ingest_event(RadrootsEventIngest::new(request_event, 1_700_400_000_000)) + .await + .expect("request ingest"); + + let mut watch = sdk + .trades() + .watch( + TradeStatusWatchRequest::new(status_request(order_id)) + .with_capacity(2) + .with_refresh_interval_ms(1) + .with_refresh_limit(2), + ) + .await + .expect("watch"); + + let first = watch.next().await.expect("first").expect("first update"); + let second = watch.next().await.expect("second").expect("second update"); + let closed = watch.next().await.expect("closed"); + let cancel = watch.cancel().await; + + assert_eq!(watch.capacity(), 2); + assert_eq!(first.sequence, 1); + assert_eq!(second.sequence, 2); + assert_eq!(first.status.status, TradeStatusKind::Requested); + assert_eq!(second.status.status, TradeStatusKind::Requested); + assert!(closed.is_none()); + assert_eq!(cancel.state, TradeStatusWatchCancelState::AlreadyFinished); +} + +#[tokio::test] +async fn trade_status_watch_backpressures_slow_consumer_with_bounded_buffer() { + let (_tempdir, sdk, _store) = directory_sdk_and_store().await; + let mut watch = sdk + .trades() + .watch( + TradeStatusWatchRequest::parse("watch-slow-consumer") + .expect("watch request") + .with_capacity(1) + .with_refresh_interval_ms(1) + .with_refresh_limit(50), + ) + .await + .expect("watch"); + + tokio::time::sleep(Duration::from_millis(25)).await; + let buffered_len = watch.buffered_len(); + let cancel = watch.cancel().await; + let post_cancel = watch.next().await.expect("post cancel"); + + assert_eq!(buffered_len, 1); + assert_eq!(cancel.state, TradeStatusWatchCancelState::Cancelled); + assert_eq!(cancel.buffered_updates_dropped, 1); + assert!(post_cancel.is_none()); +} + +#[tokio::test] +async fn trade_status_watch_cancel_drains_buffer_and_closes_stream() { + let (_tempdir, sdk, _store) = directory_sdk_and_store().await; + let mut watch = sdk + .trades() + .watch( + TradeStatusWatchRequest::parse("watch-cancel-close") + .expect("watch request") + .with_capacity(4) + .with_refresh_interval_ms(1), + ) + .await + .expect("watch"); + let first = watch.next().await.expect("first").expect("first update"); + + let cancel = watch.cancel().await; + tokio::time::sleep(Duration::from_millis(5)).await; + let post_cancel = watch.next().await.expect("post cancel"); + + assert_eq!(first.sequence, 1); + assert_eq!(first.status.status, TradeStatusKind::Missing); + assert_eq!(cancel.state, TradeStatusWatchCancelState::Cancelled); + assert!(post_cancel.is_none()); +} + +#[tokio::test] +async fn trade_status_watch_closes_after_producer_error() { + let sdk = RadrootsClient::builder().build().await.expect("sdk"); + let mut watch = sdk + .trades() + .watch( + TradeStatusWatchRequest::new( + TradeStatusRequest::parse("watch-producer-error") + .expect("status request") + .with_source(SdkTradeStatusSource::ResyncThenLocal), + ) + .with_refresh_interval_ms(1) + .with_refresh_limit(2) + .with_capacity(2), + ) + .await + .expect("watch"); + + let error = watch.next().await.expect_err("producer error"); + let closed = watch.next().await.expect("closed"); + + assert!(!error.to_string().is_empty()); + assert!(closed.is_none()); +} + +#[tokio::test] +async fn trade_status_watch_rejects_unbounded_capacity() { + let sdk = RadrootsClient::builder().build().await.expect("sdk"); + let result = sdk + .trades() + .watch( + TradeStatusWatchRequest::parse("watch-invalid-capacity") + .expect("watch request") + .with_capacity(TRADE_STATUS_WATCH_MAX_CAPACITY + 1), + ) + .await; + let Err(error) = result else { + panic!("expected capacity error"); + }; + + assert!(matches!(error, RadrootsSdkError::InvalidRequest { .. })); + assert!(error.to_string().contains("capacity")); +} + +#[tokio::test] #[ignore = "manual expensive release-gate lane for the 100k local-event status target"] async fn manual_local_status_perf_gate_measures_100k_events() { let (_tempdir, sdk, store) = directory_sdk_and_store().await;