lib

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

commit 2398cc0605ad6fbb5729b9f1867cc53e3e8bbede
parent a37a06ce0dfe0c7d08e73ad324fcaddf952a173d
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 06:46:07 +0000

transport-nostr: normalize relay outcomes and status

- centralize upstream relay classification into stable delivery and fetch semantics
- redact upstream diagnostic text while retaining only safe class codes and messages
- track passive source and sink availability as available, degraded, or unavailable
- verify all outcome classes, redaction, status, package, Clippy, architecture, and boundary gates

Diffstat:
Mcrates/transport_nostr/src/client.rs | 4++++
Mcrates/transport_nostr/src/sink.rs | 86++++++++++++++++++++-----------------------------------------------------------
Mcrates/transport_nostr/src/source.rs | 51+++++++++++++++++++--------------------------------
Mcrates/transport_nostr/src/status.rs | 279++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
4 files changed, 322 insertions(+), 98 deletions(-)

diff --git a/crates/transport_nostr/src/client.rs b/crates/transport_nostr/src/client.rs @@ -128,6 +128,7 @@ pub struct NostrTransport { pub(crate) client: Arc<dyn crate::sink::RelayClient>, pub(crate) source_client: Arc<dyn crate::source::RelaySourceClient>, pub(crate) auth: Arc<crate::auth::AuthFlow>, + pub(crate) status: Arc<crate::status::StatusTracker>, } impl NostrTransport { @@ -142,6 +143,7 @@ impl NostrTransport { auth: Arc::new(crate::auth::AuthFlow::new(Arc::new( crate::auth::LiveAuthClient::new(client), ))), + status: Arc::new(crate::status::StatusTracker::default()), } } @@ -157,6 +159,7 @@ impl NostrTransport { client, source_client: Arc::new(crate::source::LiveRelaySourceClient::isolated()), auth: Arc::new(crate::auth::AuthFlow::isolated()), + status: Arc::new(crate::status::StatusTracker::default()), } } @@ -170,6 +173,7 @@ impl NostrTransport { client: Arc::new(crate::sink::LiveRelayClient::isolated()), source_client, auth: Arc::new(crate::auth::AuthFlow::isolated()), + status: Arc::new(crate::status::StatusTracker::default()), } } } diff --git a/crates/transport_nostr/src/sink.rs b/crates/transport_nostr/src/sink.rs @@ -1,12 +1,11 @@ //! Nostr implementation of the transport event sink. -use crate::{NostrTransport, RelayUrl}; +use crate::{NostrTransport, RelayUrl, status}; use core::time::Duration; use radroots_nostr::event::Event; use radroots_transport::{ BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, - capability::{Availability, Maturity, SinkCapabilities}, - outcome::{DeliveryOutcome, Retryability}, + outcome::DeliveryOutcome, sink::{DeliveryTargetReceipt, SinkStatus}, }; use std::collections::{BTreeMap, BTreeSet}; @@ -56,15 +55,15 @@ impl RelayClient for LiveRelayClient { for relay in relays { let url = relay.as_str().to_owned(); let outcome = match self.client.add_relay(url.as_str()).await { - Err(error) => connection_failure(error.to_string()), + Err(_) => status::connection_failure(), Ok(_) => match self .client .try_connect_relay(url.as_str(), connect_timeout) .await { - Err(error) => connection_failure(error.to_string()), + Err(_) => status::connection_failure(), Ok(()) => match self.client.send_event_to([url.as_str()], &event).await { - Err(error) => classify_failure(error.to_string()), + Err(error) => status::delivery_failure(error.to_string().as_str()), Ok(output) => output .success .iter() @@ -73,10 +72,12 @@ impl RelayClient for LiveRelayClient { .or_else(|| { output.failed.iter().find_map(|(failed, message)| { (failed.to_string().trim_end_matches('/') == url) - .then(|| classify_failure(message.clone())) + .then(|| status::delivery_failure(message.as_str())) }) }) - .unwrap_or_else(|| classify_failure("relay omitted result".into())), + .unwrap_or_else(|| { + status::delivery_failure("relay omitted result") + }), }, }, }; @@ -90,16 +91,7 @@ impl RelayClient for LiveRelayClient { impl EventSink for NostrTransport { fn status(&self) -> BoxFuture<'_, Result<SinkStatus, radroots_transport::Error>> { let configured = !self.config().relays().is_empty(); - Box::pin(async move { - Ok(SinkStatus::new( - radroots_transport::TransportId::NOSTR, - configured, - Maturity::Preview, - Availability::Available, - SinkCapabilities::DELIVER, - "bounded Nostr event delivery configured", - )) - }) + Box::pin(async move { Ok(status::sink_status(&self.status, configured)) }) } fn deliver( @@ -157,56 +149,20 @@ impl EventSink for NostrTransport { }); receipts.push(DeliveryTargetReceipt::attempted(target, outcome)); } + let accepted = receipts + .iter() + .filter(|receipt| status::delivery_succeeded(receipt.outcome())) + .count(); + let failed = receipts.len().saturating_sub(accepted); + let diagnostic = receipts + .iter() + .find_map(|receipt| receipt.outcome().message()); + self.status.record_sink(accepted, failed, diagnostic); DeliveryReceipt::for_request(&request, receipts) }) } } -fn connection_failure(message: String) -> DeliveryOutcome { - normalized(DeliveryOutcome::unavailable(), "connection_failed", message) -} - -fn classify_failure(message: String) -> DeliveryOutcome { - let normalized_message = message.to_ascii_lowercase(); - if normalized_message.contains("duplicate") || normalized_message.contains("already have") { - normalized(DeliveryOutcome::accepted(), "duplicate", message) - } else if normalized_message.contains("auth") { - normalized( - DeliveryOutcome::failed(Retryability::Retryable) - .expect("retryable failure classification"), - "auth_required", - message, - ) - } else if normalized_message.contains("blocked") - || normalized_message.contains("restricted") - || normalized_message.contains("invalid") - { - normalized(DeliveryOutcome::rejected(), "rejected", message) - } else if normalized_message.contains("rate") { - normalized(DeliveryOutcome::unavailable(), "rate_limited", message) - } else if normalized_message.contains("timeout") { - normalized(DeliveryOutcome::unavailable(), "timeout", message) - } else { - normalized(DeliveryOutcome::unavailable(), "relay_failure", message) - } -} - -fn normalized(outcome: DeliveryOutcome, code: &'static str, message: String) -> DeliveryOutcome { - let clean = message - .chars() - .filter(|character| !character.is_control()) - .take(1_024) - .collect::<String>(); - let clean = if clean.trim().is_empty() { - "relay operation failed".to_owned() - } else { - clean.trim().to_owned() - }; - outcome - .with_detail(code, clean) - .expect("normalized static relay outcome") -} - #[cfg(test)] mod tests { use super::*; @@ -276,7 +232,7 @@ mod tests { .expect("config"); let two = RelayUrl::parse("wss://two.example", RelayUrlPolicy::Public).expect("two"); let client = MockRelayClient { - outcomes: BTreeMap::from([(two, classify_failure("rate limited".to_owned()))]), + outcomes: BTreeMap::from([(two, status::delivery_failure("rate limited"))]), }; let transport = NostrTransport::with_client(config, Arc::new(client)); let request = request(); @@ -306,7 +262,7 @@ mod tests { ("connection failed", DeliveryOutcomeKind::Unavailable), ]; for (message, expected) in cases { - assert_eq!(classify_failure(message.to_owned()).kind(), expected); + assert_eq!(status::delivery_failure(message).kind(), expected); } } } diff --git a/crates/transport_nostr/src/source.rs b/crates/transport_nostr/src/source.rs @@ -1,12 +1,11 @@ //! Nostr implementation of the transport event source. -use crate::{NostrTransport, RelayUrl}; +use crate::{NostrTransport, RelayUrl, status}; use core::cmp::Ordering; use core::time::Duration; use nostr_sdk::prelude::{Filter, JsonUtil, Timestamp}; use radroots_transport::{ BoxFuture, EventSource, FetchPage, FetchRequest, - capability::{Availability, Maturity, SourceCapabilities}, outcome::{FetchTargetOutcome, FetchTargetState}, source::{EventProvenance, FetchCursor, NextPage, ObservedEvent, SourceStatus}, }; @@ -91,16 +90,7 @@ struct Candidate { impl EventSource for NostrTransport { fn status(&self) -> BoxFuture<'_, Result<SourceStatus, radroots_transport::Error>> { let configured = !self.config().relays().is_empty(); - Box::pin(async move { - Ok(SourceStatus::new( - radroots_transport::TransportId::NOSTR, - configured, - Maturity::Preview, - Availability::Available, - SourceCapabilities::FETCH, - "bounded Nostr event source configured", - )) - }) + Box::pin(async move { Ok(status::source_status(&self.status, configured)) }) } fn fetch( @@ -138,6 +128,7 @@ impl EventSource for NostrTransport { ) .with_message("fetch deadline elapsed before relay access") })); + self.status.record_source(0, targets.len(), Some("timeout")); return FetchPage::for_request(&request, Vec::new(), outcomes, NextPage::Complete); } @@ -152,6 +143,9 @@ impl EventSource for NostrTransport { let mut candidates = Vec::new(); let mut malformed_by_relay = BTreeMap::<RelayUrl, usize>::new(); let mut reported = BTreeSet::new(); + let mut succeeded = 0usize; + let mut failed = 0usize; + let mut diagnostic = None; for batch in batches { let Some(target) = targets.get(&batch.relay) else { return Err(radroots_transport::Error::UnexpectedFetchTargetOutcome); @@ -161,6 +155,7 @@ impl EventSource for NostrTransport { } match batch.result { Ok(raw_events) => { + succeeded += 1; for raw in raw_events { match radroots_event_codec::decode::signed_event(raw.as_str()) { Ok(event) => candidates.push(Candidate { @@ -193,17 +188,21 @@ impl EventSource for NostrTransport { }; outcomes.push(outcome); } - Err(message) => outcomes.push( - FetchTargetOutcome::new( - target.fingerprint().clone(), - FetchTargetState::FailedRetryable, - ) - .with_message(safe_message(message)), - ), + Err(message) => { + failed += 1; + let (state, safe) = status::fetch_failure(message.as_str()); + diagnostic.get_or_insert(safe); + outcomes.push( + FetchTargetOutcome::new(target.fingerprint().clone(), state) + .with_message(safe), + ); + } } } for (relay, target) in &targets { if !reported.contains(relay) { + failed += 1; + diagnostic.get_or_insert("relay returned no fetch result"); outcomes.push( FetchTargetOutcome::new( target.fingerprint().clone(), @@ -213,6 +212,7 @@ impl EventSource for NostrTransport { ); } } + self.status.record_source(succeeded, failed, diagnostic); candidates.sort_by(compare_candidate); if let Some(cursor) = &cursor { @@ -301,19 +301,6 @@ fn unix_time_ms() -> u64 { .unwrap_or_default() } -fn safe_message(message: String) -> String { - let message = message - .chars() - .filter(|character| !character.is_control()) - .take(1_024) - .collect::<String>(); - if message.trim().is_empty() { - "relay fetch failed".to_owned() - } else { - message.trim().to_owned() - } -} - #[cfg(test)] mod tests { use super::*; diff --git a/crates/transport_nostr/src/status.rs b/crates/transport_nostr/src/status.rs @@ -1 +1,278 @@ -//! Stable Nostr relay status normalization. +//! Stable Nostr relay status and outcome normalization. + +use radroots_transport::{ + SinkStatus, SourceStatus, + capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities}, + outcome::{DeliveryOutcome, DeliveryOutcomeKind, FetchTargetState, Retryability}, +}; +use std::fmt; +use std::sync::Mutex; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum FailureClass { + Duplicate, + Rejected, + AuthRequired, + RateLimited, + Timeout, + Connection, + Malformed, + Unknown, +} + +impl FailureClass { + const fn code(self) -> &'static str { + match self { + Self::Duplicate => "duplicate", + Self::Rejected => "rejected", + Self::AuthRequired => "auth_required", + Self::RateLimited => "rate_limited", + Self::Timeout => "timeout", + Self::Connection => "connection_failed", + Self::Malformed => "malformed_event", + Self::Unknown => "relay_failure", + } + } + + const fn message(self) -> &'static str { + match self { + Self::Duplicate => "relay already has the event", + Self::Rejected => "relay rejected the event", + Self::AuthRequired => "relay authentication is required", + Self::RateLimited => "relay rate limit was reached", + Self::Timeout => "relay operation timed out", + Self::Connection => "relay connection failed", + Self::Malformed => "relay returned a malformed event", + Self::Unknown => "relay operation failed", + } + } +} + +#[derive(Clone)] +struct RedactedDiagnostic { + class: FailureClass, +} + +impl fmt::Debug for RedactedDiagnostic { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("RedactedDiagnostic") + .field("class", &self.class) + .field("upstream", &"[redacted]") + .finish() + } +} + +#[derive(Clone, Debug)] +struct Snapshot { + source: Availability, + sink: Availability, + source_diagnostic: Option<RedactedDiagnostic>, + sink_diagnostic: Option<RedactedDiagnostic>, +} + +impl Default for Snapshot { + fn default() -> Self { + Self { + source: Availability::Available, + sink: Availability::Available, + source_diagnostic: None, + sink_diagnostic: None, + } + } +} + +#[derive(Debug, Default)] +pub(crate) struct StatusTracker { + snapshot: Mutex<Snapshot>, +} + +impl StatusTracker { + pub(crate) fn record_sink(&self, accepted: usize, failed: usize, diagnostic: Option<&str>) { + if let Ok(mut snapshot) = self.snapshot.lock() { + snapshot.sink = availability(accepted, failed); + snapshot.sink_diagnostic = diagnostic.map(classify); + } + } + + pub(crate) fn record_source(&self, succeeded: usize, failed: usize, diagnostic: Option<&str>) { + if let Ok(mut snapshot) = self.snapshot.lock() { + snapshot.source = availability(succeeded, failed); + snapshot.source_diagnostic = diagnostic.map(classify); + } + } + + fn source_availability(&self) -> Availability { + self.snapshot + .lock() + .map(|snapshot| snapshot.source) + .unwrap_or(Availability::Unavailable) + } + + fn sink_availability(&self) -> Availability { + self.snapshot + .lock() + .map(|snapshot| snapshot.sink) + .unwrap_or(Availability::Unavailable) + } +} + +pub(crate) fn source_status(tracker: &StatusTracker, configured: bool) -> SourceStatus { + SourceStatus::new( + radroots_transport::TransportId::NOSTR, + configured, + Maturity::Preview, + if configured { + tracker.source_availability() + } else { + Availability::Unavailable + }, + SourceCapabilities::FETCH, + if configured { + "bounded Nostr event source configured" + } else { + "Nostr event source is not configured" + }, + ) +} + +pub(crate) fn sink_status(tracker: &StatusTracker, configured: bool) -> SinkStatus { + SinkStatus::new( + radroots_transport::TransportId::NOSTR, + configured, + Maturity::Preview, + if configured { + tracker.sink_availability() + } else { + Availability::Unavailable + }, + SinkCapabilities::DELIVER, + if configured { + "bounded Nostr event delivery configured" + } else { + "Nostr event sink is not configured" + }, + ) +} + +pub(crate) fn delivery_failure(upstream: &str) -> DeliveryOutcome { + let class = classify(upstream).class; + let outcome = match class { + FailureClass::Duplicate => DeliveryOutcome::accepted(), + FailureClass::Rejected | FailureClass::Malformed => DeliveryOutcome::rejected(), + FailureClass::AuthRequired => DeliveryOutcome::failed(Retryability::Retryable) + .expect("retryable authentication outcome"), + FailureClass::RateLimited + | FailureClass::Timeout + | FailureClass::Connection + | FailureClass::Unknown => DeliveryOutcome::unavailable(), + }; + outcome + .with_detail(class.code(), class.message()) + .expect("static normalized relay outcome") +} + +pub(crate) fn connection_failure() -> DeliveryOutcome { + DeliveryOutcome::unavailable() + .with_detail( + FailureClass::Connection.code(), + FailureClass::Connection.message(), + ) + .expect("static normalized connection outcome") +} + +pub(crate) fn fetch_failure(upstream: &str) -> (FetchTargetState, &'static str) { + let class = classify(upstream).class; + let state = match class { + FailureClass::Rejected | FailureClass::Malformed => FetchTargetState::FailedTerminal, + _ => FetchTargetState::FailedRetryable, + }; + (state, class.message()) +} + +fn classify(upstream: &str) -> RedactedDiagnostic { + let message = upstream.to_ascii_lowercase(); + let class = if message.contains("duplicate") || message.contains("already have") { + FailureClass::Duplicate + } else if message.contains("auth") { + FailureClass::AuthRequired + } else if message.contains("blocked") + || message.contains("restricted") + || message.contains("invalid") + || message.contains("reject") + { + FailureClass::Rejected + } else if message.contains("rate") { + FailureClass::RateLimited + } else if message.contains("timeout") || message.contains("timed out") { + FailureClass::Timeout + } else if message.contains("connect") || message.contains("offline") { + FailureClass::Connection + } else if message.contains("malformed") || message.contains("decode") { + FailureClass::Malformed + } else { + FailureClass::Unknown + }; + RedactedDiagnostic { class } +} + +fn availability(succeeded: usize, failed: usize) -> Availability { + match (succeeded, failed) { + (0, 0) | (_, 0) => Availability::Available, + (0, _) => Availability::Unavailable, + (_, _) => Availability::Degraded, + } +} + +pub(crate) fn delivery_succeeded(outcome: &DeliveryOutcome) -> bool { + matches!( + outcome.kind(), + DeliveryOutcomeKind::Accepted | DeliveryOutcomeKind::Delivered + ) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn every_upstream_class_maps_to_stable_secret_safe_output() { + let secret = "token=very-secret-value"; + let cases = [ + ("duplicate: already have", "duplicate"), + ("blocked by policy", "rejected"), + ("auth required", "auth_required"), + ("rate limited", "rate_limited"), + ("connection timeout", "timeout"), + ("connection offline", "connection_failed"), + ("unknown failure", "relay_failure"), + ]; + for (message, code) in cases { + let outcome = delivery_failure(format!("{message} {secret}").as_str()); + assert_eq!(outcome.code(), Some(code)); + assert!(!outcome.message().expect("message").contains(secret)); + let diagnostic = classify(format!("{message} {secret}").as_str()); + assert!(!format!("{diagnostic:?}").contains(secret)); + } + } + + #[test] + fn status_tracks_available_degraded_and_unavailable_without_io() { + let tracker = StatusTracker::default(); + assert_eq!( + sink_status(&tracker, true).availability(), + Availability::Available + ); + tracker.record_sink(1, 1, Some("token=secret timeout")); + assert_eq!( + sink_status(&tracker, true).availability(), + Availability::Degraded + ); + tracker.record_source(0, 2, Some("token=secret offline")); + assert_eq!( + source_status(&tracker, true).availability(), + Availability::Unavailable + ); + assert!(!format!("{tracker:?}").contains("token=secret")); + } +}