lib

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

commit 0fe3881f88753fcb72be17b51274dd932fc5714a
parent 70abf6eb6bcbf07f4074b6fba237add4681c00eb
Author: triesap <tyson@radroots.org>
Date:   Mon, 21 Sep 2026 06:41:56 +0000

transport_nostr: preserve exact signed event bytes

- Retain EVENT payload bytes across sealed delivery preparation
- Serialize SDK and raw writes on the same hardened connection
- Preserve acknowledgement, lifetime, deadline and cancellation bounds
- Verify loopback delivery, coverage and complete workspace gates

Diffstat:
Acontracts/architecture/decisions/nostr_exact_delivery.v1.json | 14++++++++++++++
Mcrates/transport_nostr/README.md | 11++++++++++-
Mcrates/transport_nostr/src/client.rs | 8++++----
Acrates/transport_nostr/src/exact_delivery.rs | 142+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/transport_nostr/src/lib.rs | 2++
Mcrates/transport_nostr/src/relay.rs | 16+++++++++++++---
Mcrates/transport_nostr/src/sink.rs | 70+++++++++++++++++++++++++++++++++++++---------------------------------
Acrates/transport_nostr/src/socket_write.rs | 182+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/transport_nostr/src/socket_write_tests.rs | 129+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/transport_nostr/tests/exact_delivery.rs | 255+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/transport_nostr/tests/package_boundary.rs | 7++++++-
11 files changed, 794 insertions(+), 42 deletions(-)

