lib

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

commit 6803dac83bfeb6b528ef07f7d156643ab0084b95
parent 0da30c9e42d3ac45365bb4c306129a9543f11bd4
Author: triesap <tyson@radroots.org>
Date:   Tue,  7 Jul 2026 03:45:14 +0000

transport: adopt daemon transport publish protocol

- Replace the SDK radrootsd adapter with the transport publish v2 protocol.
- Route proxy push completion through daemon target outcomes instead of relay aggregate quorum.
- Carry redacted proxy bearer auth through ProxyProfile without serializing secrets.
- Validate with cargo fmt, cargo check, cargo test, and exact retired-symbol scans.

Diffstat:
Mcrates/sdk/Cargo.toml | 10++++++----
Mcrates/sdk/src/adapters/radrootsd.rs | 178++++++++++++++++++++++++++++++++++++++++++++-----------------------------------
Mcrates/sdk/src/lib.rs | 2+-
Mcrates/sdk/src/sync_runtime.rs | 279++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------
Mcrates/sdk/src/transport.rs | 48+++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/sdk/tests/sync_runtime.rs | 96++++++++++++++++++++++++++++++++++++++++++++-----------------------------------
Mcrates/sdk/tests/trade_product_publish_runtime.rs | 37+++++++++++++++++++++----------------
Mcrates/sdk/tests/trade_public_api.rs | 6++++++
Mcrates/sdk/tests/unit/adapters_radrootsd_tests.rs | 190++++++++++++++++++++++++++++++++++++++++++++++---------------------------------
Mcrates/sdk/tests/unit/sync_runtime_tests.rs | 138++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------------
10 files changed, 663 insertions(+), 321 deletions(-)

