lib

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

commit 7af64268eda3a7f9aa062b1a144c65beffbb27bd
parent 05530a10ed58910f9309e1560820c9091edf300e
Author: triesap <tyson@radroots.org>
Date:   Wed,  1 Jul 2026 07:57:28 +0000

sdk: resync trades from configured relays

- replace local-only trade resync with configured relay fetch and event-store ingest
- return typed relay evidence import receipts with projection refresh status
- cover empty import, duplicate replay, malformed evidence, and relay failures
- validate with SDK fmt, source-boundary, focused flow tests, and all-features check

Diffstat:
Mcrates/sdk/src/lib.rs | 11+++++++----
Mcrates/sdk/src/orders_runtime.rs | 301+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcrates/sdk/tests/orders_runtime.rs | 421++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------
Mcrates/sdk/tests/trade_public_api.rs | 6+++---
4 files changed, 662 insertions(+), 77 deletions(-)

diff --git a/crates/sdk/src/lib.rs b/crates/sdk/src/lib.rs @@ -105,10 +105,13 @@ pub use crate::market_runtime::{ pub use crate::orders_runtime::{ SdkTradeStatusIssue, SdkTradeStatusIssueKind, SdkTradeStatusSource, TRADE_STATUS_DEFAULT_LIMIT, TRADE_STATUS_MAX_LIMIT, TradeEvidenceIngestReceipt, TradeEvidenceIngestRequest, - TradeRequestEvidenceIngestReceipt, TradeRequestEvidenceIngestRequest, TradeResyncReceipt, - TradeResyncRequest, TradeSellerInboxReceipt, TradeSellerInboxRequest, - TradeStatusAmbiguityCandidate, TradeStatusEligibility, TradeStatusEvidenceSummary, - TradeStatusKind, TradeStatusNextActionKind, TradeStatusReceipt, TradeStatusRequest, + TradeRequestEvidenceIngestReceipt, TradeRequestEvidenceIngestRequest, + TradeResyncEventImportReceipt, TradeResyncEvidenceReceipt, TradeResyncReceipt, + TradeResyncRelayOutcomeKind, TradeResyncRelayOutcomeReceipt, + TradeResyncRelayTransportOutcomeKind, TradeResyncRequest, TradeSellerInboxReceipt, + TradeSellerInboxRequest, TradeStatusAmbiguityCandidate, TradeStatusEligibility, + TradeStatusEvidenceSummary, TradeStatusKind, TradeStatusNextActionKind, TradeStatusReceipt, + TradeStatusRequest, }; #[cfg(all(feature = "runtime", feature = "signer-adapters"))] pub use crate::orders_runtime::{ diff --git a/crates/sdk/src/orders_runtime.rs b/crates/sdk/src/orders_runtime.rs @@ -1,5 +1,11 @@ #[cfg(feature = "signer-adapters")] use crate::TradeBuyerClient; +#[cfg(feature = "runtime")] +use crate::runtime::sdk_now_ms; +#[cfg(feature = "runtime")] +use crate::sync_runtime::{ + SyncProjectionRefreshReceipt, SyncProjectionRefreshRequest, refresh_product_projections_for_sdk, +}; #[cfg(feature = "signer-adapters")] use crate::workflow_runtime::enqueue_configured_signed_workflow; #[cfg(any(feature = "signer-adapters", test))] @@ -29,10 +35,11 @@ use radroots_events::{ ids::{RadrootsEventId, RadrootsListingAddress, RadrootsOrderId, RadrootsPublicKey}, kinds::{ KIND_ORDER_CANCELLATION, KIND_ORDER_DECISION, KIND_ORDER_REQUEST, - KIND_ORDER_REVISION_DECISION, KIND_ORDER_REVISION_PROPOSAL, + KIND_ORDER_REVISION_DECISION, KIND_ORDER_REVISION_PROPOSAL, KIND_TRADE_VALIDATION_RECEIPT, + ORDER_EVENT_KINDS, }, order::RadrootsOrderEconomics, - tags::TAG_P, + tags::{TAG_D, TAG_P}, }; #[cfg(any(feature = "signer-adapters", test))] use radroots_events::{ @@ -52,6 +59,14 @@ use radroots_events_codec::order::{ }; #[cfg(any(feature = "signer-adapters", test))] use radroots_events_codec::wire::{WireEventParts, to_frozen_draft}; +#[cfg(all(feature = "runtime", feature = "relay-runtime"))] +use radroots_nostr::prelude::{RadrootsNostrFilter, RadrootsNostrKind, radroots_nostr_filter_tag}; +#[cfg(feature = "runtime")] +use radroots_relay_transport::{ + RadrootsNostrClientFetchAdapter, RadrootsRelayFetchAdapter, RadrootsRelayFetchEventReceipt, + RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchReceipt, RadrootsRelayFetchRelayOutcome, + RadrootsRelayFetchRequest, RadrootsRelayOutcomeKind, fetch_and_ingest_relay_events, +}; #[cfg(feature = "runtime")] use radroots_trade::identity::{RadrootsTradeLocator, RadrootsTradeLocatorCandidate}; #[cfg(any(feature = "signer-adapters", test))] @@ -1398,16 +1413,101 @@ impl TradeResyncRequest { self.limit = limit; self } + + fn validate(&self) -> Result<(), RadrootsSdkError> { + if self.limit == 0 || self.limit > TRADE_STATUS_MAX_LIMIT { + return Err(RadrootsSdkError::trade_status_limit_invalid( + self.limit, + 1, + TRADE_STATUS_MAX_LIMIT, + )); + } + Ok(()) + } } #[cfg(feature = "runtime")] #[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] pub struct TradeResyncReceipt { + pub relay_targets: Vec<String>, + pub evidence: TradeResyncEvidenceReceipt, + pub refresh: SyncProjectionRefreshReceipt, pub status: TradeStatusReceipt, } #[cfg(feature = "runtime")] #[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] +pub struct TradeResyncEvidenceReceipt { + pub inserted_count: usize, + pub duplicate_count: usize, + pub malformed_count: usize, + pub unsupported_count: usize, + pub eose_count: usize, + pub closed_count: usize, + pub notice_count: usize, + pub events: Vec<TradeResyncEventImportReceipt>, + pub relays: Vec<TradeResyncRelayOutcomeReceipt>, +} + +#[cfg(feature = "runtime")] +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] +pub struct TradeResyncEventImportReceipt { + pub relay_url: String, + pub event_id: Option<String>, + pub inserted: bool, + pub duplicate: bool, + pub unsupported: bool, + pub malformed: bool, + pub projection_eligible: bool, + pub verification_status: Option<String>, + pub message: Option<String>, +} + +#[cfg(feature = "runtime")] +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] +pub struct TradeResyncRelayOutcomeReceipt { + pub relay_url: String, + pub outcome_kind: TradeResyncRelayOutcomeKind, + pub transport_outcome_kind: Option<TradeResyncRelayTransportOutcomeKind>, + pub message: Option<String>, +} + +#[cfg(feature = "runtime")] +#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)] +#[serde(rename_all = "snake_case")] +#[non_exhaustive] +pub enum TradeResyncRelayOutcomeKind { + Eose, + Closed, + Notice, +} + +#[cfg(feature = "runtime")] +#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)] +#[serde(rename_all = "snake_case")] +#[non_exhaustive] +pub enum TradeResyncRelayTransportOutcomeKind { + Accepted, + DuplicateAccepted, + Blocked, + RateLimited, + Invalid, + PowRequired, + Restricted, + AuthRequired, + Muted, + Unsupported, + PaymentRequired, + Error, + Timeout, + ConnectionFailed, + RelayUrlRejected, + SkippedAlreadyAccepted, + Unknown, +} + +#[cfg(feature = "runtime")] +#[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] #[non_exhaustive] pub struct TradeStatusRequest { pub locator: RadrootsTradeLocator, @@ -2558,10 +2658,203 @@ impl<'sdk> TradeResyncClient<'sdk> { &self, request: TradeResyncRequest, ) -> Result<TradeResyncReceipt, RadrootsSdkError> { + #[cfg(feature = "relay-runtime")] + { + let adapter = RadrootsNostrClientFetchAdapter; + return self.resync_with_fetch_adapter(request, &adapter).await; + } + #[cfg(not(feature = "relay-runtime"))] + { + let _ = request; + Err(RadrootsSdkError::ProductSyncUnsupported { + operation: "trade.resync", + required_feature: "relay-runtime", + }) + } + } + + #[cfg(feature = "relay-runtime")] + pub async fn resync_with_fetch_adapter<A>( + &self, + request: TradeResyncRequest, + adapter: &A, + ) -> Result<TradeResyncReceipt, RadrootsSdkError> + where + A: RadrootsRelayFetchAdapter, + { + request.validate()?; + let relay_targets = self.sdk.relay_urls().to_vec(); + if relay_targets.is_empty() { + return Err(RadrootsSdkError::empty_target_relays("trade.resync")); + } + let fetch_request = trade_resync_fetch_request(self.sdk, &request, &relay_targets)?; + let fetch_receipt = + fetch_and_ingest_relay_events(adapter, &self.sdk._event_store, fetch_request).await?; + if trade_resync_total_relay_failure(&fetch_receipt, relay_targets.len()) { + return Err(RadrootsSdkError::ProductSyncRelaySetupFailure { + message: trade_resync_total_failure_message(&fetch_receipt), + }); + } + let evidence = TradeResyncEvidenceReceipt::from_fetch(fetch_receipt); + let refresh = refresh_product_projections_for_sdk( + self.sdk, + SyncProjectionRefreshRequest::new().with_limit(request.limit), + ) + .await?; let status = trades_client(self.sdk) - .status(TradeStatusRequest::new(request.locator).with_limit(request.limit)) + .status(TradeStatusRequest::new(request.locator.clone()).with_limit(request.limit)) .await?; - Ok(TradeResyncReceipt { status }) + if status.status == TradeStatusKind::Ambiguous { + return Err(RadrootsSdkError::TradeAmbiguous { + operation: "trade.resync".to_owned(), + locator: request.locator, + candidates: status + .ambiguity_candidates + .iter() + .map(|candidate| candidate.locator.clone()) + .collect(), + }); + } + Ok(TradeResyncReceipt { + relay_targets, + evidence, + refresh, + status, + }) + } +} + +#[cfg(all(feature = "runtime", feature = "relay-runtime"))] +fn trade_resync_fetch_request( + sdk: &crate::RadrootsClient, + request: &TradeResyncRequest, + relay_targets: &[String], +) -> Result<RadrootsRelayFetchRequest, RadrootsSdkError> { + Ok( + RadrootsRelayFetchRequest::fetch(sdk_now_ms(sdk)?, request.limit as usize) + .with_relay_urls(relay_targets.iter().cloned()) + .with_filters([trade_resync_filter(&request.locator, request.limit)?]), + ) +} + +#[cfg(all(feature = "runtime", feature = "relay-runtime"))] +fn trade_resync_filter( + locator: &RadrootsTradeLocator, + limit: u32, +) -> Result<RadrootsNostrFilter, RadrootsSdkError> { + let mut filter = RadrootsNostrFilter::new().limit(limit as usize); + for kind in ORDER_EVENT_KINDS + .iter() + .copied() + .chain([KIND_TRADE_VALIDATION_RECEIPT]) + { + let kind = u16::try_from(kind).map_err(|_| RadrootsSdkError::InvalidRequest { + message: format!("trade resync event kind {kind} exceeds Nostr filter range"), + })?; + filter = filter.kind(RadrootsNostrKind::Custom(kind)); + } + radroots_nostr_filter_tag(filter, TAG_D, vec![locator.order_id().as_str().to_owned()]).map_err( + |error| RadrootsSdkError::InvalidRequest { + message: format!("trade resync filter invalid: {error}"), + }, + ) +} + +#[cfg(feature = "runtime")] +fn trade_resync_total_relay_failure( + receipt: &RadrootsRelayFetchReceipt, + relay_count: usize, +) -> bool { + relay_count > 0 && receipt.eose_count == 0 && receipt.closed_count >= relay_count +} + +#[cfg(feature = "runtime")] +fn trade_resync_total_failure_message(receipt: &RadrootsRelayFetchReceipt) -> String { + format!( + "trade.resync failed for all configured relays: closed_count={}, notice_count={}, malformed_count={}", + receipt.closed_count, receipt.notice_count, receipt.malformed_count + ) +} + +#[cfg(feature = "runtime")] +impl TradeResyncEvidenceReceipt { + fn from_fetch(receipt: RadrootsRelayFetchReceipt) -> Self { + Self { + inserted_count: receipt.inserted_count, + duplicate_count: receipt.duplicate_count, + malformed_count: receipt.malformed_count, + unsupported_count: receipt.unsupported_count, + eose_count: receipt.eose_count, + closed_count: receipt.closed_count, + notice_count: receipt.notice_count, + events: receipt.events.into_iter().map(Into::into).collect(), + relays: receipt.relay_outcomes.into_iter().map(Into::into).collect(), + } + } +} + +#[cfg(feature = "runtime")] +impl From<RadrootsRelayFetchEventReceipt> for TradeResyncEventImportReceipt { + fn from(receipt: RadrootsRelayFetchEventReceipt) -> Self { + Self { + relay_url: receipt.relay_url, + event_id: receipt.event_id, + inserted: receipt.inserted, + duplicate: receipt.duplicate, + unsupported: receipt.unsupported, + malformed: receipt.malformed, + projection_eligible: receipt.projection_eligible, + verification_status: receipt.verification_status, + message: receipt.message, + } + } +} + +#[cfg(feature = "runtime")] +impl From<RadrootsRelayFetchRelayOutcome> for TradeResyncRelayOutcomeReceipt { + fn from(receipt: RadrootsRelayFetchRelayOutcome) -> Self { + Self { + relay_url: receipt.relay_url, + outcome_kind: receipt.kind.into(), + transport_outcome_kind: receipt.relay_outcome.map(|outcome| outcome.kind.into()), + message: receipt.message, + } + } +} + +#[cfg(feature = "runtime")] +impl From<RadrootsRelayFetchOutcomeKind> for TradeResyncRelayOutcomeKind { + fn from(kind: RadrootsRelayFetchOutcomeKind) -> Self { + match kind { + RadrootsRelayFetchOutcomeKind::Eose => Self::Eose, + RadrootsRelayFetchOutcomeKind::Closed => Self::Closed, + RadrootsRelayFetchOutcomeKind::Notice => Self::Notice, + } + } +} + +#[cfg(feature = "runtime")] +impl From<RadrootsRelayOutcomeKind> for TradeResyncRelayTransportOutcomeKind { + fn from(kind: RadrootsRelayOutcomeKind) -> Self { + match kind { + RadrootsRelayOutcomeKind::Accepted => Self::Accepted, + RadrootsRelayOutcomeKind::DuplicateAccepted => Self::DuplicateAccepted, + RadrootsRelayOutcomeKind::Blocked => Self::Blocked, + RadrootsRelayOutcomeKind::RateLimited => Self::RateLimited, + RadrootsRelayOutcomeKind::Invalid => Self::Invalid, + RadrootsRelayOutcomeKind::PowRequired => Self::PowRequired, + RadrootsRelayOutcomeKind::Restricted => Self::Restricted, + RadrootsRelayOutcomeKind::AuthRequired => Self::AuthRequired, + RadrootsRelayOutcomeKind::Muted => Self::Muted, + RadrootsRelayOutcomeKind::Unsupported => Self::Unsupported, + RadrootsRelayOutcomeKind::PaymentRequired => Self::PaymentRequired, + RadrootsRelayOutcomeKind::Error => Self::Error, + RadrootsRelayOutcomeKind::Timeout => Self::Timeout, + RadrootsRelayOutcomeKind::ConnectionFailed => Self::ConnectionFailed, + RadrootsRelayOutcomeKind::RelayUrlRejected => Self::RelayUrlRejected, + RadrootsRelayOutcomeKind::SkippedAlreadyAccepted => Self::SkippedAlreadyAccepted, + RadrootsRelayOutcomeKind::Unknown => Self::Unknown, + } } } diff --git a/crates/sdk/tests/orders_runtime.rs b/crates/sdk/tests/orders_runtime.rs @@ -4,6 +4,7 @@ use std::path::Path; use std::time::{Duration, Instant}; +use nostr::JsonUtil; use radroots_authority::RadrootsActorContext; use radroots_core::{ RadrootsCoreCurrency, RadrootsCoreDecimal, RadrootsCoreMoney, RadrootsCoreUnit, @@ -30,6 +31,7 @@ use radroots_nostr::prelude::{ radroots_nostr_build_event, }; use radroots_outbox::RadrootsOutbox; +use radroots_relay_transport::{RadrootsMockRelayFetchAdapter, RadrootsRelayFetchItem}; use radroots_sdk::{ AckPolicy, DvmValidationReceiptIngestRequest, PublishMode, RadrootsClient, RadrootsSdkError, RadrootsSdkRecoveryAction, RadrootsSdkTimestamp, RelayResolutionPolicy, SdkMutationState, @@ -37,20 +39,21 @@ use radroots_sdk::{ SdkTradeStatusSource, TRADE_STATUS_DEFAULT_LIMIT, TRADE_STATUS_MAX_LIMIT, TRADE_SUBMIT_OPERATION_KIND, TradeAcceptRequest, TradeCancelRequest, TradeDeclineRequest, TradeEvidenceIngestRequest, TradeMutationOutcome, TradeProposeRequest, - TradeRequestEvidenceIngestRequest, TradeResyncRequest, TradeRevisionDecisionRequest, + TradeRequestEvidenceIngestRequest, TradeResyncRelayOutcomeKind, + TradeResyncRelayTransportOutcomeKind, TradeResyncRequest, TradeRevisionDecisionRequest, TradeRevisionProposalRequest, TradeSellerInboxRequest, TradeStatusKind, TradeStatusNextActionKind, TradeStatusRequest, }; use radroots_sdk::{PrivacyPreflightConfirmation, PrivacyPreflightStatus, ProductSensitivityField}; #[cfg(all(feature = "signer-adapters", feature = "local-signer"))] use radroots_sdk::{RadrootsSdkLocalKeySigner, RadrootsSdkSignerProvider}; -use radroots_trade::order::RadrootsOrderIssue; use radroots_trade::validation_receipt::{ RadrootsTradeValidationReceipt, RadrootsValidationReceiptProof, RadrootsValidationReceiptProofSystem, RadrootsValidationReceiptResult, RadrootsValidationReceiptStatement, RadrootsValidationReceiptType, validation_receipt_event_build, validation_receipt_public_values_hash_hex, }; +use radroots_trade::{identity::RadrootsTradeLocator, order::RadrootsOrderIssue}; use serde::Serialize; use serde::ser::{self, SerializeStruct}; @@ -68,7 +71,6 @@ const RELAY: &str = "wss://relay.radroots.test"; #[cfg(any())] const OTHER_PUBLIC_KEY_HEX: &str = "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc"; -#[cfg(any())] const RELAY_B: &str = "wss://relay-b.radroots.test"; const PERF_TOTAL_LOCAL_EVENTS: i64 = 100_000; const PERF_TRADE_RELEVANT_EVENTS: i64 = 25_000; @@ -405,19 +407,47 @@ async fn directory_sdk_and_store() -> (tempfile::TempDir, RadrootsClient, Radroo (tempdir, sdk, store) } +async fn directory_sdk_and_store_with_relays( + relays: &[&str], +) -> (tempfile::TempDir, RadrootsClient, RadrootsEventStore) { + let tempdir = tempfile::tempdir().expect("tempdir"); + let mut builder = RadrootsClient::builder() + .directory_storage(tempdir.path().join("sdk")) + .fixed_clock(RadrootsSdkTimestamp::from_unix_seconds(1_700_000_000)); + for relay in relays { + builder = builder.relay_url(*relay); + } + let sdk = builder.build().await.expect("sdk"); + let store = + RadrootsEventStore::open_file(&sdk.storage_paths().expect("paths").event_store_path) + .await + .expect("event store"); + (tempdir, sdk, store) +} + #[cfg(all(feature = "signer-adapters", feature = "local-signer"))] async fn directory_sdk_with_signer(storage_root: &Path, secret_key_hex: &str) -> RadrootsClient { + directory_sdk_with_signer_and_relays(storage_root, secret_key_hex, &[]).await +} + +#[cfg(all(feature = "signer-adapters", feature = "local-signer"))] +async fn directory_sdk_with_signer_and_relays( + storage_root: &Path, + secret_key_hex: &str, + relays: &[&str], +) -> RadrootsClient { let secret_key = RadrootsNostrSecretKey::from_hex(secret_key_hex).expect("secret key"); let signer_keys = RadrootsNostrKeys::new(secret_key); - RadrootsClient::builder() + let mut builder = RadrootsClient::builder() .directory_storage(storage_root) .fixed_clock(RadrootsSdkTimestamp::from_unix_seconds(1_700_000_000)) .signer_provider(RadrootsSdkSignerProvider::LocalKey( RadrootsSdkLocalKeySigner::new(signer_keys).expect("local signer"), - )) - .build() - .await - .expect("sdk") + )); + for relay in relays { + builder = builder.relay_url(*relay); + } + builder.build().await.expect("sdk") } fn order_id(raw: &str) -> RadrootsOrderId { @@ -849,12 +879,21 @@ async fn order_submit_enqueue_stores_event_queues_outbox_and_status_sees_request ); } -#[cfg(all(feature = "signer-adapters", feature = "local-signer"))] +#[cfg(all( + feature = "signer-adapters", + feature = "local-signer", + feature = "relay-runtime" +))] #[tokio::test] 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 buyer_sdk = directory_sdk_with_signer_and_relays( + storage_root.as_path(), + BUYER_SECRET_KEY_HEX, + &[RELAY], + ) + .await; let propose_receipt = expect_enqueued( buyer_sdk .trades() @@ -900,7 +939,12 @@ async fn trade_product_clients_propose_inbox_accept_status_and_resync() { Some(SELLER_PUBLIC_KEY_HEX) ); - let seller_sdk = directory_sdk_with_signer(storage_root.as_path(), SELLER_SECRET_KEY_HEX).await; + let seller_sdk = directory_sdk_with_signer_and_relays( + storage_root.as_path(), + SELLER_SECRET_KEY_HEX, + &[RELAY], + ) + .await; let inbox = seller_sdk .trades() .seller() @@ -961,12 +1005,18 @@ async fn trade_product_clients_propose_inbox_accept_status_and_resync() { Some(accept_receipt.signed_event_id.as_str()) ); + let resync_adapter = RadrootsMockRelayFetchAdapter::new(vec![relay_eose(RELAY)]); let resync = seller_sdk .trades() .resync() - .resync(TradeResyncRequest::new(propose_receipt.locator)) + .resync_with_fetch_adapter( + TradeResyncRequest::new(propose_receipt.locator), + &resync_adapter, + ) .await .expect("facade resync"); + assert_eq!(resync.relay_targets, vec![RELAY.to_owned()]); + assert_eq!(resync.evidence.eose_count, 1); assert_eq!(resync.status.status, TradeStatusKind::AgreedPendingRhi); assert_eq!( resync.status.last_event_id, @@ -974,16 +1024,28 @@ async fn trade_product_clients_propose_inbox_accept_status_and_resync() { ); } -#[cfg(all(feature = "signer-adapters", feature = "local-signer"))] +#[cfg(all( + feature = "signer-adapters", + feature = "local-signer", + feature = "relay-runtime" +))] #[tokio::test] async fn trade_product_clients_resync_committed_after_rhi_validation_receipt() { let tempdir = tempfile::tempdir().expect("tempdir"); let buyer_storage_root = tempdir.path().join("buyer-sdk"); let seller_storage_root = tempdir.path().join("seller-sdk"); - let buyer_sdk = - directory_sdk_with_signer(buyer_storage_root.as_path(), BUYER_SECRET_KEY_HEX).await; - let seller_sdk = - directory_sdk_with_signer(seller_storage_root.as_path(), SELLER_SECRET_KEY_HEX).await; + let buyer_sdk = directory_sdk_with_signer_and_relays( + buyer_storage_root.as_path(), + BUYER_SECRET_KEY_HEX, + &[RELAY], + ) + .await; + let seller_sdk = directory_sdk_with_signer_and_relays( + seller_storage_root.as_path(), + SELLER_SECRET_KEY_HEX, + &[RELAY], + ) + .await; let buyer_store = RadrootsEventStore::open_file(&buyer_sdk.storage_paths().expect("paths").event_store_path) .await @@ -1019,13 +1081,25 @@ async fn trade_product_clients_resync_committed_after_rhi_validation_receipt() { .total_events, 0 ); - replay_stored_event( - &buyer_store, - &seller_store, - &propose_receipt.signed_event_id, - 4_000, - ) - .await; + let seller_proposal_resync_adapter = RadrootsMockRelayFetchAdapter::new(vec![ + relay_event_item_from_store(&buyer_store, &propose_receipt.signed_event_id, RELAY, 4_000) + .await, + relay_eose(RELAY), + ]); + let seller_proposal_resync = seller_sdk + .trades() + .resync() + .resync_with_fetch_adapter( + TradeResyncRequest::new(propose_receipt.locator.clone()), + &seller_proposal_resync_adapter, + ) + .await + .expect("seller proposal resync"); + assert_eq!( + seller_proposal_resync.status.status, + TradeStatusKind::Requested + ); + assert_eq!(seller_proposal_resync.evidence.inserted_count, 1); let accept_receipt = expect_enqueued( seller_sdk .trades() @@ -1048,21 +1122,16 @@ async fn trade_product_clients_resync_committed_after_rhi_validation_receipt() { .await .expect("accept trade"), ); - replay_stored_event( - &seller_store, - &buyer_store, - &accept_receipt.signed_event_id, - 4_100, - ) - .await; - let receipt_event = signed_validation_receipt_event( + let receipt_raw_event = signed_raw_validation_receipt_event( "trade-product-committed-resync", &propose_receipt.listing_event_id, &propose_receipt.signed_event_id, &accept_receipt.signed_event_id, 33, ); - let receipt_event_id = RadrootsEventId::parse(receipt_event.id.as_str()).expect("receipt id"); + let receipt_event = radroots_event_from_nostr(&receipt_raw_event); + let receipt_event_id = + RadrootsEventId::parse(receipt_raw_event.id.to_hex().as_str()).expect("receipt id"); let ingest = seller_sdk .dvm() @@ -1077,12 +1146,14 @@ async fn trade_product_clients_resync_committed_after_rhi_validation_receipt() { .expect("ingest validation receipt"); assert!(ingest.inserted); assert_eq!(ingest.receipt_event_id, receipt_event_id); - replay_stored_event(&seller_store, &buyer_store, &receipt_event_id, 4_200).await; let seller_resync = seller_sdk .trades() .resync() - .resync(TradeResyncRequest::new(propose_receipt.locator.clone())) + .resync_with_fetch_adapter( + TradeResyncRequest::new(propose_receipt.locator.clone()), + &RadrootsMockRelayFetchAdapter::new(vec![relay_eose(RELAY)]), + ) .await .expect("seller resync"); assert_eq!(seller_resync.status.status, TradeStatusKind::Committed); @@ -1095,10 +1166,19 @@ async fn trade_product_clients_resync_committed_after_rhi_validation_receipt() { Some(receipt_event_id.clone()) ); + let buyer_committed_resync_adapter = RadrootsMockRelayFetchAdapter::new(vec![ + relay_event_item_from_store(&seller_store, &accept_receipt.signed_event_id, RELAY, 4_100) + .await, + relay_raw_event_item(&receipt_raw_event, RELAY, 4_200), + relay_eose(RELAY), + ]); let buyer_resync = buyer_sdk .trades() .resync() - .resync(TradeResyncRequest::new(propose_receipt.locator)) + .resync_with_fetch_adapter( + TradeResyncRequest::new(propose_receipt.locator), + &buyer_committed_resync_adapter, + ) .await .expect("buyer resync"); assert_eq!(buyer_resync.status.status, TradeStatusKind::Committed); @@ -1106,41 +1186,220 @@ async fn trade_product_clients_resync_committed_after_rhi_validation_receipt() { buyer_resync.status.rhi_receipt_event_id, Some(receipt_event_id) ); + assert_eq!(buyer_resync.evidence.inserted_count, 2); +} + +#[cfg(feature = "relay-runtime")] +#[tokio::test] +async fn trade_resync_imports_relay_evidence_into_empty_local_store() { + let (_tempdir, sdk, store) = directory_sdk_and_store_with_relays(&[RELAY]).await; + let request_event = signed_raw_order_request_event("resync-empty-local-import", 41); + let request_event_id = + RadrootsEventId::parse(request_event.id.to_hex().as_str()).expect("event id"); + let adapter = RadrootsMockRelayFetchAdapter::new(vec![ + relay_raw_event_item(&request_event, RELAY, 5_000), + relay_eose(RELAY), + ]); + + let resync = sdk + .trades() + .resync() + .resync_with_fetch_adapter( + TradeResyncRequest::new(RadrootsTradeLocator::from_order_id(order_id( + "resync-empty-local-import", + ))), + &adapter, + ) + .await + .expect("resync"); + + assert_eq!(resync.status.status, TradeStatusKind::Requested); + assert_eq!(resync.evidence.inserted_count, 1); + assert_eq!(resync.evidence.duplicate_count, 0); + assert_eq!( + resync.evidence.events[0].event_id.as_deref(), + Some(request_event_id.as_str()) + ); + assert!( + store + .get_event(request_event_id.as_str()) + .await + .expect("stored event") + .is_some() + ); } -async fn replay_stored_event( +#[cfg(feature = "relay-runtime")] +#[tokio::test] +async fn trade_resync_duplicate_replay_is_idempotent() { + let (_tempdir, sdk, _store) = directory_sdk_and_store_with_relays(&[RELAY]).await; + let request_event = signed_raw_order_request_event("resync-duplicate-replay", 42); + let adapter = RadrootsMockRelayFetchAdapter::new(vec![ + relay_raw_event_item(&request_event, RELAY, 5_100), + relay_eose(RELAY), + ]); + let locator = RadrootsTradeLocator::from_order_id(order_id("resync-duplicate-replay")); + + let first = sdk + .trades() + .resync() + .resync_with_fetch_adapter(TradeResyncRequest::new(locator.clone()), &adapter) + .await + .expect("first resync"); + let second = sdk + .trades() + .resync() + .resync_with_fetch_adapter(TradeResyncRequest::new(locator), &adapter) + .await + .expect("second resync"); + + assert_eq!(first.evidence.inserted_count, 1); + assert_eq!(second.evidence.inserted_count, 0); + assert_eq!(second.evidence.duplicate_count, 1); + assert_eq!(second.status.status, TradeStatusKind::Requested); +} + +#[cfg(feature = "relay-runtime")] +#[tokio::test] +async fn trade_resync_reports_malformed_evidence_without_poisoning_store() { + let (_tempdir, sdk, store) = directory_sdk_and_store_with_relays(&[RELAY]).await; + let adapter = + RadrootsMockRelayFetchAdapter::new(vec![relay_malformed(RELAY), relay_eose(RELAY)]); + + let resync = sdk + .trades() + .resync() + .resync_with_fetch_adapter( + TradeResyncRequest::new(RadrootsTradeLocator::from_order_id(order_id( + "resync-malformed-evidence", + ))), + &adapter, + ) + .await + .expect("resync"); + + assert_eq!(resync.status.status, TradeStatusKind::Missing); + assert_eq!(resync.evidence.malformed_count, 1); + assert_eq!(resync.evidence.inserted_count, 0); + assert_eq!( + store + .status_summary() + .await + .expect("store summary") + .total_events, + 0 + ); +} + +#[cfg(feature = "relay-runtime")] +#[tokio::test] +async fn trade_resync_errors_on_total_relay_failure() { + let (_tempdir, sdk, _store) = directory_sdk_and_store_with_relays(&[RELAY, RELAY_B]).await; + let adapter = RadrootsMockRelayFetchAdapter::new(vec![ + relay_closed(RELAY, "timeout: relay offline"), + relay_closed(RELAY_B, "error: relay unavailable"), + ]); + + let error = sdk + .trades() + .resync() + .resync_with_fetch_adapter( + TradeResyncRequest::new(RadrootsTradeLocator::from_order_id(order_id( + "resync-total-relay-failure", + ))), + &adapter, + ) + .await + .expect_err("total relay failure"); + + assert_eq!(error.code(), "product_sync_relay_setup_failure"); +} + +#[cfg(feature = "relay-runtime")] +#[tokio::test] +async fn trade_resync_reports_partial_relay_failure() { + let (_tempdir, sdk, _store) = directory_sdk_and_store_with_relays(&[RELAY, RELAY_B]).await; + let request_event = signed_raw_order_request_event("resync-partial-relay-failure", 43); + let adapter = RadrootsMockRelayFetchAdapter::new(vec![ + relay_closed(RELAY_B, "timeout: relay unavailable"), + relay_raw_event_item(&request_event, RELAY, 5_200), + relay_eose(RELAY), + ]); + + let resync = sdk + .trades() + .resync() + .resync_with_fetch_adapter( + TradeResyncRequest::new(RadrootsTradeLocator::from_order_id(order_id( + "resync-partial-relay-failure", + ))), + &adapter, + ) + .await + .expect("partial failure resync"); + + assert_eq!(resync.status.status, TradeStatusKind::Requested); + assert_eq!(resync.evidence.closed_count, 1); + assert_eq!(resync.evidence.eose_count, 1); + assert_eq!( + resync.evidence.relays[0].outcome_kind, + TradeResyncRelayOutcomeKind::Closed + ); + assert_eq!( + resync.evidence.relays[0].transport_outcome_kind, + Some(TradeResyncRelayTransportOutcomeKind::Timeout) + ); +} + +async fn relay_event_item_from_store( source: &RadrootsEventStore, - target: &RadrootsEventStore, event_id: &RadrootsEventId, + relay_url: &str, observed_at_ms: i64, -) { +) -> RadrootsRelayFetchItem { let stored = source .get_event(event_id.as_str()) .await .expect("source event lookup") .expect("source event"); - let event = stored_event_from_raw_json(stored.raw_json.as_str()); - let receipt = target - .ingest_event( - RadrootsEventIngest::new(event, observed_at_ms).with_raw_json(stored.raw_json), - ) - .await - .expect("target event ingest"); - assert!(receipt.inserted); - assert!( - target - .get_event(event_id.as_str()) - .await - .expect("target event lookup") - .is_some() - ); + RadrootsRelayFetchItem::Event { + relay_url: relay_url.to_owned(), + raw_json: stored.raw_json, + observed_at_ms, + } } -fn stored_event_from_raw_json(raw_json: &str) -> RadrootsNostrEvent { - serde_json::from_str::<nostr::Event>(raw_json) - .map(|event| radroots_event_from_nostr(&event)) - .or_else(|_| serde_json::from_str::<RadrootsNostrEvent>(raw_json)) - .expect("stored raw event json") +fn relay_raw_event_item( + event: &nostr::Event, + relay_url: &str, + observed_at_ms: i64, +) -> RadrootsRelayFetchItem { + RadrootsRelayFetchItem::Event { + relay_url: relay_url.to_owned(), + raw_json: event.as_json(), + observed_at_ms, + } +} + +fn relay_eose(relay_url: &str) -> RadrootsRelayFetchItem { + RadrootsRelayFetchItem::Eose { + relay_url: relay_url.to_owned(), + } +} + +fn relay_closed(relay_url: &str, message: &str) -> RadrootsRelayFetchItem { + RadrootsRelayFetchItem::Closed { + relay_url: relay_url.to_owned(), + message: message.to_owned(), + } +} + +fn relay_malformed(relay_url: &str) -> RadrootsRelayFetchItem { + RadrootsRelayFetchItem::Event { + relay_url: relay_url.to_owned(), + raw_json: "{".to_owned(), + observed_at_ms: 4_999, + } } #[cfg(all(feature = "signer-adapters", feature = "local-signer"))] @@ -2154,13 +2413,31 @@ fn revision_economics() -> RadrootsOrderEconomics { } } -fn signed_validation_receipt_event( +fn signed_raw_validation_receipt_event( raw_order_id: &str, listing_event_id: &RadrootsEventId, root_event_id: &RadrootsEventId, target_event_id: &RadrootsEventId, created_at: u32, -) -> RadrootsNostrEvent { +) -> nostr::Event { + signed_raw_event( + SERVICE_SECRET_KEY_HEX, + created_at, + validation_receipt_wire_parts( + raw_order_id, + listing_event_id, + root_event_id, + target_event_id, + ), + ) +} + +fn validation_receipt_wire_parts( + raw_order_id: &str, + listing_event_id: &RadrootsEventId, + root_event_id: &RadrootsEventId, + target_event_id: &RadrootsEventId, +) -> WireEventParts { let receipt = RadrootsTradeValidationReceipt { changed_records_root: hash32('6'), domain: "radroots.receipt".to_owned(), @@ -2187,8 +2464,7 @@ fn signed_validation_receipt_event( }, version: 1, }; - let parts = validation_receipt_event_build(raw_order_id, &receipt).expect("receipt event"); - signed_event(SERVICE_SECRET_KEY_HEX, created_at, parts) + validation_receipt_event_build(raw_order_id, &receipt).expect("receipt event") } async fn insert_perf_non_trade_events(store: &RadrootsEventStore, base: i64, count: i64) { @@ -2285,14 +2561,18 @@ fn signed_event( created_at: u32, parts: WireEventParts, ) -> RadrootsNostrEvent { + let event = signed_raw_event(secret_key_hex, created_at, parts); + radroots_event_from_nostr(&event) +} + +fn signed_raw_event(secret_key_hex: &str, created_at: u32, parts: WireEventParts) -> nostr::Event { let secret_key = RadrootsNostrSecretKey::from_hex(secret_key_hex).expect("secret key"); let keys = RadrootsNostrKeys::new(secret_key); - let event = radroots_nostr_build_event(parts.kind, parts.content, parts.tags) + radroots_nostr_build_event(parts.kind, parts.content, parts.tags) .expect("event builder") .custom_created_at(RadrootsNostrTimestamp::from_secs(u64::from(created_at))) .sign_with_keys(&keys) - .expect("signed event"); - radroots_event_from_nostr(&event) + .expect("signed event") } fn signed_order_request_event(raw_order_id: &str, created_at: u32) -> RadrootsNostrEvent { @@ -2304,6 +2584,15 @@ fn signed_order_request_event(raw_order_id: &str, created_at: u32) -> RadrootsNo signed_event(BUYER_SECRET_KEY_HEX, created_at, draft) } +fn signed_raw_order_request_event(raw_order_id: &str, created_at: u32) -> nostr::Event { + let draft = radroots_events_codec::order::order_request_event_build( + &listing_event_ptr(), + &order_request(raw_order_id), + ) + .expect("request draft"); + signed_raw_event(BUYER_SECRET_KEY_HEX, created_at, draft) +} + #[cfg(any())] fn request_event_ptr(event: &RadrootsNostrEvent) -> RadrootsNostrEventPtr { RadrootsNostrEventPtr { diff --git a/crates/sdk/tests/trade_public_api.rs b/crates/sdk/tests/trade_public_api.rs @@ -25,13 +25,13 @@ async fn grouped_trade_surface_is_the_public_product_entrypoint() { assert_eq!(status.status, TradeStatusKind::Missing); - let resync = trades + let resync_error = trades .resync() .resync(TradeResyncRequest::new(status.locator)) .await - .expect("resync"); + .expect_err("resync requires configured relays"); - assert_eq!(resync.status.status, TradeStatusKind::Missing); + assert_eq!(resync_error.code(), "empty_target_relays"); let seller_actor = RadrootsActorContext::test(SELLER_PUBLIC_KEY_HEX, [RadrootsActorRole::Seller])