diff --git a/contracts/architecture/decisions/nostr_exact_delivery.v1.json b/contracts/architecture/decisions/nostr_exact_delivery.v1.json @@ -0,0 +1,14 @@ +{ + "schema": "radroots.nostr-exact-delivery.v1", + "owner": "radroots_transport_nostr", + "status": "implemented", + "scope": "Existing sealed PreparedDelivery and EventSink publication of a validated SignedEvent.", + "wire": "The EVENT frame contains the byte-identical SignedEvent.raw_json(), including whitespace, key order and admitted extensions. Validate the event conversion and complete frame bound before constructing prepared authority. Never replace retained JSON with an SDK serialization.", + "max_wire_message_bytes": 524288, + "connection": "SDK authentication, subscription requests and exact publication share one serialized writer on the existing DNS-pinned and TLS-verified connection. No second socket engine, executor, adapter task, event rewrite registry, signer or durable journal is introduced.", + "lifetime": "An instance-local registry contains only the configured relay keys and weak connection references. SDK sink destruction or write failure revokes raw send authority. Reconnection installs a distinct writer; an old handle cannot switch to the replacement connection.", + "receipt": "Subscribe to the same relay notification stream before sending. Only a matching event-ID OK may accept or reject this attempt. Other IDs cannot satisfy it. Timeout, disconnect, shutdown, notification loss and missing results cannot invent acceptance.", + "cancellation": "Futures are caller-polled. Cancelling a write releases its serialization guard but does not prove that no frame was buffered or published. Durable retry, settlement and late facts remain caller-owned.", + "bounds": "Preserve configured target access, bounded relay inventory, connection concurrency, frame bounds, timeouts and reconnect suppression. Queued relay attempts consume the same frozen operation deadline.", + "compatibility": "Public API, generic transport SPI, dependency versions, persisted event schemas, verification thresholds and release authority remain unchanged." +} diff --git a/crates/transport_nostr/README.md b/crates/transport_nostr/README.md @@ -66,6 +66,15 @@ request to durable Submitted state and then pass the capability to only half of this boundary that may contact relays. Executing a capability through a differently configured transport fails closed. +The prepared event retains its signed raw JSON. Delivery wraps those exact +bytes in the Nostr `EVENT` frame without parsing and serializing them again; +whitespace, field order, and admitted extension fields remain unchanged. The +frame must fit the existing 512 KiB wire bound. Publication shares the existing +hardened connection writer with SDK authentication and subscription messages. +It requires a matching event-ID `OK` from that relay before recording acceptance. +Closed or replaced connections do not grant fresh send authority to a retained +writer. Missing acknowledgements remain unknown or unavailable evidence. + Callers cannot forge or mutate prepared authority: ```compile_fail @@ -181,7 +190,7 @@ exact relay provenance and that current checkpoint. Event limits, absolute deadlines, explicit cancellation, source closure, and stable repeated terminal results follow the generic subscription contract. -Delivery converts an already validated signed Radroots event to Nostr, attempts +Delivery validates an already signed Radroots event and sends its retained JSON, attempts each configured writable target once, and returns one normalized receipt entry per requested target. Relay rejection, authentication requirements, rate limits, timeouts, connection failures, missing results, and partial acceptance remain diff --git a/crates/transport_nostr/src/client.rs b/crates/transport_nostr/src/client.rs @@ -228,16 +228,16 @@ pub struct NostrTransport { impl NostrTransport { /// Creates an inert transport from validated explicit configuration. pub fn new(config: Config) -> Self { + let connector = crate::relay::HardenedWebsocketTransport::new(config.endpoints()); + let writers = connector.writers.clone(); let client = nostr_sdk::Client::builder() - .websocket_transport(crate::relay::HardenedWebsocketTransport::new( - config.endpoints(), - )) + .websocket_transport(connector) .build(); client.automatic_authentication(false); let status = Arc::new(crate::status::StatusTracker::new(&config)); Self { config, - client: Arc::new(crate::sink::LiveRelayClient::new(client.clone())), + client: Arc::new(crate::sink::LiveRelayClient::new(client.clone(), writers)), source_client: Arc::new(crate::source::LiveRelaySourceClient::new(client.clone())), subscription_client: Arc::new(crate::subscription::LiveRelaySubscriptionClient::new( client.clone(), diff --git a/crates/transport_nostr/src/exact_delivery.rs b/crates/transport_nostr/src/exact_delivery.rs @@ -0,0 +1,142 @@ +//! Exact retained EVENT framing and connection-local acknowledgement handling. + +use crate::{relay::MAX_WIRE_MESSAGE_BYTES, socket_write::WriterRegistry, status}; +use nostr_relay_pool::RelayNotification; +use nostr_sdk::{EventId, RelayMessage}; +use radroots_transport::{DeliveryRequest, outcome::DeliveryOutcome}; +use std::sync::Arc; + +#[derive(Clone)] +pub(crate) struct ExactEvent { + id: EventId, + frame: Arc<str>, +} + +impl ExactEvent { + pub(crate) fn from_request(request: &DeliveryRequest) -> Option<Self> { + let signed = request.payload().event(); + let raw = signed.raw_json(); + let event = radroots_nostr::event::to_nostr(signed.envelope()).ok()?; + Some(Self { + id: event.id, + frame: frame(raw)?, + }) + } + + pub(crate) async fn publish( + &self, + client: &nostr_sdk::Client, + writers: &WriterRegistry, + url: &str, + ) -> Result<DeliveryOutcome, String> { + let relay = client.relay(url).await.map_err(|error| error.to_string())?; + // Subscribe before sending on the same connection; no success from enqueue alone. + let mut notifications = relay.notifications(); + let writer = writers.get(url).map_err(|error| error.to_string())?; + writer + .send(async_wsocket::Message::Text(self.frame.to_string())) + .await + .map_err(|error| error.to_string())?; + loop { + let notification = notifications + .recv() + .await + .map_err(|_| "relay acknowledgement unavailable".to_owned())?; + if let Some(outcome) = acknowledgement(self.id, notification) { + return Ok(outcome); + } + } + } +} + +fn frame(raw: &str) -> Option<Arc<str>> { + // Bound before allocation; preserve whitespace, field order and extensions. + (raw.len() <= MAX_WIRE_MESSAGE_BYTES - "[\"EVENT\",]".len()) + .then(|| format!("[\"EVENT\",{raw}]").into()) +} + +fn acknowledgement(id: EventId, notification: RelayNotification) -> Option<DeliveryOutcome> { + match notification { + RelayNotification::Message { + message: + RelayMessage::Ok { + event_id, + status: accepted, + message, + }, + } if event_id == id => Some(if accepted { + DeliveryOutcome::accepted() + } else { + status::delivery_failure(&message) + }), + RelayNotification::RelayStatus { + status: + nostr_sdk::RelayStatus::Disconnected + | nostr_sdk::RelayStatus::Terminated + | nostr_sdk::RelayStatus::Banned, + } => Some(status::delivery_failure("relay disconnected")), + RelayNotification::Shutdown => Some(status::delivery_failure("relay shutdown")), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn framing_bounds_complete_message_bytes_before_allocation() { + let limit = MAX_WIRE_MESSAGE_BYTES - "[\"EVENT\",]".len(); + let exact = " ".repeat(limit); + assert_eq!(frame(&exact).unwrap().len(), MAX_WIRE_MESSAGE_BYTES); + assert!(frame(&(exact + " ")).is_none()); + let unicode = "🌱".repeat(limit / 4 + 1); + assert!(frame(&unicode).is_none()); + assert_eq!( + &*frame(" \n{\"extension\":true} \t").unwrap(), + "[\"EVENT\", \n{\"extension\":true} \t]" + ); + } + + #[test] + fn only_matching_ok_can_accept_and_disconnects_remain_failures() { + let id = EventId::from_byte_array([1; 32]); + let other = EventId::from_byte_array([2; 32]); + let ok = |event_id, status, message: &'static str| RelayNotification::Message { + message: RelayMessage::Ok { + event_id, + status, + message: message.into(), + }, + }; + assert!(acknowledgement(id, ok(other, true, "")).is_none()); + assert_eq!( + acknowledgement(id, ok(id, true, "")).unwrap(), + DeliveryOutcome::accepted() + ); + assert!(!status::delivery_succeeded( + &acknowledgement(id, ok(id, false, "blocked: denied")).unwrap() + )); + for status in [ + nostr_sdk::RelayStatus::Disconnected, + nostr_sdk::RelayStatus::Terminated, + nostr_sdk::RelayStatus::Banned, + ] { + assert!(!super::status::delivery_succeeded( + &acknowledgement(id, RelayNotification::RelayStatus { status }).unwrap() + )); + } + assert!( + acknowledgement( + id, + RelayNotification::RelayStatus { + status: nostr_sdk::RelayStatus::Connected + } + ) + .is_none() + ); + assert!(!status::delivery_succeeded( + &acknowledgement(id, RelayNotification::Shutdown).unwrap() + )); + } +} diff --git a/crates/transport_nostr/src/lib.rs b/crates/transport_nostr/src/lib.rs @@ -6,9 +6,11 @@ mod auth; mod client; mod cursor; mod error; +mod exact_delivery; mod profile; mod relay; mod sink; +mod socket_write; mod source; mod status; mod subscription; diff --git a/crates/transport_nostr/src/relay.rs b/crates/transport_nostr/src/relay.rs @@ -19,7 +19,7 @@ use tokio::net::TcpStream; use url::Url; const MAX_RESOLVED_ADDRESSES: usize = 32; -const MAX_WIRE_MESSAGE_BYTES: usize = 512 * 1024; +pub(crate) const MAX_WIRE_MESSAGE_BYTES: usize = 512 * 1024; /// Validated canonical Nostr relay URL. #[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] @@ -97,11 +97,17 @@ impl fmt::Display for RelayUrl { #[derive(Clone, Debug)] pub(crate) struct HardenedWebsocketTransport { policies: Arc<BTreeMap<String, RelayUrlPolicy>>, + pub(crate) writers: crate::socket_write::WriterRegistry, } impl HardenedWebsocketTransport { pub(crate) fn new(endpoints: &[RelayEndpoint]) -> Self { Self { + writers: crate::socket_write::WriterRegistry::new( + endpoints + .iter() + .map(|endpoint| endpoint.url().as_str().to_owned()), + ), policies: Arc::new( endpoints .iter() @@ -167,7 +173,11 @@ impl WebSocketTransport for HardenedWebsocketTransport { .map_err(TransportError::backend)?; let socket = WebSocket::Tokio(stream); let (tx, rx) = socket.split(); - let sink: WebSocketSink = Box::new(HardenedTransportSink(tx)); + let writer = + crate::socket_write::SocketWriter::new(Box::new(HardenedTransportSink(tx))); + self.writers.install(relay.as_str(), &writer)?; + let sink: WebSocketSink = + Box::new(crate::socket_write::SharedSocketSink::new(writer)); let stream: WebSocketStream = Box::pin(rx.map_err(TransportError::backend)) as WebSocketStream; Ok((sink, stream)) @@ -229,7 +239,7 @@ impl fmt::Display for NetworkPolicyError { impl std::error::Error for NetworkPolicyError {} -fn policy_error(message: &'static str) -> TransportError { +pub(crate) fn policy_error(message: &'static str) -> TransportError { TransportError::backend(NetworkPolicyError(message)) } diff --git a/crates/transport_nostr/src/sink.rs b/crates/transport_nostr/src/sink.rs @@ -1,9 +1,9 @@ //! Nostr implementation of the transport event sink. +use crate::exact_delivery::ExactEvent; use crate::{NostrTransport, RelayUrl, status}; use core::{fmt, time::Duration}; use futures::{StreamExt, stream}; -use radroots_nostr::event::Event; use radroots_transport::{ BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure, Target, outcome::DeliveryOutcome, @@ -19,7 +19,7 @@ pub(crate) struct RelayPublishResult { /// Sealed, no-I/O result of validating one delivery against this adapter. /// -/// The value retains the exact request and converted signed event. It is +/// The value retains the exact request and signed event bytes. It is /// constructed only by [`NostrTransport::prepare_delivery`] and is consumed by /// [`NostrTransport::execute_prepared_delivery`]. Ordinary `Debug` never /// exposes event bytes, request identities, or relay destinations. @@ -27,7 +27,7 @@ pub(crate) struct RelayPublishResult { pub struct PreparedDelivery { request: DeliveryRequest, config: crate::Config, - event: Event, + event: ExactEvent, authorized: Vec<(RelayUrl, Target)>, skipped: Vec<DeliveryTargetReceipt>, } @@ -50,7 +50,7 @@ pub(crate) trait RelayClient: Send + Sync { fn publish<'a>( &'a self, relays: Vec<RelayUrl>, - event: Event, + event: ExactEvent, max_connections: usize, connect_timeout: Duration, operation_timeout: Duration, @@ -60,60 +60,64 @@ pub(crate) trait RelayClient: Send + Sync { #[derive(Clone, Debug)] pub(crate) struct LiveRelayClient { client: nostr_sdk::Client, + writers: crate::socket_write::WriterRegistry, } impl LiveRelayClient { - pub(crate) const fn new(client: nostr_sdk::Client) -> Self { - Self { client } + pub(crate) const fn new( + client: nostr_sdk::Client, + writers: crate::socket_write::WriterRegistry, + ) -> Self { + Self { client, writers } } #[cfg(test)] pub(crate) fn isolated() -> Self { let client = nostr_sdk::Client::default(); client.automatic_authentication(false); - Self::new(client) + Self::new( + client, + crate::socket_write::WriterRegistry::new(std::iter::empty()), + ) } } impl RelayClient for LiveRelayClient { - // The live SDK loop requires external relays. Its result normalization is - // covered through the injected RelayClient boundary below. + // Socket scheduling is exercised by the real loopback suite. Deterministic + // coverage owns framing, writer lifecycle and acknowledgement normalization. #[cfg_attr(coverage_nightly, coverage(off))] fn publish<'a>( &'a self, relays: Vec<RelayUrl>, - event: Event, + event: ExactEvent, max_connections: usize, connect_timeout: Duration, operation_timeout: Duration, ) -> BoxFuture<'a, Vec<RelayPublishResult>> { Box::pin(async move { + let deadline = tokio::time::Instant::now() + operation_timeout; stream::iter(relays.into_iter().map(|relay| { let event = event.clone(); async move { let url = relay.as_str().to_owned(); let attempt = async { - self.client.add_relay(url.as_str()).await?; + if tokio::time::Instant::now() >= deadline { + return Err("timeout".to_owned()); + } + self.client + .add_relay(url.as_str()) + .await + .map_err(|error| error.to_string())?; self.client .try_connect_relay(url.as_str(), connect_timeout) - .await?; - self.client.send_event_to([url.as_str()], &event).await + .await + .map_err(|error| error.to_string())?; + event.publish(&self.client, &self.writers, &url).await }; - let outcome = match tokio::time::timeout(operation_timeout, attempt).await { + let outcome = match tokio::time::timeout_at(deadline, attempt).await { Err(_) => status::delivery_failure("timeout"), - Ok(Err(error)) => status::delivery_failure(error.to_string().as_str()), - Ok(Ok(output)) => output - .success - .iter() - .any(|accepted| accepted.to_string().trim_end_matches('/') == url) - .then(DeliveryOutcome::accepted) - .or_else(|| { - output.failed.iter().find_map(|(failed, message)| { - (failed.to_string().trim_end_matches('/') == url) - .then(|| status::delivery_failure(message.as_str())) - }) - }) - .unwrap_or_else(|| status::delivery_failure("relay omitted result")), + Ok(Err(error)) => status::delivery_failure(&error), + Ok(Ok(outcome)) => outcome, }; RelayPublishResult { relay, outcome } } @@ -149,8 +153,8 @@ impl NostrTransport { ), } } - let event = radroots_nostr::event::to_nostr(request.payload().event().envelope()) - .map_err(|_| Box::new(SinkFailure::invalid_contract(&request)))?; + let event = ExactEvent::from_request(&request) + .ok_or_else(|| Box::new(SinkFailure::invalid_contract(&request)))?; Ok(PreparedDelivery { request, config: self.config().clone(), @@ -305,7 +309,7 @@ mod tests { fn publish<'a>( &'a self, relays: Vec<RelayUrl>, - _event: Event, + _event: ExactEvent, _max_connections: usize, _connect_timeout: Duration, _operation_timeout: Duration, @@ -333,7 +337,7 @@ mod tests { fn publish<'a>( &'a self, _relays: Vec<RelayUrl>, - _event: Event, + _event: ExactEvent, _max_connections: usize, _connect_timeout: Duration, _operation_timeout: Duration, @@ -507,7 +511,7 @@ mod tests { fn publish<'a>( &'a self, _relays: Vec<RelayUrl>, - _event: Event, + _event: ExactEvent, _max_connections: usize, _connect_timeout: Duration, _operation_timeout: Duration, @@ -592,7 +596,7 @@ mod tests { let client = LiveRelayClient::isolated(); let results = futures::executor::block_on(client.publish( vec![], - radroots_nostr::event::to_nostr(payload().event().envelope()).expect("nostr event"), + ExactEvent::from_request(&request()).expect("exact event"), 1, Duration::from_millis(1), Duration::from_millis(1), diff --git a/crates/transport_nostr/src/socket_write.rs b/crates/transport_nostr/src/socket_write.rs @@ -0,0 +1,182 @@ +//! One serialized writer for SDK messages and exact signed-event publication. + +use async_wsocket::Message; +use core::{ + fmt, + pin::Pin, + task::{Context, Poll}, +}; +use futures::{Sink, SinkExt, lock::Mutex}; +use nostr_relay_pool::transport::{error::TransportError, websocket::WebSocketSink}; +use radroots_transport::BoxFuture; +use std::{ + collections::BTreeMap, + sync::{ + Arc, Mutex as RegistryMutex, Weak, + atomic::{AtomicBool, Ordering}, + }, +}; + +use crate::relay::policy_error; + +pub(crate) struct SocketWriter { + sink: Mutex<WebSocketSink>, + open: AtomicBool, +} + +impl SocketWriter { + pub(crate) fn new(sink: WebSocketSink) -> Arc<Self> { + Arc::new(Self { + sink: Mutex::new(sink), + open: AtomicBool::new(true), + }) + } + + pub(crate) async fn send(&self, message: Message) -> Result<(), TransportError> { + let mut sink = self.sink.lock().await; + if !self.open.load(Ordering::Acquire) { + return Err(policy_error("relay writer is closed")); + } + let result = sink.send(message).await; + if result.is_err() { + self.invalidate(); + } + result + } + + fn invalidate(&self) { + self.open.store(false, Ordering::Release); + } + + async fn close(&self) -> Result<(), TransportError> { + self.invalidate(); + self.sink.lock().await.close().await + } +} + +/// Configured keys only; the registry never owns a connection or event bytes. +#[derive(Clone)] +pub(crate) struct WriterRegistry(Arc<RegistryMutex<BTreeMap<String, Weak<SocketWriter>>>>); + +impl WriterRegistry { + pub(crate) fn new(keys: impl Iterator<Item = String>) -> Self { + Self(Arc::new(RegistryMutex::new( + keys.map(|key| (key, Weak::new())).collect(), + ))) + } + + pub(crate) fn install( + &self, + key: &str, + writer: &Arc<SocketWriter>, + ) -> Result<(), TransportError> { + let mut entries = self + .0 + .lock() + .map_err(|_| policy_error("relay writer registry unavailable"))?; + let slot = entries + .get_mut(key) + .ok_or_else(|| policy_error("relay writer is not configured"))?; + if let Some(previous) = slot.upgrade() { + previous.invalidate(); + } + *slot = Arc::downgrade(writer); + Ok(()) + } + + pub(crate) fn get(&self, key: &str) -> Result<Arc<SocketWriter>, TransportError> { + self.0 + .lock() + .map_err(|_| policy_error("relay writer registry unavailable"))? + .get(key) + .and_then(Weak::upgrade) + .filter(|writer| writer.open.load(Ordering::Acquire)) + .ok_or_else(|| policy_error("relay writer is unavailable")) + } +} + +impl fmt::Debug for WriterRegistry { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("WriterRegistry([redacted])") + } +} + +/// The SDK retains connection ownership; dropping its sink revokes raw writes. +pub(crate) struct SharedSocketSink { + writer: Arc<SocketWriter>, + pending: Option<BoxFuture<'static, Result<(), TransportError>>>, + closing: bool, +} + +impl SharedSocketSink { + pub(crate) fn new(writer: Arc<SocketWriter>) -> Self { + Self { + writer, + pending: None, + closing: false, + } + } + + fn poll_pending(&mut self, context: &mut Context<'_>) -> Poll<Result<(), TransportError>> { + if let Some(pending) = &mut self.pending { + let result = futures::ready!(pending.as_mut().poll(context)); + self.pending = None; + return Poll::Ready(result); + } + Poll::Ready(Ok(())) + } +} + +impl Sink<Message> for SharedSocketSink { + type Error = TransportError; + + fn poll_ready( + mut self: Pin<&mut Self>, + context: &mut Context<'_>, + ) -> Poll<Result<(), Self::Error>> { + if self.closing { + return Poll::Ready(Err(policy_error("relay writer is closing"))); + } + self.poll_pending(context) + } + + fn start_send(mut self: Pin<&mut Self>, message: Message) -> Result<(), Self::Error> { + if self.closing || self.pending.is_some() { + return Err(policy_error("relay writer is not ready")); + } + let writer = Arc::clone(&self.writer); + self.pending = Some(Box::pin(async move { writer.send(message).await })); + Ok(()) + } + + fn poll_flush( + mut self: Pin<&mut Self>, + context: &mut Context<'_>, + ) -> Poll<Result<(), Self::Error>> { + self.poll_pending(context) + } + + fn poll_close( + mut self: Pin<&mut Self>, + context: &mut Context<'_>, + ) -> Poll<Result<(), Self::Error>> { + futures::ready!(self.poll_pending(context))?; + if !self.closing { + self.closing = true; + self.writer.invalidate(); + let writer = Arc::clone(&self.writer); + self.pending = Some(Box::pin(async move { writer.close().await })); + } + self.poll_pending(context) + } +} + +impl Drop for SharedSocketSink { + fn drop(&mut self) { + self.writer.invalidate(); + } +} + +#[cfg(test)] +#[path = "socket_write_tests.rs"] +mod tests; diff --git a/crates/transport_nostr/src/socket_write_tests.rs b/crates/transport_nostr/src/socket_write_tests.rs @@ -0,0 +1,129 @@ +use super::*; +use futures::{FutureExt, channel::mpsc, task::noop_waker_ref}; + +fn message(text: &str) -> Message { + Message::Text(text.to_owned()) +} + +fn channel_writer() -> (Arc<SocketWriter>, mpsc::Receiver<Message>) { + let (tx, rx) = mpsc::channel(0); + ( + SocketWriter::new(Box::new(tx.sink_map_err(TransportError::backend))), + rx, + ) +} + +#[test] +fn concurrent_sdk_and_exact_messages_share_backpressure_and_order() { + use futures::StreamExt; + futures::executor::block_on(async { + let (writer, mut rx) = channel_writer(); + let mut sdk = SharedSocketSink::new(Arc::clone(&writer)); + let mut first = Box::pin(sdk.send(message("AUTH"))); + assert!(first.as_mut().now_or_never().is_none()); + let mut raw = Box::pin(writer.send(message("EVENT"))); + assert!(raw.as_mut().now_or_never().is_none()); + assert_eq!(rx.next().await, Some(message("AUTH"))); + first.await.unwrap(); + assert!(raw.as_mut().now_or_never().is_none()); + assert_eq!(rx.next().await, Some(message("EVENT"))); + raw.await.unwrap(); + sdk.close().await.unwrap(); + sdk.close().await.unwrap(); + assert!(sdk.send(message("REQ")).await.is_err()); + assert!(writer.send(message("EVENT")).await.is_err()); + assert!(rx.next().await.is_none()); + }); +} + +#[test] +fn cancelling_a_waiter_releases_its_lock_without_claiming_no_effect() { + use futures::StreamExt; + futures::executor::block_on(async { + let (writer, mut rx) = channel_writer(); + let mut sdk = SharedSocketSink::new(Arc::clone(&writer)); + let mut pending = Box::pin(writer.send(message("first"))); + assert!(pending.as_mut().now_or_never().is_none()); + drop(pending); + // It was already buffered: cancelling a send is not proof of no publication. + assert_eq!(rx.next().await, Some(message("first"))); + let mut next = Box::pin(sdk.send(message("second"))); + assert!(next.as_mut().now_or_never().is_none()); + assert_eq!(rx.next().await, Some(message("second"))); + next.await.unwrap(); + }); +} + +#[test] +fn pending_sdk_message_must_flush_before_next_send_or_close() { + use futures::StreamExt; + futures::executor::block_on(async { + let (writer, mut rx) = channel_writer(); + let mut sdk = SharedSocketSink::new(writer); + let mut context = Context::from_waker(noop_waker_ref()); + assert!(Pin::new(&mut sdk).poll_ready(&mut context).is_ready()); + Pin::new(&mut sdk).start_send(message("first")).unwrap(); + assert!(Pin::new(&mut sdk).start_send(message("unready")).is_err()); + assert!(Pin::new(&mut sdk).poll_close(&mut context).is_pending()); + assert_eq!(rx.next().await, Some(message("first"))); + sdk.close().await.unwrap(); + assert!(Pin::new(&mut sdk).start_send(message("closed")).is_err()); + }); +} + +#[test] +fn sdk_drop_revokes_existing_raw_handles_and_failed_io_revokes_writer() { + futures::executor::block_on(async { + let (writer, rx) = channel_writer(); + drop(SharedSocketSink::new(Arc::clone(&writer))); + assert!(writer.send(message("after drop")).await.is_err()); + drop(rx); + let (writer, rx) = channel_writer(); + drop(rx); + let mut sdk = SharedSocketSink::new(Arc::clone(&writer)); + assert!(sdk.send(message("failure")).await.is_err()); + assert!(!writer.open.load(Ordering::Acquire)); + assert!(writer.send(message("after error")).await.is_err()); + }); +} + +#[test] +fn registry_is_configured_bounded_weak_and_reconnect_isolated() { + let registry = WriterRegistry::new(["relay".to_owned()].into_iter()); + assert_eq!(format!("{registry:?}"), "WriterRegistry([redacted])"); + assert!(registry.get("relay").is_err()); + assert!(registry.get("unknown").is_err()); + let (first, _rx) = channel_writer(); + assert!(registry.install("unknown", &first).is_err()); + registry.install("relay", &first).unwrap(); + let retained = registry.get("relay").unwrap(); + assert!(Arc::ptr_eq(&retained, &first)); + let first_sdk = SharedSocketSink::new(Arc::clone(&first)); + let (second, _rx) = channel_writer(); + registry.install("relay", &second).unwrap(); + assert!(!retained.open.load(Ordering::Acquire)); + drop(first_sdk); + assert!(!retained.open.load(Ordering::Acquire)); + assert!(Arc::ptr_eq(&registry.get("relay").unwrap(), &second)); + second.invalidate(); + assert!(registry.get("relay").is_err()); + drop(second); + assert!(registry.get("relay").is_err()); + assert_eq!(registry.0.lock().unwrap().len(), 1); +} + +#[test] +fn poisoned_registry_fails_closed() { + let registry = WriterRegistry::new(["relay".to_owned()].into_iter()); + let other = registry.clone(); + assert!( + std::panic::catch_unwind(move || { + let _guard = other.0.lock().unwrap(); + panic!("fixture registry poison"); + }) + .is_err() + ); + let (writer, _rx) = channel_writer(); + assert!(registry.install("relay", &writer).is_err()); + assert!(registry.get("relay").is_err()); +} diff --git a/crates/transport_nostr/tests/exact_delivery.rs b/crates/transport_nostr/tests/exact_delivery.rs @@ -0,0 +1,255 @@ +use futures::{FutureExt, SinkExt, StreamExt}; +use nostr_sdk::prelude::{EventBuilder, JsonUtil, Keys, Kind, Tag, Timestamp}; +use radroots_transport::{ + DeliveryRequest, EventSink, EventSource, FetchRequest, TargetSet, + policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, + sink::DeliveryPayload, + source::FetchBounds, +}; +use radroots_transport_nostr::{ + Config, NostrTransport, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind, + RelayUrlPolicy, +}; +use serde_json::{Value, json}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; +use tokio::net::{TcpListener, TcpStream}; +use tokio_tungstenite::{WebSocketStream, accept_async, tungstenite::Message}; + +const KEY: &str = "0000000000000000000000000000000000000000000000000000000000000001"; + +fn now() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_millis() + .try_into() + .unwrap() +} + +fn raw_event() -> String { + let event = EventBuilder::text_note("harvest café\n🌱") + .custom_created_at(Timestamp::from_secs(1_800_000_000)) + .sign_with_keys(&Keys::parse(KEY).unwrap()) + .unwrap(); + let mut value: Value = serde_json::from_str(&event.as_json()).unwrap(); + value["extension"] = json!({"retained": true}); + format!(" \n{} \t", serde_json::to_string_pretty(&value).unwrap()) +} + +fn config(url: &str, timeout: u64) -> Config { + Config::from_profile( + RelayProfile::explicit( + RelayProfileKind::Simulator, + [RelayEndpoint::new(url, RelayUrlPolicy::Local, RelayAccess::ReadWrite).unwrap()], + ) + .unwrap(), + ) + .with_timeouts(1_000, timeout, 500) + .unwrap() +} + +fn targets(config: &Config) -> TargetSet { + TargetSet::new( + config + .relays() + .iter() + .map(|relay| relay.to_target().unwrap()) + .collect(), + ) + .unwrap() +} + +fn request(config: &Config, raw: &str) -> DeliveryRequest { + DeliveryRequest::new( + "exact-wire", + DeliveryPayload::new(radroots_event_codec::decode::signed_event(raw).unwrap()), + targets(config), + SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), + now() + 5_000, + ) + .unwrap() +} + +async fn text(socket: &mut WebSocketStream<TcpStream>) -> String { + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if let Message::Text(text) = socket.next().await.unwrap().unwrap() { + return text.to_string(); + } + } + }) + .await + .unwrap() +} + +async fn reply(socket: &mut WebSocketStream<TcpStream>, value: Value) { + socket + .send(Message::Text(value.to_string().into())) + .await + .unwrap(); +} + +#[tokio::test] +async fn exact_signed_bytes_read_requests_and_auth_use_the_same_connection() { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let config = config(&format!("ws://{}", listener.local_addr().unwrap()), 2_000); + let raw = raw_event(); + let expected = format!("[\"EVENT\",{raw}]"); + let (auth_sent, auth_seen) = tokio::sync::oneshot::channel(); + let server = tokio::spawn(async move { + let (tcp, _) = listener.accept().await.unwrap(); + let mut socket = accept_async(tcp).await.unwrap(); + let req: Value = serde_json::from_str(&text(&mut socket).await).unwrap(); + assert_eq!(req[0], "REQ"); + reply(&mut socket, json!(["EOSE", req[1]])).await; + loop { + let message: Value = serde_json::from_str(&text(&mut socket).await).unwrap(); + if message[0] == "AUTH" { + assert_eq!(message[1]["kind"], 22242); + auth_sent.send(()).unwrap(); + break; + } + assert_eq!(message[0], "CLOSE"); + } + let wire = text(&mut socket).await; + assert_eq!(wire, expected); + let event: Value = serde_json::from_str(&wire).unwrap(); + reply( + &mut socket, + json!(["OK", "11".repeat(32), true, "unrelated event"]), + ) + .await; + reply(&mut socket, json!(["OK", event[1]["id"], true, ""])).await; + }); + let transport = NostrTransport::new(config.clone()); + let page = transport + .fetch( + FetchRequest::new( + "same-socket-read", + targets(&config), + FetchBounds::new(1, now() + 5_000).unwrap(), + ) + .unwrap(), + ) + .await + .unwrap(); + assert!(page.events().is_empty()); + let relay = &config.relays()[0]; + let at = now() / 1_000 * 1_000; + transport + .begin_authentication(relay, "challenge", at, at + 5_000) + .unwrap(); + let auth = EventBuilder::new(Kind::Authentication, "") + .tags([ + Tag::parse(["relay", relay.as_str()]).unwrap(), + Tag::parse(["challenge", "challenge"]).unwrap(), + ]) + .custom_created_at(Timestamp::from_secs(at / 1_000)) + .sign_with_keys(&Keys::parse(KEY).unwrap()) + .unwrap(); + transport + .complete_authentication(relay, "challenge", Some(&auth.as_json()), at) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(5), auth_seen) + .await + .unwrap() + .unwrap(); + let receipt = transport.deliver(request(&config, &raw)).await.unwrap(); + assert!( + receipt.target_receipts()[0] + .outcome() + .satisfies(SatisfactionClass::Accepted) + ); + server.await.unwrap(); +} + +#[tokio::test] +async fn rejection_lost_ack_and_disconnect_never_invent_acceptance() { + for mode in ["reject", "lost", "disconnect"] { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let config = config(&format!("ws://{}", listener.local_addr().unwrap()), 250); + let raw = raw_event(); + let expected = format!("[\"EVENT\",{raw}]"); + let server = tokio::spawn(async move { + let (tcp, _) = listener.accept().await.unwrap(); + let mut socket = accept_async(tcp).await.unwrap(); + let wire = text(&mut socket).await; + assert_eq!(wire, expected); + let event: Value = serde_json::from_str(&wire).unwrap(); + match mode { + "reject" => { + reply( + &mut socket, + json!(["OK", event[1]["id"], false, "blocked: fixture"]), + ) + .await + } + "lost" => { + reply( + &mut socket, + json!(["OK", "22".repeat(32), true, "wrong ID"]), + ) + .await; + tokio::time::sleep(Duration::from_millis(500)).await; + } + "disconnect" => socket.close(None).await.unwrap(), + _ => unreachable!(), + } + }); + let receipt = NostrTransport::new(config.clone()) + .deliver(request(&config, &raw)) + .await + .unwrap(); + assert!( + !receipt.target_receipts()[0] + .outcome() + .satisfies(SatisfactionClass::Accepted), + "{mode}" + ); + assert!(receipt.target_receipts()[0].was_attempted()); + server.await.unwrap(); + } +} + +#[tokio::test] +async fn queued_delivery_targets_cannot_start_after_the_shared_deadline() { + let first = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let second = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let endpoints = [&first, &second].map(|listener| { + RelayEndpoint::new( + format!("ws://{}", listener.local_addr().unwrap()), + RelayUrlPolicy::Local, + RelayAccess::ReadWrite, + ) + .unwrap() + }); + let config = Config::from_profile( + RelayProfile::explicit(RelayProfileKind::Simulator, endpoints).unwrap(), + ) + .with_timeouts(1_000, 250, 500) + .unwrap() + .with_max_connections(1) + .unwrap(); + let server = tokio::spawn(async move { + let (tcp, _) = first.accept().await.unwrap(); + let mut socket = accept_async(tcp).await.unwrap(); + let _published = text(&mut socket).await; + futures::future::pending::<()>().await; + }); + let transport = NostrTransport::new(config.clone()); + let receipt = transport + .deliver(request(&config, &raw_event())) + .await + .unwrap(); + assert_eq!(receipt.target_receipts().len(), 2); + assert!( + receipt + .target_receipts() + .iter() + .all(|target| !target.outcome().satisfies(SatisfactionClass::Accepted)) + ); + assert!(second.accept().now_or_never().is_none()); + server.abort(); + assert!(server.await.unwrap_err().is_cancelled()); +} diff --git a/crates/transport_nostr/tests/package_boundary.rs b/crates/transport_nostr/tests/package_boundary.rs @@ -43,9 +43,11 @@ fn manifest_and_root_match_the_governed_transport_boundary() { "client", "cursor", "error", + "exact_delivery", "profile", "relay", "sink", + "socket_write", "source", "status", "subscription" @@ -203,7 +205,7 @@ fn preparation_is_sealed_and_separated_from_execution_io() { "pub struct PreparedDelivery", "#[must_use = \"prepared delivery must be durably bound before execution or deliberately discarded\"]", "request: DeliveryRequest", - "event: Event", + "event: ExactEvent", "config: crate::Config", "formatter.write_str(\"PreparedDelivery([redacted])\")", "pub const fn request(&self) -> &DeliveryRequest", @@ -315,10 +317,13 @@ fn adapter_owns_no_storage_outbox_or_orchestration_surface() { "client.rs".to_owned(), "cursor.rs".to_owned(), "error.rs".to_owned(), + "exact_delivery.rs".to_owned(), "lib.rs".to_owned(), "profile.rs".to_owned(), "relay.rs".to_owned(), "sink.rs".to_owned(), + "socket_write.rs".to_owned(), + "socket_write_tests.rs".to_owned(), "source.rs".to_owned(), "source_budget.rs".to_owned(), "source_paging_tests.rs".to_owned(),