commit 7ef2e7dc401aa2019583e480f8e589d18ef3d558
parent 4d03f12208b77e506bf4b655d5e0a4f8bdfbf264
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 06:25:38 +0000
transport-nostr: implement the Nostr event sink
- implement EventSink over a signer-free live Nostr client with one bounded attempt per configured relay
- convert canonical signed events and generic targets explicitly without storage or outbox ownership
- normalize acceptance, duplicate, rejection, auth, rate-limit, timeout, and connection receipts per relay
- verify package tests, strict Clippy, architecture, dependency, API, and source-maintenance gates
Diffstat:
4 files changed, 511 insertions(+), 8 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -274,6 +274,37 @@ dependencies = [
]
[[package]]
+name = "async-utility"
+version = "0.3.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "188f83b9a198af8c336e505611edb00d6d2ac5c694241c5a4f9a12316938cfe9"
+dependencies = [
+ "futures-util",
+ "gloo-timers",
+ "tokio",
+ "wasm-bindgen-futures",
+]
+
+[[package]]
+name = "async-wsocket"
+version = "0.13.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "1c92385c7c8b3eb2de1b78aeca225212e4c9a69a78b802832759b108681a5069"
+dependencies = [
+ "async-utility",
+ "futures",
+ "futures-util",
+ "js-sys",
+ "tokio",
+ "tokio-rustls",
+ "tokio-socks",
+ "tokio-tungstenite",
+ "url",
+ "wasm-bindgen",
+ "web-sys",
+]
+
+[[package]]
name = "atoi"
version = "2.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -292,6 +323,12 @@ dependencies = [
]
[[package]]
+name = "atomic-destructor"
+version = "0.3.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "ef49f5882e4b6afaac09ad239a4f8c70a24b8f2b0897edb1f706008efd109cf4"
+
+[[package]]
name = "atomic-waker"
version = "1.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2190,6 +2227,18 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0cc23270f6e1808e30a928bdc84dea0b9b4136a8bc82338574f23baf47bbd280"
[[package]]
+name = "gloo-timers"
+version = "0.3.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "bbb143cf96099802033e0d4f4963b19fd2e0b728bcf076cd9cf7f6634f092994"
+dependencies = [
+ "futures-channel",
+ "futures-core",
+ "js-sys",
+ "wasm-bindgen",
+]
+
+[[package]]
name = "group"
version = "0.13.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -2501,7 +2550,7 @@ dependencies = [
"tokio",
"tokio-rustls",
"tower-service",
- "webpki-roots",
+ "webpki-roots 1.0.6",
]
[[package]]
@@ -3130,6 +3179,12 @@ dependencies = [
]
[[package]]
+name = "lru"
+version = "0.16.4"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "7f66e8d5d03f609abc3a39e6f08e4164ebf1447a732906d39eb9b99b7919ef39"
+
+[[package]]
name = "lru-slab"
version = "0.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -3278,6 +3333,12 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1d87ecb2933e8aeadb3e3a02b828fed80a7528047e68b4f424523a0981a3a084"
[[package]]
+name = "negentropy"
+version = "0.5.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "f0efe882e02d206d8d279c20eb40e03baf7cb5136a1476dc084a324fbc3ec42d"
+
+[[package]]
name = "nom"
version = "7.1.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -3314,6 +3375,59 @@ dependencies = [
]
[[package]]
+name = "nostr-database"
+version = "0.44.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "7462c9d8ae5ef6a28d66a192d399ad2530f1f2130b13186296dbb11bdef5b3d1"
+dependencies = [
+ "lru 0.16.4",
+ "nostr",
+ "tokio",
+]
+
+[[package]]
+name = "nostr-gossip"
+version = "0.44.0"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "ade30de16869618919c6b5efc8258f47b654a98b51541eb77f85e8ec5e3c83a6"
+dependencies = [
+ "nostr",
+]
+
+[[package]]
+name = "nostr-relay-pool"
+version = "0.44.3"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "c85c54d6ca9aae4ae2bf19a7663ba9db5f45f783f1d24aff55f006386b8b99a1"
+dependencies = [
+ "async-utility",
+ "async-wsocket",
+ "atomic-destructor",
+ "hex",
+ "lru 0.16.4",
+ "negentropy",
+ "nostr",
+ "nostr-database",
+ "tokio",
+ "tracing",
+]
+
+[[package]]
+name = "nostr-sdk"
+version = "0.44.1"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "471732576710e779b64f04c55e3f8b5292f865fea228436daf19694f0bf70393"
+dependencies = [
+ "async-utility",
+ "nostr",
+ "nostr-database",
+ "nostr-gossip",
+ "nostr-relay-pool",
+ "tokio",
+ "tracing",
+]
+
+[[package]]
name = "nostrdb"
version = "0.9.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -5124,6 +5238,8 @@ dependencies = [
name = "radroots_transport_nostr"
version = "0.1.0-alpha"
dependencies = [
+ "futures",
+ "nostr-sdk",
"radroots_event_codec",
"radroots_nostr",
"radroots_protocol",
@@ -5438,7 +5554,7 @@ dependencies = [
"wasm-bindgen-futures",
"wasm-streams",
"web-sys",
- "webpki-roots",
+ "webpki-roots 1.0.6",
]
[[package]]
@@ -5948,6 +6064,17 @@ dependencies = [
]
[[package]]
+name = "sha1"
+version = "0.10.7"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "a978451301f4db1d02937a4ab3ccce137717b81826e79b7d49ffe3244a13c3b8"
+dependencies = [
+ "cfg-if",
+ "cpufeatures 0.2.17",
+ "digest 0.10.7",
+]
+
+[[package]]
name = "sha1_smol"
version = "1.0.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -6843,7 +6970,7 @@ dependencies = [
"hex",
"indicatif",
"itertools 0.14.0",
- "lru",
+ "lru 0.12.5",
"mti",
"num-bigint 0.4.6",
"opentelemetry",
@@ -7640,6 +7767,18 @@ dependencies = [
]
[[package]]
+name = "tokio-socks"
+version = "0.5.3"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "a7e2948f60dbe26b35f2c7fb74ac2854c1fddded0fe9d7548fcc674a246f7615"
+dependencies = [
+ "either",
+ "futures-util",
+ "thiserror 1.0.69",
+ "tokio",
+]
+
+[[package]]
name = "tokio-stream"
version = "0.1.18"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -7651,6 +7790,22 @@ dependencies = [
]
[[package]]
+name = "tokio-tungstenite"
+version = "0.26.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "7a9daff607c6d2bf6c16fd681ccb7eecc83e4e2cdc1ca067ffaadfca5de7f084"
+dependencies = [
+ "futures-util",
+ "log",
+ "rustls",
+ "rustls-pki-types",
+ "tokio",
+ "tokio-rustls",
+ "tungstenite",
+ "webpki-roots 0.26.11",
+]
+
+[[package]]
name = "tokio-util"
version = "0.7.18"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -7960,6 +8115,25 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b"
[[package]]
+name = "tungstenite"
+version = "0.26.2"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "4793cb5e56680ecbb1d843515b23b6de9a75eb04b66643e256a396d43be33c13"
+dependencies = [
+ "bytes",
+ "data-encoding",
+ "http",
+ "httparse",
+ "log",
+ "rand 0.9.2",
+ "rustls",
+ "rustls-pki-types",
+ "sha1",
+ "thiserror 2.0.18",
+ "utf-8",
+]
+
+[[package]]
name = "twox-hash"
version = "2.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
@@ -8351,6 +8525,15 @@ dependencies = [
[[package]]
name = "webpki-roots"
+version = "0.26.11"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "521bc38abb08001b01866da9f51eb7c5d647a19260e00054a8c7fd5f9e57f7a9"
+dependencies = [
+ "webpki-roots 1.0.6",
+]
+
+[[package]]
+name = "webpki-roots"
version = "1.0.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "22cfaf3c063993ff62e73cb4311efde4db1efb31ab78a3e5c457939ad5cc0bed"
diff --git a/crates/transport_nostr/Cargo.toml b/crates/transport_nostr/Cargo.toml
@@ -15,11 +15,15 @@ readme = "README.md"
name = "radroots_transport_nostr"
[dependencies]
-radroots_event_codec = { workspace = true, default-features = false }
-radroots_nostr = { workspace = true, default-features = false }
+radroots_event_codec = { workspace = true, default-features = false, features = ["json"] }
+radroots_nostr = { workspace = true, default-features = false, features = ["events"] }
radroots_protocol = { workspace = true, default-features = false }
radroots_transport = { workspace = true, default-features = false }
+nostr-sdk = { workspace = true }
url = { workspace = true }
+[dev-dependencies]
+futures = { workspace = true }
+
[lints]
workspace = true
diff --git a/crates/transport_nostr/src/client.rs b/crates/transport_nostr/src/client.rs
@@ -1,7 +1,9 @@
//! Concrete Nostr transport composition.
use crate::{Error, RelayUrl, RelayUrlPolicy};
+use core::fmt;
use std::collections::BTreeSet;
+use std::sync::Arc;
/// Maximum relay targets accepted by one transport instance.
pub(crate) const MAX_RELAYS: usize = 64;
@@ -120,21 +122,39 @@ fn validate_timeout(field: &'static str, value_ms: u64) -> Result<(), Error> {
}
/// Concrete Nostr implementation of the transport source and sink SPIs.
-#[derive(Clone, Debug)]
+#[derive(Clone)]
pub struct NostrTransport {
config: Config,
+ pub(crate) client: Arc<dyn crate::sink::RelayClient>,
}
impl NostrTransport {
/// Creates an inert transport from validated explicit configuration.
- pub const fn new(config: Config) -> Self {
- Self { config }
+ pub fn new(config: Config) -> Self {
+ Self {
+ config,
+ client: Arc::new(crate::sink::LiveRelayClient),
+ }
}
/// Returns the transport configuration.
pub const fn config(&self) -> &Config {
&self.config
}
+
+ #[cfg(test)]
+ pub(crate) fn with_client(config: Config, client: Arc<dyn crate::sink::RelayClient>) -> Self {
+ Self { config, client }
+ }
+}
+
+impl fmt::Debug for NostrTransport {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("NostrTransport")
+ .field("config", &self.config)
+ .finish_non_exhaustive()
+ }
}
#[cfg(test)]
diff --git a/crates/transport_nostr/src/sink.rs b/crates/transport_nostr/src/sink.rs
@@ -1 +1,297 @@
//! Nostr implementation of the transport event sink.
+
+use crate::{NostrTransport, RelayUrl};
+use core::time::Duration;
+use radroots_nostr::event::Event;
+use radroots_transport::{
+ BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink,
+ capability::{Availability, Maturity, SinkCapabilities},
+ outcome::{DeliveryOutcome, Retryability},
+ sink::{DeliveryTargetReceipt, SinkStatus},
+};
+use std::collections::{BTreeMap, BTreeSet};
+
+#[derive(Clone, Debug)]
+pub(crate) struct RelayPublishResult {
+ relay: RelayUrl,
+ outcome: DeliveryOutcome,
+}
+
+pub(crate) trait RelayClient: Send + Sync {
+ fn publish<'a>(
+ &'a self,
+ relays: Vec<RelayUrl>,
+ event: Event,
+ connect_timeout: Duration,
+ ) -> BoxFuture<'a, Vec<RelayPublishResult>>;
+}
+
+#[derive(Debug)]
+pub(crate) struct LiveRelayClient;
+
+impl RelayClient for LiveRelayClient {
+ fn publish<'a>(
+ &'a self,
+ relays: Vec<RelayUrl>,
+ event: Event,
+ connect_timeout: Duration,
+ ) -> BoxFuture<'a, Vec<RelayPublishResult>> {
+ Box::pin(async move {
+ let client = nostr_sdk::Client::default();
+ let mut results = Vec::with_capacity(relays.len());
+ for relay in relays {
+ let url = relay.as_str().to_owned();
+ let outcome = match client.add_relay(url.as_str()).await {
+ Err(error) => connection_failure(error.to_string()),
+ Ok(_) => match client
+ .try_connect_relay(url.as_str(), connect_timeout)
+ .await
+ {
+ Err(error) => connection_failure(error.to_string()),
+ Ok(()) => match client.send_event_to([url.as_str()], &event).await {
+ Err(error) => classify_failure(error.to_string()),
+ Ok(output) => output
+ .success
+ .iter()
+ .any(|accepted| accepted.to_string().trim_end_matches('/') == url)
+ .then(DeliveryOutcome::accepted)
+ .or_else(|| {
+ output.failed.iter().find_map(|(failed, message)| {
+ (failed.to_string().trim_end_matches('/') == url)
+ .then(|| classify_failure(message.clone()))
+ })
+ })
+ .unwrap_or_else(|| classify_failure("relay omitted result".into())),
+ },
+ },
+ };
+ results.push(RelayPublishResult { relay, outcome });
+ }
+ results
+ })
+ }
+}
+
+impl EventSink for NostrTransport {
+ fn status(&self) -> BoxFuture<'_, Result<SinkStatus, radroots_transport::Error>> {
+ let configured = !self.config().relays().is_empty();
+ Box::pin(async move {
+ Ok(SinkStatus::new(
+ radroots_transport::TransportId::NOSTR,
+ configured,
+ Maturity::Preview,
+ Availability::Available,
+ SinkCapabilities::DELIVER,
+ "bounded Nostr event delivery configured",
+ ))
+ })
+ }
+
+ fn deliver(
+ &self,
+ request: DeliveryRequest,
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::Error>> {
+ Box::pin(async move {
+ let mut requested = Vec::new();
+ let mut skipped = Vec::new();
+ for target in request.target_set().targets() {
+ match RelayUrl::from_target(target, self.config().relay_url_policy()) {
+ Ok(relay) if self.config().relays().contains(&relay) => {
+ requested.push((relay, target.clone()));
+ }
+ _ => skipped.push(DeliveryTargetReceipt::skipped(
+ target.clone(),
+ DeliveryOutcome::rejected().with_detail(
+ "target_denied",
+ "target is not configured for this sink",
+ )?,
+ )?),
+ }
+ }
+
+ let event = match radroots_nostr::event::to_nostr(request.payload().event().envelope())
+ {
+ Ok(event) => event,
+ Err(_) => return Err(radroots_transport::Error::InvalidDeliveryOutcome),
+ };
+ let expected: BTreeSet<_> = requested.iter().map(|(relay, _)| relay.clone()).collect();
+ let results = self
+ .client
+ .publish(
+ requested.iter().map(|(relay, _)| relay.clone()).collect(),
+ event,
+ Duration::from_millis(self.config().connect_timeout_ms()),
+ )
+ .await;
+ let mut by_relay = BTreeMap::new();
+ for result in results {
+ if expected.contains(&result.relay)
+ && by_relay.insert(result.relay, result.outcome).is_none()
+ {
+ continue;
+ }
+ return Err(radroots_transport::Error::InvalidDeliveryOutcome);
+ }
+
+ let mut receipts = skipped;
+ for (relay, target) in requested {
+ let outcome = by_relay.remove(&relay).unwrap_or_else(|| {
+ DeliveryOutcome::unavailable()
+ .with_detail("missing_result", "relay returned no result")
+ .expect("static normalized outcome")
+ });
+ receipts.push(DeliveryTargetReceipt::attempted(target, outcome));
+ }
+ DeliveryReceipt::for_request(&request, receipts)
+ })
+ }
+}
+
+fn connection_failure(message: String) -> DeliveryOutcome {
+ normalized(DeliveryOutcome::unavailable(), "connection_failed", message)
+}
+
+fn classify_failure(message: String) -> DeliveryOutcome {
+ let normalized_message = message.to_ascii_lowercase();
+ if normalized_message.contains("duplicate") || normalized_message.contains("already have") {
+ normalized(DeliveryOutcome::accepted(), "duplicate", message)
+ } else if normalized_message.contains("auth") {
+ normalized(
+ DeliveryOutcome::failed(Retryability::Retryable)
+ .expect("retryable failure classification"),
+ "auth_required",
+ message,
+ )
+ } else if normalized_message.contains("blocked")
+ || normalized_message.contains("restricted")
+ || normalized_message.contains("invalid")
+ {
+ normalized(DeliveryOutcome::rejected(), "rejected", message)
+ } else if normalized_message.contains("rate") {
+ normalized(DeliveryOutcome::unavailable(), "rate_limited", message)
+ } else if normalized_message.contains("timeout") {
+ normalized(DeliveryOutcome::unavailable(), "timeout", message)
+ } else {
+ normalized(DeliveryOutcome::unavailable(), "relay_failure", message)
+ }
+}
+
+fn normalized(outcome: DeliveryOutcome, code: &'static str, message: String) -> DeliveryOutcome {
+ let clean = message
+ .chars()
+ .filter(|character| !character.is_control())
+ .take(1_024)
+ .collect::<String>();
+ let clean = if clean.trim().is_empty() {
+ "relay operation failed".to_owned()
+ } else {
+ clean.trim().to_owned()
+ };
+ outcome
+ .with_detail(code, clean)
+ .expect("normalized static relay outcome")
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use crate::{Config, RelayUrlPolicy};
+ use radroots_transport::{
+ Target, TargetSet,
+ outcome::DeliveryOutcomeKind,
+ policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ sink::DeliveryPayload,
+ };
+ use std::sync::Arc;
+
+ #[derive(Debug)]
+ struct MockRelayClient {
+ outcomes: BTreeMap<RelayUrl, DeliveryOutcome>,
+ }
+
+ impl RelayClient for MockRelayClient {
+ fn publish<'a>(
+ &'a self,
+ relays: Vec<RelayUrl>,
+ _event: Event,
+ _connect_timeout: Duration,
+ ) -> BoxFuture<'a, Vec<RelayPublishResult>> {
+ Box::pin(async move {
+ relays
+ .into_iter()
+ .map(|relay| RelayPublishResult {
+ outcome: self
+ .outcomes
+ .get(&relay)
+ .cloned()
+ .unwrap_or_else(DeliveryOutcome::accepted),
+ relay,
+ })
+ .collect()
+ })
+ }
+ }
+
+ fn payload() -> DeliveryPayload {
+ let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#;
+ DeliveryPayload::new(radroots_event_codec::decode::signed_event(raw).expect("signed event"))
+ }
+
+ fn request() -> DeliveryRequest {
+ DeliveryRequest::new(
+ "nostr-delivery",
+ payload(),
+ TargetSet::new(vec![
+ Target::nostr_relay("wss://one.example").expect("one"),
+ Target::nostr_relay("wss://two.example").expect("two"),
+ ])
+ .expect("targets"),
+ SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
+ 1_800_000_000_000,
+ )
+ .expect("request")
+ }
+
+ #[test]
+ fn sink_returns_normalized_per_relay_partial_success() {
+ let config = Config::new(
+ RelayUrlPolicy::Public,
+ ["wss://one.example", "wss://two.example"],
+ )
+ .expect("config");
+ let two = RelayUrl::parse("wss://two.example", RelayUrlPolicy::Public).expect("two");
+ let client = MockRelayClient {
+ outcomes: BTreeMap::from([(two, classify_failure("rate limited".to_owned()))]),
+ };
+ let transport = NostrTransport::with_client(config, Arc::new(client));
+ let request = request();
+ let receipt = futures::executor::block_on(transport.deliver(request.clone()))
+ .expect("delivery receipt");
+
+ assert_eq!(receipt.target_receipts().len(), 2);
+ assert_eq!(
+ receipt.target_receipts()[0].outcome().kind(),
+ DeliveryOutcomeKind::Accepted
+ );
+ assert_eq!(
+ receipt.target_receipts()[1].outcome().kind(),
+ DeliveryOutcomeKind::Unavailable
+ );
+ assert!(!receipt.is_satisfied(&request).expect("satisfaction"));
+ }
+
+ #[test]
+ fn upstream_messages_map_to_stable_outcomes() {
+ let cases = [
+ ("duplicate: already have", DeliveryOutcomeKind::Accepted),
+ ("blocked by policy", DeliveryOutcomeKind::Rejected),
+ ("rate limited", DeliveryOutcomeKind::Unavailable),
+ ("auth required", DeliveryOutcomeKind::Failed),
+ ("connection timeout", DeliveryOutcomeKind::Unavailable),
+ ("connection failed", DeliveryOutcomeKind::Unavailable),
+ ];
+ for (message, expected) in cases {
+ assert_eq!(classify_failure(message.to_owned()).kind(), expected);
+ }
+ }
+}