diff --git a/crates/sdk/Cargo.toml b/crates/sdk/Cargo.toml @@ -42,11 +42,13 @@ radrootsd-proxy = [ "std", "serde_json", "dep:futures", - "dep:radroots_publish_proxy_protocol", + "dep:radroots_transport", + "dep:radroots_transport_publish_protocol", "dep:radroots_transport_nostr", "dep:reqwest", - "radroots_publish_proxy_protocol/serde", - "radroots_publish_proxy_protocol/std", + "radroots_transport/serde", + "radroots_transport_publish_protocol/serde", + "radroots_transport_publish_protocol/std", "radroots_transport_nostr/std", ] signer-adapters = [ @@ -123,8 +125,8 @@ radroots_events = { workspace = true, default-features = false } radroots_events_codec = { workspace = true, default-features = false } radroots_geocoder = { workspace = true, optional = true } radroots_outbox = { workspace = true, optional = true, default-features = false } -radroots_publish_proxy_protocol = { workspace = true, optional = true, default-features = false } radroots_transport = { workspace = true, optional = true, default-features = false } +radroots_transport_publish_protocol = { workspace = true, optional = true, default-features = false } radroots_transport_nostr = { workspace = true, optional = true, default-features = false } radroots_transport_reticulum = { workspace = true, optional = true, default-features = false } radroots_runtime_paths = { workspace = true, optional = true, default-features = false } diff --git a/crates/sdk/src/adapters/radrootsd.rs b/crates/sdk/src/adapters/radrootsd.rs @@ -2,22 +2,26 @@ use core::fmt; use core::time::Duration; use radroots_events::draft::RadrootsSignedNostrEvent; -use radroots_publish_proxy_protocol::{ - METHOD_EVENT, PublishDeliveryPolicy, PublishEventRequest, PublishEventResponse, - PublishProxyProtocolError, PublishRelayOutcomeKind, PublishRelayPolicy, SignedNostrEventWire, +use radroots_transport::{ + RadrootsTransportKind, RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget, }; -use radroots_transport::RadrootsTransportSatisfactionPolicy; use radroots_transport_nostr::{ RadrootsRelayOutcome, RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishReceipt, RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, RadrootsRelayTransportError, }; +use radroots_transport_publish_protocol::{ + METHOD_EVENT, SignedNostrEventWire, TransportPublishDeliveryPolicy, + TransportPublishEventRequest, TransportPublishEventResponse, TransportPublishOutcomeKind, + TransportPublishPreviewBehavior, TransportPublishProtocolError, TransportPublishTarget, + TransportPublishTargetOutcome, TransportPublishTargetPolicy, +}; use reqwest::header::{AUTHORIZATION, CONTENT_TYPE, HeaderMap, HeaderValue}; use serde::{Deserialize, Serialize, de::DeserializeOwned}; use serde_json::{Value, json}; -pub const SDK_RADROOTSD_PROXY_REQUEST_ID: &str = "radroots-sdk-publish-event"; -pub const SDK_RADROOTSD_PROXY_MAX_RELAYS: usize = 20; +pub const SDK_RADROOTSD_PROXY_REQUEST_ID: &str = "radroots-sdk-transport-publish-event"; +pub const SDK_RADROOTSD_PROXY_MAX_TARGETS: usize = 20; #[derive(Clone, PartialEq, Eq, Default, Serialize, Deserialize)] pub enum RadrootsdAuth { @@ -39,7 +43,6 @@ impl fmt::Debug for RadrootsdAuth { pub struct RadrootsdProxyConfig { pub endpoint: String, pub auth: RadrootsdAuth, - pub relay_policy: PublishRelayPolicy, pub timeout: Duration, pub request_timeout_ms: Option<u64>, } @@ -49,7 +52,6 @@ impl RadrootsdProxyConfig { Self { endpoint: endpoint.into(), auth: RadrootsdAuth::None, - relay_policy: PublishRelayPolicy::RequestThenAuthorWriteThenDaemonDefault, timeout: Duration::from_secs(10), request_timeout_ms: None, } @@ -60,11 +62,6 @@ impl RadrootsdProxyConfig { self } - pub fn with_relay_policy(mut self, relay_policy: PublishRelayPolicy) -> Self { - self.relay_policy = relay_policy; - self - } - pub fn with_timeout(mut self, timeout: Duration) -> Self { self.timeout = timeout; self @@ -93,19 +90,18 @@ impl RadrootsdProxyPublishAdapter { pub async fn publish_signed_event( &self, request: RadrootsdProxyPublishRequest, - ) -> Result<RadrootsRelayPublishReceipt, RadrootsdError> { - let request = request.into_protocol_request(self.config.relay_policy)?; + ) -> Result<TransportPublishEventResponse, RadrootsdError> { + let request = request.into_protocol_request(); request - .validate(SDK_RADROOTSD_PROXY_MAX_RELAYS) + .validate(SDK_RADROOTSD_PROXY_MAX_TARGETS) .map_err(RadrootsdError::from_protocol)?; - let response = publish_event( + publish_event( self.config.endpoint.as_str(), &self.config.auth, &request, self.config.timeout, ) - .await?; - proxy_receipt_from_response(response) + .await } } @@ -118,20 +114,30 @@ impl RadrootsRelayPublishAdapter for RadrootsdProxyPublishAdapter { Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>, > { Box::pin(async move { + let targets = request + .targets + .relay_strings() + .into_iter() + .map(|relay| RadrootsTransportTarget::new(RadrootsTransportKind::Nostr, relay)) + .collect::<Result<Vec<_>, _>>()?; let request = RadrootsdProxyPublishRequest { delivery_policy: delivery_policy_from_relay_request( - request.targets.len(), + targets.len(), &request.satisfaction_policy, )?, signed_event: request.signed_event, - relays: request.targets.relay_strings(), + target_policy: TransportPublishTargetPolicy::explicit_targets( + targets.iter().map(transport_publish_target).collect(), + ), idempotency_key: None, timeout_ms: self.config.request_timeout_ms, }; - let receipt = self + let response = self .publish_signed_event(request) .await .map_err(|error| RadrootsRelayTransportError::Transport(error.to_string()))?; + let receipt = proxy_relay_receipt_from_response(response) + .map_err(|error| RadrootsRelayTransportError::Transport(error.to_string()))?; Ok(receipt.relays) }) } @@ -140,25 +146,21 @@ impl RadrootsRelayPublishAdapter for RadrootsdProxyPublishAdapter { #[derive(Clone, Debug, PartialEq, Eq)] pub struct RadrootsdProxyPublishRequest { pub signed_event: RadrootsSignedNostrEvent, - pub relays: Vec<String>, - pub delivery_policy: PublishDeliveryPolicy, + pub target_policy: TransportPublishTargetPolicy, + pub delivery_policy: TransportPublishDeliveryPolicy, pub idempotency_key: Option<String>, pub timeout_ms: Option<u64>, } impl RadrootsdProxyPublishRequest { - fn into_protocol_request( - self, - relay_policy: PublishRelayPolicy, - ) -> Result<PublishEventRequest, RadrootsdError> { - Ok(PublishEventRequest { + fn into_protocol_request(self) -> TransportPublishEventRequest { + TransportPublishEventRequest { event: signed_event_wire(&self.signed_event), - relays: self.relays, - relay_policy, + target_policy: self.target_policy, delivery_policy: self.delivery_policy, idempotency_key: self.idempotency_key, timeout_ms: self.timeout_ms, - }) + } } } @@ -172,7 +174,7 @@ pub enum RadrootsdError { } impl RadrootsdError { - fn from_protocol(error: PublishProxyProtocolError) -> Self { + fn from_protocol(error: TransportPublishProtocolError) -> Self { Self::InvalidRequest(error.to_string()) } } @@ -212,9 +214,9 @@ struct JsonRpcError { pub async fn publish_event( endpoint: &str, auth: &RadrootsdAuth, - request: &PublishEventRequest, + request: &TransportPublishEventRequest, timeout: Duration, -) -> Result<PublishEventResponse, RadrootsdError> { +) -> Result<TransportPublishEventResponse, RadrootsdError> { jsonrpc_call( endpoint, auth, @@ -240,8 +242,10 @@ fn auth_headers(auth: &RadrootsdAuth) -> Result<HeaderMap, RadrootsdError> { } } -pub fn publish_event_request_json(request: &PublishEventRequest) -> Result<Value, RadrootsdError> { - Ok(serde_json::to_value(request).expect("radrootsd publish.event request serializes")) +pub fn publish_event_request_json( + request: &TransportPublishEventRequest, +) -> Result<Value, RadrootsdError> { + Ok(serde_json::to_value(request).expect("radrootsd transport publish request serializes")) } fn http_status_error(status: reqwest::StatusCode, body: &str) -> RadrootsdError { @@ -352,23 +356,35 @@ fn signed_event_wire(event: &RadrootsSignedNostrEvent) -> SignedNostrEventWire { } } +fn transport_publish_target(target: &RadrootsTransportTarget) -> TransportPublishTarget { + TransportPublishTarget { + transport_kind: target.kind.canonical_label(), + endpoint_uri: target.uri.as_str().to_owned(), + preview_behavior: if target.kind == RadrootsTransportKind::Reticulum { + Some(TransportPublishPreviewBehavior::RejectDeliveryAttempts) + } else { + None + }, + } +} + fn delivery_policy_from_relay_request( target_count: usize, satisfaction_policy: &RadrootsTransportSatisfactionPolicy, -) -> Result<PublishDeliveryPolicy, RadrootsRelayTransportError> { +) -> Result<TransportPublishDeliveryPolicy, RadrootsRelayTransportError> { let required = satisfaction_policy.required_target_count(target_count)?; let delivery_policy = if required >= target_count { - PublishDeliveryPolicy::All + TransportPublishDeliveryPolicy::All } else if required <= 1 { - PublishDeliveryPolicy::Any + TransportPublishDeliveryPolicy::Any } else { - PublishDeliveryPolicy::Quorum { quorum: required } + TransportPublishDeliveryPolicy::Quorum { quorum: required } }; Ok(delivery_policy) } -fn proxy_receipt_from_response( - response: PublishEventResponse, +fn proxy_relay_receipt_from_response( + response: TransportPublishEventResponse, ) -> Result<RadrootsRelayPublishReceipt, RadrootsdError> { response .job @@ -377,26 +393,15 @@ fn proxy_receipt_from_response( let quorum = response .job .delivery_policy - .required_ack_count(response.job.relay_count); - let attempted_count = response - .job - .relays - .iter() - .filter(|relay| relay.attempted) - .count(); + .required_target_count(response.job.target_count); let relays = response .job - .relays + .targets .into_iter() - .map(|relay| RadrootsRelayPublishRelayReceipt { - relay_url: relay.relay_url, - attempted: relay.attempted, - outcome: RadrootsRelayOutcome { - kind: relay_outcome_kind(relay.outcome_kind), - message: relay.message, - }, - }) - .collect(); + .filter(|target| target.transport_kind == "nostr") + .map(relay_receipt_from_target_outcome) + .collect::<Vec<_>>(); + let attempted_count = relays.iter().filter(|relay| relay.attempted).count(); Ok(RadrootsRelayPublishReceipt { event_id: response.job.event_id, attempted_count, @@ -409,27 +414,44 @@ fn proxy_receipt_from_response( }) } -fn relay_outcome_kind(kind: PublishRelayOutcomeKind) -> RadrootsRelayOutcomeKind { +fn relay_receipt_from_target_outcome( + target: TransportPublishTargetOutcome, +) -> RadrootsRelayPublishRelayReceipt { + RadrootsRelayPublishRelayReceipt { + relay_url: target.endpoint_uri, + attempted: target.attempted, + outcome: RadrootsRelayOutcome { + kind: relay_outcome_kind(target.outcome_kind), + message: target.message, + }, + } +} + +fn relay_outcome_kind(kind: TransportPublishOutcomeKind) -> RadrootsRelayOutcomeKind { match kind { - PublishRelayOutcomeKind::Accepted => RadrootsRelayOutcomeKind::Accepted, - PublishRelayOutcomeKind::DuplicateAccepted => RadrootsRelayOutcomeKind::DuplicateAccepted, - PublishRelayOutcomeKind::Blocked => RadrootsRelayOutcomeKind::Blocked, - PublishRelayOutcomeKind::RateLimited => RadrootsRelayOutcomeKind::RateLimited, - PublishRelayOutcomeKind::Invalid => RadrootsRelayOutcomeKind::Invalid, - PublishRelayOutcomeKind::PowRequired => RadrootsRelayOutcomeKind::PowRequired, - PublishRelayOutcomeKind::Restricted => RadrootsRelayOutcomeKind::Restricted, - PublishRelayOutcomeKind::AuthRequired => RadrootsRelayOutcomeKind::AuthRequired, - PublishRelayOutcomeKind::Muted => RadrootsRelayOutcomeKind::Muted, - PublishRelayOutcomeKind::Unsupported => RadrootsRelayOutcomeKind::Unsupported, - PublishRelayOutcomeKind::PaymentRequired => RadrootsRelayOutcomeKind::PaymentRequired, - PublishRelayOutcomeKind::Error => RadrootsRelayOutcomeKind::Error, - PublishRelayOutcomeKind::Timeout => RadrootsRelayOutcomeKind::Timeout, - PublishRelayOutcomeKind::ConnectionFailed => RadrootsRelayOutcomeKind::ConnectionFailed, - PublishRelayOutcomeKind::RelayUrlRejected => RadrootsRelayOutcomeKind::RelayUrlRejected, - PublishRelayOutcomeKind::SkippedAlreadyAccepted => { + TransportPublishOutcomeKind::Accepted => RadrootsRelayOutcomeKind::Accepted, + TransportPublishOutcomeKind::DuplicateAccepted => { + RadrootsRelayOutcomeKind::DuplicateAccepted + } + TransportPublishOutcomeKind::Blocked => RadrootsRelayOutcomeKind::Blocked, + TransportPublishOutcomeKind::RateLimited => RadrootsRelayOutcomeKind::RateLimited, + TransportPublishOutcomeKind::Invalid => RadrootsRelayOutcomeKind::Invalid, + TransportPublishOutcomeKind::PowRequired => RadrootsRelayOutcomeKind::PowRequired, + TransportPublishOutcomeKind::Restricted => RadrootsRelayOutcomeKind::Restricted, + TransportPublishOutcomeKind::AuthRequired => RadrootsRelayOutcomeKind::AuthRequired, + TransportPublishOutcomeKind::Muted => RadrootsRelayOutcomeKind::Muted, + TransportPublishOutcomeKind::Unsupported => RadrootsRelayOutcomeKind::Unsupported, + TransportPublishOutcomeKind::PaymentRequired => RadrootsRelayOutcomeKind::PaymentRequired, + TransportPublishOutcomeKind::Error => RadrootsRelayOutcomeKind::Error, + TransportPublishOutcomeKind::Timeout => RadrootsRelayOutcomeKind::Timeout, + TransportPublishOutcomeKind::ConnectionFailed => RadrootsRelayOutcomeKind::ConnectionFailed, + TransportPublishOutcomeKind::TargetRejected => RadrootsRelayOutcomeKind::RelayUrlRejected, + TransportPublishOutcomeKind::SkippedAlreadyAccepted => { RadrootsRelayOutcomeKind::SkippedAlreadyAccepted } - PublishRelayOutcomeKind::Unknown => RadrootsRelayOutcomeKind::Unknown, + TransportPublishOutcomeKind::Deferred + | TransportPublishOutcomeKind::Unavailable + | TransportPublishOutcomeKind::Unknown => RadrootsRelayOutcomeKind::Unknown, } } diff --git a/crates/sdk/src/lib.rs b/crates/sdk/src/lib.rs @@ -230,7 +230,7 @@ pub use crate::trade_storage::{ }; #[cfg(feature = "runtime")] pub use crate::transport::{ - HybridProfile, NostrProfile, NostrRelayUrlPolicy, ProxyProfile, PublishMode, + HybridProfile, NostrProfile, NostrRelayUrlPolicy, ProxyAuth, ProxyProfile, PublishMode, ReticulumPreviewBehavior, ReticulumPreviewProfile, SDK_TRANSPORT_TARGET_MAX_COUNT, SatisfactionPolicy, TargetPolicy, TargetSet, TransportDeliveryReceipt, TransportDeliveryTargetStatus, TransportKind, TransportOutcome, TransportProfile, diff --git a/crates/sdk/src/sync_runtime.rs b/crates/sdk/src/sync_runtime.rs @@ -1,11 +1,11 @@ #[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] use crate::adapters::radrootsd::{ - RadrootsdError, RadrootsdProxyConfig, RadrootsdProxyPublishAdapter, + RadrootsdAuth, RadrootsdError, RadrootsdProxyConfig, RadrootsdProxyPublishAdapter, RadrootsdProxyPublishRequest, }; #[cfg(feature = "runtime")] use crate::{ - NostrRelayUrlPolicy, RadrootsSdkError, SyncClient, + NostrRelayUrlPolicy, ProxyAuth, ProxyProfile, RadrootsSdkError, SyncClient, runtime::{RadrootsClient, sdk_now_ms}, transport::TransportProfile, }; @@ -22,8 +22,6 @@ use radroots_outbox::{ }; #[cfg(feature = "runtime")] use radroots_outbox::{RadrootsOutboxEventState, RadrootsOutboxStatusSummary}; -#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] -use radroots_publish_proxy_protocol::PublishDeliveryPolicy; #[cfg(feature = "runtime")] use radroots_trade::projection::{ RADROOTS_PRODUCT_PROJECTION_ID, RADROOTS_PRODUCT_PROJECTION_VERSION, @@ -41,6 +39,12 @@ use radroots_transport_nostr::{ RadrootsOutboxPublishPolicy, RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishReceipt, RadrootsRelayPublishRelayReceipt, publish_claimed_outbox_event, }; +#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] +use radroots_transport_publish_protocol::{ + NostrPublishTargetSourcePolicy, TransportPublishDeliveryPolicy, TransportPublishJobStatus, + TransportPublishJobView, TransportPublishOutcomeKind, TransportPublishTarget, + TransportPublishTargetOutcome, TransportPublishTargetPolicy, +}; #[cfg(feature = "runtime")] pub const PUSH_OUTBOX_DEFAULT_LIMIT: usize = 20; @@ -429,6 +433,8 @@ pub enum PushOutboxTargetOutcomeKind { ConnectionFailed, TargetUriRejected, SkippedAlreadyAccepted, + Deferred, + Unavailable, Unknown, } @@ -515,9 +521,8 @@ impl<'sdk> SyncClient<'sdk> { } #[cfg(feature = "radrootsd-proxy")] TransportProfile::Proxy { profile } => { - let adapter = RadrootsdProxyPublishAdapter::new(RadrootsdProxyConfig::new( - profile.endpoint_url().to_owned(), - )); + let adapter = + RadrootsdProxyPublishAdapter::new(radrootsd_proxy_config_from_profile(profile)); self.push_outbox_with_proxy_adapter(&adapter, request).await } TransportProfile::LocalOnly | TransportProfile::ReticulumPreview { .. } => { @@ -619,16 +624,23 @@ impl<'sdk> SyncClient<'sdk> { publish_now_ms, ) .await?; - receipt.push_event(push_event_receipt( - claimed.outbox_event_id, - push_event_final_state(&publish), - publish, - )); + receipt.push_event(push_proxy_event_receipt(claimed.outbox_event_id, publish)); } Ok(receipt) } } +#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] +fn radrootsd_proxy_config_from_profile(profile: &ProxyProfile) -> RadrootsdProxyConfig { + let config = RadrootsdProxyConfig::new(profile.endpoint_url().to_owned()); + match profile.auth() { + ProxyAuth::None => config, + ProxyAuth::BearerToken(token) => { + config.with_auth(RadrootsdAuth::BearerToken(token.to_owned())) + } + } +} + #[cfg(feature = "runtime")] async fn claim_ready_signed_event_for_push( sdk: &RadrootsClient, @@ -685,7 +697,7 @@ async fn push_proxy_claimed_outbox_event( claimed: &RadrootsOutboxClaimedEvent, next_attempt_delay_ms: i64, now_ms: i64, -) -> Result<RadrootsRelayPublishReceipt, RadrootsSdkError> { +) -> Result<TransportPublishJobView, RadrootsSdkError> { let signed_event = claimed.signed_event.clone().ok_or( radroots_transport_nostr::RadrootsRelayTransportError::MissingSignedOutboxEvent( claimed.outbox_event_id, @@ -700,14 +712,11 @@ async fn push_proxy_claimed_outbox_event( now_ms, ) .await?; - let relays = proxy_nostr_relay_targets(claimed) - .iter() - .map(|target| target.endpoint_uri.as_str().to_owned()) - .collect::<Vec<_>>(); + let target_policy = proxy_transport_publish_target_policy(claimed); let request = RadrootsdProxyPublishRequest { signed_event: signed_event.clone(), - delivery_policy: proxy_delivery_policy(sync, claimed, relays.len()).await?, - relays, + delivery_policy: proxy_delivery_policy(sync, claimed, &target_policy).await?, + target_policy, idempotency_key: Some(proxy_outbox_idempotency_key( claimed.outbox_event_id, claimed.attempt_count, @@ -716,7 +725,7 @@ async fn push_proxy_claimed_outbox_event( timeout_ms: adapter.config().request_timeout_ms, }; let publish = match adapter.publish_signed_event(request).await { - Ok(publish) => publish, + Ok(response) => response.job, Err(error) => { let message = proxy_error_message(&error); sync.sdk @@ -729,7 +738,7 @@ async fn push_proxy_claimed_outbox_event( now_ms, ) .await?; - return Ok(proxy_transport_error_receipt(signed_event.id)); + return Ok(proxy_transport_error_job(&signed_event)); } }; complete_proxy_publish_attempt(sync, claimed, &publish, next_attempt_delay_ms, now_ms).await?; @@ -740,8 +749,8 @@ async fn push_proxy_claimed_outbox_event( async fn proxy_delivery_policy( sync: &SyncClient<'_>, claimed: &RadrootsOutboxClaimedEvent, - target_count: usize, -) -> Result<PublishDeliveryPolicy, RadrootsSdkError> { + target_policy: &TransportPublishTargetPolicy, +) -> Result<TransportPublishDeliveryPolicy, RadrootsSdkError> { let plans = sync .sdk ._outbox @@ -762,23 +771,26 @@ async fn proxy_delivery_policy( claimed.outbox_event_id ), })?; - proxy_delivery_policy_from_satisfaction(target_count, &plan.satisfaction_policy) + proxy_delivery_policy_from_satisfaction( + target_policy.request_target_count(), + &plan.satisfaction_policy, + ) } #[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] fn proxy_delivery_policy_from_satisfaction( target_count: usize, satisfaction_policy: &RadrootsTransportSatisfactionPolicy, -) -> Result<PublishDeliveryPolicy, RadrootsSdkError> { +) -> Result<TransportPublishDeliveryPolicy, RadrootsSdkError> { if target_count == 0 { - return Ok(PublishDeliveryPolicy::Any); + return Ok(TransportPublishDeliveryPolicy::Any); } let required = satisfaction_policy.required_target_count(target_count)?; Ok(match satisfaction_policy { - RadrootsTransportSatisfactionPolicy::AnyTarget => PublishDeliveryPolicy::Any, - RadrootsTransportSatisfactionPolicy::AllTargets => PublishDeliveryPolicy::All, + RadrootsTransportSatisfactionPolicy::AnyTarget => TransportPublishDeliveryPolicy::Any, + RadrootsTransportSatisfactionPolicy::AllTargets => TransportPublishDeliveryPolicy::All, RadrootsTransportSatisfactionPolicy::AtLeast(_) => { - PublishDeliveryPolicy::Quorum { quorum: required } + TransportPublishDeliveryPolicy::Quorum { quorum: required } } }) } @@ -796,17 +808,19 @@ fn proxy_outbox_idempotency_key( async fn complete_proxy_publish_attempt( sync: &SyncClient<'_>, claimed: &RadrootsOutboxClaimedEvent, - publish: &RadrootsRelayPublishReceipt, + publish: &TransportPublishJobView, next_attempt_delay_ms: i64, now_ms: i64, ) -> Result<(), RadrootsSdkError> { let mut completed_target_ids = std::collections::BTreeSet::new(); - for relay in &publish.relays { - if let Some(target) = proxy_nostr_relay_targets(claimed) - .into_iter() - .find(|target| target.endpoint_uri.as_str() == relay.relay_url.as_str()) + for outcome in &publish.targets { + if let Some(target) = claimed + .delivery_targets + .iter() + .filter(|target| target.status.is_ready_for_attempt()) + .find(|target| proxy_target_matches_outcome(target, outcome)) { - complete_proxy_delivery_target(sync, claimed, target, relay, now_ms).await?; + complete_proxy_delivery_target(sync, claimed, target, outcome, now_ms).await?; completed_target_ids.insert(target.delivery_target_id); } } @@ -816,7 +830,7 @@ async fn complete_proxy_publish_attempt( .filter(|target| target.status.is_ready_for_attempt()) .filter(|target| !completed_target_ids.contains(&target.delivery_target_id)) { - complete_unmatched_proxy_delivery_target(sync, claimed, target, publish, now_ms).await?; + complete_missing_proxy_delivery_target(sync, claimed, target, publish, now_ms).await?; } sync.sdk ._outbox @@ -833,15 +847,58 @@ async fn complete_proxy_publish_attempt( } #[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] -fn proxy_nostr_relay_targets( +fn proxy_transport_publish_target_policy( claimed: &RadrootsOutboxClaimedEvent, -) -> Vec<&RadrootsOutboxDeliveryTargetRecord> { - claimed +) -> TransportPublishTargetPolicy { + let ready_targets = claimed .delivery_targets .iter() - .filter(|target| target.transport_kind == RadrootsTransportKind::Nostr) .filter(|target| target.status.is_ready_for_attempt()) - .collect() + .collect::<Vec<_>>(); + if ready_targets.len() == 1 && is_proxy_delegate_target(ready_targets[0]) { + TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + Vec::new(), + ) + } else { + TransportPublishTargetPolicy::explicit_targets( + ready_targets + .into_iter() + .map(transport_publish_target_from_outbox_target) + .collect(), + ) + } +} + +#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] +fn transport_publish_target_from_outbox_target( + target: &RadrootsOutboxDeliveryTargetRecord, +) -> TransportPublishTarget { + TransportPublishTarget { + transport_kind: target.transport_kind.canonical_label(), + endpoint_uri: target.endpoint_uri.as_str().to_owned(), + preview_behavior: if target.transport_kind == RadrootsTransportKind::Reticulum { + Some( + radroots_transport_publish_protocol::TransportPublishPreviewBehavior::RejectDeliveryAttempts, + ) + } else { + None + }, + } +} + +#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] +fn proxy_target_matches_outcome( + target: &RadrootsOutboxDeliveryTargetRecord, + outcome: &TransportPublishTargetOutcome, +) -> bool { + target.transport_kind.canonical_label() == outcome.transport_kind + && target.endpoint_uri.as_str() == outcome.endpoint_uri +} + +#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] +fn is_proxy_delegate_target(target: &RadrootsOutboxDeliveryTargetRecord) -> bool { + target.transport_kind.canonical_label() == "radrootsd_proxy" } #[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] @@ -849,10 +906,10 @@ async fn complete_proxy_delivery_target( sync: &SyncClient<'_>, claimed: &RadrootsOutboxClaimedEvent, target: &RadrootsOutboxDeliveryTargetRecord, - relay: &RadrootsRelayPublishRelayReceipt, + outcome: &TransportPublishTargetOutcome, now_ms: i64, ) -> Result<(), RadrootsSdkError> { - if relay.outcome.counts_toward_quorum() { + if outcome.outcome_kind.counts_toward_satisfaction() { sync.sdk ._outbox .mark_delivery_target_accepted( @@ -862,15 +919,14 @@ async fn complete_proxy_delivery_target( now_ms, ) .await?; - } else if relay.outcome.is_retryable() { + } else if outcome.outcome_kind.is_retryable() { sync.sdk ._outbox .mark_delivery_target_failed_retryable( claimed.outbox_event_id, claimed.claim_token.as_str(), target.delivery_target_id, - relay - .outcome + outcome .message .as_deref() .unwrap_or("radrootsd proxy publish retryable"), @@ -884,8 +940,7 @@ async fn complete_proxy_delivery_target( claimed.outbox_event_id, claimed.claim_token.as_str(), target.delivery_target_id, - relay - .outcome + outcome .message .as_deref() .unwrap_or("radrootsd proxy publish terminal"), @@ -897,14 +952,14 @@ async fn complete_proxy_delivery_target( } #[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] -async fn complete_unmatched_proxy_delivery_target( +async fn complete_missing_proxy_delivery_target( sync: &SyncClient<'_>, claimed: &RadrootsOutboxClaimedEvent, target: &RadrootsOutboxDeliveryTargetRecord, - publish: &RadrootsRelayPublishReceipt, + publish: &TransportPublishJobView, now_ms: i64, ) -> Result<(), RadrootsSdkError> { - if publish.quorum_met { + if is_proxy_delegate_target(target) && publish.delivery_satisfied { sync.sdk ._outbox .mark_delivery_target_accepted( @@ -915,6 +970,7 @@ async fn complete_unmatched_proxy_delivery_target( ) .await?; } else if publish.retryable_count > 0 + || !publish.terminal || target.status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable { sync.sdk @@ -943,16 +999,30 @@ async fn complete_unmatched_proxy_delivery_target( } #[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] -fn proxy_transport_error_receipt(event_id: String) -> RadrootsRelayPublishReceipt { - RadrootsRelayPublishReceipt { - event_id, - attempted_count: 1, - accepted_count: 0, +fn proxy_transport_error_job( + event: &radroots_events::draft::RadrootsSignedNostrEvent, +) -> TransportPublishJobView { + TransportPublishJobView { + job_id: "radroots-sdk-transport-error".to_owned(), + status: TransportPublishJobStatus::DeliveryUnsatisfiedRetryable, + terminal: false, + delivery_satisfied: false, + event_id: event.id.clone(), + pubkey: event.pubkey.clone(), + event_kind: event.kind, + target_policy: TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + Vec::new(), + ), + delivery_policy: TransportPublishDeliveryPolicy::Any, + target_count: 1, + acknowledged_count: 0, retryable_count: 1, terminal_count: 0, - quorum: 1, - quorum_met: false, - relays: Vec::new(), + requested_at_ms: 0, + completed_at_ms: None, + last_error: Some("radrootsd proxy publish failed".to_owned()), + targets: Vec::new(), } } @@ -961,6 +1031,97 @@ fn proxy_error_message(error: &RadrootsdError) -> String { format!("radrootsd proxy publish failed: {error}") } +#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] +fn proxy_push_event_final_state(publish: &TransportPublishJobView) -> PushOutboxEventState { + if publish.delivery_satisfied { + PushOutboxEventState::Published + } else if publish.retryable_count > 0 || !publish.terminal { + PushOutboxEventState::PublishRetryable + } else { + PushOutboxEventState::FailedTerminal + } +} + +#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] +fn push_proxy_event_receipt( + outbox_event_id: i64, + publish: TransportPublishJobView, +) -> PushOutboxEventReceipt { + let event_id = RadrootsEventId::parse(publish.event_id.as_str()) + .expect("transport publish daemon job uses signed event id"); + let quorum = publish + .delivery_policy + .required_target_count(publish.target_count); + PushOutboxEventReceipt { + event_id, + outbox_event_id, + final_state: proxy_push_event_final_state(&publish), + attempted_count: publish + .targets + .iter() + .filter(|target| target.attempted) + .count(), + accepted_count: publish.acknowledged_count, + retryable_count: publish.retryable_count, + terminal_count: publish.terminal_count, + quorum, + quorum_met: publish.delivery_satisfied, + targets: publish + .targets + .into_iter() + .map(push_proxy_target_receipt) + .collect(), + } +} + +#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] +fn push_proxy_target_receipt(outcome: TransportPublishTargetOutcome) -> PushOutboxTargetReceipt { + PushOutboxTargetReceipt { + transport_kind: outcome.transport_kind, + endpoint_uri: outcome.endpoint_uri, + outcome_kind: push_proxy_target_outcome_kind(outcome.outcome_kind), + attempted: outcome.attempted, + message: outcome.message, + } +} + +#[cfg(all(feature = "runtime", feature = "radrootsd-proxy"))] +fn push_proxy_target_outcome_kind( + outcome_kind: TransportPublishOutcomeKind, +) -> PushOutboxTargetOutcomeKind { + match outcome_kind { + TransportPublishOutcomeKind::Accepted => PushOutboxTargetOutcomeKind::Accepted, + TransportPublishOutcomeKind::DuplicateAccepted => { + PushOutboxTargetOutcomeKind::DuplicateAccepted + } + TransportPublishOutcomeKind::Blocked => PushOutboxTargetOutcomeKind::Blocked, + TransportPublishOutcomeKind::RateLimited => PushOutboxTargetOutcomeKind::RateLimited, + TransportPublishOutcomeKind::Invalid => PushOutboxTargetOutcomeKind::Invalid, + TransportPublishOutcomeKind::PowRequired => PushOutboxTargetOutcomeKind::PowRequired, + TransportPublishOutcomeKind::Restricted => PushOutboxTargetOutcomeKind::Restricted, + TransportPublishOutcomeKind::AuthRequired => PushOutboxTargetOutcomeKind::AuthRequired, + TransportPublishOutcomeKind::Muted => PushOutboxTargetOutcomeKind::Muted, + TransportPublishOutcomeKind::Unsupported => PushOutboxTargetOutcomeKind::Unsupported, + TransportPublishOutcomeKind::PaymentRequired => { + PushOutboxTargetOutcomeKind::PaymentRequired + } + TransportPublishOutcomeKind::Error => PushOutboxTargetOutcomeKind::Error, + TransportPublishOutcomeKind::Timeout => PushOutboxTargetOutcomeKind::Timeout, + TransportPublishOutcomeKind::ConnectionFailed => { + PushOutboxTargetOutcomeKind::ConnectionFailed + } + TransportPublishOutcomeKind::TargetRejected => { + PushOutboxTargetOutcomeKind::TargetUriRejected + } + TransportPublishOutcomeKind::SkippedAlreadyAccepted => { + PushOutboxTargetOutcomeKind::SkippedAlreadyAccepted + } + TransportPublishOutcomeKind::Deferred => PushOutboxTargetOutcomeKind::Deferred, + TransportPublishOutcomeKind::Unavailable => PushOutboxTargetOutcomeKind::Unavailable, + TransportPublishOutcomeKind::Unknown => PushOutboxTargetOutcomeKind::Unknown, + } +} + #[cfg(feature = "runtime")] fn push_outbox_claim_token() -> String { format!("radroots-sdk-sync-{}", uuid::Uuid::now_v7()) diff --git a/crates/sdk/src/transport.rs b/crates/sdk/src/transport.rs @@ -4,7 +4,7 @@ use radroots_transport::{ RadrootsTransportTargetFingerprint, RadrootsTransportTargetReceipt, RadrootsTransportTargetSet, }; use radroots_transport_nostr::{RadrootsRelayUrl, RadrootsRelayUrlPolicy}; -use serde::ser::SerializeStruct; +use serde::ser::{SerializeStruct, Serializer}; use std::collections::BTreeSet; pub use radroots_transport::{ @@ -358,22 +358,68 @@ impl HybridProfile { } } +#[derive(Clone, PartialEq, Eq)] +pub enum ProxyAuth { + None, + BearerToken(String), +} + +impl Default for ProxyAuth { + fn default() -> Self { + Self::None + } +} + +impl core::fmt::Debug for ProxyAuth { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + match self { + Self::None => f.write_str("None"), + Self::BearerToken(_) => f.write_str("BearerToken(<redacted>)"), + } + } +} + +impl serde::Serialize for ProxyAuth { + fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error> + where + S: Serializer, + { + let mut state = serializer.serialize_struct("ProxyAuth", 1)?; + match self { + Self::None => state.serialize_field("kind", "none")?, + Self::BearerToken(_) => state.serialize_field("kind", "bearer_token")?, + } + state.end() + } +} + #[derive(Clone, Debug, PartialEq, Eq, serde::Serialize)] pub struct ProxyProfile { endpoint_url: String, + auth: ProxyAuth, } impl ProxyProfile { pub fn new(endpoint_url: impl Into<String>) -> Self { Self { endpoint_url: endpoint_url.into(), + auth: ProxyAuth::None, } } + pub fn with_bearer_token(mut self, token: impl Into<String>) -> Self { + self.auth = ProxyAuth::BearerToken(token.into()); + self + } + pub fn endpoint_url(&self) -> &str { self.endpoint_url.as_str() } + pub fn auth(&self) -> &ProxyAuth { + &self.auth + } + pub(crate) fn target_set(&self) -> Result<TargetSet, RadrootsSdkError> { TargetSet::transport_targets(vec![RadrootsTransportTarget::new( RadrootsTransportKind::custom("radrootsd_proxy")?, diff --git a/crates/sdk/tests/sync_runtime.rs b/crates/sdk/tests/sync_runtime.rs @@ -65,13 +65,13 @@ struct FixtureSigner { struct TransportFailurePublishAdapter; #[cfg(feature = "radrootsd-proxy")] -struct RecordedProxyRequest { +struct RecordedTransportPublishRequest { body: String, } #[cfg(feature = "radrootsd-proxy")] #[derive(Clone, Copy)] -enum ProxyResponseMode { +enum TransportPublishResponseMode { Accepted, Retryable, Terminal, @@ -85,23 +85,28 @@ struct RecordingPublishAdapter { } #[cfg(feature = "radrootsd-proxy")] -fn spawn_publish_proxy_server() -> (String, JoinHandle<RecordedProxyRequest>) { - let listener = TcpListener::bind("127.0.0.1:0").expect("bind proxy server"); +fn spawn_transport_publish_server() -> (String, JoinHandle<RecordedTransportPublishRequest>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind transport publish server"); let endpoint = format!("http://{}/rpc", listener.local_addr().expect("addr")); let handle = std::thread::spawn(move || { let (mut stream, _) = listener.accept().expect("accept"); - let body = read_proxy_request_body(&mut stream); - write_proxy_response(&mut stream, body.as_str(), ProxyResponseMode::Accepted, 1); - RecordedProxyRequest { body } + let body = read_transport_publish_request_body(&mut stream); + write_transport_publish_response( + &mut stream, + body.as_str(), + TransportPublishResponseMode::Accepted, + 1, + ); + RecordedTransportPublishRequest { body } }); (endpoint, handle) } #[cfg(feature = "radrootsd-proxy")] -fn spawn_publish_proxy_sequence_server( - responses: Vec<ProxyResponseMode>, -) -> (String, JoinHandle<Vec<RecordedProxyRequest>>) { - let listener = TcpListener::bind("127.0.0.1:0").expect("bind proxy server"); +fn spawn_transport_publish_sequence_server( + responses: Vec<TransportPublishResponseMode>, +) -> (String, JoinHandle<Vec<RecordedTransportPublishRequest>>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind transport publish server"); let endpoint = format!("http://{}/rpc", listener.local_addr().expect("addr")); let handle = std::thread::spawn(move || { responses @@ -109,9 +114,9 @@ fn spawn_publish_proxy_sequence_server( .enumerate() .map(|(index, mode)| { let (mut stream, _) = listener.accept().expect("accept"); - let body = read_proxy_request_body(&mut stream); - write_proxy_response(&mut stream, body.as_str(), mode, index + 1); - RecordedProxyRequest { body } + let body = read_transport_publish_request_body(&mut stream); + write_transport_publish_response(&mut stream, body.as_str(), mode, index + 1); + RecordedTransportPublishRequest { body } }) .collect() }); @@ -119,7 +124,7 @@ fn spawn_publish_proxy_sequence_server( } #[cfg(feature = "radrootsd-proxy")] -fn read_proxy_request_body(stream: &mut TcpStream) -> String { +fn read_transport_publish_request_body(stream: &mut TcpStream) -> String { let mut request = Vec::new(); let mut buffer = [0u8; 1024]; loop { @@ -159,10 +164,10 @@ fn read_proxy_request_body(stream: &mut TcpStream) -> String { } #[cfg(feature = "radrootsd-proxy")] -fn write_proxy_response( +fn write_transport_publish_response( stream: &mut TcpStream, body: &str, - mode: ProxyResponseMode, + mode: TransportPublishResponseMode, job_number: usize, ) { let body_json: serde_json::Value = serde_json::from_str(body).expect("body json"); @@ -174,9 +179,9 @@ fn write_proxy_response( acknowledged_count, retryable_count, terminal_count, - relay, + target, ) = match mode { - ProxyResponseMode::Accepted => ( + TransportPublishResponseMode::Accepted => ( "delivery_satisfied", true, true, @@ -184,14 +189,15 @@ fn write_proxy_response( 0, 0, serde_json::json!({ - "relay_url": "wss://daemon-resolved.example.com", + "transport_kind": "nostr", + "endpoint_uri": "wss://daemon-resolved.example.com", "source": "daemon_default", "attempted": true, "outcome_kind": "accepted", "message": "accepted" }), ), - ProxyResponseMode::Retryable => ( + TransportPublishResponseMode::Retryable => ( "delivery_unsatisfied_retryable", false, false, @@ -199,14 +205,15 @@ fn write_proxy_response( 1, 0, serde_json::json!({ - "relay_url": "wss://daemon-resolved.example.com", + "transport_kind": "nostr", + "endpoint_uri": "wss://daemon-resolved.example.com", "source": "daemon_default", "attempted": false, "outcome_kind": "connection_failed", "message": "dns lookup failed" }), ), - ProxyResponseMode::Terminal => ( + TransportPublishResponseMode::Terminal => ( "delivery_unsatisfied_terminal", true, false, @@ -214,7 +221,8 @@ fn write_proxy_response( 0, 1, serde_json::json!({ - "relay_url": "wss://daemon-resolved.example.com", + "transport_kind": "nostr", + "endpoint_uri": "wss://daemon-resolved.example.com", "source": "daemon_default", "attempted": true, "outcome_kind": "invalid", @@ -235,15 +243,15 @@ fn write_proxy_response( "event_id": event["id"], "pubkey": event["pubkey"], "event_kind": event["kind"], - "relay_policy": body_json["params"]["relay_policy"], + "target_policy": body_json["params"]["target_policy"], "delivery_policy": body_json["params"]["delivery_policy"], - "relay_count": 1, + "target_count": 1, "acknowledged_count": acknowledged_count, "retryable_count": retryable_count, "terminal_count": terminal_count, "requested_at_ms": 1700000000000i64, "completed_at_ms": 1700000000100i64, - "relays": [relay] + "targets": [target] } } }) @@ -1446,7 +1454,7 @@ async fn push_outbox_empty_queue_returns_zero_counts() { #[cfg(feature = "radrootsd-proxy")] #[tokio::test] async fn product_push_outbox_uses_radrootsd_proxy_transport_with_daemon_resolved_relays() { - let (endpoint, handle) = spawn_publish_proxy_server(); + let (endpoint, handle) = spawn_transport_publish_server(); let tempdir = tempfile::tempdir().expect("tempdir"); let sdk = RadrootsClient::builder() .directory_storage(tempdir.path().join("sdk")) @@ -1473,7 +1481,7 @@ async fn product_push_outbox_uses_radrootsd_proxy_transport_with_daemon_resolved .sync() .push_outbox(PushOutboxRequest::new().with_limit(1)) .await - .expect("proxy push"); + .expect("transport publish push"); assert_eq!(receipt.attempted_events, 1); assert_eq!(receipt.published_events, 1); @@ -1492,12 +1500,16 @@ async fn product_push_outbox_uses_radrootsd_proxy_transport_with_daemon_resolved PushOutboxTargetOutcomeKind::Accepted ); - let recorded = handle.join().expect("proxy request"); + let recorded = handle.join().expect("transport publish request"); let body: serde_json::Value = serde_json::from_str(recorded.body.as_str()).expect("body"); - assert_eq!(body["method"], "publish.event"); - assert_eq!(body["params"]["relays"], serde_json::json!([])); + assert_eq!(body["method"], "transport.publish.event"); + assert_eq!(body["params"]["target_policy"]["kind"], "nostr"); + assert_eq!( + body["params"]["target_policy"]["relay_urls"], + serde_json::json!([]) + ); assert_eq!( - body["params"]["relay_policy"], + body["params"]["target_policy"]["source_policy"], "request_then_author_write_then_daemon_default" ); assert_eq!(body["params"]["delivery_policy"]["mode"], "any"); @@ -1517,9 +1529,9 @@ async fn product_push_outbox_uses_radrootsd_proxy_transport_with_daemon_resolved #[cfg(feature = "radrootsd-proxy")] #[tokio::test] async fn product_push_outbox_radrootsd_proxy_idempotency_is_attempt_scoped() { - let (endpoint, handle) = spawn_publish_proxy_sequence_server(vec![ - ProxyResponseMode::Retryable, - ProxyResponseMode::Accepted, + let (endpoint, handle) = spawn_transport_publish_sequence_server(vec![ + TransportPublishResponseMode::Retryable, + TransportPublishResponseMode::Accepted, ]); let tempdir = tempfile::tempdir().expect("tempdir"); let storage = tempdir.path().join("sdk"); @@ -1553,7 +1565,7 @@ async fn product_push_outbox_radrootsd_proxy_idempotency_is_attempt_scoped() { .with_next_attempt_delay_ms(1), ) .await - .expect("first proxy push"); + .expect("first transport publish push"); assert_eq!(first.attempted_events, 1); assert_eq!(first.retryable_events, 1); @@ -1575,7 +1587,7 @@ async fn product_push_outbox_radrootsd_proxy_idempotency_is_attempt_scoped() { .sync() .push_outbox(PushOutboxRequest::new().with_limit(1)) .await - .expect("second proxy push"); + .expect("second transport publish push"); assert_eq!(second.attempted_events, 1); assert_eq!(second.published_events, 1); @@ -1585,7 +1597,7 @@ async fn product_push_outbox_radrootsd_proxy_idempotency_is_attempt_scoped() { PushOutboxEventState::Published ); - let recorded = handle.join().expect("proxy requests"); + let recorded = handle.join().expect("transport publish requests"); assert_eq!(recorded.len(), 2); let first_body: serde_json::Value = serde_json::from_str(recorded[0].body.as_str()).expect("first body"); @@ -1647,7 +1659,7 @@ async fn product_push_outbox_radrootsd_proxy_error_and_terminal_paths_update_out .with_next_attempt_delay_ms(1), ) .await - .expect("retryable proxy push"); + .expect("retryable transport publish push"); assert_eq!(retryable.retryable_events, 1); assert_eq!( retryable.events[0].final_state, @@ -1669,7 +1681,7 @@ async fn product_push_outbox_radrootsd_proxy_error_and_terminal_paths_update_out ); let (terminal_endpoint, terminal_handle) = - spawn_publish_proxy_sequence_server(vec![ProxyResponseMode::Terminal]); + spawn_transport_publish_sequence_server(vec![TransportPublishResponseMode::Terminal]); let terminal_sdk = RadrootsClient::builder() .directory_storage(tempdir.path().join("terminal-sdk")) .fixed_clock(RadrootsSdkTimestamp::from_unix_seconds(1_700_000_000)) @@ -1696,7 +1708,7 @@ async fn product_push_outbox_radrootsd_proxy_error_and_terminal_paths_update_out .sync() .push_outbox(PushOutboxRequest::new().with_limit(1)) .await - .expect("terminal proxy push"); + .expect("terminal transport publish push"); assert_eq!(terminal.terminal_events, 1); assert_eq!(terminal.events[0].outbox_event_id, enqueue.outbox_event_id); assert_eq!( diff --git a/crates/sdk/tests/trade_product_publish_runtime.rs b/crates/sdk/tests/trade_product_publish_runtime.rs @@ -39,23 +39,23 @@ const SELLER_PUBLIC_KEY_HEX: &str = "e0266e3cfb0d2886f91c73f5f868f3b98273713e5fcd97c081663f5518a4b3af"; const RELAY: &str = "wss://relay.radroots.test"; -struct RecordedProxyRequest { +struct RecordedTransportPublishRequest { body: String, } -fn spawn_trade_publish_proxy_server() -> (String, JoinHandle<RecordedProxyRequest>) { - let listener = TcpListener::bind("127.0.0.1:0").expect("bind proxy server"); +fn spawn_trade_transport_publish_server() -> (String, JoinHandle<RecordedTransportPublishRequest>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind transport publish server"); let endpoint = format!("http://{}/rpc", listener.local_addr().expect("addr")); let handle = std::thread::spawn(move || { let (mut stream, _) = listener.accept().expect("accept"); - let body = read_proxy_request_body(&mut stream); - write_proxy_accept_response(&mut stream, body.as_str()); - RecordedProxyRequest { body } + let body = read_transport_publish_request_body(&mut stream); + write_transport_publish_accept_response(&mut stream, body.as_str()); + RecordedTransportPublishRequest { body } }); (endpoint, handle) } -fn read_proxy_request_body(stream: &mut TcpStream) -> String { +fn read_transport_publish_request_body(stream: &mut TcpStream) -> String { let mut request = Vec::new(); let mut buffer = [0u8; 1024]; loop { @@ -94,7 +94,7 @@ fn read_proxy_request_body(stream: &mut TcpStream) -> String { body.to_owned() } -fn write_proxy_accept_response(stream: &mut TcpStream, body: &str) { +fn write_transport_publish_accept_response(stream: &mut TcpStream, body: &str) { let body_json: serde_json::Value = serde_json::from_str(body).expect("body json"); let event = &body_json["params"]["event"]; let response_body = serde_json::json!({ @@ -110,16 +110,17 @@ fn write_proxy_accept_response(stream: &mut TcpStream, body: &str) { "event_id": event["id"], "pubkey": event["pubkey"], "event_kind": event["kind"], - "relay_policy": body_json["params"]["relay_policy"], + "target_policy": body_json["params"]["target_policy"], "delivery_policy": body_json["params"]["delivery_policy"], - "relay_count": 1, + "target_count": 1, "acknowledged_count": 1, "retryable_count": 0, "terminal_count": 0, "requested_at_ms": 1700000000000i64, "completed_at_ms": 1700000000100i64, - "relays": [{ - "relay_url": RELAY, + "targets": [{ + "transport_kind": "nostr", + "endpoint_uri": RELAY, "source": "request", "attempted": true, "outcome_kind": "accepted", @@ -231,7 +232,7 @@ fn explicit_trade_relays() -> TargetPolicy { #[tokio::test] async fn trade_product_propose_enqueue_and_publish_uses_ack_policy() { - let (endpoint, handle) = spawn_trade_publish_proxy_server(); + let (endpoint, handle) = spawn_trade_transport_publish_server(); let tempdir = tempfile::tempdir().expect("tempdir"); let secret_key = RadrootsNostrSecretKey::from_hex(BUYER_SECRET_KEY_HEX).expect("secret key"); let signer_keys = RadrootsNostrKeys::new(secret_key); @@ -278,9 +279,13 @@ async fn trade_product_propose_enqueue_and_publish_uses_ack_policy() { PushOutboxTargetOutcomeKind::Accepted ); - let recorded = handle.join().expect("proxy request"); + let recorded = handle.join().expect("transport publish request"); let body: serde_json::Value = serde_json::from_str(recorded.body.as_str()).expect("body"); - assert_eq!(body["method"], "publish.event"); + assert_eq!(body["method"], "transport.publish.event"); assert_eq!(body["params"]["delivery_policy"]["mode"], "any"); - assert_eq!(body["params"]["relays"], serde_json::json!([RELAY])); + assert_eq!(body["params"]["target_policy"]["kind"], "explicit_targets"); + assert_eq!( + body["params"]["target_policy"]["targets"][0]["endpoint_uri"], + RELAY + ); } diff --git a/crates/sdk/tests/trade_public_api.rs b/crates/sdk/tests/trade_public_api.rs @@ -32,7 +32,10 @@ async fn grouped_trade_surface_is_the_public_product_entrypoint() { .await .expect_err("resync requires configured relays"); + #[cfg(feature = "relay-runtime")] assert_eq!(resync_error.code(), "empty_transport_targets"); + #[cfg(not(feature = "relay-runtime"))] + assert_eq!(resync_error.code(), "product_sync_unsupported"); let validation_receipt_error = validation_receipts .list( @@ -42,7 +45,10 @@ async fn grouped_trade_surface_is_the_public_product_entrypoint() { .await .expect_err("validation receipt list requires configured relays"); + #[cfg(feature = "relay-runtime")] assert_eq!(validation_receipt_error.code(), "empty_transport_targets"); + #[cfg(not(feature = "relay-runtime"))] + assert_eq!(validation_receipt_error.code(), "product_sync_unsupported"); let seller_actor = RadrootsActorContext::test(SELLER_PUBLIC_KEY_HEX, [RadrootsActorRole::Seller]) diff --git a/crates/sdk/tests/unit/adapters_radrootsd_tests.rs b/crates/sdk/tests/unit/adapters_radrootsd_tests.rs @@ -1,11 +1,13 @@ use super::*; -use radroots_publish_proxy_protocol::{ - PublishJobStatus, PublishJobView, PublishRelayOutcome, PublishRelayOutcomeKind, - PublishRelaySource, -}; use radroots_transport_nostr::{ RadrootsRelayPublishRequest, RadrootsRelayTargetSet, RadrootsRelayUrlPolicy, }; +use radroots_transport_publish_protocol::{ + NostrPublishTargetSourcePolicy, TransportPublishDeliveryPolicy, TransportPublishEventRequest, + TransportPublishEventResponse, TransportPublishJobStatus, TransportPublishJobView, + TransportPublishOutcomeKind, TransportPublishTarget, TransportPublishTargetOutcome, + TransportPublishTargetPolicy, TransportPublishTargetSource, +}; use std::io::{Read, Write}; use std::net::TcpListener; use std::thread::JoinHandle; @@ -107,38 +109,44 @@ fn signed_event() -> RadrootsSignedNostrEvent { } } -fn publish_request() -> PublishEventRequest { - PublishEventRequest { +fn publish_request() -> TransportPublishEventRequest { + TransportPublishEventRequest { event: signed_event_wire(&signed_event()), - relays: vec!["wss://relay.example.com".to_owned()], - relay_policy: PublishRelayPolicy::RequestThenAuthorWriteThenDaemonDefault, - delivery_policy: PublishDeliveryPolicy::Any, + target_policy: TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + vec!["wss://relay.example.com".to_owned()], + ), + delivery_policy: TransportPublishDeliveryPolicy::Any, idempotency_key: Some("idem-1".to_owned()), timeout_ms: Some(10_000), } } -fn job(outcome_kind: PublishRelayOutcomeKind) -> PublishJobView { - PublishJobView { +fn job(outcome_kind: TransportPublishOutcomeKind) -> TransportPublishJobView { + TransportPublishJobView { job_id: "job-1".to_owned(), - status: PublishJobStatus::DeliverySatisfied, + status: TransportPublishJobStatus::DeliverySatisfied, terminal: true, delivery_satisfied: true, event_id: "a".repeat(64), pubkey: "b".repeat(64), event_kind: 30_402, - relay_policy: PublishRelayPolicy::RequestThenAuthorWriteThenDaemonDefault, - delivery_policy: PublishDeliveryPolicy::Any, - relay_count: 1, - acknowledged_count: usize::from(outcome_kind.counts_toward_quorum()), + target_policy: TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + vec!["wss://relay.example.com".to_owned()], + ), + delivery_policy: TransportPublishDeliveryPolicy::Any, + target_count: 1, + acknowledged_count: usize::from(outcome_kind.counts_toward_satisfaction()), retryable_count: usize::from(outcome_kind.is_retryable()), terminal_count: usize::from(outcome_kind.is_terminal_failure()), requested_at_ms: 1_700_000_000_000, completed_at_ms: Some(1_700_000_000_100), last_error: None, - relays: vec![PublishRelayOutcome { - relay_url: "wss://relay.example.com".to_owned(), - source: PublishRelaySource::Request, + targets: vec![TransportPublishTargetOutcome { + transport_kind: "nostr".to_owned(), + endpoint_uri: "wss://relay.example.com".to_owned(), + source: TransportPublishTargetSource::Request, attempted: true, outcome_kind, message: Some("relay outcome".to_owned()), @@ -153,7 +161,7 @@ fn publish_response_json() -> String { "id": SDK_RADROOTSD_PROXY_REQUEST_ID, "result": { "deduplicated": false, - "job": job(PublishRelayOutcomeKind::Accepted) + "job": job(TransportPublishOutcomeKind::Accepted) } }) .to_string() @@ -212,7 +220,6 @@ fn auth_headers_omit_or_redact_bearer_authorization() { fn proxy_config_builders_preserve_typed_runtime_options() { let config = RadrootsdProxyConfig::new("http://127.0.0.1:8080/rpc") .with_auth(RadrootsdAuth::BearerToken("sdk-token".to_owned())) - .with_relay_policy(PublishRelayPolicy::ExplicitOnly) .with_timeout(Duration::from_millis(250)) .with_request_timeout_ms(1_500); let adapter = RadrootsdProxyPublishAdapter::new(config.clone()); @@ -223,10 +230,6 @@ fn proxy_config_builders_preserve_typed_runtime_options() { adapter.config().auth, RadrootsdAuth::BearerToken("sdk-token".to_owned()) ); - assert_eq!( - adapter.config().relay_policy, - PublishRelayPolicy::ExplicitOnly - ); assert_eq!(adapter.config().timeout, Duration::from_millis(250)); assert_eq!(adapter.config().request_timeout_ms, Some(1_500)); } @@ -238,11 +241,15 @@ fn publish_event_request_json_uses_signed_event_contract() { assert_eq!(value["event"]["id"], "a".repeat(64)); assert_eq!(value["event"]["pubkey"], "b".repeat(64)); assert_eq!(value["event"]["kind"], 30_402); - assert_eq!(value["relays"][0], "wss://relay.example.com"); + assert_eq!(value["target_policy"]["kind"], "nostr"); assert_eq!( - value["relay_policy"], + value["target_policy"]["source_policy"], "request_then_author_write_then_daemon_default" ); + assert_eq!( + value["target_policy"]["relay_urls"][0], + "wss://relay.example.com" + ); assert_eq!(value["delivery_policy"]["mode"], "any"); assert_eq!(value["idempotency_key"], "idem-1"); let rendered = value.to_string(); @@ -252,7 +259,7 @@ fn publish_event_request_json_uses_signed_event_contract() { #[test] fn decode_jsonrpc_response_validates_envelope_and_errors() { - let response: PublishEventResponse = decode_jsonrpc_response( + let response: TransportPublishEventResponse = decode_jsonrpc_response( METHOD_EVENT, SDK_RADROOTSD_PROXY_REQUEST_ID, publish_response_json().as_str(), @@ -260,10 +267,10 @@ fn decode_jsonrpc_response_validates_envelope_and_errors() { .expect("response"); assert_eq!(response.job.event_id, "a".repeat(64)); - let error = decode_jsonrpc_response::<PublishEventResponse>( + let error = decode_jsonrpc_response::<TransportPublishEventResponse>( METHOD_EVENT, SDK_RADROOTSD_PROXY_REQUEST_ID, - r#"{"jsonrpc":"2.0","id":"radroots-sdk-publish-event","error":{"code":-32001,"message":"principal unauthorized"}}"#, + r#"{"jsonrpc":"2.0","id":"radroots-sdk-transport-publish-event","error":{"code":-32001,"message":"principal unauthorized"}}"#, ) .expect_err("jsonrpc error"); assert!(matches!( @@ -312,9 +319,9 @@ fn decode_jsonrpc_response_validates_envelope_and_errors() { #[test] fn daemon_outcomes_map_to_relay_transport_receipts() { - let payment = proxy_receipt_from_response(PublishEventResponse { + let payment = proxy_relay_receipt_from_response(TransportPublishEventResponse { deduplicated: false, - job: job(PublishRelayOutcomeKind::PaymentRequired), + job: job(TransportPublishOutcomeKind::PaymentRequired), }) .expect("payment receipt"); assert_eq!( @@ -323,9 +330,9 @@ fn daemon_outcomes_map_to_relay_transport_receipts() { ); assert_eq!(payment.terminal_count, 1); - let skipped = proxy_receipt_from_response(PublishEventResponse { + let skipped = proxy_relay_receipt_from_response(TransportPublishEventResponse { deduplicated: true, - job: job(PublishRelayOutcomeKind::SkippedAlreadyAccepted), + job: job(TransportPublishOutcomeKind::SkippedAlreadyAccepted), }) .expect("skipped receipt"); assert_eq!( @@ -336,68 +343,68 @@ fn daemon_outcomes_map_to_relay_transport_receipts() { let cases = [ ( - PublishRelayOutcomeKind::Accepted, + TransportPublishOutcomeKind::Accepted, RadrootsRelayOutcomeKind::Accepted, ), ( - PublishRelayOutcomeKind::DuplicateAccepted, + TransportPublishOutcomeKind::DuplicateAccepted, RadrootsRelayOutcomeKind::DuplicateAccepted, ), ( - PublishRelayOutcomeKind::Blocked, + TransportPublishOutcomeKind::Blocked, RadrootsRelayOutcomeKind::Blocked, ), ( - PublishRelayOutcomeKind::RateLimited, + TransportPublishOutcomeKind::RateLimited, RadrootsRelayOutcomeKind::RateLimited, ), ( - PublishRelayOutcomeKind::Invalid, + TransportPublishOutcomeKind::Invalid, RadrootsRelayOutcomeKind::Invalid, ), ( - PublishRelayOutcomeKind::PowRequired, + TransportPublishOutcomeKind::PowRequired, RadrootsRelayOutcomeKind::PowRequired, ), ( - PublishRelayOutcomeKind::Restricted, + TransportPublishOutcomeKind::Restricted, RadrootsRelayOutcomeKind::Restricted, ), ( - PublishRelayOutcomeKind::AuthRequired, + TransportPublishOutcomeKind::AuthRequired, RadrootsRelayOutcomeKind::AuthRequired, ), ( - PublishRelayOutcomeKind::Muted, + TransportPublishOutcomeKind::Muted, RadrootsRelayOutcomeKind::Muted, ), ( - PublishRelayOutcomeKind::Unsupported, + TransportPublishOutcomeKind::Unsupported, RadrootsRelayOutcomeKind::Unsupported, ), ( - PublishRelayOutcomeKind::Error, + TransportPublishOutcomeKind::Error, RadrootsRelayOutcomeKind::Error, ), ( - PublishRelayOutcomeKind::Timeout, + TransportPublishOutcomeKind::Timeout, RadrootsRelayOutcomeKind::Timeout, ), ( - PublishRelayOutcomeKind::ConnectionFailed, + TransportPublishOutcomeKind::ConnectionFailed, RadrootsRelayOutcomeKind::ConnectionFailed, ), ( - PublishRelayOutcomeKind::RelayUrlRejected, + TransportPublishOutcomeKind::TargetRejected, RadrootsRelayOutcomeKind::RelayUrlRejected, ), ( - PublishRelayOutcomeKind::Unknown, + TransportPublishOutcomeKind::Unknown, RadrootsRelayOutcomeKind::Unknown, ), ]; for (proxy_kind, relay_kind) in cases { - let receipt = proxy_receipt_from_response(PublishEventResponse { + let receipt = proxy_relay_receipt_from_response(TransportPublishEventResponse { deduplicated: false, job: job(proxy_kind), }) @@ -407,7 +414,7 @@ fn daemon_outcomes_map_to_relay_transport_receipts() { } #[tokio::test] -async fn publish_event_posts_publish_proxy_jsonrpc() { +async fn publish_event_posts_transport_publish_jsonrpc() { let (endpoint, handle) = spawn_http_server("200 OK", publish_response_json().as_str()); let receipt = publish_event( @@ -432,7 +439,11 @@ async fn publish_event_posts_publish_proxy_jsonrpc() { assert_eq!(body["method"], METHOD_EVENT); assert_eq!(body["id"], SDK_RADROOTSD_PROXY_REQUEST_ID); assert_eq!(body["params"]["event"]["content"], "{\"name\":\"carrots\"}"); - assert_eq!(body["params"]["relays"][0], "wss://relay.example.com"); + assert_eq!(body["params"]["target_policy"]["kind"], "nostr"); + assert_eq!( + body["params"]["target_policy"]["relay_urls"][0], + "wss://relay.example.com" + ); } #[tokio::test] @@ -447,15 +458,17 @@ async fn publish_signed_event_posts_typed_proxy_request() { let receipt = adapter .publish_signed_event(RadrootsdProxyPublishRequest { signed_event: signed_event(), - relays: vec!["wss://relay.example.com".to_owned()], - delivery_policy: PublishDeliveryPolicy::All, + target_policy: TransportPublishTargetPolicy::explicit_targets(vec![ + TransportPublishTarget::nostr("wss://relay.example.com"), + ]), + delivery_policy: TransportPublishDeliveryPolicy::All, idempotency_key: Some("idem-typed".to_owned()), timeout_ms: adapter.config().request_timeout_ms, }) .await .expect("typed publish"); - assert!(receipt.quorum_met); + assert!(receipt.job.delivery_satisfied); let recorded = handle.join().expect("server thread"); assert!( recorded @@ -465,6 +478,11 @@ async fn publish_signed_event_posts_typed_proxy_request() { ); let body: serde_json::Value = serde_json::from_str(recorded.body.as_str()).expect("body"); assert_eq!(body["params"]["delivery_policy"]["mode"], "all"); + assert_eq!(body["params"]["target_policy"]["kind"], "explicit_targets"); + assert_eq!( + body["params"]["target_policy"]["targets"][0]["endpoint_uri"], + "wss://relay.example.com" + ); assert_eq!(body["params"]["idempotency_key"], "idem-typed"); assert_eq!(body["params"]["timeout_ms"], 7_000); } @@ -512,17 +530,17 @@ async fn relay_publish_adapter_derives_delivery_policy_and_timeout() { ( 2, radroots_transport::RadrootsTransportSatisfactionPolicy::AllTargets, - PublishDeliveryPolicy::All, + TransportPublishDeliveryPolicy::All, ), ( 2, radroots_transport::RadrootsTransportSatisfactionPolicy::AnyTarget, - PublishDeliveryPolicy::Any, + TransportPublishDeliveryPolicy::Any, ), ( 3, radroots_transport::RadrootsTransportSatisfactionPolicy::AtLeast(2), - PublishDeliveryPolicy::Quorum { quorum: 2 }, + TransportPublishDeliveryPolicy::Quorum { quorum: 2 }, ), ] { let response_body = publish_response_json(); @@ -550,7 +568,7 @@ async fn relay_publish_adapter_derives_delivery_policy_and_timeout() { serde_json::from_str(recorded.body.as_str()).expect("request body"); assert_eq!(body["params"]["timeout_ms"], 4_000); assert_eq!( - serde_json::from_value::<PublishDeliveryPolicy>( + serde_json::from_value::<TransportPublishDeliveryPolicy>( body["params"]["delivery_policy"].clone() ) .expect("delivery policy"), @@ -592,8 +610,11 @@ async fn publish_signed_event_rejects_invalid_protocol_requests_before_http() { RadrootsdProxyPublishAdapter::new(RadrootsdProxyConfig::new("http://127.0.0.1:9/rpc")); let base = RadrootsdProxyPublishRequest { signed_event: signed_event(), - relays: vec!["wss://relay.example.com".to_owned()], - delivery_policy: PublishDeliveryPolicy::Any, + target_policy: TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + vec!["wss://relay.example.com".to_owned()], + ), + delivery_policy: TransportPublishDeliveryPolicy::Any, idempotency_key: Some("idem-1".to_owned()), timeout_ms: Some(1_000), }; @@ -603,13 +624,19 @@ async fn publish_signed_event_rejects_invalid_protocol_requests_before_http() { let mut empty_event_tag = base.clone(); empty_event_tag.signed_event.tags = vec![Vec::new()]; let mut invalid_quorum = base.clone(); - invalid_quorum.delivery_policy = PublishDeliveryPolicy::Quorum { quorum: 0 }; - let mut too_many_relays = base.clone(); - too_many_relays.relays = (0..=SDK_RADROOTSD_PROXY_MAX_RELAYS) - .map(|index| format!("wss://relay-{index}.example.com")) - .collect(); - let mut empty_relay = base.clone(); - empty_relay.relays = vec![" ".to_owned()]; + invalid_quorum.delivery_policy = TransportPublishDeliveryPolicy::Quorum { quorum: 0 }; + let mut too_many_targets = base.clone(); + too_many_targets.target_policy = TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + (0..=SDK_RADROOTSD_PROXY_MAX_TARGETS) + .map(|index| format!("wss://relay-{index}.example.com")) + .collect(), + ); + let mut empty_endpoint_uri = base.clone(); + empty_endpoint_uri.target_policy = TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + vec![" ".to_owned()], + ); let mut empty_idempotency = base; empty_idempotency.idempotency_key = Some(" ".to_owned()); @@ -617,8 +644,8 @@ async fn publish_signed_event_rejects_invalid_protocol_requests_before_http() { invalid_event_kind, empty_event_tag, invalid_quorum, - too_many_relays, - empty_relay, + too_many_targets, + empty_endpoint_uri, empty_idempotency, ] { assert!(matches!( @@ -629,17 +656,17 @@ async fn publish_signed_event_rejects_invalid_protocol_requests_before_http() { } #[test] -fn proxy_receipt_from_response_rejects_invalid_daemon_job_contracts() { - let mut empty_job_id = job(PublishRelayOutcomeKind::Accepted); +fn proxy_relay_receipt_from_response_rejects_invalid_daemon_job_contracts() { + let mut empty_job_id = job(TransportPublishOutcomeKind::Accepted); empty_job_id.job_id = " ".to_owned(); - let mut invalid_event_id = job(PublishRelayOutcomeKind::Accepted); + let mut invalid_event_id = job(TransportPublishOutcomeKind::Accepted); invalid_event_id.event_id = "not-an-event-id".to_owned(); - let mut invalid_pubkey = job(PublishRelayOutcomeKind::Accepted); + let mut invalid_pubkey = job(TransportPublishOutcomeKind::Accepted); invalid_pubkey.pubkey = "not-a-pubkey".to_owned(); - let mut invalid_kind = job(PublishRelayOutcomeKind::Accepted); + let mut invalid_kind = job(TransportPublishOutcomeKind::Accepted); invalid_kind.event_kind = 70_000; - let mut invalid_quorum = job(PublishRelayOutcomeKind::Accepted); - invalid_quorum.delivery_policy = PublishDeliveryPolicy::Quorum { quorum: 0 }; + let mut invalid_quorum = job(TransportPublishOutcomeKind::Accepted); + invalid_quorum.delivery_policy = TransportPublishDeliveryPolicy::Quorum { quorum: 0 }; for job in [ empty_job_id, @@ -649,7 +676,7 @@ fn proxy_receipt_from_response_rejects_invalid_daemon_job_contracts() { invalid_quorum, ] { assert!(matches!( - proxy_receipt_from_response(PublishEventResponse { + proxy_relay_receipt_from_response(TransportPublishEventResponse { deduplicated: false, job, }), @@ -664,8 +691,11 @@ async fn adapter_rejects_invalid_request_before_transport() { RadrootsdProxyPublishAdapter::new(RadrootsdProxyConfig::new("http://127.0.0.1:9/rpc")); let mut request = RadrootsdProxyPublishRequest { signed_event: signed_event(), - relays: Vec::new(), - delivery_policy: PublishDeliveryPolicy::Quorum { quorum: 0 }, + target_policy: TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + Vec::new(), + ), + delivery_policy: TransportPublishDeliveryPolicy::Quorum { quorum: 0 }, idempotency_key: None, timeout_ms: None, }; diff --git a/crates/sdk/tests/unit/sync_runtime_tests.rs b/crates/sdk/tests/unit/sync_runtime_tests.rs @@ -1,7 +1,7 @@ #[cfg(feature = "radrootsd-proxy")] use super::{ CLAIM_OWNER, complete_proxy_publish_attempt, proxy_delivery_policy_from_satisfaction, - proxy_error_message, proxy_outbox_idempotency_key, proxy_transport_error_receipt, + proxy_error_message, proxy_outbox_idempotency_key, proxy_transport_error_job, push_proxy_claimed_outbox_event, }; use super::{ @@ -39,13 +39,17 @@ use radroots_nostr::prelude::{ use radroots_outbox::RadrootsOutboxClaimedEvent; use radroots_outbox::{RadrootsOutboxEventState, RadrootsOutboxStatusSummary}; #[cfg(feature = "radrootsd-proxy")] -use radroots_publish_proxy_protocol::PublishDeliveryPolicy; -#[cfg(feature = "radrootsd-proxy")] use radroots_transport::RadrootsTransportSatisfactionPolicy; use radroots_transport_nostr::{ RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishReceipt, RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, RadrootsRelayTransportError, }; +#[cfg(feature = "radrootsd-proxy")] +use radroots_transport_publish_protocol::{ + NostrPublishTargetSourcePolicy, TransportPublishDeliveryPolicy, TransportPublishJobStatus, + TransportPublishJobView, TransportPublishOutcomeKind, TransportPublishTargetOutcome, + TransportPublishTargetPolicy, TransportPublishTargetSource, +}; use std::collections::BTreeSet; #[cfg(feature = "radrootsd-proxy")] @@ -168,6 +172,49 @@ async fn claimed_proxy_event(d_tag: &str) -> (crate::RadrootsClient, RadrootsOut (sdk, claimed) } +#[cfg(feature = "radrootsd-proxy")] +fn proxy_job(event_id: &str, outcome_kind: TransportPublishOutcomeKind) -> TransportPublishJobView { + let delivery_satisfied = outcome_kind.counts_toward_satisfaction(); + let retryable = outcome_kind.is_retryable(); + let terminal_failure = outcome_kind.is_terminal_failure(); + TransportPublishJobView { + job_id: "proxy-unit-job".to_owned(), + status: if delivery_satisfied { + TransportPublishJobStatus::DeliverySatisfied + } else if retryable { + TransportPublishJobStatus::DeliveryUnsatisfiedRetryable + } else { + TransportPublishJobStatus::DeliveryUnsatisfiedTerminal + }, + terminal: !retryable, + delivery_satisfied, + event_id: event_id.to_owned(), + pubkey: PROXY_SIGNER_PUBLIC_KEY_HEX.to_owned(), + event_kind: KIND_FARM, + target_policy: TransportPublishTargetPolicy::nostr( + NostrPublishTargetSourcePolicy::RequestThenAuthorWriteThenDaemonDefault, + vec!["wss://relay.example.com".to_owned()], + ), + delivery_policy: TransportPublishDeliveryPolicy::Any, + target_count: 1, + acknowledged_count: usize::from(delivery_satisfied), + retryable_count: usize::from(retryable), + terminal_count: usize::from(terminal_failure), + requested_at_ms: 1_700_000_000_000, + completed_at_ms: Some(1_700_000_000_100), + last_error: None, + targets: vec![TransportPublishTargetOutcome { + transport_kind: "nostr".to_owned(), + endpoint_uri: "wss://relay.example.com".to_owned(), + source: TransportPublishTargetSource::Request, + attempted: true, + outcome_kind, + message: Some("daemon outcome".to_owned()), + latency_ms: Some(4), + }], + } +} + #[test] fn push_outbox_claim_tokens_are_unique_under_immediate_generation() { let mut tokens = BTreeSet::new(); @@ -501,7 +548,7 @@ async fn proxy_push_empty_queue_and_private_helpers_are_deterministic() { .sync() .push_outbox_with_proxy_adapter(&adapter, super::PushOutboxRequest::new()) .await - .expect("empty proxy push"); + .expect("empty transport publish push"); assert_eq!(receipt.attempted_events, 0); assert_eq!( @@ -510,7 +557,7 @@ async fn proxy_push_empty_queue_and_private_helpers_are_deterministic() { &RadrootsTransportSatisfactionPolicy::AllTargets ) .expect("zero-target proxy policy"), - PublishDeliveryPolicy::Any + TransportPublishDeliveryPolicy::Any ); assert_eq!( proxy_delivery_policy_from_satisfaction( @@ -518,24 +565,27 @@ async fn proxy_push_empty_queue_and_private_helpers_are_deterministic() { &RadrootsTransportSatisfactionPolicy::AllTargets ) .expect("all-target proxy policy"), - PublishDeliveryPolicy::All + TransportPublishDeliveryPolicy::All ); assert_eq!( proxy_delivery_policy_from_satisfaction(2, &RadrootsTransportSatisfactionPolicy::AnyTarget) .expect("any-target proxy policy"), - PublishDeliveryPolicy::Any + TransportPublishDeliveryPolicy::Any ); assert_eq!( proxy_outbox_idempotency_key(7, 3, "event-id"), "radroots-sdk-outbox-7-3-event-id" ); - let proxy_receipt = proxy_transport_error_receipt("a".repeat(64)); - assert_eq!(proxy_receipt.attempted_count, 1); - assert_eq!(proxy_receipt.retryable_count, 1); - assert_eq!(proxy_receipt.quorum, 1); - assert!(!proxy_receipt.quorum_met); - assert!(proxy_receipt.relays.is_empty()); + let signed_event = ProxyFixtureSigner::new() + .sign_frozen_draft(&proxy_frozen_draft("proxy-transport-error-job")) + .expect("signed event"); + let proxy_job = proxy_transport_error_job(&signed_event); + assert_eq!(proxy_job.event_id, signed_event.id); + assert_eq!(proxy_job.target_count, 1); + assert_eq!(proxy_job.retryable_count, 1); + assert!(!proxy_job.delivery_satisfied); + assert!(proxy_job.targets.is_empty()); assert_eq!( proxy_error_message(&RadrootsdError::Http("connection refused".to_owned())), "radrootsd proxy publish failed: connection refused" @@ -628,7 +678,7 @@ async fn proxy_claim_publish_marks_retryable_transport_errors() { let receipt = push_proxy_claimed_outbox_event(&sync, &adapter, &claimed, 60_000, 1_700_000_000_000) .await - .expect("transport error receipt"); + .expect("transport error job"); assert_eq!(receipt.retryable_count, 1); let stored = sdk @@ -645,43 +695,51 @@ async fn proxy_claim_publish_marks_retryable_transport_errors() { #[tokio::test] async fn proxy_completion_updates_outbox_for_success_retryable_and_terminal_receipts() { let cases = [ - ("proxy-complete-success", PushOutboxEventState::Published, { - let mut receipt = relay_publish_receipt("a".repeat(64).as_str()); - receipt.quorum_met = true; - receipt.quorum = 1; - receipt.accepted_count = 1; - receipt - }), + ( + "proxy-complete-success", + PushOutboxEventState::Published, + TransportPublishOutcomeKind::Accepted, + ), ( "proxy-complete-retryable", PushOutboxEventState::PublishRetryable, - { - let mut receipt = relay_publish_receipt("b".repeat(64).as_str()); - receipt.retryable_count = 1; - receipt.quorum = 1; - receipt - }, + TransportPublishOutcomeKind::Timeout, ), ( "proxy-complete-terminal", PushOutboxEventState::FailedTerminal, - { - let mut receipt = relay_publish_receipt("c".repeat(64).as_str()); - receipt.terminal_count = 1; - receipt.quorum = 1; - receipt - }, + TransportPublishOutcomeKind::Blocked, ), ]; - for (d_tag, expected_state, mut publish) in cases { + for (d_tag, expected_state, outcome_kind) in cases { let (sdk, claimed) = claimed_proxy_event(d_tag).await; - publish.event_id = claimed - .signed_event - .as_ref() - .expect("signed event") - .id - .clone(); + let publish = proxy_job( + claimed + .signed_event + .as_ref() + .expect("signed event") + .id + .as_str(), + outcome_kind, + ); + assert_eq!( + publish.event_id, + claimed + .signed_event + .as_ref() + .expect("signed event") + .id + .clone() + ); + assert_eq!( + publish.pubkey, + claimed.signed_event.as_ref().expect("signed event").pubkey + ); + assert_eq!( + publish.event_kind, + claimed.signed_event.as_ref().expect("signed event").kind + ); let sync = sdk.sync(); complete_proxy_publish_attempt(&sync, &claimed, &publish, 60_000, 1_700_000_000_000) .await