commit eb827e0a4b4dd97e444ab138c1377b3884855fb7
parent f5e7b656ce9d3dac3bc8281ee09088cc5c7a1a81
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:
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;