lib

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

commit f86dede27f1905cb7135085957b61ba1825c069b
parent 9ad72b6c05b72ddbab7d1f51728bd2abcbdd07f0
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 06:12:03 +0000

transport-nostr: align package manifest and module root

- align the package with its final pre-release identity and exact Radroots dependency boundary
- replace the legacy storage-owning surface with the approved private module skeleton and curated exports
- remove obsolete feature requests from transitional workspace consumers so Cargo metadata resolves
- verify format, package check and tests, strict Clippy, architecture, dependency, API, and source-maintenance gates

Diffstat:
MCargo.lock | 201++-----------------------------------------------------------------------------
Mcrates/net/Cargo.toml | 4+---
Mcrates/nostr_runtime/Cargo.toml | 4+---
Mcrates/transport_nostr/Cargo.toml | 53+++--------------------------------------------------
Acrates/transport_nostr/src/auth.rs | 1+
Mcrates/transport_nostr/src/client.rs | 616++-----------------------------------------------------------------------------
Mcrates/transport_nostr/src/error.rs | 189++++++-------------------------------------------------------------------------
Mcrates/transport_nostr/src/lib.rs | 62++++++++------------------------------------------------------
Mcrates/transport_nostr/src/relay.rs | 448+++----------------------------------------------------------------------------
Acrates/transport_nostr/src/sink.rs | 1+
Acrates/transport_nostr/src/source.rs | 1+
Acrates/transport_nostr/src/status.rs | 1+
Acrates/transport_nostr/tests/package_boundary.rs | 69+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Dcrates/transport_nostr/tests/transport.rs | 4985-------------------------------------------------------------------------------
14 files changed, 132 insertions(+), 6503 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -274,37 +274,6 @@ dependencies = [ ] [[package]] -name = "async-utility" -version = "0.3.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a34a3b57207a7a1007832416c3e4862378c8451b4e8e093e436f48c2d3d2c151" -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" @@ -323,12 +292,6 @@ 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" @@ -2227,18 +2190,6 @@ 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" @@ -2550,7 +2501,7 @@ dependencies = [ "tokio", "tokio-rustls", "tower-service", - "webpki-roots 1.0.6", + "webpki-roots", ] [[package]] @@ -3179,12 +3130,6 @@ dependencies = [ ] [[package]] -name = "lru" -version = "0.16.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a1dc47f592c06f33f8e3aea9591776ec7c9f9e4124778ff8a3c3b87159f7e593" - -[[package]] name = "lru-slab" version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -3333,12 +3278,6 @@ 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" @@ -3375,59 +3314,6 @@ 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.3", - "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.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4b1073ccfbaea5549fb914a9d52c68dab2aecda61535e5143dd73e95445a804b" -dependencies = [ - "async-utility", - "async-wsocket", - "atomic-destructor", - "hex", - "lru 0.16.3", - "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" @@ -5238,20 +5124,10 @@ dependencies = [ name = "radroots_transport_nostr" version = "0.1.0-alpha" dependencies = [ - "futures", - "nostr", - "nostr-sdk", - "radroots_event", - "radroots_event_store", + "radroots_event_codec", "radroots_nostr", - "radroots_outbox", + "radroots_protocol", "radroots_transport", - "reqwest", - "serde", - "serde_json", - "thiserror 1.0.69", - "tokio", - "url", ] [[package]] @@ -5561,7 +5437,7 @@ dependencies = [ "wasm-bindgen-futures", "wasm-streams", "web-sys", - "webpki-roots 1.0.6", + "webpki-roots", ] [[package]] @@ -6071,17 +5947,6 @@ dependencies = [ ] [[package]] -name = "sha1" -version = "0.10.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e3bf829a2d51ab4a5ddf1352d8470c140cadc8301b2ae1789db023f01cedd6ba" -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" @@ -6977,7 +6842,7 @@ dependencies = [ "hex", "indicatif", "itertools 0.14.0", - "lru 0.12.5", + "lru", "mti", "num-bigint 0.4.6", "opentelemetry", @@ -7774,18 +7639,6 @@ dependencies = [ ] [[package]] -name = "tokio-socks" -version = "0.5.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0d4770b8024672c1101b3f6733eab95b18007dbe0847a8afe341fcf79e06043f" -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" @@ -7797,22 +7650,6 @@ 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" @@ -8122,25 +7959,6 @@ 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" @@ -8532,15 +8350,6 @@ 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/net/Cargo.toml b/crates/net/Cargo.toml @@ -48,9 +48,7 @@ radroots_nostr = { workspace = true, optional = true, default-features = true, f "events", ] } nostr = { workspace = true, optional = true, features = ["std"] } -radroots_transport_nostr = { workspace = true, optional = true, default-features = false, features = [ - "client", -] } +radroots_transport_nostr = { workspace = true, optional = true, default-features = false } directories = { workspace = true, optional = true } hex = { workspace = true, optional = true } radroots_runtime_paths = { workspace = true, optional = true } diff --git a/crates/nostr_runtime/Cargo.toml b/crates/nostr_runtime/Cargo.toml @@ -20,9 +20,7 @@ nostrdb = ["nostr-client"] [dependencies] radroots_nostr = { workspace = true, optional = true, default-features = true } -radroots_transport_nostr = { workspace = true, optional = true, default-features = false, features = [ - "client", -] } +radroots_transport_nostr = { workspace = true, optional = true, default-features = false } nostr = { workspace = true, optional = true, features = ["std"] } futures = { workspace = true, optional = true } thiserror = { workspace = true } diff --git a/crates/transport_nostr/Cargo.toml b/crates/transport_nostr/Cargo.toml @@ -14,58 +14,11 @@ readme = "README.md" [lib] name = "radroots_transport_nostr" -[features] -default = ["std", "client", "storage", "runtime-tokio"] -std = [] -client = [ - "dep:radroots_nostr", - "dep:nostr-sdk", - "dep:tokio", - "radroots_nostr/std", - "radroots_nostr/events", -] -nip11 = ["client", "dep:reqwest"] -storage = ["dep:radroots_event_store", "dep:radroots_outbox", "client"] -runtime-tokio = [ - "dep:tokio", - "storage", - "radroots_event_store/runtime-tokio", - "radroots_outbox/runtime-tokio", -] - [dependencies] -radroots_event = { workspace = true, default-features = false, features = [ - "std", - "serde", -] } -radroots_event_store = { workspace = true, optional = true, default-features = false, features = [ - "sqlite", - "runtime-tokio", -] } -radroots_nostr = { workspace = true, optional = true, default-features = false, features = [ - "std", - "events", -] } -radroots_outbox = { workspace = true, optional = true, default-features = false, features = [ - "sqlite", - "runtime-tokio", -] } +radroots_event_codec = { workspace = true, default-features = false } +radroots_nostr = { workspace = true, default-features = false } +radroots_protocol = { workspace = true, default-features = false } radroots_transport = { workspace = true, default-features = false } -futures = { workspace = true } -nostr = { workspace = true, features = ["std"] } -nostr-sdk = { workspace = true, optional = true } -reqwest = { workspace = true, optional = true, default-features = false, features = [ - "json", - "rustls-tls", -] } -serde = { workspace = true, features = ["derive", "std"] } -serde_json = { workspace = true, features = ["std"] } -thiserror = { workspace = true } -tokio = { workspace = true, optional = true, features = ["rt", "sync"] } -url = { workspace = true } - -[dev-dependencies] -tokio = { workspace = true, features = ["macros", "rt"] } [lints] workspace = true diff --git a/crates/transport_nostr/src/auth.rs b/crates/transport_nostr/src/auth.rs @@ -0,0 +1 @@ +//! Explicit relay authentication state. diff --git a/crates/transport_nostr/src/client.rs b/crates/transport_nostr/src/client.rs @@ -1,609 +1,23 @@ -#![forbid(unsafe_code)] +//! Concrete Nostr transport composition. -use core::time::Duration; -use core::{ - fmt::Debug, - pin::Pin, - task::{Context, Poll}, -}; -use std::collections::{HashMap, HashSet}; -#[cfg(not(target_arch = "wasm32"))] -use std::net::SocketAddr; +/// Configuration for a concrete Nostr transport. +#[derive(Clone, Debug, Default, Eq, PartialEq)] +pub struct Config; -use futures::Stream; -use nostr_sdk::{Client, ClientBuilder, ClientOptions}; - -use crate::{RadrootsRelayTransportError, RelayUrl}; -use nostr::{Keys, SecretKey, SubscriptionId}; -use radroots_nostr::event::Event as RadrootsNostrEvent; -use radroots_nostr::event::EventId as RadrootsNostrEventId; -use radroots_nostr::filter::Filter as RadrootsNostrFilter; - -#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] -pub enum RadrootsNostrRelayStatus { - Initialized, - Pending, - Connecting, - Connected, - Disconnected, - Terminated, - Banned, - Sleeping, -} - -fn normalize_relay_status(value: nostr_sdk::RelayStatus) -> RadrootsNostrRelayStatus { - match value { - nostr_sdk::RelayStatus::Initialized => RadrootsNostrRelayStatus::Initialized, - nostr_sdk::RelayStatus::Pending => RadrootsNostrRelayStatus::Pending, - nostr_sdk::RelayStatus::Connecting => RadrootsNostrRelayStatus::Connecting, - nostr_sdk::RelayStatus::Connected => RadrootsNostrRelayStatus::Connected, - nostr_sdk::RelayStatus::Disconnected => RadrootsNostrRelayStatus::Disconnected, - nostr_sdk::RelayStatus::Terminated => RadrootsNostrRelayStatus::Terminated, - nostr_sdk::RelayStatus::Banned => RadrootsNostrRelayStatus::Banned, - nostr_sdk::RelayStatus::Sleeping => RadrootsNostrRelayStatus::Sleeping, - } -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub enum RadrootsNostrMonitorNotification { - StatusChanged { - relay_url: RelayUrl, - status: RadrootsNostrRelayStatus, - }, -} - -#[derive(Debug, Clone)] -pub struct RadrootsNostrMonitor { - inner: nostr_sdk::prelude::Monitor, -} - -impl RadrootsNostrMonitor { - pub fn new(channel_size: usize) -> Self { - Self { - inner: nostr_sdk::prelude::Monitor::new(channel_size), - } - } - - pub fn subscribe(&self) -> RadrootsNostrMonitorReceiver { - RadrootsNostrMonitorReceiver { - inner: self.inner.subscribe(), - } - } -} - -pub struct RadrootsNostrMonitorReceiver { - inner: tokio::sync::broadcast::Receiver<nostr_sdk::prelude::MonitorNotification>, -} - -impl RadrootsNostrMonitorReceiver { - pub async fn recv( - &mut self, - ) -> Result<RadrootsNostrMonitorNotification, RadrootsNostrMonitorReceiveError> { - self.inner - .recv() - .await - .map(normalize_monitor_notification) - .map_err(|error| match error { - tokio::sync::broadcast::error::RecvError::Closed => { - RadrootsNostrMonitorReceiveError::Closed - } - tokio::sync::broadcast::error::RecvError::Lagged(skipped) => { - RadrootsNostrMonitorReceiveError::Lagged { skipped } - } - }) - } -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum RadrootsNostrMonitorReceiveError { - Closed, - Lagged { skipped: u64 }, -} - -fn normalize_monitor_notification( - notification: nostr_sdk::prelude::MonitorNotification, -) -> RadrootsNostrMonitorNotification { - match notification { - nostr_sdk::prelude::MonitorNotification::StatusChanged { relay_url, status } => { - RadrootsNostrMonitorNotification::StatusChanged { - relay_url: RelayUrl::from_normalized_transport(relay_url.to_string()), - status: normalize_relay_status(status), - } - } - } -} - -#[derive(Debug, Clone)] -pub struct RadrootsNostrOutput<T> { - pub val: T, - pub success: HashSet<RelayUrl>, - pub failed: HashMap<RelayUrl, String>, -} - -fn normalize_output<T>(output: nostr_sdk::prelude::Output<T>) -> RadrootsNostrOutput<T> -where - T: Debug, -{ - RadrootsNostrOutput { - val: output.val, - success: output - .success - .into_iter() - .map(|url| RelayUrl::from_normalized_transport(url.to_string())) - .collect(), - failed: output - .failed - .into_iter() - .map(|(url, error)| (RelayUrl::from_normalized_transport(url.to_string()), error)) - .collect(), - } -} - -fn normalize_subscription_output( - output: nostr_sdk::prelude::Output<SubscriptionId>, -) -> RadrootsNostrOutput<RadrootsNostrSubscriptionId> { - let output = normalize_output(output); - RadrootsNostrOutput { - val: RadrootsNostrSubscriptionId(output.val.to_string()), - success: output.success, - failed: output.failed, - } -} - -/// An opaque local credential for the compatibility Nostr transport client. -/// -/// Secret material cannot be cloned, formatted, or serialized through this -/// boundary. Hosts should retain their authoritative credential and create a -/// short-lived transport credential only when constructing a client. -/// -/// ```compile_fail -/// use radroots_transport_nostr::RadrootsNostrClientKey; -/// -/// let key = RadrootsNostrClientKey::generate(); -/// let duplicated = key.clone(); -/// ``` -pub struct RadrootsNostrClientKey { - inner: Keys, -} - -impl RadrootsNostrClientKey { - pub fn generate() -> Self { - Self { - inner: Keys::generate(), - } - } - - pub fn from_secret_key_bytes( - secret_key: [u8; 32], - ) -> Result<Self, RadrootsRelayTransportError> { - let secret_key = SecretKey::from_slice(&secret_key) - .map_err(|error| RadrootsRelayTransportError::ClientConfig(error.to_string()))?; - Ok(Self { - inner: Keys::new(secret_key), - }) - } - - pub fn public_key_hex(&self) -> String { - self.inner.public_key().to_hex() - } - - fn into_inner(self) -> Keys { - self.inner - } -} - -impl Debug for RadrootsNostrClientKey { - fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { - formatter.write_str("RadrootsNostrClientKey([REDACTED])") - } -} - -#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] -pub struct RadrootsNostrSubscriptionId(String); - -impl RadrootsNostrSubscriptionId { - pub fn as_str(&self) -> &str { - self.0.as_str() - } -} - -impl core::fmt::Display for RadrootsNostrSubscriptionId { - fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { - formatter.write_str(self.as_str()) - } -} - -pub struct RadrootsNostrEventStream { - inner: nostr_sdk::pool::stream::BoxedStream<RadrootsNostrEvent>, -} - -impl Stream for RadrootsNostrEventStream { - type Item = RadrootsNostrEvent; - - fn poll_next(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<Self::Item>> { - self.get_mut().inner.as_mut().poll_next(context) - } -} - -#[derive(Clone)] -pub struct RadrootsNostrRelay { - inner: nostr_sdk::Relay, - url: RelayUrl, -} - -impl RadrootsNostrRelay { - pub fn is_connected(&self) -> bool { - self.inner.is_connected() - } - - pub fn url(&self) -> &RelayUrl { - &self.url - } +/// Concrete Nostr implementation of the transport source and sink SPIs. +#[derive(Clone, Debug, Default)] +pub struct NostrTransport { + config: Config, } -#[derive(Debug, Clone, Copy, Default)] -pub struct RadrootsNostrSubscribeAutoCloseOptions { - timeout: Option<Duration>, - idle_timeout: Option<Duration>, -} - -impl RadrootsNostrSubscribeAutoCloseOptions { - pub fn timeout(mut self, timeout: Option<Duration>) -> Self { - self.timeout = timeout; - self - } - - pub fn idle_timeout(mut self, timeout: Option<Duration>) -> Self { - self.idle_timeout = timeout; - self - } - - fn into_sdk(self) -> nostr_sdk::SubscribeAutoCloseOptions { - nostr_sdk::SubscribeAutoCloseOptions::default() - .timeout(self.timeout) - .idle_timeout(self.idle_timeout) - } -} - -#[derive(Clone)] -pub struct RadrootsNostrClient { - inner: Client, - monitor: Option<RadrootsNostrMonitor>, -} - -#[derive(Debug, Clone, Default)] -pub struct RadrootsNostrClientOptions { - automatic_authentication: Option<bool>, - max_avg_latency_ms: Option<u64>, - verify_subscriptions: Option<bool>, - ban_relay_on_mismatch: Option<bool>, - #[cfg(not(target_arch = "wasm32"))] - proxy: Option<SocketAddr>, -} - -impl RadrootsNostrClientOptions { - pub fn new() -> Self { - Self::default() - } - - pub fn automatic_authentication(mut self, enabled: bool) -> Self { - self.automatic_authentication = Some(enabled); - self - } - - pub fn max_avg_latency_ms(mut self, max_ms: u64) -> Self { - self.max_avg_latency_ms = Some(max_ms); - self - } - - pub fn verify_subscriptions(mut self, enabled: bool) -> Self { - self.verify_subscriptions = Some(enabled); - self - } - - pub fn ban_relay_on_mismatch(mut self, enabled: bool) -> Self { - self.ban_relay_on_mismatch = Some(enabled); - self - } - - #[cfg(not(target_arch = "wasm32"))] - pub fn proxy_addr(mut self, addr: SocketAddr) -> Self { - self.proxy = Some(addr); - self - } - - #[cfg(not(target_arch = "wasm32"))] - pub fn proxy_str(mut self, addr: &str) -> Result<Self, RadrootsRelayTransportError> { - let parsed: SocketAddr = addr.parse().map_err(|err: std::net::AddrParseError| { - RadrootsRelayTransportError::ClientConfig(err.to_string()) - })?; - self.proxy = Some(parsed); - Ok(self) - } - - fn to_client_options(&self) -> ClientOptions { - let mut options = ClientOptions::new(); - if let Some(enabled) = self.automatic_authentication { - options = options.automatic_authentication(enabled); - } - if let Some(max_ms) = self.max_avg_latency_ms { - options = options.max_avg_latency(Duration::from_millis(max_ms)); - } - if let Some(enabled) = self.verify_subscriptions { - options = options.verify_subscriptions(enabled); - } - if let Some(enabled) = self.ban_relay_on_mismatch { - options = options.ban_relay_on_mismatch(enabled); - } - #[cfg(not(target_arch = "wasm32"))] - if let Some(proxy) = self.proxy { - let connection = nostr_sdk::client::options::Connection::new().proxy(proxy); - options = options.connection(connection); - } - options - } -} - -impl RadrootsNostrClient { - pub fn new_signerless() -> Self { - Self { - inner: Client::default(), - monitor: None, - } - } - - pub fn new_signerless_with_options(options: RadrootsNostrClientOptions) -> Self { - let inner = ClientBuilder::new() - .opts(options.to_client_options()) - .build(); - Self { - inner, - monitor: None, - } - } - - pub fn new(keys: RadrootsNostrClientKey) -> Self { - Self { - inner: Client::new(keys.into_inner()), - monitor: None, - } - } - - pub fn from_keys_with_options( - keys: RadrootsNostrClientKey, - options: RadrootsNostrClientOptions, - ) -> Self { - let inner = ClientBuilder::new() - .signer(keys.into_inner()) - .opts(options.to_client_options()) - .build(); - Self { - inner, - monitor: None, - } - } - - pub fn new_with_monitor(keys: RadrootsNostrClientKey, monitor: RadrootsNostrMonitor) -> Self { - let inner = Client::builder() - .signer(keys.into_inner()) - .monitor(monitor.inner.clone()) - .build(); - Self { - inner, - monitor: Some(monitor), - } - } - - pub async fn has_signer(&self) -> bool { - self.inner.has_signer().await - } - - pub async fn public_key_hex(&self) -> Result<String, RadrootsRelayTransportError> { - self.inner - .public_key() - .await - .map(|public_key| public_key.to_hex()) - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub fn monitor(&self) -> Option<&RadrootsNostrMonitor> { - self.monitor.as_ref() - } - - pub async fn connect(&self) { - self.inner.connect().await; - } - - pub async fn wait_for_connection(&self, timeout: Duration) { - self.inner.wait_for_connection(timeout).await; - } - - pub async fn try_connect(&self, timeout: Duration) -> RadrootsNostrOutput<()> { - normalize_output(self.inner.try_connect(timeout).await) - } - - pub async fn add_relay(&self, url: &str) -> Result<bool, RadrootsRelayTransportError> { - self.inner - .add_relay(url) - .await - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn add_write_relay(&self, url: &str) -> Result<bool, RadrootsRelayTransportError> { - self.inner - .add_write_relay(url) - .await - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn add_read_relay(&self, url: &str) -> Result<bool, RadrootsRelayTransportError> { - self.inner - .add_read_relay(url) - .await - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn remove_relay(&self, url: &str) -> Result<(), RadrootsRelayTransportError> { - self.inner - .force_remove_relay(url) - .await - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn relays(&self) -> HashMap<RelayUrl, RadrootsNostrRelay> { - self.inner - .relays() - .await - .into_iter() - .map(|(url, inner)| { - let url = RelayUrl::from_normalized_transport(url.to_string()); - (url.clone(), RadrootsNostrRelay { inner, url }) - }) - .collect() - } - - pub async fn fetch_events( - &self, - filter: RadrootsNostrFilter, - timeout: Duration, - ) -> Result<Vec<RadrootsNostrEvent>, RadrootsRelayTransportError> { - self.inner - .fetch_events(filter, timeout) - .await - .map(|events| events.to_vec()) - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn query_database( - &self, - filter: RadrootsNostrFilter, - ) -> Result<Vec<RadrootsNostrEvent>, RadrootsRelayTransportError> { - self.inner - .database() - .query(filter) - .await - .map(|events| events.to_vec()) - .map_err(|error| RadrootsRelayTransportError::ClientDatabase(error.to_string())) - } - - pub async fn stream_events( - &self, - filter: RadrootsNostrFilter, - timeout: Duration, - ) -> Result<RadrootsNostrEventStream, RadrootsRelayTransportError> { - self.inner - .stream_events(filter, timeout) - .await - .map(|inner| RadrootsNostrEventStream { inner }) - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn subscribe( - &self, - filter: RadrootsNostrFilter, - options: Option<RadrootsNostrSubscribeAutoCloseOptions>, - ) -> Result<RadrootsNostrOutput<RadrootsNostrSubscriptionId>, RadrootsRelayTransportError> { - self.inner - .subscribe( - filter, - options.map(RadrootsNostrSubscribeAutoCloseOptions::into_sdk), - ) - .await - .map(normalize_subscription_output) - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn subscribe_to_relays( - &self, - relays: &[RelayUrl], - filter: RadrootsNostrFilter, - options: Option<RadrootsNostrSubscribeAutoCloseOptions>, - ) -> Result<RadrootsNostrOutput<RadrootsNostrSubscriptionId>, RadrootsRelayTransportError> { - self.inner - .subscribe_to( - relays.iter().map(RelayUrl::as_str), - filter, - options.map(RadrootsNostrSubscribeAutoCloseOptions::into_sdk), - ) - .await - .map(normalize_subscription_output) - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn unsubscribe(&self, subscription_id: &RadrootsNostrSubscriptionId) { - self.inner - .unsubscribe(&SubscriptionId::new(subscription_id.as_str())) - .await; - } - - /// Relays a caller-supplied signed event. - /// - /// This is a transport boundary, not an authored-builder boundary. The - /// caller is responsible for the event's authoring policy and signature. - pub async fn send_event( - &self, - event: &RadrootsNostrEvent, - ) -> Result<RadrootsNostrOutput<RadrootsNostrEventId>, RadrootsRelayTransportError> { - self.inner - .send_event(event) - .await - .map(normalize_output) - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn send_event_to_relays( - &self, - relays: &[RelayUrl], - event: &RadrootsNostrEvent, - ) -> Result<RadrootsNostrOutput<RadrootsNostrEventId>, RadrootsRelayTransportError> { - self.inner - .send_event_to(relays.iter().map(RelayUrl::as_str), event) - .await - .map(normalize_output) - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } - - pub async fn send_event_to( - &self, - relays: Vec<String>, - event: &RadrootsNostrEvent, - ) -> Result<RadrootsNostrOutput<RadrootsNostrEventId>, RadrootsRelayTransportError> { - self.inner - .send_event_to(relays, event) - .await - .map(normalize_output) - .map_err(|error| RadrootsRelayTransportError::Client(error.to_string())) - } -} - -pub async fn radroots_nostr_fetch_event_by_id( - client: &RadrootsNostrClient, - id: &str, -) -> Result<RadrootsNostrEvent, RadrootsRelayTransportError> { - let event_id = RadrootsNostrEventId::parse(id) - .map_err(|error| RadrootsRelayTransportError::NostrEvent(error.to_string()))?; - let filter = RadrootsNostrFilter::new().id(event_id); - let events = client.fetch_events(filter, Duration::from_secs(10)).await?; - let event = events - .first() - .ok_or_else(|| RadrootsRelayTransportError::EventNotFound(event_id.to_hex()))?; - Ok(event.clone()) -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn client_key_debug_output_is_redacted() { - let key = RadrootsNostrClientKey::generate(); - - assert_eq!(format!("{key:?}"), "RadrootsNostrClientKey([REDACTED])"); - assert!(!format!("{key:?}").contains(&key.public_key_hex())); +impl NostrTransport { + /// Creates an inert transport from explicit configuration. + pub const fn new(config: Config) -> Self { + Self { config } } - #[test] - fn client_key_rejects_invalid_secret_scalar() { - assert!(RadrootsNostrClientKey::from_secret_key_bytes([0; 32]).is_err()); + /// Returns the transport configuration. + pub const fn config(&self) -> &Config { + &self.config } } diff --git a/crates/transport_nostr/src/error.rs b/crates/transport_nostr/src/error.rs @@ -1,182 +1,21 @@ -#![forbid(unsafe_code)] +//! Stable Nostr transport failures. -use thiserror::Error; +use core::fmt; -#[derive(Debug, Error)] -pub enum RadrootsRelayTransportError { - #[cfg(feature = "client")] - #[error("Nostr client error: {0}")] - Client(String), - - #[cfg(feature = "client")] - #[error("Nostr client database error: {0}")] - ClientDatabase(String), - - #[cfg(feature = "client")] - #[error("Nostr protocol event error: {0}")] - NostrEvent(String), - - #[cfg(feature = "client")] - #[error("Nostr client configuration error: {0}")] - ClientConfig(String), - - #[cfg(feature = "client")] - #[error("Nostr event not found: {0}")] - EventNotFound(String), - - #[error("Relay URL parse failed for `{url}`: {reason}")] - RelayUrlParse { url: String, reason: String }, - - #[error("Relay URL `{url}` uses ws outside localhost relay policy")] - WsRequiresLocalhostPolicy { url: String }, - - #[error("Relay URL `{url}` has unsupported scheme `{scheme}`")] - UnsupportedRelayScheme { url: String, scheme: String }, - - #[error("Relay URL `{url}` must include a host")] - EmptyRelayHost { url: String }, - - #[error("Relay URL `{url}` must not include userinfo")] - RelayUrlUserinfo { url: String }, - - #[error("Relay URL `{url}` must not include query or fragment")] - RelayUrlQueryOrFragment { url: String }, - - #[error("Relay URL `{url}` targets forbidden destination: {reason}")] - RelayUrlForbiddenDestination { url: String, reason: String }, - - #[error("Relay URL `{url}` resolved to forbidden address `{address}`: {reason}")] - RelayUrlResolvedForbiddenDestination { - url: String, - address: String, - reason: String, - }, - - #[error("Relay URL `{url}` did not resolve to any addresses")] - RelayUrlResolvedNoAddresses { url: String }, - - #[error("Relay target set must not be empty")] - EmptyTargetSet, - - #[error("Relay target set contains duplicate URL `{url}`")] - DuplicateRelayUrl { url: String }, - - #[error("Relay fetch item contains invalid relay URL `{url}`: {reason}")] - InvalidFetchItemRelayUrl { url: String, reason: String }, - - #[error("Relay fetch item came from unrequested relay URL `{url}`")] - UnexpectedFetchItemRelayUrl { url: String }, - - #[error("Relay fetch adapter returned duplicate terminal outcome for relay URL `{url}`")] - DuplicateFetchTerminalRelayUrl { url: String }, - - #[error( - "Relay fetch adapter returned conflicting terminal outcomes for relay URL `{url}`: first={first}, next={next}" - )] - ConflictingFetchTerminalRelayUrl { - url: String, - first: &'static str, - next: &'static str, - }, - - #[error("Relay publish receipt contains invalid relay URL `{url}`: {reason}")] - InvalidPublishReceiptRelayUrl { url: String, reason: String }, - - #[error("Relay publish receipt came from unrequested relay URL `{url}`")] - UnexpectedPublishReceiptRelayUrl { url: String }, - - #[error("Relay publish adapter returned duplicate receipts for relay URL `{url}`")] - DuplicatePublishReceiptRelayUrl { url: String }, - - #[error("Relay publish adapter returned incoherent attempt state for relay URL `{url}`")] - InvalidPublishReceiptAttemptState { url: String }, - - #[error("Transport returned conflicting publish receipts for relay URL `{url}`")] - ConflictingTransportReceiptRelayUrl { url: String }, - - #[error("Expected transport kind `{expected}`, received `{actual}`")] - UnexpectedTransportKind { - expected: &'static str, - actual: String, - }, - - #[error("Relay fetch filters must not be empty")] - EmptyFetchFilters, - - #[error("Relay fetch {field} must be greater than zero")] - InvalidFetchLimit { field: &'static str }, - - #[error("Relay fetch {field} {actual} exceeds maximum {max}")] - FetchLimitTooLarge { - field: &'static str, - max: usize, - actual: usize, - }, - - #[error("Relay transport {field} cannot be negative: {value}")] - InvalidTimestamp { field: &'static str, value: i64 }, - - #[error( - "Relay publish idempotency key is invalid: {reason}; bytes={actual_bytes}, max={max_bytes}" - )] - InvalidIdempotencyKey { - reason: &'static str, - actual_bytes: usize, - max_bytes: usize, - }, - - #[error("Relay publish required target `{fingerprint}` is not in the requested relay set")] - RequiredTargetNotRequested { fingerprint: String }, - - #[error("Transport contract error: {0}")] - TransportContract(String), - - #[error("JSON error: {0}")] - Json(#[from] serde_json::Error), - - #[error("Nostr event JSON error: {0}")] - NostrEventJson(String), - - #[cfg(feature = "storage")] - #[error("Event store error: {0}")] - EventStore(#[from] radroots_event_store::RadrootsEventStoreError), - - #[cfg(feature = "storage")] - #[error("Event store returned no current visibility for persisted event `{event_id}`")] - MissingStoredEventVisibility { event_id: String }, - - #[cfg(feature = "storage")] - #[error("Persisted relay fetch event receipt is missing its event id")] - MissingPersistedFetchReceiptEventId, - - #[cfg(feature = "storage")] - #[error("Event store returned an unsupported current visibility for event `{event_id}`")] - UnsupportedStoredEventVisibility { event_id: String }, - - #[cfg(feature = "storage")] - #[error("Outbox error: {0}")] - Outbox(#[from] radroots_outbox::RadrootsOutboxError), - - #[cfg(feature = "storage")] - #[error("Outbox claim {0} does not contain a signed event")] - MissingSignedOutboxEvent(i64), - - #[error("Relay transport error: {0}")] - Transport(String), +/// Error returned by the Nostr transport adapter. +#[derive(Clone, Debug, Eq, PartialEq)] +#[non_exhaustive] +pub enum Error { + /// The adapter has not yet been configured for an operation. + NotConfigured, } -pub(crate) fn ensure_nonnegative_timestamp( - field: &'static str, - value: i64, -) -> Result<(), RadrootsRelayTransportError> { - if value < 0 { - return Err(RadrootsRelayTransportError::InvalidTimestamp { field, value }); +impl fmt::Display for Error { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::NotConfigured => formatter.write_str("Nostr transport is not configured"), + } } - Ok(()) } -impl From<radroots_transport::RadrootsTransportError> for RadrootsRelayTransportError { - fn from(value: radroots_transport::RadrootsTransportError) -> Self { - Self::TransportContract(value.to_string()) - } -} +impl std::error::Error for Error {} diff --git a/crates/transport_nostr/src/lib.rs b/crates/transport_nostr/src/lib.rs @@ -1,60 +1,14 @@ -#![cfg_attr(coverage_nightly, feature(coverage_attribute))] +#![doc = include_str!("../README.md")] #![forbid(unsafe_code)] -#[cfg(feature = "client")] +mod auth; mod client; mod error; -#[cfg(feature = "storage")] -mod fetch; -#[cfg(feature = "nip11")] -mod nip11; -#[cfg(feature = "storage")] -mod outbox; -mod outcome; -mod publish; mod relay; -#[cfg(feature = "client")] -mod relays; +mod sink; +mod source; +mod status; -#[cfg(feature = "client")] -pub use client::{ - RadrootsNostrClient, RadrootsNostrClientKey, RadrootsNostrClientOptions, - RadrootsNostrEventStream, RadrootsNostrMonitor, RadrootsNostrMonitorNotification, - RadrootsNostrMonitorReceiveError, RadrootsNostrMonitorReceiver, RadrootsNostrOutput, - RadrootsNostrRelay, RadrootsNostrRelayStatus, RadrootsNostrSubscribeAutoCloseOptions, - RadrootsNostrSubscriptionId, radroots_nostr_fetch_event_by_id, -}; -pub use error::RadrootsRelayTransportError; -#[cfg(all(feature = "storage", feature = "runtime-tokio"))] -pub use fetch::fetch_relay_events_blocking; -#[cfg(feature = "storage")] -pub use fetch::{ - RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX, RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX, - RADROOTS_RELAY_FETCH_RAW_JSON_BYTE_LIMIT_MAX, RadrootsMockRelayFetchAdapter, - RadrootsNostrClientFetchAdapter, RadrootsRelayFetchAdapter, RadrootsRelayFetchEventAdmission, - RadrootsRelayFetchEventReceipt, RadrootsRelayFetchEventValidStream, - RadrootsRelayFetchEventVerification, RadrootsRelayFetchEventVisibility, - RadrootsRelayFetchFailure, RadrootsRelayFetchFilters, RadrootsRelayFetchItem, - RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchReceipt, - RadrootsRelayFetchRelayOutcome, RadrootsRelayFetchRequest, RadrootsRelayFetchedEvent, - RadrootsRelayFetchedEventsReceipt, fetch_and_ingest_relay_events, fetch_relay_events, -}; -#[cfg(feature = "nip11")] -pub use nip11::fetch_nip11; -#[cfg(feature = "storage")] -pub use outbox::{ - RadrootsOutboxPublishPolicy, RadrootsOutboxPublishReceipt, RadrootsOutboxPublishTargetReceipt, - publish_claimed_outbox_event, publish_claimed_outbox_event_with_transport, -}; -pub use outcome::{RadrootsRelayOutcome, RadrootsRelayOutcomeKind}; -#[cfg(feature = "client")] -pub use publish::RadrootsNostrClientPublishAdapter; -pub use publish::{ - RADROOTS_RELAY_PUBLISH_IDEMPOTENCY_KEY_MAX_BYTES, RadrootsMockRelayPublishAdapter, - RadrootsNostrTransport, RadrootsRelayPublishAdapter, RadrootsRelayPublishReceipt, - RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, publish_signed_event, - verified_signed_event_payload, -}; -pub use relay::{RadrootsRelayTargetSet, RadrootsRelayUrlPolicy, RelayUrl}; -#[cfg(feature = "client")] -pub use relays::{radroots_nostr_add_relay, radroots_nostr_connect, radroots_nostr_remove_relay}; +pub use client::{Config, NostrTransport}; +pub use error::Error; +pub use relay::{RelayUrl, RelayUrlPolicy}; diff --git a/crates/transport_nostr/src/relay.rs b/crates/transport_nostr/src/relay.rs @@ -1,446 +1,22 @@ -#![forbid(unsafe_code)] +//! Nostr relay identifiers and network policy. -use crate::RadrootsRelayTransportError; -use radroots_transport::{Target, TransportId}; -use std::fmt; -use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; -use url::Url; - -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub enum RadrootsRelayUrlPolicy { - /// Allows `wss` endpoints from trusted configuration. - /// - /// This performs canonical syntax, literal-address, and local-hostname - /// checks. It is not an SSRF boundary for attacker-controlled hostnames - /// because the default SDK connector does not pin DNS resolution. - Public, - /// Allows `ws` or `wss` only for exact loopback hosts. - Localhost, -} - -impl RadrootsRelayUrlPolicy { - fn accepts_ws_host(self, host: &str) -> bool { - matches!(self, Self::Localhost) - && matches!(host, "localhost" | "127.0.0.1" | "::1" | "[::1]") - } -} - -#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +/// Validated Nostr relay URL. +#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] pub struct RelayUrl(String); impl RelayUrl { - pub(crate) fn from_normalized_transport(value: impl Into<String>) -> Self { - Self(value.into()) - } - - pub fn parse( - value: impl AsRef<str>, - policy: RadrootsRelayUrlPolicy, - ) -> Result<Self, RadrootsRelayTransportError> { - let original = value.as_ref(); - let parsed = - Url::parse(original).map_err(|error| RadrootsRelayTransportError::RelayUrlParse { - url: original.to_owned(), - reason: error.to_string(), - })?; - if !parsed.username().is_empty() || parsed.password().is_some() { - return Err(RadrootsRelayTransportError::RelayUrlUserinfo { - url: original.to_owned(), - }); - } - let Some(host) = parsed.host_str().filter(|host| !host.is_empty()) else { - return Err(RadrootsRelayTransportError::EmptyRelayHost { - url: original.to_owned(), - }); - }; - validate_host_destination(original, host, policy)?; - if parsed.query().is_some() || parsed.fragment().is_some() { - return Err(RadrootsRelayTransportError::RelayUrlQueryOrFragment { - url: original.to_owned(), - }); - } - let scheme = parsed.scheme(); - match scheme { - "wss" => {} - "ws" if policy.accepts_ws_host(host) => {} - "ws" => { - return Err(RadrootsRelayTransportError::WsRequiresLocalhostPolicy { - url: original.to_owned(), - }); - } - other => { - return Err(RadrootsRelayTransportError::UnsupportedRelayScheme { - url: original.to_owned(), - scheme: other.to_owned(), - }); - } - } - let target = Target::new(TransportId::NOSTR, original).map_err(|error| { - RadrootsRelayTransportError::RelayUrlParse { - url: original.to_owned(), - reason: error.to_string(), - } - })?; - Ok(Self(target.uri().as_str().to_owned())) - } - - pub fn validate_public_resolved_ip_addrs<I>( - &self, - addrs: I, - ) -> Result<(), RadrootsRelayTransportError> - where - I: IntoIterator<Item = IpAddr>, - { - let mut resolved_any = false; - for address in addrs { - resolved_any = true; - if let Some(reason) = forbidden_public_ip_reason(address) { - return Err( - RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination { - url: self.0.clone(), - address: address.to_string(), - reason: reason.to_owned(), - }, - ); - } - } - if !resolved_any { - return Err(RadrootsRelayTransportError::RelayUrlResolvedNoAddresses { - url: self.0.clone(), - }); - } - Ok(()) - } - + /// Returns the validated URL representation. pub fn as_str(&self) -> &str { self.0.as_str() } - - pub fn into_string(self) -> String { - self.0 - } } -fn validate_host_destination( - original: &str, - host: &str, - policy: RadrootsRelayUrlPolicy, -) -> Result<(), RadrootsRelayTransportError> { - let host = host - .strip_prefix('[') - .and_then(|value| value.strip_suffix(']')) - .unwrap_or(host); - if matches!(policy, RadrootsRelayUrlPolicy::Localhost) { - if !policy.accepts_ws_host(host) { - return Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { - url: original.to_owned(), - reason: "localhost policy permits only exact loopback hosts".to_owned(), - }); - } - return Ok(()); - } - if let Ok(address) = host.parse::<IpAddr>() { - if let Some(reason) = forbidden_public_ip_reason(address) { - return Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { - url: original.to_owned(), - reason: reason.to_owned(), - }); - } - } else if let Some(reason) = forbidden_public_hostname_reason(host) { - return Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { - url: original.to_owned(), - reason: reason.to_owned(), - }); - } - Ok(()) -} - -fn forbidden_public_hostname_reason(host: &str) -> Option<&'static str> { - let host = host.to_ascii_lowercase(); - if host == "localhost" || host.ends_with(".localhost") { - Some("localhost DNS name") - } else if host.ends_with(".local") { - Some("link-local multicast DNS name") - } else if host.ends_with(".home.arpa") { - Some("special-use home network DNS name") - } else if !host.contains('.') { - Some("single-label DNS name") - } else { - None - } -} - -fn forbidden_public_ip_reason(address: IpAddr) -> Option<&'static str> { - match address { - IpAddr::V4(address) => forbidden_public_ipv4_reason(address), - IpAddr::V6(address) => forbidden_public_ipv6_reason(address), - } -} - -fn forbidden_public_ipv4_reason(address: Ipv4Addr) -> Option<&'static str> { - let octets = address.octets(); - if address.is_unspecified() || octets[0] == 0 { - Some("unspecified or this-network IPv4 address") - } else if address.is_loopback() { - Some("loopback IPv4 address") - } else if address.is_private() { - Some("private IPv4 address") - } else if address.is_link_local() { - Some("link-local IPv4 address") - } else if address.is_multicast() { - Some("multicast IPv4 address") - } else if address.is_broadcast() { - Some("broadcast IPv4 address") - } else if address.is_documentation() { - Some("documentation IPv4 address") - } else if octets[0] == 100 && (64..=127).contains(&octets[1]) { - Some("shared IPv4 address space") - } else if octets[0] == 192 && octets[1] == 0 && octets[2] == 0 { - Some("IETF protocol-assignment IPv4 address") - } else if octets[0] == 192 && octets[1] == 88 && octets[2] == 99 { - Some("deprecated or local-use relay-anycast IPv4 address") - } else if octets[0] == 198 && matches!(octets[1], 18 | 19) { - Some("benchmark IPv4 address") - } else if octets[0] >= 240 { - Some("reserved IPv4 address") - } else { - None - } -} - -fn forbidden_public_ipv6_reason(address: Ipv6Addr) -> Option<&'static str> { - let segments = address.segments(); - if let Some(mapped) = address.to_ipv4_mapped() { - return forbidden_public_ipv4_reason(mapped); - } - if address.is_unspecified() { - Some("unspecified IPv6 address") - } else if address.is_loopback() { - Some("loopback IPv6 address") - } else if address.is_multicast() { - Some("multicast IPv6 address") - } else if (segments[0] & 0xfe00) == 0xfc00 { - Some("unique-local IPv6 address") - } else if (segments[0] & 0xffc0) == 0xfe80 { - Some("link-local IPv6 address") - } else if segments[0] == 0x0064 - && segments[1] == 0xff9b - && segments[2..6].iter().all(|segment| *segment == 0) - { - Some("IPv4/IPv6 translation address") - } else if !is_supported_global_ipv6_unicast(segments) { - Some("non-global or reserved IPv6 unicast address") - } else if segments[0] == 0x2001 && segments[1] == 0x0db8 { - Some("documentation IPv6 address") - } else if segments[0] == 0x2001 && segments[1] < 0x0200 { - Some("IETF protocol-assignment IPv6 address") - } else if segments[0] == 0x2002 { - Some("6to4 IPv6 address") - } else if segments[0] == 0x3fff && (segments[1] & 0xf000) == 0 { - Some("documentation IPv6 address") - } else { - None - } -} - -fn is_supported_global_ipv6_unicast(segments: [u16; 8]) -> bool { - (segments[0] & 0xe000) == 0x2000 -} - -impl fmt::Display for RelayUrl { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - f.write_str(self.0.as_str()) - } -} - -#[derive(Clone, Debug, PartialEq, Eq)] -pub struct RadrootsRelayTargetSet { - relays: Vec<RelayUrl>, -} - -impl RadrootsRelayTargetSet { - pub fn new<I, S>( - relays: I, - policy: RadrootsRelayUrlPolicy, - ) -> Result<Self, RadrootsRelayTransportError> - where - I: IntoIterator<Item = S>, - S: AsRef<str>, - { - let mut ordered_relays = Vec::new(); - for relay in relays { - let relay = RelayUrl::parse(relay, policy)?; - if ordered_relays.iter().any(|existing| existing == &relay) { - return Err(RadrootsRelayTransportError::DuplicateRelayUrl { - url: relay.into_string(), - }); - } - ordered_relays.push(relay); - } - let relays = ordered_relays; - if relays.is_empty() { - return Err(RadrootsRelayTransportError::EmptyTargetSet); - } - Ok(Self { relays }) - } - - pub fn from_urls(relays: Vec<RelayUrl>) -> Result<Self, RadrootsRelayTransportError> { - let mut ordered_relays = Vec::new(); - for relay in relays { - if ordered_relays.iter().any(|existing| existing == &relay) { - return Err(RadrootsRelayTransportError::DuplicateRelayUrl { - url: relay.into_string(), - }); - } - ordered_relays.push(relay); - } - let relays = ordered_relays; - if relays.is_empty() { - return Err(RadrootsRelayTransportError::EmptyTargetSet); - } - Ok(Self { relays }) - } - - pub fn relays(&self) -> &[RelayUrl] { - &self.relays - } - - pub fn relay_strings(&self) -> Vec<String> { - self.relays - .iter() - .map(|relay| relay.as_str().to_owned()) - .collect() - } - - pub fn len(&self) -> usize { - self.relays.len() - } - - pub fn is_empty(&self) -> bool { - self.relays.is_empty() - } -} - -#[cfg(test)] -mod tests { - use super::{ - RadrootsRelayUrlPolicy, forbidden_public_ipv4_reason, forbidden_public_ipv6_reason, - validate_host_destination, - }; - use std::net::{Ipv4Addr, Ipv6Addr}; - - #[test] - fn host_destination_validation_covers_public_and_local_policy_edges() { - assert!(!RadrootsRelayUrlPolicy::Public.accepts_ws_host("localhost")); - assert!(RadrootsRelayUrlPolicy::Localhost.accepts_ws_host("localhost")); - validate_host_destination( - "wss://93.184.216.34", - "93.184.216.34", - RadrootsRelayUrlPolicy::Public, - ) - .expect("public ipv4 host"); - validate_host_destination( - "wss://relay.example.com", - "relay.example.com", - RadrootsRelayUrlPolicy::Public, - ) - .expect("public dns host"); - validate_host_destination( - "ws://127.0.0.1", - "127.0.0.1", - RadrootsRelayUrlPolicy::Localhost, - ) - .expect("localhost policy host"); - } - - #[test] - fn public_ipv4_classifier_covers_forbidden_ranges_and_global_addresses() { - let cases = [ - Ipv4Addr::new(0, 0, 0, 0), - Ipv4Addr::new(0, 1, 2, 3), - Ipv4Addr::new(127, 0, 0, 1), - Ipv4Addr::new(10, 1, 2, 3), - Ipv4Addr::new(169, 254, 1, 2), - Ipv4Addr::new(224, 0, 0, 1), - Ipv4Addr::new(255, 255, 255, 255), - Ipv4Addr::new(192, 0, 2, 1), - Ipv4Addr::new(100, 64, 0, 1), - Ipv4Addr::new(192, 0, 0, 8), - Ipv4Addr::new(192, 88, 99, 2), - Ipv4Addr::new(198, 18, 0, 1), - Ipv4Addr::new(240, 0, 0, 1), - ]; - for address in cases { - assert!(forbidden_public_ipv4_reason(address).is_some()); - } - assert_eq!( - forbidden_public_ipv4_reason(Ipv4Addr::new(93, 184, 216, 34)), - None - ); - assert_eq!( - forbidden_public_ipv4_reason(Ipv4Addr::new(100, 128, 0, 1)), - None - ); - assert_eq!( - forbidden_public_ipv4_reason(Ipv4Addr::new(193, 0, 0, 8)), - None - ); - assert_eq!( - forbidden_public_ipv4_reason(Ipv4Addr::new(192, 1, 0, 8)), - None - ); - assert_eq!( - forbidden_public_ipv4_reason(Ipv4Addr::new(192, 0, 1, 8)), - None - ); - assert_eq!( - forbidden_public_ipv4_reason(Ipv4Addr::new(198, 20, 0, 1)), - None - ); - } - - #[test] - fn public_ipv6_classifier_covers_forbidden_ranges_and_global_addresses() { - let cases = [ - "::ffff:192.168.1.10", - "::", - "::1", - "ff02::1", - "fd00::1", - "fe80::1", - "64:ff9b::7f00:1", - "64:ff9b::a00:1", - "64:ff9b::5db8:d822", - "64:ff9b:1::1", - "100::1", - "100:0:0:1::1", - "2001:db8::1", - "2001:1::1", - "2002:db8::1", - "2002:1::1", - "3fff::1", - "5f00::1", - ]; - for address in cases { - assert!( - forbidden_public_ipv6_reason(address.parse::<Ipv6Addr>().expect("ipv6")).is_some() - ); - } - assert_eq!( - forbidden_public_ipv6_reason( - "2001:4860:4860::8888" - .parse::<Ipv6Addr>() - .expect("public ipv6") - ), - None - ); - assert_eq!( - forbidden_public_ipv6_reason("2001:db9::1".parse::<Ipv6Addr>().expect("ipv6")), - None - ); - assert_eq!( - forbidden_public_ipv6_reason("2001:200::1".parse::<Ipv6Addr>().expect("ipv6")), - None - ); - } +/// Network destinations accepted for a relay URL. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +#[non_exhaustive] +pub enum RelayUrlPolicy { + /// Public TLS relay endpoints only. + Public, + /// Exact loopback relay endpoints only. + Local, } diff --git a/crates/transport_nostr/src/sink.rs b/crates/transport_nostr/src/sink.rs @@ -0,0 +1 @@ +//! Nostr implementation of the transport event sink. diff --git a/crates/transport_nostr/src/source.rs b/crates/transport_nostr/src/source.rs @@ -0,0 +1 @@ +//! Nostr implementation of the transport event source. diff --git a/crates/transport_nostr/src/status.rs b/crates/transport_nostr/src/status.rs @@ -0,0 +1 @@ +//! Stable Nostr relay status normalization. diff --git a/crates/transport_nostr/tests/package_boundary.rs b/crates/transport_nostr/tests/package_boundary.rs @@ -0,0 +1,69 @@ +use std::collections::BTreeSet; + +const MANIFEST: &str = include_str!("../Cargo.toml"); +const ROOT: &str = include_str!("../src/lib.rs"); + +#[test] +fn manifest_and_root_match_the_governed_transport_boundary() { + for required in [ + "name = \"radroots_transport_nostr\"", + "version = \"0.1.0-alpha\"", + "publish = false", + "[lib]\nname = \"radroots_transport_nostr\"", + ] { + assert!( + MANIFEST.contains(required), + "manifest is missing `{required}`" + ); + } + assert!(!MANIFEST.contains("[features]")); + assert_eq!( + radroots_dependency_keys(MANIFEST), + BTreeSet::from([ + "radroots_event_codec", + "radroots_nostr", + "radroots_protocol", + "radroots_transport", + ]) + ); + assert_eq!( + private_modules(ROOT), + BTreeSet::from([ + "auth", "client", "error", "relay", "sink", "source", "status" + ]) + ); + for export in [ + "pub use client::{Config, NostrTransport};", + "pub use error::Error;", + "pub use relay::{RelayUrl, RelayUrlPolicy};", + ] { + assert!(ROOT.contains(export), "crate root is missing `{export}`"); + } +} + +fn radroots_dependency_keys(manifest: &str) -> BTreeSet<&str> { + dependency_keys(manifest) + .into_iter() + .filter(|key| key.starts_with("radroots_")) + .collect() +} + +fn dependency_keys(manifest: &str) -> BTreeSet<&str> { + manifest + .split_once("[dependencies]") + .map(|(_, dependencies)| dependencies) + .unwrap_or_default() + .lines() + .skip(1) + .take_while(|line| !line.starts_with('[')) + .filter_map(|line| line.split_once('=').map(|(key, _)| key.trim())) + .filter(|key| !key.is_empty()) + .collect() +} + +fn private_modules(root: &str) -> BTreeSet<&str> { + root.lines() + .filter_map(|line| line.trim().strip_prefix("mod ")) + .filter_map(|module| module.strip_suffix(';')) + .collect() +} diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -1,4985 +0,0 @@ -use futures::future::BoxFuture; -use nostr::{EventBuilder, JsonUtil}; -use nostr::{Keys as RadrootsNostrKeys, SecretKey as RadrootsNostrSecretKey}; -use radroots_event::draft::{EventDraft, SignedEvent}; -use radroots_event::envelope::kind::{ - KIND_DELETION_REQUEST, KIND_FOLLOW, KIND_GEOCHAT, KIND_POST, KIND_PROFILE, -}; -use radroots_event::wire::v1::DEFAULT_RAW_JSON_MAX_BYTES; -use radroots_event_store::{ - RadrootsEventStore, RadrootsTransportObservationRow, RadrootsTransportObservationType, -}; -use radroots_nostr::event::Kind as RadrootsNostrKind; -use radroots_nostr::event::Timestamp as RadrootsNostrTimestamp; -use radroots_nostr::filter::Filter as RadrootsNostrFilter; -use radroots_nostr::filter::with_tag; -use radroots_nostr::signing::sign_frozen_draft; -use radroots_nostr::tag::Tag as RadrootsNostrTag; -use radroots_nostr::tag::TagKind as RadrootsNostrTagKind; -use radroots_outbox::{ - RadrootsOutbox, RadrootsOutboxClaimedEvent, RadrootsOutboxDeliveryPlanInput, - RadrootsOutboxDeliveryTargetStatus, RadrootsOutboxEventState, RadrootsOutboxOperationInput, - RadrootsOutboxOperationStatus, -}; -use radroots_transport::capability::{Availability, Maturity, SinkCapabilities}; -use radroots_transport::outcome::{DeliveryOutcome, DeliveryOutcomeKind}; -use radroots_transport::policy::{ - SatisfactionClass as SinkSatisfactionClass, SatisfactionPolicy as SinkSatisfactionPolicy, - TargetPolicy as SinkTargetPolicy, -}; -use radroots_transport::sink::{ - DeliveryPayload, DeliveryReceipt, DeliveryRequest, DeliveryTargetReceipt, EventSink, SinkStatus, -}; -use radroots_transport::target::{TargetLabel, TargetScope}; -use radroots_transport::{ - RadrootsTransportDeliveryTargetStatus, RadrootsTransportError, RadrootsTransportOutcome, - RadrootsTransportOutcomeKind, RadrootsTransportPayload, RadrootsTransportSatisfactionClass, - RadrootsTransportSatisfactionPolicy, Target, TargetSet, TransportId, -}; -use radroots_transport_nostr::{ - RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX, RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX, - RADROOTS_RELAY_FETCH_RAW_JSON_BYTE_LIMIT_MAX, RadrootsMockRelayFetchAdapter, - RadrootsMockRelayPublishAdapter, RadrootsNostrTransport, RadrootsOutboxPublishPolicy, - RadrootsRelayFetchEventAdmission, RadrootsRelayFetchEventValidStream, - RadrootsRelayFetchEventVerification, RadrootsRelayFetchEventVisibility, - RadrootsRelayFetchFilters, RadrootsRelayFetchItem, RadrootsRelayFetchMode, - RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchRequest, RadrootsRelayOutcome, - RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt, - RadrootsRelayPublishRequest, RadrootsRelayTargetSet, RadrootsRelayTransportError, - RadrootsRelayUrlPolicy, RelayUrl, fetch_and_ingest_relay_events, fetch_relay_events, - fetch_relay_events_blocking, publish_claimed_outbox_event, - publish_claimed_outbox_event_with_transport, publish_signed_event, - verified_signed_event_payload, -}; -use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; - -const FIXTURE_ALICE_SECRET_KEY_HEX: &str = - "10c5304d6c9ae3a1a16f7860f1cc8f5e3a76225a2663b3a989a0d775919b7df5"; -const FIXTURE_ALICE_PUBLIC_KEY_HEX: &str = - "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; -const RELAY_PRIMARY_WSS: &str = "wss://relay.example.com"; -const RELAY_SECONDARY_WSS: &str = "wss://relay-2.example.com"; -const RELAY_TERTIARY_WSS: &str = "wss://relay-3.example.com"; - -struct TransportFailurePublishAdapter; - -impl RadrootsRelayPublishAdapter for TransportFailurePublishAdapter { - fn publish<'a>( - &'a self, - _request: RadrootsRelayPublishRequest, - ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> - { - Box::pin(async { - Err(RadrootsRelayTransportError::Transport( - "adapter boundary unavailable".to_owned(), - )) - }) - } -} - -struct PartialPublishAdapter; - -impl RadrootsRelayPublishAdapter for PartialPublishAdapter { - fn publish<'a>( - &'a self, - _request: RadrootsRelayPublishRequest, - ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> - { - Box::pin(async { - Ok(vec![RadrootsRelayPublishRelayReceipt::attempted( - RELAY_PRIMARY_WSS, - RadrootsRelayOutcome::accepted(), - )]) - }) - } -} - -struct SlashSpelledRelayReceiptPublishAdapter; - -impl RadrootsRelayPublishAdapter for SlashSpelledRelayReceiptPublishAdapter { - fn publish<'a>( - &'a self, - _request: RadrootsRelayPublishRequest, - ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> - { - Box::pin(async { - Ok(vec![RadrootsRelayPublishRelayReceipt::attempted( - format!("{RELAY_PRIMARY_WSS}/"), - RadrootsRelayOutcome::accepted(), - )]) - }) - } -} - -struct NostrJsonFailurePublishAdapter; - -impl RadrootsRelayPublishAdapter for NostrJsonFailurePublishAdapter { - fn publish<'a>( - &'a self, - _request: RadrootsRelayPublishRequest, - ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> - { - Box::pin(async { - Err(RadrootsRelayTransportError::NostrEventJson( - "adapter rejected raw event".to_owned(), - )) - }) - } -} - -struct UnknownRelayReceiptPublishAdapter; - -impl RadrootsRelayPublishAdapter for UnknownRelayReceiptPublishAdapter { - fn publish<'a>( - &'a self, - request: RadrootsRelayPublishRequest, - ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> - { - Box::pin(async move { - let relay = request - .targets() - .relays() - .first() - .expect("fixture target") - .as_str() - .to_owned(); - Ok(vec![ - RadrootsRelayPublishRelayReceipt::attempted( - relay, - RadrootsRelayOutcome::accepted(), - ), - RadrootsRelayPublishRelayReceipt::attempted( - RELAY_TERTIARY_WSS, - RadrootsRelayOutcome::accepted(), - ), - ]) - }) - } -} - -struct DuplicateRelayReceiptPublishAdapter; - -impl RadrootsRelayPublishAdapter for DuplicateRelayReceiptPublishAdapter { - fn publish<'a>( - &'a self, - request: RadrootsRelayPublishRequest, - ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> - { - Box::pin(async move { - let relay = request - .targets() - .relays() - .first() - .expect("fixture target") - .as_str() - .to_owned(); - Ok(vec![ - RadrootsRelayPublishRelayReceipt::attempted( - relay.clone(), - RadrootsRelayOutcome::accepted(), - ), - RadrootsRelayPublishRelayReceipt::attempted( - format!("{relay}/"), - RadrootsRelayOutcome::accepted(), - ), - ]) - }) - } -} - -struct InvalidRelayReceiptPublishAdapter; - -impl RadrootsRelayPublishAdapter for InvalidRelayReceiptPublishAdapter { - fn publish<'a>( - &'a self, - _request: RadrootsRelayPublishRequest, - ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> - { - Box::pin(async { - Ok(vec![RadrootsRelayPublishRelayReceipt::attempted( - "not a relay URL", - RadrootsRelayOutcome::accepted(), - )]) - }) - } -} - -struct SkippedAcceptedRelayReceiptPublishAdapter; - -impl RadrootsRelayPublishAdapter for SkippedAcceptedRelayReceiptPublishAdapter { - fn publish<'a>( - &'a self, - request: RadrootsRelayPublishRequest, - ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> - { - Box::pin(async move { - let relay = request - .targets() - .relays() - .first() - .expect("fixture target") - .as_str(); - Ok(vec![RadrootsRelayPublishRelayReceipt::skipped( - relay, - RadrootsRelayOutcome::accepted(), - )]) - }) - } -} - -struct AttemptedSkippedRelayReceiptPublishAdapter; - -impl RadrootsRelayPublishAdapter for AttemptedSkippedRelayReceiptPublishAdapter { - fn publish<'a>( - &'a self, - request: RadrootsRelayPublishRequest, - ) -> BoxFuture<'a, Result<Vec<RadrootsRelayPublishRelayReceipt>, RadrootsRelayTransportError>> - { - Box::pin(async move { - let relay = request - .targets() - .relays() - .first() - .expect("fixture target") - .as_str(); - Ok(vec![RadrootsRelayPublishRelayReceipt::attempted( - relay, - RadrootsRelayOutcome::skipped_already_accepted("already accepted"), - )]) - }) - } -} - -#[derive(Clone)] -struct ScriptedTransport { - kind: TransportId, - outcomes: Vec<RadrootsTransportOutcome>, -} - -impl ScriptedTransport { - fn new(outcomes: Vec<RadrootsTransportOutcome>) -> Self { - Self { - kind: TransportId::NOSTR, - outcomes, - } - } - - fn with_kind(mut self, kind: TransportId) -> Self { - self.kind = kind; - self - } -} - -impl EventSink for ScriptedTransport { - fn status( - &self, - ) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, RadrootsTransportError>> { - Box::pin(async move { - Ok(SinkStatus::new( - self.kind, - true, - Maturity::Stable, - Availability::Available, - SinkCapabilities::DELIVER, - "scripted", - )) - }) - } - - fn deliver( - &self, - request: DeliveryRequest, - ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, RadrootsTransportError>> { - Box::pin(async move { - DeliveryReceipt::for_request( - &request, - request - .target_set() - .targets() - .iter() - .cloned() - .zip(self.outcomes.iter()) - .map(|(target, outcome)| { - DeliveryTargetReceipt::attempted( - target, - legacy_outcome_to_delivery_outcome(outcome), - ) - }) - .collect(), - ) - }) - } -} - -fn legacy_outcome_to_delivery_outcome(outcome: &RadrootsTransportOutcome) -> DeliveryOutcome { - match outcome.kind { - RadrootsTransportOutcomeKind::Delivered => DeliveryOutcome::delivered(), - RadrootsTransportOutcomeKind::Accepted - | RadrootsTransportOutcomeKind::DuplicateAccepted - | RadrootsTransportOutcomeKind::Forwarded - | RadrootsTransportOutcomeKind::StoredByGateway - | RadrootsTransportOutcomeKind::Seen => DeliveryOutcome::accepted(), - RadrootsTransportOutcomeKind::Timeout - | RadrootsTransportOutcomeKind::ConnectionFailed - | RadrootsTransportOutcomeKind::TransportUnavailable => DeliveryOutcome::unavailable(), - RadrootsTransportOutcomeKind::DeferredUntilImplemented - | RadrootsTransportOutcomeKind::Rejected - | RadrootsTransportOutcomeKind::RouteUnavailable - | RadrootsTransportOutcomeKind::PayloadTooLarge - | RadrootsTransportOutcomeKind::PolicyDenied => DeliveryOutcome::rejected(), - } -} - -#[derive(Clone, Copy)] -enum ForgedDeliveryReceipt { - RequestId, - TargetSet, -} - -struct ForgedReceiptTransport { - forged: ForgedDeliveryReceipt, -} - -impl EventSink for ForgedReceiptTransport { - fn status( - &self, - ) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, RadrootsTransportError>> { - Box::pin(async { - Ok(SinkStatus::new( - TransportId::NOSTR, - true, - Maturity::Stable, - Availability::Available, - SinkCapabilities::DELIVER, - "forged receipt fixture", - )) - }) - } - - fn deliver( - &self, - _request: DeliveryRequest, - ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, RadrootsTransportError>> { - let error = match self.forged { - ForgedDeliveryReceipt::RequestId => { - RadrootsTransportError::DeliveryReceiptRequestIdMismatch - } - ForgedDeliveryReceipt::TargetSet => { - RadrootsTransportError::DeliveryReceiptTargetSetMismatch - } - }; - Box::pin(async move { Err(error) }) - } -} - -fn fixture_keys() -> RadrootsNostrKeys { - let secret_key = - RadrootsNostrSecretKey::from_hex(FIXTURE_ALICE_SECRET_KEY_HEX).expect("secret key"); - RadrootsNostrKeys::new(secret_key) -} - -fn test_event_builder( - kind: u32, - content: impl Into<String>, - tags: Vec<Vec<String>>, -) -> EventBuilder { - let tags: Vec<_> = tags - .into_iter() - .filter(|tag| !tag.is_empty()) - .map(|mut tag| { - let key = tag.remove(0); - RadrootsNostrTag::custom(RadrootsNostrTagKind::Custom(key.into()), tag) - }) - .collect(); - EventBuilder::new( - RadrootsNostrKind::Custom(u16::try_from(kind).expect("test kind must fit NIP-01")), - content.into(), - ) - .tags(tags) - .allow_self_tagging() -} - -fn signed_post(content: &str) -> SignedEvent { - signed_event_with_kind_and_hashtag(content, KIND_POST, "soil") -} - -fn signed_ephemeral(content: &str) -> SignedEvent { - let raw_event = test_event_builder(KIND_GEOCHAT, content, Vec::new()) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) - .sign_with_keys(&fixture_keys()) - .expect("signed ephemeral event"); - let raw_json = raw_event.as_json(); - let wire = radroots_event::wire::Nip01EventWire::parse_json(&raw_json).expect("wire"); - SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") -} - -fn generic_draft(content: &str) -> EventDraft { - EventDraft::new( - "radroots.social.follow_list.v1", - KIND_FOLLOW, - 1_700_000_000, - Vec::new(), - serde_json::json!({ "label": content }).to_string(), - FIXTURE_ALICE_PUBLIC_KEY_HEX, - ) - .expect("generic draft") -} - -fn assert_outbox_publish_observations( - observations: &[RadrootsTransportObservationRow], - publish_ack_count: usize, -) { - assert_eq!(observations.len(), publish_ack_count + 1); - assert_eq!( - observations - .iter() - .filter(|observation| observation.observation_type - == RadrootsTransportObservationType::LocalImport - && observation.endpoint_uri.as_str() == "local:outbox") - .count(), - 1 - ); - assert_eq!( - observations - .iter() - .filter(|observation| observation.observation_type - == RadrootsTransportObservationType::PublishAck) - .count(), - publish_ack_count - ); -} - -fn signed_event_with_kind_and_hashtag(content: &str, kind: u32, hashtag: &str) -> SignedEvent { - let raw_event = signed_raw_event_with_kind_and_hashtag(content, kind, hashtag); - let raw_json = raw_event.as_json(); - let wire = radroots_event::wire::Nip01EventWire::parse_json(&raw_json).expect("wire"); - SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") -} - -fn signed_raw_event_with_kind_and_hashtag(content: &str, kind: u32, hashtag: &str) -> nostr::Event { - test_event_builder( - kind, - content, - vec![vec!["t".to_owned(), hashtag.to_owned()]], - ) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) - .sign_with_keys(&fixture_keys()) - .expect("signed event") -} - -async fn complete_claimed_signing( - outbox: &RadrootsOutbox, - claimed: &RadrootsOutboxClaimedEvent, - now_ms: i64, -) -> SignedEvent { - if let Some(signed_event) = claimed.signed_event.clone() { - return signed_event; - } - let signed_event = sign_frozen_draft(&fixture_keys(), &claimed.draft).expect("signed event"); - outbox - .complete_signing( - claimed.outbox_event_id, - claimed.claim_token.as_str(), - signed_event, - now_ms, - ) - .await - .expect("complete signing") -} - -fn nostr_target(relay_url: &str) -> Target { - Target::new(TransportId::NOSTR, relay_url).expect("nostr target") -} - -fn scoped_nostr_target(relay_url: &str, scope: &str, label: &str) -> Target { - Target::new_with_metadata( - TransportId::NOSTR, - relay_url, - Some(TargetScope::parse(scope).expect("target scope")), - Some(TargetLabel::parse(label).expect("target label")), - ) - .expect("scoped nostr target") -} - -fn reticulum_target() -> Target { - Target::new_with_metadata( - TransportId::RETICULUM, - "reticulum:local", - Some(TargetScope::parse("local").expect("Reticulum scope")), - None, - ) - .expect("Reticulum target") -} - -fn sink_delivery_request( - request_id: &str, - signed_event: &SignedEvent, - targets: Vec<Target>, - target_policy: SinkTargetPolicy, -) -> DeliveryRequest { - DeliveryRequest::new( - request_id, - DeliveryPayload::new(signed_event.clone()), - TargetSet::new(targets).expect("target set"), - SinkSatisfactionPolicy::new(SinkSatisfactionClass::Accepted, target_policy), - 10_000, - ) - .expect("delivery request") -} - -fn outbox_operation_input<I, S>( - draft: EventDraft, - relays: I, - satisfaction_policy: RadrootsTransportSatisfactionPolicy, -) -> RadrootsOutboxOperationInput -where - I: IntoIterator<Item = S>, - S: AsRef<str>, -{ - let targets = relays - .into_iter() - .map(|relay| nostr_target(relay.as_ref())) - .collect::<Vec<_>>(); - RadrootsOutboxOperationInput::new( - "publish_post", - draft, - RadrootsOutboxDeliveryPlanInput::new( - "transport.nostr.local", - 1, - satisfaction_policy, - targets, - ), - 1_000, - ) -} - -fn all_accepted_outbox_operation_input<I, S>( - draft: EventDraft, - relays: I, -) -> RadrootsOutboxOperationInput -where - I: IntoIterator<Item = S>, - S: AsRef<str>, -{ - outbox_operation_input( - draft, - relays, - RadrootsTransportSatisfactionPolicy::all_accepted(), - ) -} - -fn unsupported_raw_event() -> String { - let event = test_event_builder(999, "unsupported", Vec::new()) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_001)) - .sign_with_keys(&fixture_keys()) - .expect("signed unsupported event"); - event.as_json() -} - -fn invalid_contract_shape_raw_event() -> String { - let event = test_event_builder( - KIND_POST, - "invalid reply", - vec![ - vec!["t".to_owned(), "soil".to_owned()], - vec![ - "e".to_owned(), - "invalid-event-id".to_owned(), - String::new(), - "root".to_owned(), - ], - ], - ) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_002)) - .sign_with_keys(&fixture_keys()) - .expect("signed contract-invalid event"); - event.as_json() -} - -fn post_relay_fetch_filter(limit: usize) -> RadrootsNostrFilter { - with_tag( - RadrootsNostrFilter::new() - .kind(RadrootsNostrKind::Custom( - u16::try_from(KIND_POST).expect("post kind must fit NIP-01"), - )) - .limit(limit), - "t", - vec!["soil".to_owned()], - ) - .expect("post relay fetch filter") -} - -fn unsupported_relay_fetch_filter(limit: usize) -> RadrootsNostrFilter { - RadrootsNostrFilter::new() - .kind(RadrootsNostrKind::Custom(999)) - .limit(limit) -} - -fn fixture_relay_targets() -> RadrootsRelayTargetSet { - RadrootsRelayTargetSet::new( - [RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS, RELAY_TERTIARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("fixture relay targets") -} - -fn primary_relay_target() -> RadrootsRelayTargetSet { - RadrootsRelayTargetSet::new([RELAY_PRIMARY_WSS], RadrootsRelayUrlPolicy::Public) - .expect("primary relay target") -} - -fn fixture_relay_fetch_request( - observed_at_ms: i64, - max_events: usize, -) -> RadrootsRelayFetchRequest { - RadrootsRelayFetchRequest::fetch( - observed_at_ms, - max_events, - fixture_relay_targets(), - [ - post_relay_fetch_filter(max_events), - unsupported_relay_fetch_filter(max_events), - ], - ) - .expect("fixture relay fetch request") -} - -fn post_relay_fetch_request(observed_at_ms: i64, max_events: usize) -> RadrootsRelayFetchRequest { - RadrootsRelayFetchRequest::fetch( - observed_at_ms, - max_events, - fixture_relay_targets(), - [post_relay_fetch_filter(max_events)], - ) - .expect("post relay fetch request") -} - -fn tampered_raw_event() -> String { - let signed = signed_post("trusted"); - let mut value = serde_json::from_str::<serde_json::Value>(signed.raw_json()).expect("raw json"); - value["content"] = serde_json::Value::String("tampered".to_owned()); - serde_json::to_string(&value).expect("tampered json") -} - -#[test] -fn relay_url_validation_and_target_normalization() { - let relay = - RelayUrl::parse("wss://Relay.Example.com", RadrootsRelayUrlPolicy::Public).expect("relay"); - assert_eq!(relay.as_str(), RELAY_PRIMARY_WSS); - assert_eq!(relay.clone().into_string(), RELAY_PRIMARY_WSS); - let relay_path = RelayUrl::parse( - "wss://Relay.Example.com/nostr", - RadrootsRelayUrlPolicy::Public, - ) - .expect("relay path"); - assert_eq!(relay_path.as_str(), "wss://relay.example.com/nostr"); - - assert!(RelayUrl::parse("ws://127.0.0.1:7777", RadrootsRelayUrlPolicy::Public).is_err()); - let local = RelayUrl::parse("ws://localhost:7777", RadrootsRelayUrlPolicy::Localhost) - .expect("local relay"); - assert_eq!(local.as_str(), "ws://localhost:7777"); - let local_ipv4 = RelayUrl::parse("ws://127.0.0.1:7777", RadrootsRelayUrlPolicy::Localhost) - .expect("local ipv4 relay"); - assert_eq!(local_ipv4.as_str(), "ws://127.0.0.1:7777"); - let local_ipv6 = RelayUrl::parse("ws://[::1]:7777", RadrootsRelayUrlPolicy::Localhost) - .expect("local ipv6 relay"); - assert_eq!(local_ipv6.as_str(), "ws://[::1]:7777"); - assert!(RelayUrl::parse("ws://example.com", RadrootsRelayUrlPolicy::Localhost).is_err()); - assert!(RelayUrl::parse("ws://192.168.1.10:7777", RadrootsRelayUrlPolicy::Localhost).is_err()); - assert!(RelayUrl::parse("wss://192.168.1.10", RadrootsRelayUrlPolicy::Localhost).is_err()); - assert!(RelayUrl::parse("wss://relay.example.com", RadrootsRelayUrlPolicy::Localhost).is_err()); - assert!(RelayUrl::parse("wss://localhost", RadrootsRelayUrlPolicy::Localhost).is_ok()); - assert!(matches!( - RelayUrl::parse("wss://127.0.0.1", RadrootsRelayUrlPolicy::Public), - Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. }) - )); - assert!(matches!( - RelayUrl::parse("wss://10.1.2.3", RadrootsRelayUrlPolicy::Public), - Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. }) - )); - assert!(matches!( - RelayUrl::parse("wss://[::1]", RadrootsRelayUrlPolicy::Public), - Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. }) - )); - assert!(matches!( - RelayUrl::parse("wss://[fd00::1]", RadrootsRelayUrlPolicy::Public), - Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. }) - )); - for relay_url in [ - "wss://0.0.0.0", - "wss://169.254.1.2", - "wss://224.0.0.1", - "wss://255.255.255.255", - "wss://100.64.0.1", - "wss://192.0.0.8", - "wss://192.88.99.2", - "wss://198.18.0.1", - "wss://240.0.0.1", - "wss://[::]", - "wss://[64:ff9b::7f00:1]", - "wss://[64:ff9b::a00:1]", - "wss://[64:ff9b::5db8:d822]", - "wss://[64:ff9b:1::1]", - "wss://[100::1]", - "wss://[100:0:0:1::1]", - "wss://[ff02::1]", - "wss://[fe80::1]", - "wss://[2001:db8::1]", - "wss://[2001:1::1]", - "wss://[2002::1]", - "wss://[3fff::1]", - "wss://[5f00::1]", - "wss://[::ffff:192.168.1.10]", - "wss://localhost", - "wss://relay.local", - "wss://relay.home.arpa", - "wss://relay", - ] { - assert!(matches!( - RelayUrl::parse(relay_url, RadrootsRelayUrlPolicy::Public), - Err(RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. }) - )); - } - let public_relay = RelayUrl::parse("wss://relay.example.com", RadrootsRelayUrlPolicy::Public) - .expect("public relay"); - public_relay - .validate_public_resolved_ip_addrs([IpAddr::V4(Ipv4Addr::new(93, 184, 216, 34))]) - .expect("public resolved ip"); - assert!(matches!( - public_relay - .validate_public_resolved_ip_addrs([IpAddr::V4(Ipv4Addr::new(192, 168, 1, 10))]), - Err(RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination { .. }) - )); - assert!(matches!( - public_relay.validate_public_resolved_ip_addrs([IpAddr::V6( - "::ffff:192.168.1.10" - .parse::<Ipv6Addr>() - .expect("mapped ipv6") - )]), - Err(RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination { .. }) - )); - public_relay - .validate_public_resolved_ip_addrs([IpAddr::V6( - "2001:4860:4860::8888" - .parse::<Ipv6Addr>() - .expect("public ipv6"), - )]) - .expect("public resolved ipv6"); - assert!(matches!( - public_relay.validate_public_resolved_ip_addrs([IpAddr::V6( - "64:ff9b::5db8:d822" - .parse::<Ipv6Addr>() - .expect("translation prefix"), - )]), - Err(RadrootsRelayTransportError::RelayUrlResolvedForbiddenDestination { .. }) - )); - assert!(matches!( - public_relay.validate_public_resolved_ip_addrs(Vec::<IpAddr>::new()), - Err(RadrootsRelayTransportError::RelayUrlResolvedNoAddresses { .. }) - )); - - assert!(RelayUrl::parse("https://relay.example.com", RadrootsRelayUrlPolicy::Public).is_err()); - assert!( - RelayUrl::parse( - "wss://user@relay.example.com", - RadrootsRelayUrlPolicy::Public - ) - .is_err() - ); - assert!(matches!( - RelayUrl::parse( - "wss://user:password@relay.example.com", - RadrootsRelayUrlPolicy::Public - ), - Err(RadrootsRelayTransportError::RelayUrlUserinfo { .. }) - )); - assert!(matches!( - RelayUrl::parse( - "wss://:password@relay.example.com", - RadrootsRelayUrlPolicy::Public - ), - Err(RadrootsRelayTransportError::RelayUrlUserinfo { .. }) - )); - assert!( - RelayUrl::parse( - "wss://relay.example.com:bad", - RadrootsRelayUrlPolicy::Public - ) - .is_err() - ); - assert!(RelayUrl::parse("wss://", RadrootsRelayUrlPolicy::Public).is_err()); - assert!(matches!( - RelayUrl::parse("radroots:relay", RadrootsRelayUrlPolicy::Public), - Err(RadrootsRelayTransportError::EmptyRelayHost { .. }) - )); - assert!(matches!( - RelayUrl::parse("relay.example.com", RadrootsRelayUrlPolicy::Public), - Err(RadrootsRelayTransportError::RelayUrlParse { .. }) - )); - assert!( - RelayUrl::parse( - "wss://relay.example.com?subscription=1", - RadrootsRelayUrlPolicy::Public - ) - .is_err() - ); - assert!( - RelayUrl::parse( - "wss://relay.example.com#fragment", - RadrootsRelayUrlPolicy::Public - ) - .is_err() - ); - - let targets = RadrootsRelayTargetSet::new( - vec![RELAY_TERTIARY_WSS, RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("targets"); - assert_eq!( - targets.relay_strings(), - vec![ - RELAY_TERTIARY_WSS.to_owned(), - RELAY_PRIMARY_WSS.to_owned(), - RELAY_SECONDARY_WSS.to_owned() - ] - ); - - let from_urls = RadrootsRelayTargetSet::from_urls(vec![ - relay_path.clone(), - RelayUrl::parse(RELAY_SECONDARY_WSS, RadrootsRelayUrlPolicy::Public).expect("secondary"), - ]) - .expect("from urls"); - assert_eq!(from_urls.len(), 2); - assert!(!from_urls.is_empty()); - assert_eq!(from_urls.relays()[0], relay_path); - assert_eq!( - from_urls.relays()[0].to_string(), - "wss://relay.example.com/nostr" - ); - assert!(matches!( - RadrootsRelayTargetSet::new(Vec::<&str>::new(), RadrootsRelayUrlPolicy::Public), - Err(RadrootsRelayTransportError::EmptyTargetSet) - )); - assert!(matches!( - RadrootsRelayTargetSet::from_urls(Vec::new()), - Err(RadrootsRelayTransportError::EmptyTargetSet) - )); - assert!(matches!( - RadrootsRelayTargetSet::new( - [RELAY_PRIMARY_WSS, "WSS://Relay.Example.com/"], - RadrootsRelayUrlPolicy::Public, - ), - Err(RadrootsRelayTransportError::DuplicateRelayUrl { .. }) - )); - assert!(matches!( - RadrootsRelayTargetSet::from_urls(vec![relay_path.clone(), relay_path]), - Err(RadrootsRelayTransportError::DuplicateRelayUrl { .. }) - )); -} - -#[test] -fn transport_target_and_relay_adapter_share_canonical_url_identity() { - for (raw, policy) in [ - ( - "WSS://Relay.Example.com:443/", - RadrootsRelayUrlPolicy::Public, - ), - ( - "wss://relay.example.com/nostr/%2Ffeed", - RadrootsRelayUrlPolicy::Public, - ), - ( - "WSS://[2001:4860:4860:0:0:0:0:8888]:443/", - RadrootsRelayUrlPolicy::Public, - ), - ("ws://LOCALHOST:80/", RadrootsRelayUrlPolicy::Localhost), - ( - "ws://[0:0:0:0:0:0:0:1]:7777/", - RadrootsRelayUrlPolicy::Localhost, - ), - ] { - let target = Target::new(TransportId::NOSTR, raw).expect("transport target"); - let relay = RelayUrl::parse(raw, policy).expect("relay URL"); - assert_eq!(target.uri().as_str(), relay.as_str(), "{raw}"); - } - - for raw in [ - " wss://relay.example.com", - "wss://relay.example.com ", - "wss://relay.example.com/a/./b", - "wss://relay.example.com/a/../b", - "wss://relay.example.com/a/%2E/b", - "wss://relay.example.com\\path", - "wss://relay.example.com:0", - "wss://relay.example.com:01", - "wss://relay.example.com.", - "wss://xn--fa-hia.example.com", - "wss://faß.example.com", - "wss://%65xample.com", - "wss://relay.example.com/%2f", - "wss://relay.example.com?subscription=1", - ] { - assert!( - Target::new(TransportId::NOSTR, raw).is_err(), - "transport target accepted {raw}" - ); - assert!( - RelayUrl::parse(raw, RadrootsRelayUrlPolicy::Public).is_err(), - "relay adapter accepted {raw}" - ); - } -} - -#[test] -fn outcome_prefix_classification_covers_required_kinds() { - let cases = [ - ("blocked: policy", RadrootsRelayOutcomeKind::Blocked), - ( - "rate-limited: slow down", - RadrootsRelayOutcomeKind::RateLimited, - ), - ("invalid: bad event", RadrootsRelayOutcomeKind::Invalid), - ("pow: difficulty 24", RadrootsRelayOutcomeKind::PowRequired), - ( - "restricted: group write denied", - RadrootsRelayOutcomeKind::Restricted, - ), - ( - "auth-required: challenge", - RadrootsRelayOutcomeKind::AuthRequired, - ), - ("mute: pubkey muted", RadrootsRelayOutcomeKind::Muted), - ( - "unsupported: event kind", - RadrootsRelayOutcomeKind::Unsupported, - ), - ( - "payment-required: paid relay", - RadrootsRelayOutcomeKind::PaymentRequired, - ), - ( - "duplicate: already have it", - RadrootsRelayOutcomeKind::DuplicateAccepted, - ), - ("error: relay failed", RadrootsRelayOutcomeKind::Error), - ("timeout: no OK", RadrootsRelayOutcomeKind::Timeout), - ("strange relay text", RadrootsRelayOutcomeKind::Unknown), - ]; - - for (message, kind) in cases { - let outcome = RadrootsRelayOutcome::classify(message); - assert_eq!(outcome.kind, kind); - } - let labels = [ - (RadrootsRelayOutcomeKind::Accepted, "accepted"), - ( - RadrootsRelayOutcomeKind::DuplicateAccepted, - "duplicate_accepted", - ), - (RadrootsRelayOutcomeKind::Blocked, "blocked"), - (RadrootsRelayOutcomeKind::RateLimited, "rate_limited"), - (RadrootsRelayOutcomeKind::Invalid, "invalid"), - (RadrootsRelayOutcomeKind::PowRequired, "pow_required"), - (RadrootsRelayOutcomeKind::Restricted, "restricted"), - (RadrootsRelayOutcomeKind::AuthRequired, "auth_required"), - (RadrootsRelayOutcomeKind::Muted, "muted"), - (RadrootsRelayOutcomeKind::Unsupported, "unsupported"), - ( - RadrootsRelayOutcomeKind::PaymentRequired, - "payment_required", - ), - (RadrootsRelayOutcomeKind::Error, "error"), - (RadrootsRelayOutcomeKind::Timeout, "timeout"), - ( - RadrootsRelayOutcomeKind::ConnectionFailed, - "connection_failed", - ), - ( - RadrootsRelayOutcomeKind::RelayUrlRejected, - "relay_url_rejected", - ), - ( - RadrootsRelayOutcomeKind::SkippedAlreadyAccepted, - "skipped_already_accepted", - ), - (RadrootsRelayOutcomeKind::Unknown, "unknown"), - ]; - for (kind, label) in labels { - assert_eq!(kind.as_str(), label); - } - - assert!(RadrootsRelayOutcome::classify("duplicate: already have it").counts_toward_quorum()); - assert!( - RadrootsRelayOutcome::skipped_already_accepted("already accepted").counts_toward_quorum() - ); - assert!(RadrootsRelayOutcome::classify("auth-required: challenge").is_retryable()); - assert!(RadrootsRelayOutcome::classify("restricted: denied").is_terminal_failure()); - assert!(RadrootsRelayOutcome::relay_url_rejected("unsafe relay").is_terminal_failure()); - assert!(RadrootsRelayOutcome::classify("mute: pubkey muted").is_terminal_failure()); - assert_eq!( - RadrootsRelayOutcome::accepted().to_transport_outcome().kind, - radroots_transport::RadrootsTransportOutcomeKind::Accepted - ); - assert_eq!( - RadrootsRelayOutcome::accepted() - .to_transport_outcome() - .status, - radroots_transport::RadrootsTransportDeliveryTargetStatus::Accepted - ); - assert_eq!( - RadrootsRelayOutcome::timeout("timeout: no OK") - .to_transport_outcome() - .kind, - radroots_transport::RadrootsTransportOutcomeKind::Timeout - ); - assert_eq!( - RadrootsRelayOutcome::timeout("timeout: no OK") - .to_transport_outcome() - .status, - radroots_transport::RadrootsTransportDeliveryTargetStatus::FailedRetryable - ); - assert_eq!( - RadrootsRelayOutcome::classify("restricted: denied") - .to_transport_outcome() - .kind, - radroots_transport::RadrootsTransportOutcomeKind::Rejected - ); - assert_eq!( - RadrootsRelayOutcome::classify("restricted: denied") - .to_transport_outcome() - .status, - radroots_transport::RadrootsTransportDeliveryTargetStatus::FailedTerminal - ); - assert_eq!( - RadrootsRelayOutcome::relay_url_rejected("unsafe") - .to_transport_outcome() - .kind, - radroots_transport::RadrootsTransportOutcomeKind::RouteUnavailable - ); - assert_eq!( - RadrootsRelayOutcome::connection_failed("offline") - .kind - .as_str(), - "connection_failed" - ); - assert_eq!( - RadrootsRelayOutcome::unknown("adapter omitted receipt") - .to_transport_outcome() - .kind, - radroots_transport::RadrootsTransportOutcomeKind::TransportUnavailable - ); - assert_eq!( - RadrootsRelayOutcome::relay_url_rejected("unsafe") - .kind - .as_str(), - "relay_url_rejected" - ); -} - -#[test] -fn relay_transport_error_wraps_transport_contract_errors() { - let error = RadrootsRelayTransportError::from( - radroots_transport::RadrootsTransportError::EmptyTargetSet, - ); - - assert_eq!( - error.to_string(), - "Transport contract error: transport target set is empty" - ); -} - -#[tokio::test] -async fn mock_publish_preserves_exact_raw_json_and_counts_outcomes() { - let signed = signed_post("hello"); - let targets = RadrootsRelayTargetSet::new( - vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS, RELAY_TERTIARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("targets"); - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome( - RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::classify("duplicate: already have it"), - ) - .with_outcome( - RELAY_TERTIARY_WSS, - RadrootsRelayOutcome::classify("auth-required: challenge"), - ); - - let receipt = publish_signed_event( - &adapter, - radroots_transport_nostr::RadrootsRelayPublishRequest::new(signed.clone(), targets, 1_000) - .expect("publish request") - .with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::quorum_accepted(2)), - ) - .await - .expect("publish"); - - assert_eq!( - adapter.captured_raw_events(), - vec![signed.raw_json().to_owned()] - ); - assert_eq!(receipt.attempted_count, 3); - assert_eq!(receipt.accepted_count, 2); - assert_eq!(receipt.retryable_count, 1); - assert!(receipt.quorum_met); - serde_json::to_string(&receipt).expect("receipt json"); -} - -#[tokio::test] -async fn nostr_transport_facade_delivers_signed_event_payloads() { - let signed = signed_post("facade payload"); - let adapter = RadrootsMockRelayPublishAdapter::new(); - let expected_status = SinkStatus::new( - TransportId::NOSTR, - true, - Maturity::Stable, - Availability::Available, - SinkCapabilities::DELIVER, - "fixture ready", - ); - let transport = RadrootsNostrTransport::new(&adapter).with_status(expected_status.clone()); - assert!(transport.adapter().captured_raw_events().is_empty()); - let target = nostr_target(RELAY_PRIMARY_WSS); - let request = sink_delivery_request( - "facade-request-1", - &signed, - vec![target.clone()], - SinkTargetPolicy::all(), - ); - - let receipt = transport.deliver(request.clone()).await.expect("delivery"); - let status = transport.status().await.expect("status"); - - assert_eq!( - adapter.captured_raw_events(), - vec![signed.raw_json().to_owned()] - ); - assert_eq!(status, expected_status); - assert_eq!(receipt.request_id().as_str(), "facade-request-1"); - assert_eq!(receipt.target_receipts().len(), 1); - assert_eq!(receipt.target_receipts()[0].target(), &target); - assert_eq!( - receipt.target_receipts()[0].outcome().kind(), - DeliveryOutcomeKind::Accepted - ); - assert!(receipt.is_satisfied(&request).expect("satisfaction")); -} - -#[test] -fn verified_signed_event_payload_preserves_transport_payload_identity() { - let signed = signed_post("verified payload"); - let payload = verified_signed_event_payload(&signed).expect("verified payload"); - let RadrootsTransportPayload::SignedEventJson { - event_id, - raw_json, - digest, - } = payload - else { - panic!("signed event payload expected"); - }; - - assert_eq!(event_id, signed.id_str()); - assert_eq!(raw_json, signed.raw_json().to_owned()); - assert_eq!(digest.len(), 64); -} - -#[tokio::test] -async fn nostr_transport_sink_rejects_non_nostr_targets() { - let signed = signed_post("facade rejected target"); - let transport = RadrootsNostrTransport::new(RadrootsMockRelayPublishAdapter::new()); - let request = sink_delivery_request( - "facade-request-target", - &signed, - vec![reticulum_target()], - SinkTargetPolicy::all(), - ); - let error = transport - .deliver(request) - .await - .expect_err("target rejected"); - assert_eq!(error, RadrootsTransportError::InvalidTargetUri); -} - -#[tokio::test] -async fn nostr_transport_facade_preserves_adapter_failure_and_omission_evidence() { - let signed = signed_post("facade failures"); - let targets = vec![ - nostr_target(RELAY_PRIMARY_WSS), - nostr_target(RELAY_SECONDARY_WSS), - ]; - - let transport = RadrootsNostrTransport::new(TransportFailurePublishAdapter); - let failed = transport - .deliver(sink_delivery_request( - "facade-transport-failure", - &signed, - targets.clone(), - SinkTargetPolicy::all(), - )) - .await - .expect("failure receipts"); - assert_eq!(failed.target_receipts().len(), 2); - assert!(failed.target_receipts().iter().all(|receipt| { - receipt.was_attempted() - && receipt.outcome().kind() == DeliveryOutcomeKind::Unavailable - && receipt.outcome().is_retryable() - })); - - let partial = RadrootsNostrTransport::new(PartialPublishAdapter) - .deliver(sink_delivery_request( - "facade-partial", - &signed, - targets.clone(), - SinkTargetPolicy::all(), - )) - .await - .expect("partial receipts"); - assert_eq!(partial.target_receipts().len(), 2); - assert_eq!( - partial.target_receipts()[1].outcome().kind(), - DeliveryOutcomeKind::Unavailable - ); - - let error = RadrootsNostrTransport::new(NostrJsonFailurePublishAdapter) - .deliver(sink_delivery_request( - "facade-json-failure", - &signed, - targets, - SinkTargetPolicy::all(), - )) - .await - .expect_err("adapter JSON error"); - assert_eq!(error, RadrootsTransportError::InvalidPayloadBytes); -} - -#[tokio::test] -async fn nostr_transport_facade_matches_canonical_equivalent_relay_receipts() { - let signed = signed_post("facade canonical receipt"); - let target = nostr_target(RELAY_PRIMARY_WSS); - let sink_policy = SinkTargetPolicy::required(vec![target.fingerprint().clone()]) - .expect("required target policy"); - let legacy_policy = RadrootsTransportSatisfactionPolicy::required_targets( - RadrootsTransportSatisfactionClass::Accepted, - vec![target.fingerprint().clone()], - ) - .expect("required target policy"); - let transport = RadrootsNostrTransport::new(SlashSpelledRelayReceiptPublishAdapter); - let receipt = transport - .deliver(sink_delivery_request( - "facade-canonical-receipt", - &signed, - vec![target.clone()], - sink_policy, - )) - .await - .expect("delivery"); - - assert_eq!(receipt.target_receipts().len(), 1); - assert_eq!(receipt.target_receipts()[0].target(), &target); - assert_eq!( - receipt.target_receipts()[0].outcome().kind(), - DeliveryOutcomeKind::Accepted - ); - - let relay_receipt = publish_signed_event( - &SlashSpelledRelayReceiptPublishAdapter, - RadrootsRelayPublishRequest::new( - signed, - RadrootsRelayTargetSet::new(vec![RELAY_PRIMARY_WSS], RadrootsRelayUrlPolicy::Public) - .expect("targets"), - 1_070, - ) - .expect("publish request") - .with_satisfaction_policy(legacy_policy), - ) - .await - .expect("relay publish"); - assert!(relay_receipt.quorum_met); -} - -#[tokio::test] -async fn nostr_transport_facade_preserves_scoped_duplicate_target_metadata() { - let signed = signed_post("facade scoped duplicate"); - let adapter = RadrootsMockRelayPublishAdapter::new(); - let transport = RadrootsNostrTransport::new(&adapter); - let first = scoped_nostr_target(RELAY_PRIMARY_WSS, "local_food_buyers", "buyers"); - let second = scoped_nostr_target(RELAY_PRIMARY_WSS, "local_food_farmers", "farmers"); - let policy = SinkTargetPolicy::required(vec![ - first.fingerprint().clone(), - second.fingerprint().clone(), - ]) - .expect("required targets"); - let request = sink_delivery_request( - "facade-request-scoped", - &signed, - vec![first.clone(), second.clone()], - policy, - ); - - let receipt = transport.deliver(request.clone()).await.expect("delivery"); - - assert_eq!(receipt.target_receipts().len(), 2); - assert_eq!(receipt.target_receipts()[0].target(), &first); - assert_eq!(receipt.target_receipts()[1].target(), &second); - assert!(receipt.is_satisfied(&request).expect("satisfaction")); - assert_eq!(adapter.captured_raw_events().len(), 1); -} - -#[tokio::test] -async fn publish_receipts_track_terminal_skipped_and_adapter_errors() { - let signed = signed_post("terminal"); - let targets = RadrootsRelayTargetSet::new( - vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("targets"); - let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome( - RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::classify("restricted: group write denied"), - ); - - let receipt = publish_signed_event( - &adapter, - RadrootsRelayPublishRequest::new(signed.clone(), targets, 1_050) - .expect("publish request") - .with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::all_accepted()), - ) - .await - .expect("publish"); - - assert_eq!(receipt.event_id, signed.id_str()); - assert_eq!(receipt.attempted_count, 2); - assert_eq!(receipt.accepted_count, 1); - assert_eq!(receipt.retryable_count, 0); - assert_eq!(receipt.terminal_count, 1); - assert_eq!(receipt.quorum, 2); - assert!(!receipt.quorum_met); - - let skipped = RadrootsRelayPublishRelayReceipt::skipped( - RELAY_TERTIARY_WSS, - RadrootsRelayOutcome::timeout("timeout: no OK"), - ); - assert_eq!(skipped.relay_url, RELAY_TERTIARY_WSS); - assert!(!skipped.attempted); - assert_eq!(skipped.outcome.kind, RadrootsRelayOutcomeKind::Timeout); - - let error = publish_signed_event( - &TransportFailurePublishAdapter, - RadrootsRelayPublishRequest::new( - signed, - RadrootsRelayTargetSet::new(vec![RELAY_PRIMARY_WSS], RadrootsRelayUrlPolicy::Public) - .expect("targets"), - 1_060, - ) - .expect("publish request"), - ) - .await - .expect_err("transport failure"); - assert!(matches!(error, RadrootsRelayTransportError::Transport(_))); -} - -#[tokio::test] -async fn publish_required_target_policy_uses_relay_fingerprints() { - let signed = signed_post("required relay"); - let required_target = - Target::new(TransportId::NOSTR, RELAY_PRIMARY_WSS).expect("required target"); - let targets = RadrootsRelayTargetSet::new( - vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("targets"); - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome( - RELAY_PRIMARY_WSS, - RadrootsRelayOutcome::classify("restricted: required relay rejected"), - ) - .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted()); - - let receipt = publish_signed_event( - &adapter, - RadrootsRelayPublishRequest::new(signed, targets, 1_070) - .expect("publish request") - .with_satisfaction_policy( - RadrootsTransportSatisfactionPolicy::required_targets( - RadrootsTransportSatisfactionClass::Accepted, - vec![required_target.fingerprint().clone()], - ) - .expect("required relay policy"), - ), - ) - .await - .expect("publish"); - - assert_eq!(receipt.accepted_count, 1); - assert_eq!(receipt.quorum, 1); - assert!(!receipt.quorum_met); -} - -#[tokio::test] -async fn publish_all_policy_uses_requested_target_count() { - let signed = signed_post("partial adapter"); - let targets = RadrootsRelayTargetSet::new( - vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("targets"); - - let receipt = publish_signed_event( - &PartialPublishAdapter, - RadrootsRelayPublishRequest::new(signed.clone(), targets.clone(), 1_080) - .expect("publish request") - .with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::all_accepted()), - ) - .await - .expect("publish"); - - assert_eq!(receipt.attempted_count, 1); - assert_eq!(receipt.accepted_count, 1); - assert_eq!(receipt.retryable_count, 1); - assert_eq!(receipt.quorum, 2); - assert!(!receipt.quorum_met); - assert_eq!(receipt.relays.len(), 2); - assert!(!receipt.relays[1].attempted); - assert_eq!( - receipt.relays[1].outcome.kind, - RadrootsRelayOutcomeKind::Unknown - ); - - let no_wait = publish_signed_event( - &PartialPublishAdapter, - RadrootsRelayPublishRequest::new(signed, targets, 1_081) - .expect("publish request") - .with_satisfaction_policy(RadrootsTransportSatisfactionPolicy::NoWait), - ) - .await - .expect("no-wait publish"); - assert_eq!(no_wait.quorum, 0); - assert!(no_wait.quorum_met); -} - -#[tokio::test] -async fn publish_rejects_untrusted_adapter_receipt_provenance() { - let signed = signed_post("adapter provenance"); - let request = || { - RadrootsRelayPublishRequest::new(signed.clone(), primary_relay_target(), 1_090) - .expect("publish request") - }; - - assert!(matches!( - publish_signed_event(&UnknownRelayReceiptPublishAdapter, request()).await, - Err(RadrootsRelayTransportError::UnexpectedPublishReceiptRelayUrl { url }) - if url == RELAY_TERTIARY_WSS - )); - assert!(matches!( - publish_signed_event(&DuplicateRelayReceiptPublishAdapter, request()).await, - Err(RadrootsRelayTransportError::DuplicatePublishReceiptRelayUrl { url }) - if url == RELAY_PRIMARY_WSS - )); - assert!(matches!( - publish_signed_event(&InvalidRelayReceiptPublishAdapter, request()).await, - Err(RadrootsRelayTransportError::InvalidPublishReceiptRelayUrl { url, .. }) - if url == "not a relay URL" - )); - assert!(matches!( - publish_signed_event(&SkippedAcceptedRelayReceiptPublishAdapter, request()).await, - Err(RadrootsRelayTransportError::InvalidPublishReceiptAttemptState { url }) - if url == RELAY_PRIMARY_WSS - )); - assert!(matches!( - publish_signed_event(&AttemptedSkippedRelayReceiptPublishAdapter, request()).await, - Err(RadrootsRelayTransportError::InvalidPublishReceiptAttemptState { url }) - if url == RELAY_PRIMARY_WSS - )); -} - -#[test] -fn relay_publish_request_rejects_negative_time() { - let signed = signed_post("negative publish time"); - - assert!(matches!( - RadrootsRelayPublishRequest::new(signed, primary_relay_target(), -1), - Err(RadrootsRelayTransportError::InvalidTimestamp { - field: "now_ms", - value: -1, - }) - )); -} - -#[test] -fn relay_publish_request_seals_fields_and_validates_idempotency_keys() { - let signed = signed_post("sealed publish request"); - let request = RadrootsRelayPublishRequest::new(signed.clone(), primary_relay_target(), 7) - .expect("publish request"); - assert_eq!(request.signed_event(), &signed); - assert_eq!(request.targets().len(), 1); - assert_eq!( - request.satisfaction_policy(), - &RadrootsTransportSatisfactionPolicy::all_accepted() - ); - assert_eq!(request.idempotency_key(), None); - assert_eq!(request.now_ms(), 7); - - for invalid in ["", " ", " leading", "trailing ", "line\nbreak"] { - assert!(matches!( - request.clone().try_with_idempotency_key(invalid), - Err(RadrootsRelayTransportError::InvalidIdempotencyKey { .. }) - )); - } - assert!(matches!( - request.clone().try_with_idempotency_key("x".repeat( - radroots_transport_nostr::RADROOTS_RELAY_PUBLISH_IDEMPOTENCY_KEY_MAX_BYTES + 1, - ),), - Err(RadrootsRelayTransportError::InvalidIdempotencyKey { .. }) - )); - let request = request - .try_with_idempotency_key("publish-7") - .expect("idempotency key"); - assert_eq!(request.idempotency_key(), Some("publish-7")); -} - -#[tokio::test] -async fn relay_publish_request_rejects_unrequested_required_target_before_adapter() { - let signed = signed_post("missing required target"); - let required = Target::new(TransportId::NOSTR, RELAY_SECONDARY_WSS) - .expect("required target") - .fingerprint() - .clone(); - let request = RadrootsRelayPublishRequest::new(signed, primary_relay_target(), 8) - .expect("publish request") - .with_satisfaction_policy( - RadrootsTransportSatisfactionPolicy::required_targets( - RadrootsTransportSatisfactionClass::Accepted, - vec![required.clone()], - ) - .expect("required policy"), - ); - let adapter = RadrootsMockRelayPublishAdapter::new(); - - assert!(matches!( - publish_signed_event(&adapter, request).await, - Err(RadrootsRelayTransportError::RequiredTargetNotRequested { fingerprint }) - if fingerprint == required.as_str() - )); - assert!(adapter.captured_raw_events().is_empty()); -} - -#[test] -fn fetch_requests_reject_empty_filter_sets() { - assert!(matches!( - RadrootsRelayFetchRequest::fetch( - 1_000, - 10, - primary_relay_target(), - Vec::<RadrootsNostrFilter>::new(), - ), - Err(RadrootsRelayTransportError::EmptyFetchFilters) - )); - assert!(matches!( - RadrootsRelayFetchRequest::subscription( - 1_000, - 10, - primary_relay_target(), - Vec::<RadrootsNostrFilter>::new(), - ), - Err(RadrootsRelayTransportError::EmptyFetchFilters) - )); -} - -#[test] -fn fetch_requests_reject_zero_limits_and_timeouts() { - let filter = post_relay_fetch_filter(1); - let filters = RadrootsRelayFetchFilters::new([filter.clone()]).expect("filters"); - let as_ref_filters: &[RadrootsNostrFilter] = filters.as_ref(); - assert_eq!(as_ref_filters.len(), 1); - - assert!(matches!( - RadrootsRelayFetchRequest::fetch(-1, 1, primary_relay_target(), [filter.clone()]), - Err(RadrootsRelayTransportError::InvalidTimestamp { - field: "observed_at_ms", - value: -1, - }) - )); - assert!(matches!( - RadrootsRelayFetchRequest::fetch(1_000, 0, primary_relay_target(), [filter.clone()]), - Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_events" - )); - assert!(matches!( - RadrootsRelayFetchRequest::subscription( - 1_000, - 0, - primary_relay_target(), - [filter.clone()], - ), - Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_events" - )); - - let request = RadrootsRelayFetchRequest::fetch(1_000, 1, primary_relay_target(), [filter]) - .expect("valid fetch request"); - assert!(matches!( - request.clone().with_timeout_ms(0), - Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "timeout_ms" - )); - assert!(matches!( - request.clone().with_raw_event_scan_limit(0), - Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_raw_events" - )); - assert!(matches!( - request.clone().with_raw_json_byte_limit(0), - Err(RadrootsRelayTransportError::InvalidFetchLimit { field }) if field == "max_raw_json_bytes" - )); - - let request = request - .with_timeout_ms(1) - .expect("minimum timeout") - .with_raw_event_scan_limit(1) - .expect("minimum raw scan limit") - .with_raw_json_byte_limit(1) - .expect("minimum raw JSON byte limit"); - assert_eq!(request.timeout_ms(), 1); - assert_eq!(request.max_raw_events(), 1); - assert_eq!(request.max_raw_json_bytes(), 1); - - let request = RadrootsRelayFetchRequest::subscription( - 1_005, - 2, - RadrootsRelayTargetSet::new( - [RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("relay targets"), - [post_relay_fetch_filter(2)], - ) - .expect("subscription request") - .with_timeout_ms(25) - .expect("timeout") - .with_raw_event_scan_limit(3) - .expect("raw limit") - .with_raw_json_byte_limit(4_096) - .expect("raw JSON byte limit"); - assert_eq!(request.mode(), RadrootsRelayFetchMode::Subscription); - assert_eq!(request.observed_at_ms(), 1_005); - assert_eq!(request.max_events(), 2); - assert_eq!(request.max_raw_events(), 3); - assert_eq!(request.max_raw_json_bytes(), 4_096); - assert_eq!( - request.relay_targets().relay_strings(), - vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()] - ); - assert_eq!(request.filters().len(), 1); - assert_eq!(request.timeout_ms(), 25); -} - -#[test] -fn fetch_requests_enforce_the_coherent_visibility_batch_limit() { - let filter = post_relay_fetch_filter(RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX); - let fetch = RadrootsRelayFetchRequest::fetch( - 1_006, - RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX, - primary_relay_target(), - [filter.clone()], - ) - .expect("maximum fetch request"); - assert_eq!(fetch.max_events(), RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX); - assert_eq!( - fetch.max_raw_events(), - RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX - ); - assert_eq!( - fetch.max_raw_json_bytes(), - RADROOTS_RELAY_FETCH_RAW_JSON_BYTE_LIMIT_MAX - ); - let exact_raw_limits = fetch - .clone() - .with_raw_event_scan_limit(RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX) - .expect("maximum raw event limit") - .with_raw_json_byte_limit(RADROOTS_RELAY_FETCH_RAW_JSON_BYTE_LIMIT_MAX) - .expect("maximum raw JSON byte limit"); - assert_eq!( - exact_raw_limits.max_raw_events(), - RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX - ); - assert_eq!( - exact_raw_limits.max_raw_json_bytes(), - RADROOTS_RELAY_FETCH_RAW_JSON_BYTE_LIMIT_MAX - ); - let subscription = RadrootsRelayFetchRequest::subscription( - 1_006, - RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX, - primary_relay_target(), - [filter.clone()], - ) - .expect("maximum subscription request"); - assert_eq!( - subscription.max_events(), - RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX - ); - - let above_max = RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX + 1; - assert!(matches!( - RadrootsRelayFetchRequest::fetch( - 1_006, - above_max, - primary_relay_target(), - [filter.clone()], - ), - Err(RadrootsRelayTransportError::FetchLimitTooLarge { - field: "max_events", - max, - actual, - }) if max == RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX && actual == above_max - )); - assert!(matches!( - RadrootsRelayFetchRequest::subscription( - 1_006, - above_max, - primary_relay_target(), - [filter], - ), - Err(RadrootsRelayTransportError::FetchLimitTooLarge { - field: "max_events", - max, - actual, - }) if max == RADROOTS_RELAY_FETCH_EVENT_LIMIT_MAX && actual == above_max - )); - assert!(matches!( - fetch - .clone() - .with_raw_event_scan_limit(RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX + 1), - Err(RadrootsRelayTransportError::FetchLimitTooLarge { - field: "max_raw_events", - max: RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX, - actual, - }) if actual == RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX + 1 - )); - assert!(matches!( - fetch.with_raw_json_byte_limit(RADROOTS_RELAY_FETCH_RAW_JSON_BYTE_LIMIT_MAX + 1), - Err(RadrootsRelayTransportError::FetchLimitTooLarge { - field: "max_raw_json_bytes", - max: RADROOTS_RELAY_FETCH_RAW_JSON_BYTE_LIMIT_MAX, - actual, - }) if actual == RADROOTS_RELAY_FETCH_RAW_JSON_BYTE_LIMIT_MAX + 1 - )); -} - -#[test] -fn fetch_blocking_facade_runs_mock_adapter() { - let signed = signed_post("blocking fetch"); - let accepted_id = signed.id_str().to_owned(); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: signed.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - ]); - - let receipt = fetch_relay_events_blocking(&adapter, post_relay_fetch_request(1_090, 10)) - .expect("blocking fetch"); - - assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id); - assert_eq!(receipt.events[0].observed_at_ms, 1_090); - assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]); -} - -#[tokio::test] -async fn fetch_canonicalizes_adapter_relay_spelling_and_uses_request_observation_time() { - let signed = signed_post("canonical fetch relay"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: "wss://RELAY.EXAMPLE.COM/".to_owned(), - raw_json: signed.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: "wss://RELAY.EXAMPLE.COM/".to_owned(), - }, - ]); - - let receipt = fetch_relay_events(&adapter, post_relay_fetch_request(1_091, 10)) - .await - .expect("canonical fetch"); - - assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].relay_url, RELAY_PRIMARY_WSS); - assert_eq!(receipt.events[0].observed_at_ms, 1_091); - assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]); -} - -#[tokio::test] -async fn fetch_verifies_events_before_acceptance_budgeting() { - let accepted = signed_post("verified after tampered"); - let accepted_id = accepted.id_str().to_owned(); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: tampered_raw_event(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: accepted.raw_json().to_owned(), - }, - ]); - - let receipt = fetch_relay_events(&adapter, post_relay_fetch_request(1_091, 1)) - .await - .expect("verified fetch"); - - assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id); - assert_eq!(receipt.verification_failed_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 0); - assert_eq!( - receipt.event_receipts[0].verification, - RadrootsRelayFetchEventVerification::Failed - ); - assert_eq!( - receipt.event_receipts[1].verification, - RadrootsRelayFetchEventVerification::Verified - ); -} - -#[tokio::test] -async fn fetch_deduplicates_event_ids_without_starving_unique_events() { - let first = signed_post("first unique"); - let second = signed_post("second unique"); - let first_id = first.id_str().to_owned(); - let second_id = second.id_str().to_owned(); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: first.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: first.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: second.raw_json().to_owned(), - }, - ]); - - let receipt = fetch_relay_events(&adapter, post_relay_fetch_request(1_091, 2)) - .await - .expect("deduplicated fetch"); - - assert_eq!( - receipt - .events - .iter() - .map(|event| event.event.id.to_hex()) - .collect::<Vec<_>>(), - vec![first_id.clone(), second_id] - ); - assert_eq!(receipt.duplicate_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 0); - assert_eq!(receipt.event_receipts.len(), 3); - assert_eq!( - receipt.event_receipts[1].event_id.as_deref(), - Some(first_id.as_str()) - ); - assert!(receipt.event_receipts[1].duplicate); -} - -#[tokio::test] -async fn fetch_reports_local_truncation_without_claiming_eose() { - let adapter = RadrootsMockRelayFetchAdapter::new(vec![RadrootsRelayFetchItem::Truncated { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - message: "local budget reached".to_owned(), - }]); - - let receipt = fetch_relay_events(&adapter, post_relay_fetch_request(1_091, 1)) - .await - .expect("truncated fetch"); - - assert!(receipt.events.is_empty()); - assert!(receipt.connected_relays.is_empty()); - assert_eq!(receipt.eose_count, 0); - assert_eq!(receipt.truncated_count, 1); - assert_eq!(receipt.relay_outcomes.len(), 1); - assert_eq!( - receipt.relay_outcomes[0].kind, - RadrootsRelayFetchOutcomeKind::Truncated - ); -} - -#[tokio::test] -async fn fetch_rejects_duplicate_and_conflicting_terminal_outcomes() { - let duplicate = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - ]); - assert!(matches!( - fetch_relay_events(&duplicate, post_relay_fetch_request(1_091, 1)).await, - Err(RadrootsRelayTransportError::DuplicateFetchTerminalRelayUrl { url }) - if url == RELAY_PRIMARY_WSS - )); - - let conflicting = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - RadrootsRelayFetchItem::Closed { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - message: "closed after EOSE".to_owned(), - }, - ]); - assert!(matches!( - fetch_relay_events(&conflicting, post_relay_fetch_request(1_091, 1)).await, - Err(RadrootsRelayTransportError::ConflictingFetchTerminalRelayUrl { - url, - first: "eose", - next: "closed", - }) if url == RELAY_PRIMARY_WSS - )); -} - -#[tokio::test] -async fn fetch_rejects_unrequested_adapter_relay_for_every_item_before_store_mutation() { - let signed = signed_post("forged fetch relay"); - let unexpected_relay = "wss://unexpected.example.com"; - let forged_items = [ - RadrootsRelayFetchItem::Event { - relay_url: unexpected_relay.to_owned(), - raw_json: signed.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: unexpected_relay.to_owned(), - }, - RadrootsRelayFetchItem::Truncated { - relay_url: unexpected_relay.to_owned(), - message: "truncated".to_owned(), - }, - RadrootsRelayFetchItem::Closed { - relay_url: unexpected_relay.to_owned(), - message: "closed".to_owned(), - }, - RadrootsRelayFetchItem::Notice { - relay_url: unexpected_relay.to_owned(), - message: "notice".to_owned(), - }, - ]; - - for forged_item in forged_items { - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: signed.raw_json().to_owned(), - }, - forged_item, - ]); - let error = - fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(1_092, 10)) - .await - .expect_err("unrequested relay must fail the whole fetch"); - - assert!(matches!( - error, - RadrootsRelayTransportError::UnexpectedFetchItemRelayUrl { ref url } - if url == unexpected_relay - )); - assert!( - store - .raw_event(signed.id_str()) - .await - .expect("raw event") - .is_none() - ); - } -} - -#[tokio::test] -async fn fetch_ingests_events_and_records_transport_observations() { - let signed = signed_post("hello"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: signed.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: signed.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: unsupported_raw_event(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: invalid_contract_shape_raw_event(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_TERTIARY_WSS.to_owned(), - raw_json: tampered_raw_event(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_TERTIARY_WSS.to_owned(), - raw_json: "{not json".to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - RadrootsRelayFetchItem::Closed { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - message: "auth-required: challenge".to_owned(), - }, - RadrootsRelayFetchItem::Closed { - relay_url: RELAY_TERTIARY_WSS.to_owned(), - message: "restricted: group write denied".to_owned(), - }, - RadrootsRelayFetchItem::Notice { - relay_url: RELAY_TERTIARY_WSS.to_owned(), - message: "notice: test".to_owned(), - }, - ]); - - let receipt = - fetch_and_ingest_relay_events(&adapter, &store, fixture_relay_fetch_request(1_000, 10)) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 3); - assert_eq!(receipt.duplicate_count, 1); - assert_eq!(receipt.not_persisted_count, 0); - assert_eq!(receipt.verification_failed_count, 1); - assert_eq!(receipt.admission_unsupported_count, 1); - assert_eq!(receipt.admission_invalid_count, 1); - assert_eq!(receipt.valid_stream_eligible_count, 2); - assert_eq!(receipt.visible_count, 2); - assert_eq!(receipt.not_admitted_count, 2); - assert_eq!(receipt.not_current_count, 0); - assert_eq!(receipt.suppressed_count, 0); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.eose_count, 1); - assert_eq!(receipt.closed_count, 2); - assert_eq!(receipt.notice_count, 1); - assert_eq!( - receipt.inserted_count, - receipt.events.iter().filter(|event| event.inserted).count() - ); - assert_eq!( - receipt.duplicate_count, - receipt - .events - .iter() - .filter(|event| event.duplicate) - .count() - ); - assert_eq!( - receipt.not_persisted_count, - receipt - .events - .iter() - .filter(|event| event.not_persisted) - .count() - ); - assert_eq!( - receipt.admission_unsupported_count, - receipt - .events - .iter() - .filter(|event| event.admission == RadrootsRelayFetchEventAdmission::Unsupported) - .count() - ); - assert_eq!( - receipt.admission_invalid_count, - receipt - .events - .iter() - .filter(|event| event.admission == RadrootsRelayFetchEventAdmission::Invalid) - .count() - ); - assert_eq!( - receipt.verification_failed_count, - receipt - .events - .iter() - .filter(|event| event.verification == RadrootsRelayFetchEventVerification::Failed) - .count() - ); - assert_eq!( - receipt.malformed_count, - receipt - .events - .iter() - .filter(|event| event.malformed) - .count() - ); - assert!(receipt.events.iter().all(|event| { - usize::from(event.inserted) - + usize::from(event.duplicate) - + usize::from(event.not_persisted) - <= 1 - && (!event.malformed - || event.verification == RadrootsRelayFetchEventVerification::NotEvaluated) - })); - assert_eq!(receipt.relay_outcomes.len(), 4); - assert_eq!(receipt.relay_outcomes[0].relay_url, RELAY_PRIMARY_WSS); - assert_eq!( - receipt.relay_outcomes[0].kind, - RadrootsRelayFetchOutcomeKind::Eose - ); - assert!(receipt.relay_outcomes[0].relay_outcome.is_none()); - assert_eq!(receipt.relay_outcomes[1].relay_url, RELAY_SECONDARY_WSS); - assert_eq!( - receipt.relay_outcomes[1] - .relay_outcome - .as_ref() - .expect("auth outcome") - .kind, - RadrootsRelayOutcomeKind::AuthRequired - ); - assert_eq!(receipt.relay_outcomes[2].relay_url, RELAY_TERTIARY_WSS); - assert_eq!( - receipt.relay_outcomes[2] - .relay_outcome - .as_ref() - .expect("restricted outcome") - .kind, - RadrootsRelayOutcomeKind::Restricted - ); - assert_eq!( - receipt.relay_outcomes[3].kind, - RadrootsRelayFetchOutcomeKind::Notice - ); - assert!(receipt.relay_outcomes[3].relay_outcome.is_none()); - assert_eq!( - receipt.events[0].admission, - RadrootsRelayFetchEventAdmission::Admitted - ); - assert_eq!( - receipt.events[0].valid_stream, - RadrootsRelayFetchEventValidStream::Eligible - ); - assert_eq!( - receipt.events[0].visibility, - RadrootsRelayFetchEventVisibility::Visible - ); - assert_eq!( - receipt.events[1].admission, - RadrootsRelayFetchEventAdmission::Admitted - ); - assert_eq!( - receipt.events[1].valid_stream, - RadrootsRelayFetchEventValidStream::Eligible - ); - assert_eq!( - receipt.events[1].visibility, - RadrootsRelayFetchEventVisibility::Visible - ); - assert_eq!( - receipt.events[2].admission, - RadrootsRelayFetchEventAdmission::Unsupported - ); - assert_eq!( - receipt.events[2].admission_code.as_deref(), - Some("unsupported_kind") - ); - assert_eq!( - receipt.events[2].valid_stream, - RadrootsRelayFetchEventValidStream::Ineligible - ); - assert_eq!( - receipt.events[2].visibility, - RadrootsRelayFetchEventVisibility::NotAdmitted - ); - assert_eq!( - receipt.events[3].admission, - RadrootsRelayFetchEventAdmission::Invalid - ); - assert_eq!( - receipt.events[3].admission_code.as_deref(), - Some("reply_event_id_invalid") - ); - assert_eq!( - receipt.events[3].valid_stream, - RadrootsRelayFetchEventValidStream::Ineligible - ); - assert_eq!( - receipt.events[3].visibility, - RadrootsRelayFetchEventVisibility::NotAdmitted - ); - assert_eq!( - receipt.events[4].verification, - RadrootsRelayFetchEventVerification::Failed - ); - assert_eq!( - receipt.events[4].admission, - RadrootsRelayFetchEventAdmission::NotEvaluated - ); - assert_eq!( - receipt.events[5].verification, - RadrootsRelayFetchEventVerification::NotEvaluated - ); - assert_eq!( - receipt.events[5].admission, - RadrootsRelayFetchEventAdmission::NotEvaluated - ); - - let serialized = serde_json::to_value(&receipt).expect("serialized fetch receipt"); - assert!(serialized.get("verification_failed_count").is_some()); - assert!(serialized.get("admission_invalid_count").is_some()); - assert!(serialized.get("visible_count").is_some()); - assert!(serialized.get("not_persisted_count").is_some()); - assert!(serialized.get("truncated_count").is_some()); - let serialized_event = serialized["events"][0] - .as_object() - .expect("serialized event receipt"); - assert_eq!(serialized_event.len(), 14); - for field in [ - "relay_url", - "event_id", - "inserted", - "duplicate", - "not_persisted", - "malformed", - "out_of_filter", - "skipped_over_limit", - "verification", - "admission", - "admission_code", - "valid_stream", - "visibility", - "message", - ] { - assert!( - serialized_event.contains_key(field), - "serialized receipt must contain {field}" - ); - } - assert!(!serialized_event.contains_key("projection_eligible")); - assert!(!serialized_event.contains_key("admission_status")); - - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_eq!(observations.len(), 1); - assert_eq!(observations[0].transport_kind, TransportId::NOSTR); - assert_eq!(observations[0].endpoint_uri.as_str(), RELAY_PRIMARY_WSS); - assert_eq!( - observations[0].observation_type, - RadrootsTransportObservationType::Fetch - ); - assert_eq!(observations[0].observation_count, 2); -} - -#[tokio::test] -async fn fetch_reports_final_replaceable_visibility_when_newer_arrives_first() { - let newer = test_event_builder(KIND_PROFILE, r#"{"name":"newer"}"#, Vec::new()) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_001)) - .sign_with_keys(&fixture_keys()) - .expect("signed newer profile"); - let older = test_event_builder(KIND_PROFILE, r#"{"name":"older"}"#, Vec::new()) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) - .sign_with_keys(&fixture_keys()) - .expect("signed older profile"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: newer.as_json(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: older.as_json(), - }, - ]); - let filter = RadrootsNostrFilter::new() - .kind(RadrootsNostrKind::Custom( - u16::try_from(KIND_PROFILE).expect("profile kind must fit NIP-01"), - )) - .limit(2); - let request = RadrootsRelayFetchRequest::fetch(1_001, 2, primary_relay_target(), [filter]) - .expect("profile fetch request"); - - let receipt = fetch_and_ingest_relay_events(&adapter, &store, request) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 2); - assert_eq!(receipt.valid_stream_eligible_count, 2); - assert_eq!(receipt.visible_count, 1); - assert_eq!(receipt.not_current_count, 1); - assert_eq!( - receipt.events[0].visibility, - RadrootsRelayFetchEventVisibility::Visible - ); - assert_eq!( - receipt.events[1].admission, - RadrootsRelayFetchEventAdmission::Admitted - ); - assert_eq!( - receipt.events[1].valid_stream, - RadrootsRelayFetchEventValidStream::Eligible - ); - assert_eq!( - receipt.events[1].visibility, - RadrootsRelayFetchEventVisibility::NotCurrent - ); -} - -#[tokio::test] -async fn fetch_reports_final_replaceable_visibility_when_older_arrives_first() { - let older = test_event_builder(KIND_PROFILE, r#"{"name":"older"}"#, Vec::new()) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) - .sign_with_keys(&fixture_keys()) - .expect("signed older profile"); - let newer = test_event_builder(KIND_PROFILE, r#"{"name":"newer"}"#, Vec::new()) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_001)) - .sign_with_keys(&fixture_keys()) - .expect("signed newer profile"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: older.as_json(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: newer.as_json(), - }, - ]); - let filter = RadrootsNostrFilter::new() - .kind(RadrootsNostrKind::Custom( - u16::try_from(KIND_PROFILE).expect("profile kind must fit NIP-01"), - )) - .limit(2); - let request = RadrootsRelayFetchRequest::fetch(1_001, 2, primary_relay_target(), [filter]) - .expect("profile fetch request"); - - let receipt = fetch_and_ingest_relay_events(&adapter, &store, request) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 2); - assert_eq!(receipt.valid_stream_eligible_count, 2); - assert_eq!(receipt.visible_count, 1); - assert_eq!(receipt.not_current_count, 1); - assert_eq!( - receipt.events[0].visibility, - RadrootsRelayFetchEventVisibility::NotCurrent - ); - assert_eq!( - receipt.events[1].visibility, - RadrootsRelayFetchEventVisibility::Visible - ); -} - -#[tokio::test] -async fn fetch_maps_one_final_visibility_snapshot_back_to_duplicate_receipts() { - let older = test_event_builder(KIND_PROFILE, r#"{"name":"older"}"#, Vec::new()) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) - .sign_with_keys(&fixture_keys()) - .expect("signed older profile"); - let newer = test_event_builder(KIND_PROFILE, r#"{"name":"newer"}"#, Vec::new()) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_001)) - .sign_with_keys(&fixture_keys()) - .expect("signed newer profile"); - let older_id = older.id.to_hex(); - let newer_id = newer.id.to_hex(); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: older.as_json(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: newer.as_json(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: older.as_json(), - }, - ]); - let filter = RadrootsNostrFilter::new() - .kind(RadrootsNostrKind::Custom( - u16::try_from(KIND_PROFILE).expect("profile kind must fit NIP-01"), - )) - .limit(2); - let targets = RadrootsRelayTargetSet::new( - [RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("relay targets"); - let request = RadrootsRelayFetchRequest::fetch(1_001, 2, targets, [filter]) - .expect("profile fetch request"); - - let receipt = fetch_and_ingest_relay_events(&adapter, &store, request) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 2); - assert_eq!(receipt.duplicate_count, 1); - assert_eq!(receipt.valid_stream_eligible_count, 3); - assert_eq!(receipt.visible_count, 1); - assert_eq!(receipt.not_current_count, 2); - assert_eq!(receipt.events.len(), 3); - assert_eq!( - receipt - .events - .iter() - .map(|event| event.event_id.as_deref()) - .collect::<Vec<_>>(), - vec![ - Some(older_id.as_str()), - Some(newer_id.as_str()), - Some(older_id.as_str()), - ] - ); - assert_eq!( - receipt - .events - .iter() - .map(|event| event.visibility) - .collect::<Vec<_>>(), - vec![ - RadrootsRelayFetchEventVisibility::NotCurrent, - RadrootsRelayFetchEventVisibility::Visible, - RadrootsRelayFetchEventVisibility::NotCurrent, - ] - ); - assert!(receipt.events[2].duplicate); -} - -#[tokio::test] -async fn fetch_reports_store_suppression_when_deletion_precedes_target_replay() { - let target = test_event_builder(KIND_POST, "deleted target", Vec::new()) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_000)) - .sign_with_keys(&fixture_keys()) - .expect("signed target"); - let deletion = test_event_builder( - KIND_DELETION_REQUEST, - "", - vec![vec!["e".to_owned(), target.id.to_hex()]], - ) - .custom_created_at(RadrootsNostrTimestamp::from_secs(1_700_000_001)) - .sign_with_keys(&fixture_keys()) - .expect("signed deletion"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: deletion.as_json(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: target.as_json(), - }, - ]); - let deletion_filter = RadrootsNostrFilter::new().kind(RadrootsNostrKind::Custom( - u16::try_from(KIND_DELETION_REQUEST).expect("deletion kind must fit NIP-01"), - )); - let post_filter = RadrootsNostrFilter::new().kind(RadrootsNostrKind::Custom( - u16::try_from(KIND_POST).expect("post kind must fit NIP-01"), - )); - let request = RadrootsRelayFetchRequest::fetch( - 1_002, - 2, - primary_relay_target(), - [deletion_filter, post_filter], - ) - .expect("deletion replay fetch request"); - - let receipt = fetch_and_ingest_relay_events(&adapter, &store, request) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 2); - assert_eq!(receipt.valid_stream_eligible_count, 2); - assert_eq!(receipt.visible_count, 1); - assert_eq!(receipt.suppressed_count, 1); - assert_eq!(receipt.not_current_count, 0); - assert_eq!( - receipt.events[0].visibility, - RadrootsRelayFetchEventVisibility::Visible - ); - assert_eq!( - receipt.events[1].event_id.as_deref(), - Some(target.id.to_hex().as_str()) - ); - assert_eq!( - receipt.events[1].admission, - RadrootsRelayFetchEventAdmission::Admitted - ); - assert_eq!( - receipt.events[1].valid_stream, - RadrootsRelayFetchEventValidStream::Eligible - ); - assert_eq!( - receipt.events[1].visibility, - RadrootsRelayFetchEventVisibility::Suppressed - ); -} - -#[tokio::test] -async fn fetch_reports_ephemeral_events_as_not_persisted_without_duplicate_or_store_state() { - let signed = signed_ephemeral("live geochat"); - let event_id = signed.id_str().to_owned(); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: signed.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: signed.raw_json().to_owned(), - }, - ]); - let filter = RadrootsNostrFilter::new() - .kind(RadrootsNostrKind::Custom( - u16::try_from(KIND_GEOCHAT).expect("ephemeral kind"), - )) - .limit(10); - let request = RadrootsRelayFetchRequest::fetch(1_010, 10, primary_relay_target(), [filter]) - .expect("ephemeral fetch request"); - - let receipt = fetch_and_ingest_relay_events(&adapter, &store, request) - .await - .expect("ephemeral fetch ingest"); - - assert_eq!(receipt.inserted_count, 0); - assert_eq!(receipt.duplicate_count, 0); - assert_eq!(receipt.not_persisted_count, 2); - assert_eq!(receipt.malformed_count, 0); - assert_eq!(receipt.verification_failed_count, 0); - assert_eq!(receipt.admission_unsupported_count, 0); - assert_eq!(receipt.admission_invalid_count, 0); - assert_eq!(receipt.valid_stream_eligible_count, 0); - assert_eq!(receipt.events.len(), 2); - assert!(receipt.events.iter().all(|event| { - !event.inserted - && !event.duplicate - && event.not_persisted - && event.verification == RadrootsRelayFetchEventVerification::Verified - && event.admission == RadrootsRelayFetchEventAdmission::Admitted - && event.admission_code.is_none() - && event.valid_stream == RadrootsRelayFetchEventValidStream::Ineligible - && event.visibility == RadrootsRelayFetchEventVisibility::NotPersisted - })); - assert!( - store - .raw_event(event_id.as_str()) - .await - .expect("raw event") - .is_none() - ); - assert!( - store - .observations_for_event(event_id.as_str()) - .await - .expect("observations") - .is_empty() - ); - let summary = store.status_summary().await.expect("status summary"); - assert_eq!(summary.total_events, 0); - assert_eq!(summary.valid_stream_events, 0); - assert_eq!(summary.transport_observations, 0); -} - -#[tokio::test] -async fn fetch_rejects_out_of_filter_events_before_store_mutation() { - let accepted = signed_post("filter match"); - let wrong_tag = signed_event_with_kind_and_hashtag("filter wrong tag", KIND_POST, "compost"); - let wrong_kind = signed_raw_event_with_kind_and_hashtag("filter wrong kind", 999, "soil"); - let wrong_kind_event_id = wrong_kind.id.to_hex(); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: wrong_tag.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: accepted.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: wrong_kind.as_json(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - ]); - let filter = with_tag( - RadrootsNostrFilter::new() - .kind(RadrootsNostrKind::Custom( - u16::try_from(KIND_POST).expect("post kind must fit NIP-01"), - )) - .limit(10), - "t", - vec!["soil".to_owned()], - ) - .expect("filter"); - - let receipt = fetch_and_ingest_relay_events( - &adapter, - &store, - RadrootsRelayFetchRequest::fetch(1_005, 10, fixture_relay_targets(), [filter]) - .expect("fetch request"), - ) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 1); - assert_eq!(receipt.out_of_filter_count, 2); - assert_eq!(receipt.malformed_count, 0); - assert_eq!(receipt.admission_unsupported_count, 0); - assert_eq!(receipt.events.len(), 3); - assert!(receipt.events[0].out_of_filter); - assert!(!receipt.events[1].out_of_filter); - assert!(receipt.events[2].out_of_filter); - assert!( - store - .raw_event(accepted.id_str()) - .await - .expect("accepted lookup") - .is_some() - ); - assert!( - store - .raw_event(wrong_tag.id_str()) - .await - .expect("wrong tag lookup") - .is_none() - ); - assert!( - store - .raw_event(wrong_kind_event_id.as_str()) - .await - .expect("wrong kind lookup") - .is_none() - ); -} - -#[tokio::test] -async fn fetch_event_cap_counts_accepted_in_filter_events_and_preserves_later_control_outcomes() { - let accepted = signed_post("accepted capped event"); - let skipped = signed_post("skipped capped event"); - let wrong_tag = signed_event_with_kind_and_hashtag("wrong capped tag", KIND_POST, "compost"); - let accepted_id = accepted.id_str().to_owned(); - let skipped_id = skipped.id_str().to_owned(); - let wrong_tag_id = wrong_tag.id_str().to_owned(); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: "{not json".to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: wrong_tag.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: accepted.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: skipped.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - RadrootsRelayFetchItem::Closed { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - message: "auth-required: challenge".to_owned(), - }, - RadrootsRelayFetchItem::Notice { - relay_url: RELAY_TERTIARY_WSS.to_owned(), - message: "notice: still visible".to_owned(), - }, - ]); - - let receipt = - fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(1_100, 1)) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 1); - assert_eq!(receipt.duplicate_count, 0); - assert_eq!(receipt.admission_unsupported_count, 0); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.out_of_filter_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 1); - assert_eq!(receipt.events.len(), 4); - assert!(receipt.events[0].malformed); - assert!(receipt.events[1].out_of_filter); - assert!(receipt.events[2].inserted); - assert!(receipt.events[3].skipped_over_limit); - assert_eq!(receipt.eose_count, 1); - assert_eq!(receipt.closed_count, 1); - assert_eq!(receipt.notice_count, 1); - assert_eq!(receipt.relay_outcomes.len(), 3); - assert_eq!( - receipt.relay_outcomes[0].kind, - RadrootsRelayFetchOutcomeKind::Eose - ); - assert_eq!( - receipt.relay_outcomes[1] - .relay_outcome - .as_ref() - .expect("closed outcome") - .kind, - RadrootsRelayOutcomeKind::AuthRequired - ); - assert_eq!( - receipt.relay_outcomes[2].kind, - RadrootsRelayFetchOutcomeKind::Notice - ); - assert!( - store - .raw_event(accepted_id.as_str()) - .await - .expect("accepted lookup") - .is_some() - ); - assert!( - store - .raw_event(skipped_id.as_str()) - .await - .expect("skipped lookup") - .is_none() - ); - assert!( - store - .raw_event(wrong_tag_id.as_str()) - .await - .expect("wrong tag lookup") - .is_none() - ); -} - -#[tokio::test] -async fn fetch_relay_events_applies_shared_filter_limit_and_outcome_evidence() { - let accepted = signed_event_with_kind_and_hashtag("shared fetch accepted", KIND_POST, "soil"); - let skipped = signed_event_with_kind_and_hashtag("shared fetch skipped", KIND_POST, "soil"); - let wrong_tag = - signed_event_with_kind_and_hashtag("shared fetch wrong tag", KIND_POST, "compost"); - let filter = with_tag( - RadrootsNostrFilter::new() - .kind(RadrootsNostrKind::Custom( - u16::try_from(KIND_POST).expect("post kind must fit NIP-01"), - )) - .limit(10), - "t", - vec!["soil".to_owned()], - ) - .expect("filter"); - let accepted_id = accepted.id_str().to_owned(); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: "{not json".to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: wrong_tag.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: accepted.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: skipped.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - RadrootsRelayFetchItem::Closed { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - message: "auth-required: challenge".to_owned(), - }, - RadrootsRelayFetchItem::Notice { - relay_url: RELAY_TERTIARY_WSS.to_owned(), - message: "notice: still visible".to_owned(), - }, - ]); - - let receipt = fetch_relay_events( - &adapter, - RadrootsRelayFetchRequest::fetch(2_100, 1, fixture_relay_targets(), [filter]) - .expect("fetch request"), - ) - .await - .expect("fetch events"); - - assert_eq!( - receipt.target_relays, - vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS, RELAY_TERTIARY_WSS] - ); - assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]); - assert_eq!(receipt.failed_relays.len(), 1); - assert_eq!(receipt.failed_relays[0].relay_url, RELAY_SECONDARY_WSS); - assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].event.id.to_hex(), accepted_id); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.out_of_filter_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 1); - assert_eq!(receipt.eose_count, 1); - assert_eq!(receipt.closed_count, 1); - assert_eq!(receipt.notice_count, 1); - assert_eq!(receipt.event_receipts.len(), 4); - assert!(receipt.event_receipts[0].malformed); - assert!(receipt.event_receipts[1].out_of_filter); - assert!(!receipt.event_receipts[2].malformed); - assert!(receipt.event_receipts[3].skipped_over_limit); -} - -#[tokio::test] -async fn fetch_raw_scan_limit_bounds_noisy_adapter_output() { - let accepted = signed_post("raw scan accepted event"); - let wrong_tag = signed_event_with_kind_and_hashtag("raw scan wrong tag", KIND_POST, "compost"); - let accepted_id = accepted.id_str().to_owned(); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: "{not json".to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: wrong_tag.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: accepted.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - ]); - - let receipt = fetch_and_ingest_relay_events( - &adapter, - &store, - post_relay_fetch_request(1_130, 1) - .with_raw_event_scan_limit(2) - .expect("raw scan limit"), - ) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 0); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.out_of_filter_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 1); - assert_eq!(receipt.events.len(), 2); - assert_eq!(receipt.eose_count, 1); - assert!( - store - .raw_event(accepted_id.as_str()) - .await - .expect("accepted lookup") - .is_none() - ); -} - -#[tokio::test] -async fn fetch_raw_json_byte_limit_is_exact_global_and_sticky() { - let first = signed_post("raw byte first"); - let second = signed_post("raw byte second"); - let third = signed_post("raw byte third"); - let exact_bytes = first.raw_json().len() + second.raw_json().len(); - let targets = RadrootsRelayTargetSet::new( - [RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("relay targets"); - let items = vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: first.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: second.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: third.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - }, - ]; - - let exact_receipt = fetch_relay_events( - &RadrootsMockRelayFetchAdapter::new(items.clone()), - RadrootsRelayFetchRequest::fetch(1_131, 3, targets.clone(), [post_relay_fetch_filter(3)]) - .expect("fetch request") - .with_raw_json_byte_limit(exact_bytes) - .expect("exact raw JSON byte limit"), - ) - .await - .expect("exact fetch"); - assert_eq!(exact_receipt.events.len(), 2); - assert_eq!(exact_receipt.skipped_over_limit_count, 1); - assert_eq!(exact_receipt.eose_count, 2); - - let crossing_receipt = fetch_relay_events( - &RadrootsMockRelayFetchAdapter::new(items), - RadrootsRelayFetchRequest::fetch(1_132, 3, targets, [post_relay_fetch_filter(3)]) - .expect("fetch request") - .with_raw_json_byte_limit(exact_bytes - 1) - .expect("crossing raw JSON byte limit"), - ) - .await - .expect("crossing fetch"); - assert_eq!(crossing_receipt.events.len(), 1); - assert_eq!(crossing_receipt.skipped_over_limit_count, 2); - assert_eq!(crossing_receipt.eose_count, 2); -} - -#[tokio::test] -async fn fetch_raw_json_budget_charges_every_preparse_event_class() { - let malformed = "{not json".to_owned(); - let wrong_tag = signed_event_with_kind_and_hashtag("raw bytes wrong tag", KIND_POST, "compost"); - let accepted = signed_post("raw bytes accepted"); - let over_accepted_limit = signed_post("raw bytes over accepted limit"); - let sticky_skip = signed_post("raw bytes sticky skip"); - let exact_bytes = malformed.len() - + wrong_tag.raw_json().len() - + accepted.raw_json().len() * 2 - + over_accepted_limit.raw_json().len(); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: malformed, - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: wrong_tag.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: accepted.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: accepted.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: over_accepted_limit.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - raw_json: sticky_skip.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_SECONDARY_WSS.to_owned(), - }, - ]); - let targets = RadrootsRelayTargetSet::new( - [RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - RadrootsRelayUrlPolicy::Public, - ) - .expect("relay targets"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - - let receipt = fetch_and_ingest_relay_events( - &adapter, - &store, - RadrootsRelayFetchRequest::fetch(1_133, 1, targets, [post_relay_fetch_filter(6)]) - .expect("fetch request") - .with_raw_json_byte_limit(exact_bytes) - .expect("raw JSON byte limit"), - ) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 1); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.out_of_filter_count, 1); - assert_eq!(receipt.duplicate_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 2); - assert_eq!(receipt.events.len(), 5); - assert!(receipt.events[4].skipped_over_limit); - assert_eq!(receipt.eose_count, 2); -} - -#[tokio::test] -async fn fetch_rejects_oversized_raw_json_before_radroots_adapter_parsing() { - let accepted = signed_post("after oversized raw event"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![ - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: "x".repeat(DEFAULT_RAW_JSON_MAX_BYTES + 1), - }, - RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: accepted.raw_json().to_owned(), - }, - RadrootsRelayFetchItem::Eose { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - }, - ]); - - let receipt = fetch_relay_events(&adapter, post_relay_fetch_request(1_134, 1)) - .await - .expect("fetch events"); - - assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.verification_failed_count, 1); - assert_eq!(receipt.malformed_count, 0); - assert_eq!(receipt.event_receipts.len(), 2); - assert_eq!(receipt.event_receipts[0].event_id, None); - assert_eq!( - receipt.event_receipts[0].verification, - RadrootsRelayFetchEventVerification::Failed - ); - assert_eq!(receipt.eose_count, 1); -} - -#[tokio::test] -async fn fetch_subscription_mode_and_store_errors_are_propagated() { - let signed = signed_post("subscription"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: signed.raw_json().to_owned(), - }]); - - let receipt = fetch_and_ingest_relay_events( - &adapter, - &store, - RadrootsRelayFetchRequest::subscription( - 1_200, - 10, - primary_relay_target(), - [post_relay_fetch_filter(10)], - ) - .expect("subscription request"), - ) - .await - .expect("fetch ingest"); - - assert_eq!(receipt.inserted_count, 1); - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_eq!(observations.len(), 1); - assert_eq!( - observations[0].observation_type, - RadrootsTransportObservationType::Subscription - ); - - let closed_store = RadrootsEventStore::open_memory().await.expect("store"); - closed_store.pool().close().await; - let adapter = RadrootsMockRelayFetchAdapter::new(vec![RadrootsRelayFetchItem::Event { - relay_url: RELAY_PRIMARY_WSS.to_owned(), - raw_json: signed.raw_json().to_owned(), - }]); - let error = - fetch_and_ingest_relay_events(&adapter, &closed_store, post_relay_fetch_request(1_210, 10)) - .await - .expect_err("closed local store must fail the fetch ingest"); - - assert!(matches!(error, RadrootsRelayTransportError::EventStore(_))); -} - -#[tokio::test] -async fn fetch_ingest_rejects_invalid_observation_endpoint() { - let signed = signed_post("invalid observation endpoint"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let adapter = RadrootsMockRelayFetchAdapter::new(vec![RadrootsRelayFetchItem::Event { - relay_url: " ".to_owned(), - raw_json: signed.raw_json().to_owned(), - }]); - - let error = - fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(1_300, 10)) - .await - .expect_err("invalid observation endpoint"); - - assert!(matches!( - error, - RadrootsRelayTransportError::InvalidFetchItemRelayUrl { .. } - )); -} - -#[tokio::test] -async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("hello"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - vec![ - RELAY_PRIMARY_WSS.to_owned(), - RELAY_SECONDARY_WSS.to_owned(), - RELAY_TERTIARY_WSS.to_owned(), - ], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - assert_eq!(publish_claim.state, RadrootsOutboxEventState::Publishing); - - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()) - .with_outcome( - RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::timeout("timeout: no OK"), - ) - .with_outcome( - RELAY_TERTIARY_WSS, - RadrootsRelayOutcome::duplicate_accepted("duplicate: already have it"), - ); - let first = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(first.attempted_count, 3); - assert_eq!(first.accepted_count, 2); - assert!(!first.quorum_met); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable); - - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert_eq!( - targets - .iter() - .find(|target| target.endpoint_uri.as_str() == RELAY_PRIMARY_WSS) - .expect("primary") - .status, - RadrootsOutboxDeliveryTargetStatus::Accepted - ); - assert_eq!( - targets - .iter() - .find(|target| target.endpoint_uri.as_str() == RELAY_SECONDARY_WSS) - .expect("secondary") - .status, - RadrootsOutboxDeliveryTargetStatus::FailedRetryable - ); - assert_eq!( - targets - .iter() - .find(|target| target.endpoint_uri.as_str() == RELAY_TERTIARY_WSS) - .expect("tertiary") - .status, - RadrootsOutboxDeliveryTargetStatus::Accepted - ); - - let retry_claim = outbox - .claim_next_ready_event("publisher", "publish-b", 4_000, 2_500) - .await - .expect("claim") - .expect("retry claim"); - let retry_adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted()); - let second = publish_claimed_outbox_event( - &outbox, - &store, - &retry_adapter, - &retry_claim, - RadrootsOutboxPublishPolicy::new(3_000), - 2_600, - ) - .await - .expect("retry publish"); - - assert_eq!(second.local_ingest.event_id, signed.id_str()); - assert_eq!(second.attempted_count, 1); - assert_eq!(retry_adapter.captured_raw_events().len(), 1); - - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); - let operation = outbox - .get_operation(receipt.operation_id) - .await - .expect("operation") - .expect("operation"); - assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete); - - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_outbox_publish_observations(&observations, 3); -} - -#[tokio::test] -async fn outbox_transport_facade_persists_partial_success_and_retryable_failures() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("transport facade outbox"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - [RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "transport-sign", 2_000, 1_000) - .await - .expect("sign claim") - .expect("sign claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "transport-publish", 3_000, 1_100) - .await - .expect("publish claim") - .expect("publish claim"); - - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()) - .with_outcome( - RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::timeout("timeout: transport facade"), - ); - let transport = RadrootsNostrTransport::new(adapter); - let published = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &transport, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("transport publish"); - - assert_eq!(published.event_id, signed.id_str()); - assert_eq!(published.attempted_count, 2); - assert_eq!(published.accepted_count, 1); - assert_eq!(published.retryable_count, 1); - assert_eq!(published.terminal_count, 0); - assert!(!published.quorum_met); - assert_eq!(published.relay_receipts.len(), 2); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert_eq!( - targets - .iter() - .find(|target| target.endpoint_uri.as_str() == RELAY_PRIMARY_WSS) - .expect("primary") - .status, - RadrootsOutboxDeliveryTargetStatus::Accepted - ); - assert_eq!( - targets - .iter() - .find(|target| target.endpoint_uri.as_str() == RELAY_SECONDARY_WSS) - .expect("secondary") - .status, - RadrootsOutboxDeliveryTargetStatus::FailedRetryable - ); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable); - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_outbox_publish_observations(&observations, 1); - assert_eq!( - observations - .iter() - .find(|observation| { - observation.observation_type == RadrootsTransportObservationType::PublishAck - }) - .expect("publish acknowledgement") - .observation_count, - 1 - ); -} - -#[tokio::test] -async fn outbox_transport_facade_persists_every_delivery_status() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("transport outcome matrix"); - let relays = (0..14) - .map(|index| format!("wss://relay-{index}.example.com")) - .collect::<Vec<_>>(); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input(draft, &relays)) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "matrix-sign", 2_000, 1_000) - .await - .expect("sign claim") - .expect("sign claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "matrix-publish", 3_000, 1_100) - .await - .expect("publish claim") - .expect("publish claim"); - let outcomes = [ - RadrootsTransportOutcomeKind::Accepted, - RadrootsTransportOutcomeKind::DuplicateAccepted, - RadrootsTransportOutcomeKind::Delivered, - RadrootsTransportOutcomeKind::Forwarded, - RadrootsTransportOutcomeKind::StoredByGateway, - RadrootsTransportOutcomeKind::Seen, - RadrootsTransportOutcomeKind::DeferredUntilImplemented, - RadrootsTransportOutcomeKind::Rejected, - RadrootsTransportOutcomeKind::RouteUnavailable, - RadrootsTransportOutcomeKind::PayloadTooLarge, - RadrootsTransportOutcomeKind::PolicyDenied, - RadrootsTransportOutcomeKind::Timeout, - RadrootsTransportOutcomeKind::ConnectionFailed, - RadrootsTransportOutcomeKind::TransportUnavailable, - ] - .into_iter() - .map(RadrootsTransportOutcome::new) - .collect(); - let transport = ScriptedTransport::new(outcomes); - - let published = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &transport, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("transport publish"); - - assert_eq!(published.event_id, signed.id_str()); - assert_eq!(published.attempted_count, 14); - assert_eq!(published.accepted_count, 6); - assert_eq!(published.retryable_count, 3); - assert_eq!(published.terminal_count, 5); - assert!(!published.quorum_met); - assert_eq!(published.target_receipts.len(), 14); - assert_eq!(published.relay_receipts.len(), 14); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - let expected_statuses = [ - RadrootsOutboxDeliveryTargetStatus::Accepted, - RadrootsOutboxDeliveryTargetStatus::Accepted, - RadrootsOutboxDeliveryTargetStatus::Delivered, - RadrootsOutboxDeliveryTargetStatus::Accepted, - RadrootsOutboxDeliveryTargetStatus::Accepted, - RadrootsOutboxDeliveryTargetStatus::Accepted, - RadrootsOutboxDeliveryTargetStatus::FailedTerminal, - RadrootsOutboxDeliveryTargetStatus::FailedTerminal, - RadrootsOutboxDeliveryTargetStatus::FailedTerminal, - RadrootsOutboxDeliveryTargetStatus::FailedTerminal, - RadrootsOutboxDeliveryTargetStatus::FailedTerminal, - RadrootsOutboxDeliveryTargetStatus::FailedRetryable, - RadrootsOutboxDeliveryTargetStatus::FailedRetryable, - RadrootsOutboxDeliveryTargetStatus::FailedRetryable, - ]; - assert_eq!(targets.len(), expected_statuses.len()); - for (target, expected_status) in targets.iter().zip(expected_statuses) { - assert_eq!(target.status, expected_status); - } - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable); - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_outbox_publish_observations(&observations, 6); -} - -#[tokio::test] -async fn outbox_transport_facade_normalizes_predecessor_pending_evidence() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("pending transport receipt"); - outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - [RELAY_PRIMARY_WSS], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "pending-sign", 2_000, 1_000) - .await - .expect("sign claim") - .expect("sign claim"); - complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "pending-publish", 3_000, 1_100) - .await - .expect("publish claim") - .expect("publish claim"); - let mut forged_outcome = RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Accepted); - forged_outcome.status = RadrootsTransportDeliveryTargetStatus::Pending; - let transport = ScriptedTransport::new(vec![forged_outcome]); - - let published = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &transport, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("final sink receipt"); - assert_eq!(published.target_receipts.len(), 1); - assert_eq!( - published.target_receipts[0].transport_status, - RadrootsTransportDeliveryTargetStatus::Accepted - ); -} - -#[tokio::test] -async fn outbox_transport_facade_rejects_receipts_forged_for_another_request() { - for forged in [ - ForgedDeliveryReceipt::RequestId, - ForgedDeliveryReceipt::TargetSet, - ] { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - generic_draft("forged transport receipt"), - [RELAY_PRIMARY_WSS], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "forged-sign", 2_000, 1_000) - .await - .expect("sign claim") - .expect("sign claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "forged-publish", 3_000, 1_100) - .await - .expect("publish claim") - .expect("publish claim"); - - let error = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &ForgedReceiptTransport { forged }, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect_err("forged receipt rejected"); - assert!(matches!( - error, - RadrootsRelayTransportError::TransportContract(_) - )); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Publishing); - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_outbox_publish_observations(&observations, 0); - } -} - -#[tokio::test] -async fn outbox_transport_facade_handles_empty_and_invalid_claim_plans() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("transport plan edges"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - [RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "plan-sign", 2_000, 1_000) - .await - .expect("sign claim") - .expect("sign claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "plan-publish", 3_000, 1_100) - .await - .expect("publish claim") - .expect("publish claim"); - for target in &publish_claim.delivery_targets { - outbox - .mark_delivery_target_accepted( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - target.delivery_target_id, - 2_150, - ) - .await - .expect("accepted target"); - } - let published = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &ScriptedTransport::new(Vec::new()), - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("already satisfied publish"); - assert_eq!(published.event_id, signed.id_str()); - assert_eq!(published.attempted_count, 0); - assert_eq!(published.accepted_count, 2); - assert!(published.quorum_met); - - let second_draft = generic_draft("invalid claimed plan"); - outbox - .enqueue_operation(all_accepted_outbox_operation_input( - second_draft, - [RELAY_PRIMARY_WSS], - )) - .await - .expect("second enqueue"); - let second_claimed = outbox - .claim_next_ready_event("signer", "invalid-plan-sign", 3_000, 2_200) - .await - .expect("second sign claim") - .expect("second sign claim"); - complete_claimed_signing(&outbox, &second_claimed, 2_300).await; - let second_publish_claim = outbox - .claim_next_ready_event("publisher", "invalid-plan-publish", 4_000, 2_300) - .await - .expect("second publish claim") - .expect("second publish claim"); - let mut invalid_claim = second_publish_claim.clone(); - invalid_claim.active_delivery_plan_id = None; - let error = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &ScriptedTransport::new(Vec::new()), - &invalid_claim, - RadrootsOutboxPublishPolicy::new(3_500), - 2_400, - ) - .await - .expect_err("missing plan rejected"); - assert!(matches!(error, RadrootsRelayTransportError::Transport(_))); - - invalid_claim.active_delivery_plan_id = Some(i64::MAX); - let error = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &ScriptedTransport::new(Vec::new()), - &invalid_claim, - RadrootsOutboxPublishPolicy::new(3_500), - 2_401, - ) - .await - .expect_err("unknown plan rejected"); - assert!(matches!(error, RadrootsRelayTransportError::Transport(_))); - - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); -} - -#[tokio::test] -async fn outbox_transport_facade_requires_signed_claims() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("missing transport signature"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - [RELAY_PRIMARY_WSS], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "unsigned-transport", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - - let error = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &ScriptedTransport::new(Vec::new()), - &claimed, - RadrootsOutboxPublishPolicy::new(2_500), - 1_100, - ) - .await - .expect_err("missing signature rejected"); - assert!(matches!( - error, - RadrootsRelayTransportError::MissingSignedOutboxEvent(event_id) - if event_id == receipt.outbox_event_id - )); -} - -#[tokio::test] -async fn outbox_transport_facade_rejects_non_nostr_transport_before_mutation() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - generic_draft("misrouted transport"), - [RELAY_PRIMARY_WSS], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "misrouted-sign", 2_000, 1_000) - .await - .expect("sign claim") - .expect("sign claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "misrouted-publish", 3_000, 1_100) - .await - .expect("publish claim") - .expect("publish claim"); - let transport = ScriptedTransport::new(Vec::new()).with_kind(TransportId::RETICULUM); - - let error = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &transport, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect_err("non-Nostr transport rejected"); - assert!(matches!( - error, - RadrootsRelayTransportError::UnexpectedTransportKind { - expected: "nostr", - actual, - } if actual == "reticulum" - )); - assert!( - store - .raw_event(signed.id_str()) - .await - .expect("raw event") - .is_none() - ); - assert!( - store - .observations_for_event(signed.id_str()) - .await - .expect("observations") - .is_empty() - ); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Publishing); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert!(targets.iter().all(|target| target.attempt_count == 0)); -} - -#[tokio::test] -async fn outbox_publish_fans_out_endpoint_receipts_to_scoped_logical_targets() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("scoped duplicate relay"); - let receipt = outbox - .enqueue_operation(RadrootsOutboxOperationInput::new( - "publish_post", - draft, - RadrootsOutboxDeliveryPlanInput::new( - "transport.nostr.local", - 2, - RadrootsTransportSatisfactionPolicy::all_accepted(), - vec![ - scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.west", "West foodshed"), - scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.east", "East foodshed"), - ], - ), - 1_000, - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()); - - let published = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.local_ingest.event_id, signed.id_str()); - assert_eq!(published.event_id, signed.id_str()); - assert_eq!(published.attempted_count, 2); - assert_eq!(published.accepted_count, 2); - assert_eq!(published.retryable_count, 0); - assert_eq!(published.terminal_count, 0); - assert_eq!(published.quorum, 2); - assert!(published.quorum_met); - assert_eq!(published.relay_receipts.len(), 1); - assert_eq!(published.relay_receipts[0].relay_url, RELAY_PRIMARY_WSS); - assert_eq!(published.target_receipts.len(), 2); - assert!( - published - .target_receipts - .iter() - .all(|target| target.endpoint_uri == RELAY_PRIMARY_WSS && target.attempted) - ); - assert!(published.target_receipts.iter().any(|target| { - target.target_scope.as_deref() == Some("foodshed.west") - && target.target_label.as_deref() == Some("West foodshed") - })); - assert!(published.target_receipts.iter().any(|target| { - target.target_scope.as_deref() == Some("foodshed.east") - && target.target_label.as_deref() == Some("East foodshed") - })); - assert_eq!(adapter.captured_raw_events().len(), 1); - - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert_eq!(targets.len(), 2); - assert!(targets.iter().all(|target| { - target.endpoint_uri.as_str() == RELAY_PRIMARY_WSS - && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted - && target.attempt_count == 1 - })); - assert!(targets.iter().any(|target| { - target.target_scope.as_ref().map(|scope| scope.as_str()) == Some("foodshed.west") - && target.target_label.as_ref().map(|label| label.as_str()) == Some("West foodshed") - })); - assert!(targets.iter().any(|target| { - target.target_scope.as_ref().map(|scope| scope.as_str()) == Some("foodshed.east") - && target.target_label.as_ref().map(|label| label.as_str()) == Some("East foodshed") - })); - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_outbox_publish_observations(&observations, 1); -} - -#[tokio::test] -async fn outbox_transport_facade_fans_out_endpoint_receipts_to_scoped_logical_targets() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let receipt = outbox - .enqueue_operation(RadrootsOutboxOperationInput::new( - "publish_post", - generic_draft("transport scoped duplicate relay"), - RadrootsOutboxDeliveryPlanInput::new( - "transport.nostr.local", - 2, - RadrootsTransportSatisfactionPolicy::all_accepted(), - vec![ - scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.west", "West foodshed"), - scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.east", "East foodshed"), - ], - ), - 1_000, - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()); - let transport = RadrootsNostrTransport::new(adapter.clone()); - - let published = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &transport, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("transport publish"); - - assert_eq!(published.local_ingest.event_id, signed.id_str()); - assert_eq!(published.event_id, signed.id_str()); - assert_eq!(published.attempted_count, 2); - assert_eq!(published.accepted_count, 2); - assert_eq!(published.retryable_count, 0); - assert_eq!(published.terminal_count, 0); - assert_eq!(published.quorum, 2); - assert!(published.quorum_met); - assert_eq!(published.relay_receipts.len(), 1); - assert_eq!(published.target_receipts.len(), 2); - assert!(published.target_receipts.iter().all(|target| { - target.endpoint_uri == RELAY_PRIMARY_WSS - && target.attempted - && target.transport_status == RadrootsTransportDeliveryTargetStatus::Accepted - })); - assert!(published.target_receipts.iter().any(|target| { - target.target_scope.as_deref() == Some("foodshed.west") - && target.target_label.as_deref() == Some("West foodshed") - })); - assert!(published.target_receipts.iter().any(|target| { - target.target_scope.as_deref() == Some("foodshed.east") - && target.target_label.as_deref() == Some("East foodshed") - })); - assert_eq!(adapter.captured_raw_events().len(), 1); - - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert_eq!(targets.len(), 2); - assert!(targets.iter().all(|target| { - target.endpoint_uri.as_str() == RELAY_PRIMARY_WSS - && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted - && target.attempt_count == 1 - })); - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_outbox_publish_observations(&observations, 1); - assert_eq!( - observations - .iter() - .find(|observation| { - observation.observation_type == RadrootsTransportObservationType::PublishAck - }) - .expect("scoped publish acknowledgement") - .observation_count, - 1 - ); -} - -#[tokio::test] -async fn outbox_publish_required_target_failure_is_not_satisfied_by_optional_success() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("required target optional success"); - let optional = nostr_target(RELAY_PRIMARY_WSS); - let required = nostr_target(RELAY_SECONDARY_WSS); - let receipt = outbox - .enqueue_operation(RadrootsOutboxOperationInput::new( - "publish_post", - draft, - RadrootsOutboxDeliveryPlanInput::new( - "transport.nostr.local", - 1, - RadrootsTransportSatisfactionPolicy::required_targets( - RadrootsTransportSatisfactionClass::Accepted, - vec![required.fingerprint().clone()], - ) - .expect("required target policy"), - vec![optional.clone(), required.clone()], - ), - 1_000, - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let optional_target_id = publish_claim - .delivery_targets - .iter() - .find(|target| &target.endpoint_fingerprint == optional.fingerprint()) - .expect("optional target") - .delivery_target_id; - outbox - .mark_delivery_target_accepted( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - optional_target_id, - 2_000, - ) - .await - .expect("optional accepted"); - - let adapter = RadrootsMockRelayPublishAdapter::new().with_outcome( - RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::timeout("required relay timeout"), - ); - let published = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.attempted_count, 1); - assert_eq!(published.accepted_count, 0); - assert_eq!(published.retryable_count, 1); - assert_eq!(published.quorum, 1); - assert!(!published.quorum_met); - assert_eq!(published.relay_receipts.len(), 1); - assert_eq!(published.relay_receipts[0].relay_url, RELAY_SECONDARY_WSS); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert!(targets.iter().any(|target| { - &target.endpoint_fingerprint == optional.fingerprint() - && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted - })); - assert!(targets.iter().any(|target| { - &target.endpoint_fingerprint == required.fingerprint() - && target.status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable - })); -} - -#[tokio::test] -async fn outbox_publish_required_target_success_is_not_blocked_by_optional_retryable_failure() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("required target optional failure"); - let optional = nostr_target(RELAY_PRIMARY_WSS); - let required = nostr_target(RELAY_SECONDARY_WSS); - let receipt = outbox - .enqueue_operation(RadrootsOutboxOperationInput::new( - "publish_post", - draft, - RadrootsOutboxDeliveryPlanInput::new( - "transport.nostr.local", - 1, - RadrootsTransportSatisfactionPolicy::required_targets( - RadrootsTransportSatisfactionClass::Accepted, - vec![required.fingerprint().clone()], - ) - .expect("required target policy"), - vec![optional.clone(), required.clone()], - ), - 1_000, - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let optional_target_id = publish_claim - .delivery_targets - .iter() - .find(|target| &target.endpoint_fingerprint == optional.fingerprint()) - .expect("optional target") - .delivery_target_id; - outbox - .mark_delivery_target_failed_retryable( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - optional_target_id, - "optional relay timeout", - 2_000, - ) - .await - .expect("optional retryable"); - - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted()); - let published = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.local_ingest.event_id, signed.id_str()); - assert_eq!(published.attempted_count, 1); - assert_eq!(published.accepted_count, 1); - assert_eq!(published.retryable_count, 0); - assert_eq!(published.quorum, 1); - assert!(published.quorum_met); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert!(targets.iter().any(|target| { - &target.endpoint_fingerprint == optional.fingerprint() - && target.status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable - })); - assert!(targets.iter().any(|target| { - &target.endpoint_fingerprint == required.fingerprint() - && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted - })); - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_outbox_publish_observations(&observations, 1); -} - -#[tokio::test] -async fn outbox_publish_required_targets_fan_out_same_endpoint_scoped_receipts() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("required target scoped duplicate relay"); - let required = scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.west", "West foodshed"); - let optional = scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.east", "East foodshed"); - let terminal = scoped_nostr_target(RELAY_PRIMARY_WSS, "foodshed.closed", "Closed foodshed"); - let receipt = outbox - .enqueue_operation(RadrootsOutboxOperationInput::new( - "publish_post", - draft, - RadrootsOutboxDeliveryPlanInput::new( - "transport.nostr.local", - 1, - RadrootsTransportSatisfactionPolicy::required_targets( - RadrootsTransportSatisfactionClass::Accepted, - vec![required.fingerprint().clone()], - ) - .expect("required target policy"), - vec![ - required.clone(), - optional.clone(), - terminal.clone(), - reticulum_target(), - ], - ), - 1_000, - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let optional_record = publish_claim - .delivery_targets - .iter() - .find(|target| &target.endpoint_fingerprint == optional.fingerprint()) - .expect("optional target"); - outbox - .mark_delivery_target_accepted( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - optional_record.delivery_target_id, - 2_150, - ) - .await - .expect("optional target accepted"); - let terminal_record = publish_claim - .delivery_targets - .iter() - .find(|target| &target.endpoint_fingerprint == terminal.fingerprint()) - .expect("terminal target"); - outbox - .mark_delivery_target_failed_terminal( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - terminal_record.delivery_target_id, - "terminal target", - 2_151, - ) - .await - .expect("terminal target completed"); - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()); - - let published = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500).republish_accepted_relays(true), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.attempted_count, 2); - assert_eq!(published.accepted_count, 2); - assert_eq!(published.quorum, 1); - assert!(published.quorum_met); - assert_eq!(published.relay_receipts.len(), 1); - assert_eq!(published.target_receipts.len(), 2); - assert!(published.target_receipts.iter().any(|target| { - &target.endpoint_fingerprint == required.fingerprint() - && target.target_scope.as_deref() == Some("foodshed.west") - })); - assert!(published.target_receipts.iter().any(|target| { - &target.endpoint_fingerprint == optional.fingerprint() - && target.target_scope.as_deref() == Some("foodshed.east") - })); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert_eq!(targets.len(), 4); - assert!( - targets - .iter() - .filter(|target| { target.transport_kind == TransportId::NOSTR }) - .filter(|target| &target.endpoint_fingerprint != terminal.fingerprint()) - .all(|target| { - target.endpoint_uri.as_str() == RELAY_PRIMARY_WSS - && target.status == RadrootsOutboxDeliveryTargetStatus::Accepted - }) - ); - assert!(targets.iter().any(|target| { - &target.endpoint_fingerprint == terminal.fingerprint() - && target.status == RadrootsOutboxDeliveryTargetStatus::FailedTerminal - })); - assert!(targets.iter().any(|target| { - target.transport_kind == TransportId::RETICULUM - && target.status == RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented - })); -} - -#[tokio::test] -async fn outbox_transport_publish_failure_releases_retryable_claim() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("adapter transport failure"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - - let published = publish_claimed_outbox_event( - &outbox, - &store, - &TransportFailurePublishAdapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.attempted_count, 2); - assert_eq!(published.accepted_count, 0); - assert_eq!(published.retryable_count, 2); - assert_eq!(published.terminal_count, 0); - assert!(!published.quorum_met); - assert!( - published - .relay_receipts - .iter() - .all(|relay| relay.outcome.kind == RadrootsRelayOutcomeKind::ConnectionFailed) - ); - - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::PublishRetryable); - assert!(event.claim_token.is_none()); - assert_eq!(event.next_attempt_after_ms, 2_500); - - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert_eq!(targets.len(), 2); - assert!( - targets - .iter() - .all(|target| target.status == RadrootsOutboxDeliveryTargetStatus::FailedRetryable) - ); - assert!( - outbox - .claim_next_ready_event("publisher", "publish-b", 4_000, 2_499) - .await - .expect("early claim") - .is_none() - ); - let retry_claim = outbox - .claim_next_ready_event("publisher", "publish-b", 4_000, 2_500) - .await - .expect("retry claim") - .expect("retry claim"); - assert_eq!(retry_claim.outbox_event_id, receipt.outbox_event_id); - assert_eq!(retry_claim.state, RadrootsOutboxEventState::Publishing); -} - -#[tokio::test] -async fn outbox_publish_marks_published_without_adapter_when_all_relays_already_accepted() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("already accepted"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let initial_targets = publish_claim.delivery_targets.clone(); - outbox - .mark_delivery_target_accepted( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - initial_targets[0].delivery_target_id, - 2_150, - ) - .await - .expect("primary accepted"); - outbox - .mark_delivery_target_accepted( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - initial_targets[1].delivery_target_id, - 2_151, - ) - .await - .expect("secondary accepted"); - - let adapter = RadrootsMockRelayPublishAdapter::new(); - let published = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.local_ingest.event_id, signed.id_str()); - assert_eq!(published.event_id, signed.id_str()); - assert_eq!(published.attempted_count, 0); - assert_eq!(published.accepted_count, 2); - assert_eq!(published.quorum, 0); - assert!(published.quorum_met); - assert!(published.target_receipts.is_empty()); - assert!(published.relay_receipts.is_empty()); - assert!(adapter.captured_raw_events().is_empty()); - - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); - assert!(event.claim_token.is_none()); - let operation = outbox - .get_operation(receipt.operation_id) - .await - .expect("operation") - .expect("operation"); - assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete); -} - -#[tokio::test] -async fn outbox_publish_rejects_unknown_adapter_receipts() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("unknown receipt"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - vec![RELAY_PRIMARY_WSS.to_owned()], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - - let error = publish_claimed_outbox_event( - &outbox, - &store, - &UnknownRelayReceiptPublishAdapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect_err("unknown adapter receipt"); - - assert!(matches!( - error, - RadrootsRelayTransportError::UnexpectedPublishReceiptRelayUrl { url } - if url == RELAY_TERTIARY_WSS - )); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Publishing); - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_outbox_publish_observations(&observations, 0); -} - -#[tokio::test] -async fn outbox_publish_skips_non_nostr_targets() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("mixed target"); - let receipt = outbox - .enqueue_operation(RadrootsOutboxOperationInput::new( - "publish_post", - draft, - RadrootsOutboxDeliveryPlanInput::new( - "transport.mixed.local", - 1, - RadrootsTransportSatisfactionPolicy::all_accepted(), - vec![nostr_target(RELAY_PRIMARY_WSS), reticulum_target()], - ), - 1_000, - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let adapter = RadrootsMockRelayPublishAdapter::new(); - - let published = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.attempted_count, 1); - assert_eq!(adapter.captured_raw_events().len(), 1); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert!(targets.iter().any(|target| { - target.transport_kind == TransportId::RETICULUM - && target.status == RadrootsOutboxDeliveryTargetStatus::DeferredUntilImplemented - })); -} - -#[tokio::test] -async fn outbox_publish_marks_published_when_delivery_plan_satisfaction_is_met_with_failure_diagnostics() - { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("quorum"); - let receipt = outbox - .enqueue_operation(outbox_operation_input( - draft, - vec![ - RELAY_PRIMARY_WSS.to_owned(), - RELAY_SECONDARY_WSS.to_owned(), - RELAY_TERTIARY_WSS.to_owned(), - ], - RadrootsTransportSatisfactionPolicy::quorum_accepted(2), - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()) - .with_outcome( - RELAY_SECONDARY_WSS, - RadrootsRelayOutcome::duplicate_accepted("duplicate: already have it"), - ) - .with_outcome( - RELAY_TERTIARY_WSS, - RadrootsRelayOutcome::classify("restricted: group write denied"), - ); - let published = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.quorum, 2); - assert_eq!(published.accepted_count, 2); - assert_eq!(published.terminal_count, 1); - assert!(published.quorum_met); - - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); - assert!(event.claim_token.is_none()); - let operation = outbox - .get_operation(receipt.operation_id) - .await - .expect("operation") - .expect("operation"); - assert_eq!(operation.status, RadrootsOutboxOperationStatus::Complete); - - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert_eq!( - targets - .iter() - .find(|target| target.endpoint_uri.as_str() == RELAY_TERTIARY_WSS) - .expect("tertiary") - .status, - RadrootsOutboxDeliveryTargetStatus::FailedTerminal - ); - assert!( - outbox - .claim_next_ready_event("publisher", "publish-b", 4_000, 2_300) - .await - .expect("claim") - .is_none() - ); - - let observations = store - .observations_for_event(signed.id_str()) - .await - .expect("observations"); - assert_outbox_publish_observations(&observations, 2); -} - -#[tokio::test] -async fn outbox_publish_republishes_accepted_relays_when_policy_requests_it() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("republish accepted"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - let signed = complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let initial_targets = publish_claim.delivery_targets.clone(); - outbox - .mark_delivery_target_accepted( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - initial_targets[0].delivery_target_id, - 2_150, - ) - .await - .expect("primary accepted"); - - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()) - .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted()); - let published = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500) - .republish_accepted_relays(true) - .relay_url_policy(RadrootsRelayUrlPolicy::Public), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.local_ingest.event_id, signed.id_str()); - assert_eq!(published.attempted_count, 2); - assert_eq!(published.accepted_count, 2); - assert_eq!(published.quorum, 1); - assert!(published.quorum_met); - assert_eq!(adapter.captured_raw_events().len(), 1); - - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Published); - let targets = outbox - .delivery_targets(receipt.outbox_event_id) - .await - .expect("targets"); - assert!( - targets - .iter() - .all(|target| target.status == RadrootsOutboxDeliveryTargetStatus::Accepted) - ); -} - -#[tokio::test] -async fn outbox_publish_republish_policy_keeps_terminal_targets_excluded() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("republish terminal excluded"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let initial_targets = publish_claim.delivery_targets.clone(); - outbox - .mark_delivery_target_accepted( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - initial_targets[0].delivery_target_id, - 2_150, - ) - .await - .expect("primary accepted"); - outbox - .mark_delivery_target_failed_terminal( - publish_claim.outbox_event_id, - publish_claim.claim_token.as_str(), - initial_targets[1].delivery_target_id, - "terminal", - 2_151, - ) - .await - .expect("secondary terminal"); - let adapter = RadrootsMockRelayPublishAdapter::new() - .with_outcome(RELAY_PRIMARY_WSS, RadrootsRelayOutcome::accepted()) - .with_outcome(RELAY_SECONDARY_WSS, RadrootsRelayOutcome::accepted()); - - let published = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500).republish_accepted_relays(true), - 2_200, - ) - .await - .expect("publish"); - - assert_eq!(published.attempted_count, 1); - assert_eq!(published.accepted_count, 1); - assert_eq!(published.quorum, 1); - assert!(published.quorum_met); - assert_eq!(adapter.captured_raw_events().len(), 1); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::FailedTerminal); -} - -#[tokio::test] -async fn outbox_publish_requires_claimed_signed_event() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("missing signature"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - vec![RELAY_PRIMARY_WSS.to_owned()], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - let adapter = RadrootsMockRelayPublishAdapter::new(); - - let error = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &claimed, - RadrootsOutboxPublishPolicy::new(2_500), - 1_100, - ) - .await - .expect_err("missing signed event"); - - assert!(matches!( - error, - RadrootsRelayTransportError::MissingSignedOutboxEvent(event_id) - if event_id == receipt.outbox_event_id - )); - assert!(adapter.captured_raw_events().is_empty()); -} - -#[tokio::test] -async fn outbox_publish_propagates_non_transport_adapter_errors_after_target_filtering() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("adapter non transport failure"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - vec![RELAY_PRIMARY_WSS.to_owned(), RELAY_SECONDARY_WSS.to_owned()], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - complete_claimed_signing(&outbox, &claimed, 1_100).await; - let mut publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - publish_claim.delivery_targets.truncate(1); - - let error = publish_claimed_outbox_event( - &outbox, - &store, - &NostrJsonFailurePublishAdapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect_err("adapter error"); - - assert!(matches!( - error, - RadrootsRelayTransportError::NostrEventJson(_) - )); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Publishing); -} - -#[tokio::test] -async fn outbox_publish_rejects_invalid_relay_target_uri_before_adapter_publish() { - let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); - let store = RadrootsEventStore::open_memory().await.expect("store"); - let draft = generic_draft("invalid relay target"); - let receipt = outbox - .enqueue_operation(all_accepted_outbox_operation_input( - draft, - vec!["ws://127.0.0.1:9999".to_owned()], - )) - .await - .expect("enqueue"); - let claimed = outbox - .claim_next_ready_event("signer", "sign-a", 2_000, 1_000) - .await - .expect("claim") - .expect("claim"); - complete_claimed_signing(&outbox, &claimed, 1_100).await; - let publish_claim = outbox - .claim_next_ready_event("publisher", "publish-a", 3_000, 1_100) - .await - .expect("claim") - .expect("publish claim"); - let adapter = RadrootsMockRelayPublishAdapter::new(); - - let error = publish_claimed_outbox_event( - &outbox, - &store, - &adapter, - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_200, - ) - .await - .expect_err("invalid relay target"); - - assert!(matches!( - error, - RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. } - )); - assert!(adapter.captured_raw_events().is_empty()); - - let transport_error = publish_claimed_outbox_event_with_transport( - &outbox, - &store, - &ScriptedTransport::new(Vec::new()), - &publish_claim, - RadrootsOutboxPublishPolicy::new(2_500), - 2_201, - ) - .await - .expect_err("invalid transport relay target"); - assert!(matches!( - transport_error, - RadrootsRelayTransportError::RelayUrlForbiddenDestination { .. } - )); - let event = outbox - .get_event(receipt.outbox_event_id) - .await - .expect("event") - .expect("event"); - assert_eq!(event.state, RadrootsOutboxEventState::Publishing); -} - -#[tokio::test] -async fn smoke_relay_fetch_processes_one_thousand_event_receipts() { - let store = RadrootsEventStore::open_memory().await.expect("store"); - let mut items = Vec::new(); - for index in 0..1_000 { - let signed = signed_post(format!("fetch-smoke-{index}").as_str()); - let relay_url = match index % 3 { - 0 => RELAY_PRIMARY_WSS, - 1 => RELAY_SECONDARY_WSS, - _ => RELAY_TERTIARY_WSS, - }; - items.push(RadrootsRelayFetchItem::Event { - relay_url: relay_url.to_owned(), - raw_json: signed.raw_json().to_owned(), - }); - } - let adapter = RadrootsMockRelayFetchAdapter::new(items); - let receipt = - fetch_and_ingest_relay_events(&adapter, &store, post_relay_fetch_request(10_000, 1_000)) - .await - .expect("fetch"); - - assert_eq!(receipt.inserted_count, 1_000); - assert_eq!(receipt.duplicate_count, 0); - assert_eq!(receipt.malformed_count, 0); - assert_eq!(receipt.verification_failed_count, 0); - assert_eq!(receipt.admission_unsupported_count, 0); - assert_eq!(receipt.admission_invalid_count, 0); - assert_eq!(receipt.valid_stream_eligible_count, 1_000); - assert_eq!(receipt.visible_count, 1_000); - assert_eq!(receipt.events.len(), 1_000); - assert!(receipt.events.iter().all(|event| event.valid_stream - == RadrootsRelayFetchEventValidStream::Eligible - && event.visibility == RadrootsRelayFetchEventVisibility::Visible)); - let replay = store.valid_stream_after(0, 1_000).await.expect("replay"); - assert_eq!(replay.len(), 1_000); -}