commit cd32f1cb5160e54d77bd296966e4c1050a24cf11
parent 8ebaa2445f1d6db1d691775eef6be454116524f1
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:
12 files changed, 672 insertions(+), 330 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -1966,13 +1966,6 @@ dependencies = [
]
[[package]]
-name = "radroots_publish_proxy_protocol"
-version = "0.1.0-alpha.2"
-dependencies = [
- "serde",
-]
-
-[[package]]
name = "radroots_replica_db"
version = "0.1.0-alpha.2"
dependencies = [
@@ -2102,7 +2095,6 @@ dependencies = [
"radroots_nostr_connect",
"radroots_nostr_signer",
"radroots_outbox",
- "radroots_publish_proxy_protocol",
"radroots_replica_db",
"radroots_replica_db_schema",
"radroots_replica_sync",
@@ -2111,6 +2103,7 @@ dependencies = [
"radroots_trade",
"radroots_transport",
"radroots_transport_nostr",
+ "radroots_transport_publish_protocol",
"radroots_transport_reticulum",
"reqwest",
"serde",
@@ -2227,6 +2220,13 @@ dependencies = [
]
[[package]]
+name = "radroots_transport_publish_protocol"
+version = "0.1.0-alpha.2"
+dependencies = [
+ "serde",
+]
+
+[[package]]
name = "radroots_transport_reticulum"
version = "0.1.0-alpha.2"
dependencies = [
diff --git a/Cargo.toml b/Cargo.toml
@@ -43,8 +43,8 @@ radroots_nostr = { path = "../lib/crates/nostr", version = "0.1.0-alpha.2", defa
radroots_nostr_connect = { path = "../lib/crates/nostr_connect", version = "0.1.0-alpha.2", default-features = false }
radroots_nostr_signer = { path = "../lib/crates/nostr_signer", version = "0.1.0-alpha.2", default-features = false }
radroots_outbox = { path = "../lib/crates/outbox", version = "0.1.0-alpha.2", default-features = false }
-radroots_publish_proxy_protocol = { path = "../lib/crates/publish_proxy_protocol", version = "0.1.0-alpha.2", default-features = false }
radroots_transport = { path = "../lib/crates/transport", version = "0.1.0-alpha.2", default-features = false }
+radroots_transport_publish_protocol = { path = "../lib/crates/transport_publish_protocol", version = "0.1.0-alpha.2", default-features = false }
radroots_transport_nostr = { path = "../lib/crates/transport_nostr", version = "0.1.0-alpha.2", default-features = false }
radroots_transport_reticulum = { path = "../lib/crates/transport_reticulum", version = "0.1.0-alpha.2", default-features = false }
radroots_replica_db = { path = "../lib/crates/replica_db", version = "0.1.0-alpha.2", default-features = false }
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