commit 20ec51df367ca50658f94740d06b27a22ad8ae69
parent 5342e08dab1b595275b8641a91a5d6552891e366
Author: triesap <tyson@radroots.org>
Date: Fri, 7 Aug 2026 17:33:35 +0000
transport: add evidence-based relay profiles
- validate public, simulator, and device relay authority
- bind fetch cursors and retries to exact relay evidence
- expose directional status through SDK and native mobile bindings
- prove loopback I/O, offline outbox behavior, and release gates
Diffstat:
34 files changed, 2393 insertions(+), 445 deletions(-)
diff --git a/Cargo.lock b/Cargo.lock
@@ -3579,6 +3579,7 @@ dependencies = [
"radroots_storage",
"radroots_sync",
"radroots_transport",
+ "radroots_transport_nostr",
"serde",
"serde_json",
"sha2",
@@ -4155,6 +4156,7 @@ dependencies = [
"radroots_protocol",
"radroots_transport",
"serde_json",
+ "sha2",
"tokio",
"tokio-tungstenite",
"url",
diff --git a/crates/mobile_core/Cargo.toml b/crates/mobile_core/Cargo.toml
@@ -40,6 +40,7 @@ radroots_signing = { workspace = true, default-features = false, features = ["st
radroots_storage = { workspace = true, default-features = false }
radroots_sync = { workspace = true, default-features = false }
radroots_transport = { workspace = true, default-features = false, features = ["std"] }
+radroots_transport_nostr = { workspace = true }
chrono = { workspace = true }
hex = { workspace = true }
serde = { workspace = true, features = ["derive"] }
diff --git a/crates/mobile_core/src/runtime/builder.rs b/crates/mobile_core/src/runtime/builder.rs
@@ -6,15 +6,20 @@ pub struct RuntimeBuilder {
store: MobileUserStoreConfig,
#[cfg(feature = "mobile-social")]
signer: Option<std::sync::Arc<dyn radroots_signing::Signer>>,
+ #[cfg(feature = "mobile-social")]
+ relay_profile: radroots_sdk::transport::RelayProfile,
}
impl RuntimeBuilder {
#[must_use]
- pub const fn new(store: MobileUserStoreConfig) -> Self {
+ pub fn new(store: MobileUserStoreConfig) -> Self {
Self {
store,
#[cfg(feature = "mobile-social")]
signer: None,
+ #[cfg(feature = "mobile-social")]
+ relay_profile: radroots_sdk::transport::RelayProfile::public(Vec::<String>::new())
+ .expect("bundled public relay profile is valid"),
}
}
@@ -26,6 +31,15 @@ impl RuntimeBuilder {
self
}
+ /// Replaces the bundled read-only public profile with one validated host
+ /// environment profile. Construction remains inert.
+ #[cfg(feature = "mobile-social")]
+ #[must_use]
+ pub fn relay_profile(mut self, relay_profile: radroots_sdk::transport::RelayProfile) -> Self {
+ self.relay_profile = relay_profile;
+ self
+ }
+
/// Opens the exact authenticated user's durable SQLite store.
pub async fn build(self) -> Result<RadrootsRuntime, RadrootsAppError> {
if self.store.protected_data() == ProtectedDataAvailability::Unavailable {
@@ -41,6 +55,8 @@ impl RuntimeBuilder {
Some(self.store.public_key()),
#[cfg(feature = "mobile-social")]
self.signer,
+ #[cfg(feature = "mobile-social")]
+ Some(self.relay_profile),
)
}
}
@@ -81,6 +97,55 @@ mod tests {
runtime.sdk_storage_status().await.expect("status").backend,
"sqlite"
);
+ #[cfg(feature = "mobile-social")]
+ {
+ let report = runtime
+ .sdk_relay_status()
+ .expect("relay status")
+ .expect("configured profile");
+ assert_eq!(report.profile, "public");
+ assert_eq!(report.state, "configured");
+ assert_eq!(report.relays.len(), 1);
+ assert_eq!(report.relays[0].relay_url, "wss://radroots.org");
+ assert_eq!(report.relays[0].access, "read_only");
+ assert_eq!(report.relays[0].read_state, "unobserved");
+ assert_eq!(report.relays[0].write_state, "unsupported");
+ }
+ runtime.shutdown().await.expect("shutdown");
+ }
+
+ #[cfg(feature = "mobile-social")]
+ #[tokio::test]
+ async fn runtime_reconfiguration_preserves_profile_network_boundaries() {
+ let root = tempfile::tempdir().expect("tempdir");
+ let runtime = RuntimeBuilder::new(store(root.path(), ProtectedDataAvailability::Available))
+ .build()
+ .await
+ .expect("runtime");
+ assert!(
+ runtime
+ .configure_simulator_relays(vec!["ws://127.0.0.1:8080".to_owned()])
+ .is_ok()
+ );
+ let report = runtime
+ .sdk_relay_status()
+ .expect("simulator status")
+ .expect("configured profile");
+ assert_eq!(report.profile, "simulator_local");
+ assert_eq!(report.relays.len(), 1);
+ assert_eq!(report.relays[0].access, "read_write");
+ assert!(
+ runtime
+ .configure_public_relays(vec!["ws://127.0.0.1:8080".to_owned()])
+ .is_err()
+ );
+ assert_eq!(
+ runtime
+ .sdk_relay_status()
+ .expect("unchanged status")
+ .expect("configured profile"),
+ report
+ );
runtime.shutdown().await.expect("shutdown");
}
diff --git a/crates/mobile_core/src/runtime/mod.rs b/crates/mobile_core/src/runtime/mod.rs
@@ -34,15 +34,20 @@ impl RadrootsRuntime {
#[cfg(feature = "mobile-social")] signer: Option<
std::sync::Arc<dyn radroots_signing::Signer>,
>,
+ #[cfg(feature = "mobile-social")] relay_profile: Option<
+ radroots_sdk::transport::RelayProfile,
+ >,
) -> Result<Self, RadrootsAppError> {
#[cfg(feature = "mobile-social")]
- let nostr_slot = radroots_sdk::transport::NostrSlot::new(
- radroots_sdk::transport::RelayUrlPolicy::Public,
- );
- #[cfg(feature = "mobile-social")]
let builder = {
+ let nostr_slot = radroots_sdk::transport::NostrSlot::new();
+ if let Some(profile) = relay_profile {
+ nostr_slot
+ .configure(profile)
+ .map_err(RadrootsAppError::from_sdk)?;
+ }
let builder = builder
- .nostr(nostr_slot.clone())
+ .nostr(nostr_slot)
.host_sync(radroots_sdk::sync::HostPolicy::standard());
match signer {
Some(signer) => builder.signing(radroots_sdk::signing::Provider::host(signer)),
@@ -67,6 +72,8 @@ impl RadrootsRuntime {
None,
#[cfg(feature = "mobile-social")]
None,
+ #[cfg(feature = "mobile-social")]
+ None,
)
}
diff --git a/crates/mobile_core/src/runtime/product_surface/context.rs b/crates/mobile_core/src/runtime/product_surface/context.rs
@@ -1,6 +1,6 @@
use std::collections::BTreeSet;
-use radroots_event::id::RelayUrl;
+use radroots_transport_nostr::{RelayUrl, RelayUrlPolicy};
use serde::{Deserialize, Serialize};
use thiserror::Error;
@@ -53,17 +53,17 @@ impl LocalNetwork {
return Err(LocalNetworkError::MissingRelay);
}
let mut relays = BTreeSet::new();
- for relay in &relay_urls {
- if relay.is_empty()
- || relay.len() > RELAY_URL_MAX_BYTES
- || !relay.starts_with("wss://")
- || RelayUrl::parse(relay).is_err()
- {
+ let mut canonical_relay_urls = Vec::with_capacity(relay_urls.len());
+ for relay in relay_urls {
+ if relay.is_empty() || relay.len() > RELAY_URL_MAX_BYTES {
return Err(LocalNetworkError::InvalidRelay);
}
- if !relays.insert(relay) {
+ let relay = RelayUrl::parse(relay, RelayUrlPolicy::Public)
+ .map_err(|_| LocalNetworkError::InvalidRelay)?;
+ if !relays.insert(relay.clone()) {
return Err(LocalNetworkError::DuplicateRelay);
}
+ canonical_relay_urls.push(relay.to_string());
}
let mut authors = BTreeSet::new();
for author in &followed_authors {
@@ -81,7 +81,7 @@ impl LocalNetwork {
Ok(Self {
id,
label,
- relay_urls,
+ relay_urls: canonical_relay_urls,
locality,
followed_authors,
generation,
@@ -334,5 +334,40 @@ mod tests {
] {
assert!(invalid.is_err());
}
+ let canonical = LocalNetwork::new(
+ "id".into(),
+ "label".into(),
+ vec!["WSS://RELAY.EXAMPLE:443/".into()],
+ None,
+ vec![],
+ 0,
+ )
+ .expect("canonical relay");
+ assert_eq!(canonical.relay_urls, vec!["wss://relay.example"]);
+ assert!(matches!(
+ LocalNetwork::new(
+ "id".into(),
+ "label".into(),
+ vec![
+ "wss://relay.example".into(),
+ "WSS://RELAY.EXAMPLE:443/".into(),
+ ],
+ None,
+ vec![],
+ 0,
+ ),
+ Err(LocalNetworkError::DuplicateRelay)
+ ));
+ assert!(matches!(
+ LocalNetwork::new(
+ "id".into(),
+ "label".into(),
+ vec!["wss://127.0.0.1:7447".into()],
+ None,
+ vec![],
+ 0,
+ ),
+ Err(LocalNetworkError::InvalidRelay)
+ ));
}
}
diff --git a/crates/mobile_core/src/runtime/product_surface/outbox.rs b/crates/mobile_core/src/runtime/product_surface/outbox.rs
@@ -1179,6 +1179,7 @@ mod tests {
ClientBuilder::memory_default(),
Some(PublicKey::from_hex(AUTHOR).unwrap()),
None,
+ None,
)
.unwrap()
}
@@ -1192,6 +1193,7 @@ mod tests {
ClientBuilder::memory_default(),
Some(PublicKey::from_hex(AUTHOR).unwrap()),
Some(std::sync::Arc::new(signer)),
+ None,
)
.unwrap()
}
diff --git a/crates/mobile_core/src/runtime/sdk.rs b/crates/mobile_core/src/runtime/sdk.rs
@@ -21,6 +21,27 @@ pub struct SdkStorageStatusRecord {
}
#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct SdkRelayStatusRecord {
+ pub relay_url: String,
+ pub access: String,
+ pub read_state: String,
+ pub write_state: String,
+ pub read_last_attempt_unix_ms: Option<u64>,
+ pub write_last_attempt_unix_ms: Option<u64>,
+ pub read_next_attempt_unix_ms: Option<u64>,
+ pub write_next_attempt_unix_ms: Option<u64>,
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct SdkRelayStatusReportRecord {
+ pub profile: String,
+ pub state: String,
+ pub read_availability: String,
+ pub write_availability: String,
+ pub relays: Vec<SdkRelayStatusRecord>,
+}
+
+#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SdkShutdownRecord {
pub state: String,
pub already_closed: bool,
@@ -54,6 +75,136 @@ impl RadrootsRuntime {
integrity: status.integrity().health().as_str().to_owned(),
})
}
+
+ /// Installs a validated public relay profile without probing it.
+ #[cfg(feature = "mobile-social")]
+ pub fn configure_public_relays(
+ &self,
+ writable_relays: Vec<String>,
+ ) -> Result<(), RadrootsAppError> {
+ self.configure_relay_profile(
+ radroots_sdk::transport::RelayProfile::public(writable_relays)
+ .map_err(|error| RadrootsAppError::runtime(error.to_string()))?,
+ )
+ }
+
+ /// Installs an exact-loopback simulator profile without probing it.
+ #[cfg(feature = "mobile-social")]
+ pub fn configure_simulator_relays(
+ &self,
+ loopback_relays: Vec<String>,
+ ) -> Result<(), RadrootsAppError> {
+ self.configure_relay_profile(
+ radroots_sdk::transport::RelayProfile::simulator(loopback_relays)
+ .map_err(|error| RadrootsAppError::runtime(error.to_string()))?,
+ )
+ }
+
+ /// Installs an explicit physical-device TLS relay profile without probing it.
+ #[cfg(feature = "mobile-social")]
+ pub fn configure_device_relays(
+ &self,
+ writable_relays: Vec<String>,
+ ) -> Result<(), RadrootsAppError> {
+ self.configure_relay_profile(
+ radroots_sdk::transport::RelayProfile::device(writable_relays)
+ .map_err(|error| RadrootsAppError::runtime(error.to_string()))?,
+ )
+ }
+
+ #[cfg(feature = "mobile-social")]
+ fn configure_relay_profile(
+ &self,
+ profile: radroots_sdk::transport::RelayProfile,
+ ) -> Result<(), RadrootsAppError> {
+ self.client
+ .configure_nostr(profile)
+ .map_err(RadrootsAppError::from_sdk)
+ }
+
+ /// Returns passive relay evidence without DNS, socket, or probe work.
+ #[cfg(feature = "mobile-social")]
+ pub fn sdk_relay_status(&self) -> Result<Option<SdkRelayStatusReportRecord>, RadrootsAppError> {
+ let report = self
+ .client
+ .nostr_status()
+ .map_err(RadrootsAppError::from_sdk)?;
+ Ok(report.map(|report| SdkRelayStatusReportRecord {
+ profile: relay_profile_label(report.profile_kind()).to_owned(),
+ state: relay_aggregate_label(report.state()).to_owned(),
+ read_availability: transport_availability_label(report.read_availability()).to_owned(),
+ write_availability: transport_availability_label(report.write_availability())
+ .to_owned(),
+ relays: report
+ .relays()
+ .iter()
+ .map(|relay| SdkRelayStatusRecord {
+ relay_url: relay.endpoint().url().to_string(),
+ access: if relay.endpoint().access().can_write() {
+ "read_write"
+ } else {
+ "read_only"
+ }
+ .to_owned(),
+ read_state: relay_evidence_label(relay.read().state()).to_owned(),
+ write_state: relay_evidence_label(relay.write().state()).to_owned(),
+ read_last_attempt_unix_ms: relay.read().last_attempt_unix_ms(),
+ write_last_attempt_unix_ms: relay.write().last_attempt_unix_ms(),
+ read_next_attempt_unix_ms: relay.read().next_attempt_unix_ms(),
+ write_next_attempt_unix_ms: relay.write().next_attempt_unix_ms(),
+ })
+ .collect(),
+ }))
+ }
+}
+
+#[cfg(feature = "mobile-social")]
+const fn relay_evidence_label(value: radroots_sdk::transport::RelayEvidenceState) -> &'static str {
+ match value {
+ radroots_sdk::transport::RelayEvidenceState::Unsupported => "unsupported",
+ radroots_sdk::transport::RelayEvidenceState::Unobserved => "unobserved",
+ radroots_sdk::transport::RelayEvidenceState::Connecting => "connecting",
+ radroots_sdk::transport::RelayEvidenceState::Available => "available",
+ radroots_sdk::transport::RelayEvidenceState::Unavailable => "unavailable",
+ _ => "unknown",
+ }
+}
+
+#[cfg(feature = "mobile-social")]
+const fn relay_profile_label(value: radroots_sdk::transport::RelayProfileKind) -> &'static str {
+ match value {
+ radroots_sdk::transport::RelayProfileKind::Public => "public",
+ radroots_sdk::transport::RelayProfileKind::Simulator => "simulator_local",
+ radroots_sdk::transport::RelayProfileKind::Device => "device_development",
+ _ => "unknown",
+ }
+}
+
+#[cfg(feature = "mobile-social")]
+const fn relay_aggregate_label(
+ value: radroots_sdk::transport::RelayAggregateState,
+) -> &'static str {
+ match value {
+ radroots_sdk::transport::RelayAggregateState::Configured => "configured",
+ radroots_sdk::transport::RelayAggregateState::Connecting => "connecting",
+ radroots_sdk::transport::RelayAggregateState::ReadOnly => "read_only",
+ radroots_sdk::transport::RelayAggregateState::Writable => "writable",
+ radroots_sdk::transport::RelayAggregateState::Degraded => "degraded",
+ radroots_sdk::transport::RelayAggregateState::Offline => "offline",
+ radroots_sdk::transport::RelayAggregateState::Failed => "failed",
+ _ => "unknown",
+ }
+}
+
+#[cfg(feature = "mobile-social")]
+const fn transport_availability_label(
+ value: radroots_transport::capability::Availability,
+) -> &'static str {
+ match value {
+ radroots_transport::capability::Availability::Available => "available",
+ radroots_transport::capability::Availability::Degraded => "degraded",
+ radroots_transport::capability::Availability::Unavailable => "unavailable",
+ }
}
const fn availability_label(value: Availability) -> &'static str {
diff --git a/crates/mobile_ffi/src/remote.rs b/crates/mobile_ffi/src/remote.rs
@@ -82,6 +82,27 @@ pub struct SdkStorageStatusRecord {
}
#[uniffi::remote(Record)]
+pub struct SdkRelayStatusRecord {
+ pub relay_url: String,
+ pub access: String,
+ pub read_state: String,
+ pub write_state: String,
+ pub read_last_attempt_unix_ms: Option<u64>,
+ pub write_last_attempt_unix_ms: Option<u64>,
+ pub read_next_attempt_unix_ms: Option<u64>,
+ pub write_next_attempt_unix_ms: Option<u64>,
+}
+
+#[uniffi::remote(Record)]
+pub struct SdkRelayStatusReportRecord {
+ pub profile: String,
+ pub state: String,
+ pub read_availability: String,
+ pub write_availability: String,
+ pub relays: Vec<SdkRelayStatusRecord>,
+}
+
+#[uniffi::remote(Record)]
pub struct SdkShutdownRecord {
pub state: String,
pub already_closed: bool,
diff --git a/crates/mobile_ffi/src/runtime.rs b/crates/mobile_ffi/src/runtime.rs
@@ -1,7 +1,9 @@
use radroots_mobile_core::runtime::{
info::RuntimeInfo,
product_surface::{AddCommandType, CardAddParity, LocalNetwork, TodayCardType},
- sdk::{SdkCapabilityRecord, SdkShutdownRecord, SdkStorageStatusRecord},
+ sdk::{
+ SdkCapabilityRecord, SdkRelayStatusReportRecord, SdkShutdownRecord, SdkStorageStatusRecord,
+ },
};
use crate::RadrootsAppError;
@@ -89,6 +91,37 @@ impl RadrootsRuntime {
self.inner.sdk_storage_status().await.map_err(Into::into)
}
+ pub fn sdk_relay_status(&self) -> Result<Option<SdkRelayStatusReportRecord>, RadrootsAppError> {
+ self.inner.sdk_relay_status().map_err(Into::into)
+ }
+
+ pub fn configure_public_relays(
+ &self,
+ writable_relays: Vec<String>,
+ ) -> Result<(), RadrootsAppError> {
+ self.inner
+ .configure_public_relays(writable_relays)
+ .map_err(Into::into)
+ }
+
+ pub fn configure_simulator_relays(
+ &self,
+ loopback_relays: Vec<String>,
+ ) -> Result<(), RadrootsAppError> {
+ self.inner
+ .configure_simulator_relays(loopback_relays)
+ .map_err(Into::into)
+ }
+
+ pub fn configure_device_relays(
+ &self,
+ writable_relays: Vec<String>,
+ ) -> Result<(), RadrootsAppError> {
+ self.inner
+ .configure_device_relays(writable_relays)
+ .map_err(Into::into)
+ }
+
pub fn phase1_card_types(&self) -> Vec<TodayCardType> {
self.inner.phase1_card_types()
}
diff --git a/crates/mobile_ffi/tests/runtime_delegation.rs b/crates/mobile_ffi/tests/runtime_delegation.rs
@@ -25,6 +25,65 @@ async fn native_boundary_delegates_the_complete_core_surface() {
"sqlite"
);
+ let public = runtime
+ .sdk_relay_status()
+ .expect("relay status")
+ .expect("default public profile");
+ assert_eq!(public.profile, "public");
+ assert_eq!(public.state, "configured");
+ assert_eq!(public.read_availability, "unavailable");
+ assert_eq!(public.write_availability, "unavailable");
+ assert_eq!(public.relays.len(), 1);
+ assert_eq!(public.relays[0].access, "read_only");
+ assert_eq!(public.relays[0].read_state, "unobserved");
+ assert_eq!(public.relays[0].write_state, "unsupported");
+
+ runtime
+ .configure_public_relays(vec!["wss://write.example".to_owned()])
+ .expect("public relays");
+ let public = runtime
+ .sdk_relay_status()
+ .expect("relay status")
+ .expect("public profile");
+ assert_eq!(public.relays.len(), 2);
+ assert_eq!(public.relays[1].access, "read_write");
+ assert!(
+ runtime
+ .configure_public_relays(vec!["ws://127.0.0.1:7447".to_owned()])
+ .is_err()
+ );
+
+ runtime
+ .configure_simulator_relays(vec!["ws://127.0.0.1:7447".to_owned()])
+ .expect("simulator relays");
+ let simulator = runtime
+ .sdk_relay_status()
+ .expect("relay status")
+ .expect("simulator profile");
+ assert_eq!(simulator.profile, "simulator_local");
+ assert_eq!(simulator.relays.len(), 1);
+ assert_eq!(simulator.relays[0].access, "read_write");
+ assert!(
+ runtime
+ .configure_simulator_relays(vec!["wss://relay.example".to_owned()])
+ .is_err()
+ );
+
+ runtime
+ .configure_device_relays(vec!["wss://10.0.0.5:7447".to_owned()])
+ .expect("device relays");
+ let device = runtime
+ .sdk_relay_status()
+ .expect("relay status")
+ .expect("device profile");
+ assert_eq!(device.profile, "device_development");
+ assert_eq!(device.relays.len(), 2);
+ assert!(
+ runtime
+ .configure_device_relays(vec!["wss://127.0.0.1:7447".to_owned()])
+ .is_err()
+ );
+
assert_eq!(
runtime.phase1_card_types(),
vec![
@@ -83,4 +142,12 @@ async fn native_boundary_delegates_the_complete_core_surface() {
runtime.sdk_storage_status().await,
Err(RadrootsAppError::Sdk { .. })
));
+ assert!(matches!(
+ runtime.sdk_relay_status(),
+ Err(RadrootsAppError::Sdk { .. })
+ ));
+ assert!(matches!(
+ runtime.configure_public_relays(Vec::new()),
+ Err(RadrootsAppError::Sdk { .. })
+ ));
}
diff --git a/crates/nostr/tests/package_boundary.rs b/crates/nostr/tests/package_boundary.rs
@@ -258,7 +258,7 @@ fn live_client_and_http_ownership_belongs_to_transport_nostr() {
"mod sink;",
"mod source;",
"mod status;",
- "pub use client::{Config, NostrTransport};",
+ "pub use client::{Config, NostrTransport, ReconnectBackoff};",
"pub use relay::{RelayUrl, RelayUrlPolicy};",
] {
assert!(
diff --git a/crates/sdk/src/capability.rs b/crates/sdk/src/capability.rs
@@ -209,6 +209,8 @@ pub(crate) struct Context<'a> {
pub(crate) source: bool,
pub(crate) sink: bool,
pub(crate) sync: bool,
+ pub(crate) source_availability: Availability,
+ pub(crate) sink_availability: Availability,
pub(crate) lifecycle_availability: Availability,
pub(crate) explicitly_configured: &'a BTreeSet<CapabilityId>,
pub(crate) overrides: &'a BTreeMap<CapabilityId, Availability>,
@@ -230,7 +232,7 @@ pub(crate) fn report(context: Context<'_>) -> CapabilityReport {
.overrides
.get(&definition.id)
.copied()
- .unwrap_or(context.lifecycle_availability)
+ .unwrap_or_else(|| directional_availability(definition.id, &context))
};
CapabilityStatus {
id: definition.id,
@@ -244,6 +246,18 @@ pub(crate) fn report(context: Context<'_>) -> CapabilityReport {
CapabilityReport { statuses }
}
+fn directional_availability(id: CapabilityId, context: &Context<'_>) -> Availability {
+ match id {
+ CapabilityId::NOSTR_FETCH | CapabilityId::SYNC_PULL => context.source_availability,
+ CapabilityId::NOSTR_DELIVERY
+ | CapabilityId::SYNC_PUSH
+ | CapabilityId::FARM_PUBLICATION
+ | CapabilityId::LISTING_PUBLICATION
+ | CapabilityId::TRADE_COMMANDS => context.sink_availability,
+ _ => context.lifecycle_availability,
+ }
+}
+
fn configured(id: CapabilityId, context: &Context<'_>) -> bool {
match id {
CapabilityId::CANONICAL_STORAGE
@@ -291,6 +305,8 @@ mod tests {
source: false,
sink: false,
sync: false,
+ source_availability: Availability::Unavailable,
+ sink_availability: Availability::Unavailable,
lifecycle_availability: Availability::Available,
explicitly_configured: &explicitly_configured,
overrides: &overrides,
@@ -339,6 +355,8 @@ mod tests {
source: true,
sink: true,
sync: true,
+ source_availability: Availability::Available,
+ sink_availability: Availability::Degraded,
lifecycle_availability: Availability::Available,
explicitly_configured: &explicitly_configured,
overrides: &overrides,
@@ -349,6 +367,28 @@ mod tests {
.expect("pull")
.is_configured()
);
+ assert_eq!(
+ complete
+ .get(CapabilityId::SYNC_PUSH)
+ .expect("push")
+ .availability(),
+ if cfg!(feature = "sync") {
+ Availability::Degraded
+ } else {
+ Availability::Unsupported
+ }
+ );
+ assert_eq!(
+ complete
+ .get(CapabilityId::SYNC_PULL)
+ .expect("pull")
+ .availability(),
+ if cfg!(feature = "sync") {
+ Availability::Available
+ } else {
+ Availability::Unsupported
+ }
+ );
assert!(
complete
.get(CapabilityId::SYNC_PUSH)
@@ -379,6 +419,8 @@ mod tests {
source: false,
sink: false,
sync: true,
+ source_availability: Availability::Unavailable,
+ sink_availability: Availability::Unavailable,
lifecycle_availability: Availability::Degraded,
explicitly_configured: &explicitly_configured,
overrides: &overrides,
diff --git a/crates/sdk/src/client.rs b/crates/sdk/src/client.rs
@@ -38,6 +38,8 @@ pub struct ClientBuilder {
sync: Option<radroots_sync::Engine>,
#[cfg(feature = "sync")]
host_sync: Option<crate::sync::HostPolicy>,
+ #[cfg(feature = "nostr")]
+ nostr: Option<crate::transport::NostrSlot>,
capability_availability: BTreeMap<CapabilityId, Availability>,
explicitly_configured_capabilities: BTreeSet<CapabilityId>,
}
@@ -49,6 +51,8 @@ struct ClientInner {
sink: Option<Arc<dyn EventSink>>,
#[cfg(feature = "sync")]
sync: Option<radroots_sync::Engine>,
+ #[cfg(feature = "nostr")]
+ nostr: Option<crate::transport::NostrSlot>,
capability_availability: BTreeMap<CapabilityId, Availability>,
explicitly_configured_capabilities: BTreeSet<CapabilityId>,
lifecycle: AtomicU8,
@@ -177,6 +181,7 @@ impl ClientBuilder {
pub fn nostr(mut self, slot: crate::transport::NostrSlot) -> Self {
self.source = Some(Arc::new(slot.clone()));
self.sink = Some(Arc::new(slot.clone()));
+ self.nostr = Some(slot);
self
}
@@ -242,6 +247,8 @@ impl ClientBuilder {
sink: self.sink,
#[cfg(feature = "sync")]
sync: self.sync,
+ #[cfg(feature = "nostr")]
+ nostr: self.nostr,
capability_availability: self.capability_availability,
explicitly_configured_capabilities: self.explicitly_configured_capabilities,
lifecycle: AtomicU8::new(OPEN),
@@ -254,21 +261,56 @@ impl Client {
/// Returns a deterministic capability report without probing resources or
/// performing filesystem, network, signing, or storage operations.
#[must_use]
+ #[allow(unused_mut)]
pub fn capabilities(&self) -> CapabilityReport {
let lifecycle_availability = match self.inner.lifecycle.load(Ordering::Acquire) {
OPEN => Availability::Available,
CLOSING | CLOSE_RETRY_REQUIRED => Availability::Degraded,
_ => Availability::Unavailable,
};
+ let mut explicitly_configured = self.inner.explicitly_configured_capabilities.clone();
+ let mut overrides = self.inner.capability_availability.clone();
+ let mut source_configured = self.inner.source.is_some();
+ let mut sink_configured = self.inner.sink.is_some();
+ let mut source_availability = Availability::Unavailable;
+ let mut sink_availability = Availability::Unavailable;
+ #[cfg(feature = "nostr")]
+ if let Some(slot) = &self.inner.nostr {
+ match slot.relay_status() {
+ Some(status) => {
+ source_configured = !status.relays().is_empty();
+ sink_configured = status
+ .relays()
+ .iter()
+ .any(|relay| relay.endpoint().access().can_write());
+ source_availability = map_transport_availability(status.read_availability());
+ sink_availability = map_transport_availability(status.write_availability());
+ if source_configured {
+ explicitly_configured.insert(CapabilityId::NOSTR_FETCH);
+ overrides.insert(CapabilityId::NOSTR_FETCH, source_availability);
+ }
+ if sink_configured {
+ explicitly_configured.insert(CapabilityId::NOSTR_DELIVERY);
+ overrides.insert(CapabilityId::NOSTR_DELIVERY, sink_availability);
+ }
+ }
+ None => {
+ source_configured = false;
+ sink_configured = false;
+ }
+ }
+ }
crate::capability::report(crate::capability::Context {
storage: true,
signer: self.inner.signer.is_some(),
- source: self.inner.source.is_some(),
- sink: self.inner.sink.is_some(),
+ source: source_configured,
+ sink: sink_configured,
sync: self.sync_is_configured(),
+ source_availability,
+ sink_availability,
lifecycle_availability,
- explicitly_configured: &self.inner.explicitly_configured_capabilities,
- overrides: &self.inner.capability_availability,
+ explicitly_configured: &explicitly_configured,
+ overrides: &overrides,
})
}
@@ -322,6 +364,29 @@ impl Client {
Ok(self.inner.sink.as_deref())
}
+ /// Returns passive Nostr relay evidence when the concrete adapter is selected.
+ #[cfg(feature = "nostr")]
+ pub fn nostr_status(&self) -> Result<Option<crate::transport::RelayStatusReport>> {
+ self.require_open()?;
+ Ok(self
+ .inner
+ .nostr
+ .as_ref()
+ .and_then(crate::transport::NostrSlot::relay_status))
+ }
+
+ /// Atomically replaces the active Nostr relay profile after validating it.
+ /// This does not perform DNS lookup or open a socket.
+ #[cfg(feature = "nostr")]
+ pub fn configure_nostr(&self, profile: crate::transport::RelayProfile) -> Result<()> {
+ self.require_open()?;
+ self.inner
+ .nostr
+ .as_ref()
+ .ok_or_else(Error::shared_operation_unavailable)?
+ .configure(profile)
+ }
+
/// Returns client-scoped canonical synchronization operations, when configured.
#[cfg(feature = "sync")]
pub fn sync(&self) -> Result<Option<crate::sync::Operations<'_>>> {
@@ -409,6 +474,17 @@ impl Client {
}
}
+#[cfg(feature = "nostr")]
+const fn map_transport_availability(
+ availability: radroots_transport::capability::Availability,
+) -> Availability {
+ match availability {
+ radroots_transport::capability::Availability::Available => Availability::Available,
+ radroots_transport::capability::Availability::Degraded => Availability::Degraded,
+ radroots_transport::capability::Availability::Unavailable => Availability::Unavailable,
+ }
+}
+
struct CloseAttempt {
inner: Arc<ClientInner>,
completed: bool,
@@ -574,6 +650,45 @@ mod tests {
assert!(sink.sink().expect("sink capability").is_some());
}
+ #[cfg(feature = "nostr")]
+ #[test]
+ fn nostr_capabilities_require_directional_profile_authority_and_evidence() {
+ let slot = crate::transport::NostrSlot::new();
+ slot.configure(
+ crate::transport::RelayProfile::public(Vec::<String>::new()).expect("public profile"),
+ )
+ .expect("configure slot");
+ let client = ClientBuilder::memory(generation())
+ .nostr(slot)
+ .build()
+ .expect("client");
+
+ let capabilities = client.capabilities();
+ let fetch = capabilities
+ .get(CapabilityId::NOSTR_FETCH)
+ .expect("fetch capability");
+ let delivery = capabilities
+ .get(CapabilityId::NOSTR_DELIVERY)
+ .expect("delivery capability");
+ assert!(fetch.is_configured());
+ assert_eq!(fetch.availability(), Availability::Unavailable);
+ assert!(!delivery.is_configured());
+ assert_eq!(delivery.availability(), Availability::Unavailable);
+
+ client
+ .configure_nostr(
+ crate::transport::RelayProfile::public(["wss://write.example"])
+ .expect("writable profile"),
+ )
+ .expect("reconfigure");
+ let capabilities = client.capabilities();
+ let delivery = capabilities
+ .get(CapabilityId::NOSTR_DELIVERY)
+ .expect("delivery capability");
+ assert!(delivery.is_configured());
+ assert_eq!(delivery.availability(), Availability::Unavailable);
+ }
+
#[test]
fn signer_with_sink_is_valid_and_diagnostics_are_capability_only() {
let client = ClientBuilder::memory(generation())
diff --git a/crates/sdk/src/transport.rs b/crates/sdk/src/transport.rs
@@ -13,7 +13,11 @@ use radroots_transport::{
use std::sync::{Arc, RwLock};
#[cfg(feature = "nostr")]
-pub use radroots_transport_nostr::RelayUrlPolicy;
+pub use radroots_transport_nostr::{
+ DEFAULT_PUBLIC_RELAY, ReconnectBackoff, RelayAccess, RelayAggregateState,
+ RelayCapabilityEvidence, RelayCursor, RelayEndpoint, RelayEvidenceState, RelayProfile,
+ RelayProfileKind, RelayStatus, RelayStatusReport, RelayUrl, RelayUrlPolicy,
+};
const PREVIEW_UNAVAILABLE_MESSAGE: &str = "preview transport is unavailable in this SDK release";
@@ -145,7 +149,6 @@ impl Default for Profile {
#[cfg(feature = "nostr")]
#[derive(Clone)]
pub struct NostrSlot {
- policy: RelayUrlPolicy,
state: Arc<RwLock<Option<NostrState>>>,
}
@@ -153,40 +156,50 @@ pub struct NostrSlot {
#[derive(Clone)]
struct NostrState {
transport: Arc<radroots_transport_nostr::NostrTransport>,
- targets: TargetSet,
+ read_targets: TargetSet,
+ write_targets: Option<TargetSet>,
}
#[cfg(feature = "nostr")]
impl NostrSlot {
- /// Creates an inert slot with an explicit destination policy.
+ /// Creates an inert slot with no selected host profile.
#[must_use]
- pub fn new(policy: RelayUrlPolicy) -> Self {
+ pub fn new() -> Self {
Self {
- policy,
state: Arc::new(RwLock::new(None)),
}
}
- /// Validates and atomically installs the complete relay selection.
- pub fn configure<I, S>(&self, relays: I) -> crate::Result<()>
- where
- I: IntoIterator<Item = S>,
- S: AsRef<str>,
- {
- let config = radroots_transport_nostr::Config::new(self.policy, relays)
- .map_err(crate::Error::invalid_host_configuration)?;
- let targets = TargetSet::new(
+ /// Atomically installs one completely validated relay profile.
+ pub fn configure(&self, profile: RelayProfile) -> crate::Result<()> {
+ let config = radroots_transport_nostr::Config::from_profile(profile);
+ let read_targets = TargetSet::new(
config
- .relays()
- .iter()
+ .read_relays()
.map(radroots_transport_nostr::RelayUrl::to_target)
.collect::<Result<Vec<_>, _>>()
.map_err(crate::Error::invalid_host_configuration)?,
)
.map_err(|_| crate::Error::invalid_host_configuration_without_source())?;
+ let write_targets = {
+ let targets = config
+ .write_relays()
+ .map(radroots_transport_nostr::RelayUrl::to_target)
+ .collect::<Result<Vec<_>, _>>()
+ .map_err(crate::Error::invalid_host_configuration)?;
+ if targets.is_empty() {
+ None
+ } else {
+ Some(
+ TargetSet::new(targets)
+ .map_err(|_| crate::Error::invalid_host_configuration_without_source())?,
+ )
+ }
+ };
let state = NostrState {
transport: Arc::new(radroots_transport_nostr::NostrTransport::new(config)),
- targets,
+ read_targets,
+ write_targets,
};
let mut current = self
.state
@@ -203,10 +216,23 @@ impl NostrSlot {
}
}
- /// Returns the currently selected canonical targets.
+ /// Returns the currently selected canonical read targets.
+ #[must_use]
+ pub fn read_targets(&self) -> Option<TargetSet> {
+ self.snapshot().map(|state| state.read_targets)
+ }
+
+ /// Returns writable targets, or `None` when the profile is intentionally
+ /// read-only.
#[must_use]
- pub fn targets(&self) -> Option<TargetSet> {
- self.snapshot().map(|state| state.targets)
+ pub fn write_targets(&self) -> Option<TargetSet> {
+ self.snapshot().and_then(|state| state.write_targets)
+ }
+
+ /// Returns passive per-relay evidence without probing or opening sockets.
+ #[must_use]
+ pub fn relay_status(&self) -> Option<RelayStatusReport> {
+ self.snapshot().map(|state| state.transport.relay_status())
}
fn snapshot(&self) -> Option<NostrState> {
@@ -303,12 +329,19 @@ impl std::fmt::Debug for NostrSlot {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("NostrSlot")
- .field("policy", &self.policy)
- .field("configured", &self.targets().is_some())
+ .field("configured", &self.read_targets().is_some())
+ .field("writable", &self.write_targets().is_some())
.finish()
}
}
+#[cfg(feature = "nostr")]
+impl Default for NostrSlot {
+ fn default() -> Self {
+ Self::new()
+ }
+}
+
/// Explicit daemon adapter authentication configuration.
#[cfg(feature = "radrootsd")]
#[derive(Clone, Eq, PartialEq)]
@@ -674,22 +707,30 @@ mod tests {
#[cfg(feature = "nostr")]
#[test]
- fn nostr_slot_reconfiguration_is_validated_atomic_and_inert() {
- let slot = NostrSlot::new(RelayUrlPolicy::Local);
- assert!(slot.targets().is_none());
- assert!(slot.configure(["ws://127.0.0.1:7447"]).is_ok());
- let original = slot.targets().expect("configured targets");
- assert!(slot.configure(Vec::<String>::new()).is_err());
- assert_eq!(slot.targets(), Some(original));
+ fn nostr_slot_reconfiguration_is_atomic_directional_and_inert() {
+ let slot = NostrSlot::new();
+ assert!(slot.read_targets().is_none());
+ assert!(slot.relay_status().is_none());
+ assert!(
+ slot.configure(RelayProfile::simulator(["ws://127.0.0.1:7447"]).expect("profile"))
+ .is_ok()
+ );
+ let original = slot.read_targets().expect("configured targets");
+ assert_eq!(slot.write_targets(), Some(original.clone()));
+ let status = slot.relay_status().expect("status");
+ assert_eq!(status.profile_kind(), RelayProfileKind::Simulator);
+ assert_eq!(status.read_availability(), Availability::Unavailable);
+ assert_eq!(status.write_availability(), Availability::Unavailable);
slot.clear();
- assert!(slot.targets().is_none());
+ assert!(slot.read_targets().is_none());
+ assert!(slot.write_targets().is_none());
assert!(format!("{slot:?}").contains("configured: false"));
}
#[cfg(feature = "nostr")]
#[test]
fn poisoned_nostr_slot_fails_closed_for_every_host_operation() {
- let slot = NostrSlot::new(RelayUrlPolicy::Local);
+ let slot = NostrSlot::new();
let state = Arc::clone(&slot.state);
let _ = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let _guard = state.write().expect("write lock");
@@ -697,8 +738,11 @@ mod tests {
}));
slot.clear();
- assert!(slot.targets().is_none());
- assert!(slot.configure(["ws://127.0.0.1:7447"]).is_err());
+ assert!(slot.read_targets().is_none());
+ assert!(
+ slot.configure(RelayProfile::simulator(["ws://127.0.0.1:7447"]).expect("profile"))
+ .is_err()
+ );
}
#[cfg(feature = "nostr")]
@@ -709,7 +753,7 @@ mod tests {
source::FetchBounds,
};
- let slot = NostrSlot::new(RelayUrlPolicy::Local);
+ let slot = NostrSlot::new();
let source = radroots_transport::EventSource::status(&slot)
.await
.expect("source status");
diff --git a/crates/studio_nostr/src/client.rs b/crates/studio_nostr/src/client.rs
@@ -9,7 +9,7 @@ use radroots_transport::{
outcome::FetchTargetState,
source::{FetchBounds, FetchSelector},
};
-use radroots_transport_nostr::{Config, NostrTransport, RelayUrlPolicy};
+use radroots_transport_nostr::{Config, NostrTransport, RelayProfile};
use radroots_studio_application::{
BoxFuture, MAX_CONFIGURED_RELAYS, NostrClient, ProfileFetchResult,
@@ -59,16 +59,14 @@ impl NostrClient for SdkNostrClient {
if policy_relays.is_empty() {
continue;
}
- let config = Config::new(
- canonical_policy(policy),
- policy_relays.iter().map(|relay| relay.as_str()),
- )
- .and_then(|config| {
- let timeout_ms =
- timeout_millis(deadline.saturating_duration_since(Instant::now()));
- config.with_timeouts(timeout_ms, timeout_ms, timeout_ms)
- })
- .map_err(|_| invalid_relay_configuration())?;
+ let profile =
+ relay_profile(policy, policy_relays.iter().map(|relay| relay.as_str()))
+ .map_err(|_| invalid_relay_configuration())?;
+ let config = Config::from_profile(profile);
+ let timeout_ms = timeout_millis(deadline.saturating_duration_since(Instant::now()));
+ let config = config
+ .with_timeouts(timeout_ms, timeout_ms, timeout_ms)
+ .map_err(|_| invalid_relay_configuration())?;
let targets = policy_relays
.iter()
.map(|relay| Target::nostr_relay(relay.as_str()))
@@ -146,11 +144,18 @@ fn timeout_millis(timeout: Duration) -> u64 {
.clamp(1, 120_000)
}
-const fn canonical_policy(policy: RelayDestinationPolicy) -> RelayUrlPolicy {
+fn relay_profile<I, S>(
+ policy: RelayDestinationPolicy,
+ relays: I,
+) -> Result<RelayProfile, radroots_transport_nostr::Error>
+where
+ I: IntoIterator<Item = S>,
+ S: AsRef<str>,
+{
match policy {
- RelayDestinationPolicy::Public => RelayUrlPolicy::Public,
- RelayDestinationPolicy::Local => RelayUrlPolicy::Local,
- RelayDestinationPolicy::PrivateNetwork => RelayUrlPolicy::PrivateNetwork,
+ RelayDestinationPolicy::Public => RelayProfile::public(relays),
+ RelayDestinationPolicy::Local => RelayProfile::simulator(relays),
+ RelayDestinationPolicy::PrivateNetwork => RelayProfile::device(relays),
}
}
@@ -308,18 +313,27 @@ mod tests {
}
#[test]
- fn studio_policy_maps_exactly_to_the_canonical_transport_policy() {
+ fn studio_policy_maps_exactly_to_the_canonical_transport_profile() {
+ let public = super::relay_profile(RelayDestinationPolicy::Public, ["wss://public.example"])
+ .expect("public profile");
+ let local = super::relay_profile(RelayDestinationPolicy::Local, ["ws://127.0.0.1:8080"])
+ .expect("local profile");
+ let device = super::relay_profile(
+ RelayDestinationPolicy::PrivateNetwork,
+ ["wss://10.0.0.5:7447"],
+ )
+ .expect("device profile");
assert_eq!(
- super::canonical_policy(RelayDestinationPolicy::Public),
- radroots_transport_nostr::RelayUrlPolicy::Public
+ public.kind(),
+ radroots_transport_nostr::RelayProfileKind::Public
);
assert_eq!(
- super::canonical_policy(RelayDestinationPolicy::Local),
- radroots_transport_nostr::RelayUrlPolicy::Local
+ local.kind(),
+ radroots_transport_nostr::RelayProfileKind::Simulator
);
assert_eq!(
- super::canonical_policy(RelayDestinationPolicy::PrivateNetwork),
- radroots_transport_nostr::RelayUrlPolicy::PrivateNetwork
+ device.kind(),
+ radroots_transport_nostr::RelayProfileKind::Device
);
assert_eq!(super::timeout_millis(Duration::ZERO), 1);
assert_eq!(
diff --git a/crates/transport/tests/source_boundary.rs b/crates/transport/tests/source_boundary.rs
@@ -698,10 +698,10 @@ fn transport_target_identity_sources_reject_silent_dedupe() {
);
let relay_source = read_source(crates_root.join("transport_nostr/src/relay.rs").as_path());
- let client_source = read_source(crates_root.join("transport_nostr/src/client.rs").as_path());
+ let profile_source = read_source(crates_root.join("transport_nostr/src/profile.rs").as_path());
for required in ["Target::nostr_relay(original)", "Error::DuplicateRelayUrl"] {
let source = if required.contains("Duplicate") {
- client_source.as_str()
+ profile_source.as_str()
} else {
relay_source.as_str()
};
diff --git a/crates/transport_nostr/Cargo.toml b/crates/transport_nostr/Cargo.toml
@@ -43,15 +43,17 @@ radroots_nostr = { workspace = true, default-features = false, features = [
radroots_protocol = { workspace = true, default-features = false }
radroots_transport = { workspace = true, default-features = false }
async-wsocket = { workspace = true }
+futures = { workspace = true }
nostr-sdk = { workspace = true }
nostr-relay-pool = { workspace = true }
serde_json = { workspace = true, features = ["std"] }
+sha2 = { workspace = true, default-features = false }
tokio = { workspace = true, features = ["net", "time"] }
tokio-tungstenite = { workspace = true }
url = { workspace = true }
[dev-dependencies]
-futures = { workspace = true }
+tokio = { workspace = true, features = ["macros", "net", "rt-multi-thread", "time"] }
[lints]
workspace = true
diff --git a/crates/transport_nostr/README.md b/crates/transport_nostr/README.md
@@ -8,9 +8,9 @@ explicit host-mediated NIP-42 authentication, and normalizes relay outcomes
and passive status.
The crate does not own event ingestion, persistence, outbox claiming,
-projection refresh, retry scheduling, SDK profiles, or a process runtime.
-Those policies belong to `radroots_sync` and host applications. Publication
-remains disabled during the `0.1.0-alpha` refactor.
+projection refresh, durable retry scheduling, or a process runtime. Those
+policies belong to `radroots_sync` and host applications. It does own bounded
+relay profiles, per-relay reconnect suppression, and evidence-based status.
The authoritative package charter is the
[`radroots_transport_nostr` section of the Release V1 specification](https://github.com/radrootslabs/lib/blob/master/docs/specs/radroots_crates_release_v1.md#15-radroots_transport_nostr).
@@ -24,12 +24,10 @@ Configuration is explicit, validated, and inert. Constructing
```rust
use radroots_transport::{EventSink, EventSource};
-use radroots_transport_nostr::{Config, NostrTransport, RelayUrlPolicy};
+use radroots_transport_nostr::{Config, NostrTransport, RelayProfile};
-let config = Config::new(
- RelayUrlPolicy::Public,
- ["wss://relay.example.com"],
-)?.with_timeouts(5_000, 20_000, 2_000)?;
+let profile = RelayProfile::public(["wss://relay.example.com"])?;
+let config = Config::from_profile(profile).with_timeouts(5_000, 20_000, 2_000)?;
let transport = NostrTransport::new(config);
let source: &dyn EventSource = &transport;
@@ -47,23 +45,32 @@ and applies any retry or scheduling policy outside this crate.
## Public surface
-- [`Config`] validates a non-empty, duplicate-free relay set plus bounded
- connection, request, status, and concurrency limits.
+- [`RelayProfile`] defines public, loopback-simulator, and physical-device
+ profiles with independent read-only/read-write authority per endpoint.
+- [`Config`] retains the validated profile plus bounded connection, request,
+ status, concurrency, and reconnect limits.
- [`RelayUrl`] is a canonical Nostr relay URL that converts to and from the
generic `radroots_transport::Target` model.
- [`RelayUrlPolicy`] selects public-Internet, exact-loopback, or explicitly
trusted private-network destination rules.
-- [`NostrTransport`] implements both transport SPIs and exposes explicit
- NIP-42 challenge lifecycle methods.
+- [`RelayCursor`] provides the equal-timestamp-safe event ordering primitive
+ used by scoped fetch continuation cursors.
+- [`NostrTransport`] implements both transport SPIs, exposes passive typed
+ per-relay evidence, and provides explicit NIP-42 challenge lifecycle methods.
- [`Error`] contains only package-owned validation and authentication errors;
upstream failures are normalized before crossing the public boundary.
All source modules are private implementation details. The crate root contains
-only the five reviewed exports above and does not expose an upstream client,
+only reviewed adapter-owned exports and does not expose an upstream client,
relay pool, Tokio handle, signer, storage handle, or retry worker.
## Relay and network security
+The public profile always includes `wss://radroots.org` as read-only and
+requires separately configured public TLS relays for publication. The
+simulator profile admits exact loopback WebSockets only. The device profile
+requires explicit TLS endpoints and never reuses simulator loopback.
+
`RelayUrlPolicy::Public` accepts TLS WebSocket URLs with public hostnames or
global addresses. `Local` accepts exact loopback destinations and permits
plaintext WebSocket only for that class. `PrivateNetwork` accepts explicit
@@ -81,22 +88,29 @@ check before handing control to another network boundary.
## Fetch, delivery, and outcome behavior
-Fetch accepts only configured Nostr targets, translates transport-neutral kind,
+Fetch accepts only configured readable Nostr targets, translates transport-neutral kind,
author, and event-time selectors into Nostr filters, reapplies those selectors
defensively, applies the request page bound, deduplicates events by event ID,
-preserves per-relay provenance, and emits an opaque versioned cursor when more
-results remain. Malformed relay events are ignored and reported as a partial
-target outcome rather than admitted.
+preserves per-relay provenance, and emits an opaque versioned cursor bound to
+the exact target set and selector when more results remain. Equal timestamps
+are ordered by event ID so overlap-safe reconnect pagination cannot skip peers.
+Malformed relay events are ignored and reported as a partial target outcome
+rather than admitted.
Delivery converts an already validated signed Radroots event to Nostr, attempts
-each configured target once, and returns one normalized receipt entry per
+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
explicit; this crate never retries, falls back to another transport, or
rewrites an unknown result as success.
-Source and sink status are passive in-memory observations. Reading status does
-not connect to a relay, refresh DNS, or begin fetch or delivery work.
+Source and sink status are passive in-memory observations. A configured relay
+starts unobserved and never appears available before successful read or write
+evidence. Read and write evidence, failure counters, retry classes, and
+next-attempt times are independent. Aggregate status distinguishes configured,
+connecting, read-only, writable, degraded, offline, and terminally failed
+states. Reading status does not connect to a relay, refresh DNS, or begin fetch
+or delivery work.
## Deadlines, cancellation, and commit points
diff --git a/crates/transport_nostr/examples/configure_transport.rs b/crates/transport_nostr/examples/configure_transport.rs
@@ -1,8 +1,8 @@
use radroots_transport::{EventSink, EventSource};
-use radroots_transport_nostr::{Config, NostrTransport, RelayUrlPolicy};
+use radroots_transport_nostr::{Config, NostrTransport, RelayProfile};
fn main() -> Result<(), Box<dyn std::error::Error>> {
- let config = Config::new(RelayUrlPolicy::Public, ["wss://relay.example.com"])?
+ let config = Config::from_profile(RelayProfile::public(["wss://relay.example.com"])?)
.with_timeouts(5_000, 20_000, 2_000)?;
let transport = NostrTransport::new(config);
diff --git a/crates/transport_nostr/src/auth.rs b/crates/transport_nostr/src/auth.rs
@@ -265,7 +265,7 @@ fn has_exact_tag(tags: &[Vec<String>], name: &str, value: &str) -> bool {
#[cfg(test)]
mod tests {
use super::*;
- use crate::{Config, RelayUrlPolicy};
+ use crate::{Config, RelayProfile, RelayUrlPolicy};
use nostr_sdk::prelude::{
EventBuilder, JsonUtil, Keys, RelayUrl as UpstreamRelayUrl, Timestamp,
};
@@ -288,7 +288,7 @@ mod tests {
fn transport() -> (NostrTransport, Arc<MockAuthClient>, RelayUrl) {
let relay =
RelayUrl::parse("wss://relay.example.com", RelayUrlPolicy::Public).expect("relay");
- let config = Config::new(RelayUrlPolicy::Public, [relay.as_str()]).expect("config");
+ let config = Config::from_profile(RelayProfile::public([relay.as_str()]).expect("profile"));
let client = Arc::new(MockAuthClient(AtomicUsize::new(0)));
let transport = NostrTransport::new(config).with_auth_client(client.clone());
(transport, client, relay)
diff --git a/crates/transport_nostr/src/client.rs b/crates/transport_nostr/src/client.rs
@@ -1,61 +1,103 @@
//! Concrete Nostr transport composition.
-use crate::{Error, RelayUrl, RelayUrlPolicy};
+use crate::{Error, RelayEndpoint, RelayProfile, RelayProfileKind, RelayStatusReport, RelayUrl};
use core::fmt;
-use std::collections::BTreeSet;
use std::sync::Arc;
/// Maximum relay targets accepted by one transport instance.
pub(crate) const MAX_RELAYS: usize = 64;
const MAX_TIMEOUT_MS: u64 = 120_000;
const MAX_CONNECTIONS: usize = 64;
+const MAX_RECONNECT_DELAY_MS: u64 = 15 * 60 * 1_000;
+
+/// Deterministic exponential reconnect policy applied independently per relay
+/// and capability direction.
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub struct ReconnectBackoff {
+ initial_delay_ms: u64,
+ max_delay_ms: u64,
+}
+
+impl ReconnectBackoff {
+ /// Creates a bounded reconnect policy.
+ pub fn new(initial_delay_ms: u64, max_delay_ms: u64) -> Result<Self, Error> {
+ if initial_delay_ms == 0
+ || max_delay_ms < initial_delay_ms
+ || max_delay_ms > MAX_RECONNECT_DELAY_MS
+ {
+ return Err(Error::InvalidReconnectBackoff {
+ initial_delay_ms,
+ max_delay_ms,
+ });
+ }
+ Ok(Self {
+ initial_delay_ms,
+ max_delay_ms,
+ })
+ }
+
+ /// Returns the first delay after a retryable failure.
+ #[must_use]
+ pub const fn initial_delay_ms(self) -> u64 {
+ self.initial_delay_ms
+ }
+
+ /// Returns the upper bound for any computed reconnect delay.
+ #[must_use]
+ pub const fn max_delay_ms(self) -> u64 {
+ self.max_delay_ms
+ }
+
+ pub(crate) fn delay_ms(self, consecutive_failures: u32) -> u64 {
+ let exponent = consecutive_failures.saturating_sub(1).min(63);
+ self.initial_delay_ms
+ .saturating_mul(1_u64 << exponent)
+ .min(self.max_delay_ms)
+ }
+}
+
+impl Default for ReconnectBackoff {
+ fn default() -> Self {
+ Self {
+ initial_delay_ms: 1_000,
+ max_delay_ms: 60_000,
+ }
+ }
+}
/// Validated configuration for a concrete Nostr transport.
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct Config {
+ profile_kind: RelayProfileKind,
+ endpoints: Vec<RelayEndpoint>,
relays: Vec<RelayUrl>,
- relay_url_policy: RelayUrlPolicy,
connect_timeout_ms: u64,
request_timeout_ms: u64,
status_timeout_ms: u64,
max_connections: usize,
+ reconnect_backoff: ReconnectBackoff,
}
impl Config {
- /// Builds configuration from explicit relay and network policy inputs.
- pub fn new<I, S>(relay_url_policy: RelayUrlPolicy, relays: I) -> Result<Self, Error>
- where
- I: IntoIterator<Item = S>,
- S: AsRef<str>,
- {
- let mut canonical = Vec::new();
- let mut seen = BTreeSet::new();
- for relay in relays {
- let relay = RelayUrl::parse(relay, relay_url_policy)?;
- if !seen.insert(relay.clone()) {
- return Err(Error::DuplicateRelayUrl {
- url: relay.to_string(),
- });
- }
- canonical.push(relay);
- if canonical.len() > MAX_RELAYS {
- return Err(Error::TooManyRelays {
- max: MAX_RELAYS,
- actual: canonical.len(),
- });
- }
- }
- if canonical.is_empty() {
- return Err(Error::EmptyRelaySet);
- }
- Ok(Self {
- relays: canonical,
- relay_url_policy,
+ /// Builds inert transport configuration from one validated host profile.
+ #[must_use]
+ pub fn from_profile(profile: RelayProfile) -> Self {
+ let relays: Vec<_> = profile
+ .endpoints()
+ .iter()
+ .map(|endpoint| endpoint.url().clone())
+ .collect();
+ let max_connections = 8.min(relays.len());
+ Self {
+ profile_kind: profile.kind(),
+ endpoints: profile.endpoints().to_vec(),
+ relays,
connect_timeout_ms: 10_000,
request_timeout_ms: 30_000,
status_timeout_ms: 5_000,
- max_connections: 8,
- })
+ max_connections,
+ reconnect_backoff: ReconnectBackoff::default(),
+ }
}
/// Sets explicit bounded connection, request, and status timeouts.
@@ -83,14 +125,55 @@ impl Config {
Ok(self)
}
+ /// Sets the deterministic per-relay reconnect policy.
+ #[must_use]
+ pub const fn with_reconnect_backoff(mut self, value: ReconnectBackoff) -> Self {
+ self.reconnect_backoff = value;
+ self
+ }
+
+ /// Returns the selected host profile kind.
+ #[must_use]
+ pub const fn profile_kind(&self) -> RelayProfileKind {
+ self.profile_kind
+ }
+
+ /// Returns configured endpoints with directional access and network policy.
+ #[must_use]
+ pub fn endpoints(&self) -> &[RelayEndpoint] {
+ self.endpoints.as_slice()
+ }
+
/// Returns relays in caller-specified order.
pub fn relays(&self) -> &[RelayUrl] {
self.relays.as_slice()
}
- /// Returns the network policy that must also be applied after DNS resolution.
- pub const fn relay_url_policy(&self) -> RelayUrlPolicy {
- self.relay_url_policy
+ /// Returns relays authorized for reads in deterministic profile order.
+ pub fn read_relays(&self) -> impl Iterator<Item = &RelayUrl> {
+ self.endpoints
+ .iter()
+ .filter_map(|endpoint| endpoint.access().can_read().then_some(endpoint.url()))
+ }
+
+ /// Returns relays authorized for publication in deterministic profile order.
+ pub fn write_relays(&self) -> impl Iterator<Item = &RelayUrl> {
+ self.endpoints
+ .iter()
+ .filter_map(|endpoint| endpoint.access().can_write().then_some(endpoint.url()))
+ }
+
+ pub(crate) fn endpoint_for_target(
+ &self,
+ target: &radroots_transport::Target,
+ ) -> Option<&RelayEndpoint> {
+ (*target.kind() == radroots_transport::TransportId::NOSTR)
+ .then(|| target.uri().as_str())
+ .and_then(|url| {
+ self.endpoints
+ .iter()
+ .find(|endpoint| endpoint.url().as_str() == url)
+ })
}
/// Returns the connection establishment deadline in milliseconds.
@@ -112,6 +195,12 @@ impl Config {
pub const fn max_connections(&self) -> usize {
self.max_connections
}
+
+ /// Returns the per-relay reconnect policy.
+ #[must_use]
+ pub const fn reconnect_backoff(&self) -> ReconnectBackoff {
+ self.reconnect_backoff
+ }
}
fn validate_timeout(field: &'static str, value_ms: u64) -> Result<(), Error> {
@@ -136,10 +225,11 @@ impl NostrTransport {
pub fn new(config: Config) -> Self {
let client = nostr_sdk::Client::builder()
.websocket_transport(crate::relay::HardenedWebsocketTransport::new(
- config.relay_url_policy(),
+ config.endpoints(),
))
.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())),
@@ -147,7 +237,7 @@ impl NostrTransport {
auth: Arc::new(crate::auth::AuthFlow::new(Arc::new(
crate::auth::LiveAuthClient::new(client),
))),
- status: Arc::new(crate::status::StatusTracker::default()),
+ status,
}
}
@@ -156,14 +246,21 @@ impl NostrTransport {
&self.config
}
+ /// Returns passive per-relay and aggregate evidence without network I/O.
+ #[must_use]
+ pub fn relay_status(&self) -> RelayStatusReport {
+ self.status.report()
+ }
+
#[cfg(test)]
pub(crate) fn with_client(config: Config, client: Arc<dyn crate::sink::RelayClient>) -> Self {
+ let status = Arc::new(crate::status::StatusTracker::new(&config));
Self {
config,
client,
source_client: Arc::new(crate::source::LiveRelaySourceClient::isolated()),
auth: Arc::new(crate::auth::AuthFlow::isolated()),
- status: Arc::new(crate::status::StatusTracker::default()),
+ status,
}
}
@@ -172,12 +269,13 @@ impl NostrTransport {
config: Config,
source_client: Arc<dyn crate::source::RelaySourceClient>,
) -> Self {
+ let status = Arc::new(crate::status::StatusTracker::new(&config));
Self {
config,
client: Arc::new(crate::sink::LiveRelayClient::isolated()),
source_client,
auth: Arc::new(crate::auth::AuthFlow::isolated()),
- status: Arc::new(crate::status::StatusTracker::default()),
+ status,
}
}
}
@@ -197,44 +295,51 @@ mod tests {
#[test]
fn config_rejects_empty_duplicate_and_excessive_relay_sets() {
- assert!(Config::new(RelayUrlPolicy::Public, Vec::<String>::new()).is_err());
+ assert!(RelayProfile::simulator(Vec::<String>::new()).is_err());
assert!(
- Config::new(
- RelayUrlPolicy::Public,
- ["wss://relay.example.com", "wss://RELAY.EXAMPLE.COM:443/"],
- )
- .is_err()
+ RelayProfile::public(["wss://relay.example.com", "wss://RELAY.EXAMPLE.COM:443/",])
+ .is_err()
);
- let relays = (0..=MAX_RELAYS).map(|index| format!("wss://r{index}.example.com"));
- assert!(Config::new(RelayUrlPolicy::Public, relays).is_err());
+ let relays = (0..MAX_RELAYS).map(|index| format!("wss://r{index}.example.com"));
+ assert!(RelayProfile::public(relays).is_err());
}
#[test]
fn config_rejects_unbounded_limits() {
- let config =
- Config::new(RelayUrlPolicy::Public, ["wss://relay.example.com"]).expect("config");
+ let config = Config::from_profile(
+ RelayProfile::public(["wss://relay.example.com"]).expect("profile"),
+ );
assert!(config.clone().with_timeouts(0, 1, 1).is_err());
assert!(config.clone().with_timeouts(1, 120_001, 1).is_err());
- assert!(config.with_max_connections(2).is_err());
+ assert!(config.with_max_connections(3).is_err());
}
#[test]
fn valid_configuration_accessors_and_transport_debug_are_complete() {
- let config = Config::new(
- RelayUrlPolicy::Public,
- ["wss://one.example", "wss://two.example"],
+ let config = Config::from_profile(
+ RelayProfile::public(["wss://one.example", "wss://two.example"]).expect("profile"),
)
- .expect("config")
.with_timeouts(1, 2, 3)
.expect("timeouts")
.with_max_connections(2)
.expect("connections");
- assert_eq!(config.relays().len(), 2);
- assert_eq!(config.relay_url_policy(), RelayUrlPolicy::Public);
+ assert_eq!(config.relays().len(), 3);
+ assert_eq!(config.read_relays().count(), 3);
+ assert_eq!(config.write_relays().count(), 2);
+ assert_eq!(config.profile_kind(), RelayProfileKind::Public);
assert_eq!(config.connect_timeout_ms(), 1);
assert_eq!(config.request_timeout_ms(), 2);
assert_eq!(config.status_timeout_ms(), 3);
assert_eq!(config.max_connections(), 2);
+ assert_eq!(config.reconnect_backoff(), ReconnectBackoff::default());
+ assert!(ReconnectBackoff::new(0, 1).is_err());
+ assert!(ReconnectBackoff::new(2, 1).is_err());
+ assert!(ReconnectBackoff::new(1, MAX_RECONNECT_DELAY_MS + 1).is_err());
+ let backoff = ReconnectBackoff::new(2, 5).expect("backoff");
+ assert_eq!(backoff.delay_ms(0), 2);
+ assert_eq!(backoff.delay_ms(1), 2);
+ assert_eq!(backoff.delay_ms(2), 4);
+ assert_eq!(backoff.delay_ms(3), 5);
assert!(config.clone().with_timeouts(120_001, 1, 1).is_err());
assert!(config.clone().with_timeouts(1, 1, 0).is_err());
assert!(config.clone().with_max_connections(0).is_err());
diff --git a/crates/transport_nostr/src/cursor.rs b/crates/transport_nostr/src/cursor.rs
@@ -0,0 +1,86 @@
+//! Deterministic relay event cursor shared by paging and reconnect catch-up.
+
+use crate::Error;
+
+/// Stable total-order position for one Nostr event.
+///
+/// Relay timestamps are only second-granular. The canonical lowercase event id
+/// is therefore a required tie-breaker for both descending fetch pages and
+/// inclusive reconnect catch-up after a subscription interruption.
+#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct RelayCursor {
+ created_at_unix_s: u64,
+ event_id: String,
+}
+
+impl RelayCursor {
+ /// Creates a cursor from an event timestamp and canonical lowercase id.
+ pub fn new(created_at_unix_s: u64, event_id: impl Into<String>) -> Result<Self, Error> {
+ let event_id = event_id.into();
+ if event_id.len() != 64
+ || !event_id
+ .bytes()
+ .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
+ {
+ return Err(Error::InvalidRelayCursor);
+ }
+ Ok(Self {
+ created_at_unix_s,
+ event_id,
+ })
+ }
+
+ /// Returns the second-granular Nostr timestamp.
+ #[must_use]
+ pub const fn created_at_unix_s(&self) -> u64 {
+ self.created_at_unix_s
+ }
+
+ /// Returns the canonical event-id tie breaker.
+ #[must_use]
+ pub fn event_id(&self) -> &str {
+ self.event_id.as_str()
+ }
+
+ /// Returns whether a candidate follows this cursor in ascending reconnect
+ /// order. Equal timestamps are resolved by event id, preventing loss when
+ /// a reconnect query uses an inclusive `since` timestamp.
+ #[must_use]
+ pub fn precedes(&self, created_at_unix_s: u64, event_id: &str) -> bool {
+ created_at_unix_s > self.created_at_unix_s
+ || created_at_unix_s == self.created_at_unix_s && event_id > self.event_id.as_str()
+ }
+
+ /// Returns whether a candidate follows this cursor in descending page
+ /// order. Equal timestamps are resolved by the inverse event-id order.
+ #[must_use]
+ pub(crate) fn page_precedes(&self, created_at_unix_s: u64, event_id: &str) -> bool {
+ created_at_unix_s < self.created_at_unix_s
+ || created_at_unix_s == self.created_at_unix_s && event_id < self.event_id.as_str()
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[test]
+ fn equal_timestamp_ties_are_lossless_in_both_directions() {
+ let cursor = RelayCursor::new(10, "b".repeat(64)).expect("cursor");
+ assert!(cursor.precedes(11, &"0".repeat(64)));
+ assert!(cursor.precedes(10, &"c".repeat(64)));
+ assert!(!cursor.precedes(10, &"b".repeat(64)));
+ assert!(cursor.page_precedes(9, &"f".repeat(64)));
+ assert!(cursor.page_precedes(10, &"a".repeat(64)));
+ assert!(!cursor.page_precedes(10, &"c".repeat(64)));
+ assert_eq!(cursor.created_at_unix_s(), 10);
+ assert_eq!(cursor.event_id(), "b".repeat(64));
+ }
+
+ #[test]
+ fn cursor_rejects_noncanonical_event_ids() {
+ for invalid in ["a".repeat(63), "A".repeat(64), "g".repeat(64)] {
+ assert!(RelayCursor::new(1, invalid).is_err());
+ }
+ }
+}
diff --git a/crates/transport_nostr/src/error.rs b/crates/transport_nostr/src/error.rs
@@ -26,8 +26,15 @@ pub enum Error {
InvalidTimeout { field: &'static str, value_ms: u64 },
/// The per-operation connection limit is outside its governed bounds.
InvalidConnectionLimit { value: usize },
+ /// The reconnect delay policy is empty, inverted, or exceeds its bound.
+ InvalidReconnectBackoff {
+ initial_delay_ms: u64,
+ max_delay_ms: u64,
+ },
/// A transport-neutral target is not a Nostr target.
UnexpectedTransport { actual: String },
+ /// A relay cursor contains a noncanonical event position.
+ InvalidRelayCursor,
/// The generic transport target rejected the relay URL.
Target(String),
/// The relay challenge is empty, malformed, or outside its time bounds.
@@ -85,12 +92,20 @@ impl fmt::Display for Error {
Self::InvalidConnectionLimit { value } => {
write!(formatter, "invalid connection limit: {value}")
}
+ Self::InvalidReconnectBackoff {
+ initial_delay_ms,
+ max_delay_ms,
+ } => write!(
+ formatter,
+ "invalid reconnect backoff: initial={initial_delay_ms}ms max={max_delay_ms}ms"
+ ),
Self::UnexpectedTransport { actual } => {
write!(
formatter,
"expected Nostr transport target, received `{actual}`"
)
}
+ Self::InvalidRelayCursor => formatter.write_str("invalid relay cursor"),
Self::Target(reason) => write!(formatter, "transport target error: {reason}"),
Self::InvalidAuthChallenge => formatter.write_str("invalid NIP-42 challenge"),
Self::AuthChallengeConflict => {
@@ -153,9 +168,14 @@ mod tests {
value_ms: 0,
},
Error::InvalidConnectionLimit { value: 0 },
+ Error::InvalidReconnectBackoff {
+ initial_delay_ms: 0,
+ max_delay_ms: 1,
+ },
Error::UnexpectedTransport {
actual: "local".into(),
},
+ Error::InvalidRelayCursor,
Error::Target("invalid".into()),
Error::InvalidAuthChallenge,
Error::AuthChallengeConflict,
diff --git a/crates/transport_nostr/src/lib.rs b/crates/transport_nostr/src/lib.rs
@@ -4,12 +4,22 @@
mod auth;
mod client;
+mod cursor;
mod error;
+mod profile;
mod relay;
mod sink;
mod source;
mod status;
-pub use client::{Config, NostrTransport};
+pub use client::{Config, NostrTransport, ReconnectBackoff};
+pub use cursor::RelayCursor;
pub use error::Error;
+pub use profile::{
+ DEFAULT_PUBLIC_RELAY, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind,
+};
pub use relay::{RelayUrl, RelayUrlPolicy};
+pub use status::{
+ RelayAggregateState, RelayCapabilityEvidence, RelayEvidenceState, RelayStatus,
+ RelayStatusReport,
+};
diff --git a/crates/transport_nostr/src/profile.rs b/crates/transport_nostr/src/profile.rs
@@ -0,0 +1,242 @@
+//! Validated product relay profiles and directional access policy.
+
+use crate::{Error, RelayUrl, RelayUrlPolicy};
+use std::collections::BTreeSet;
+
+/// Bundled public relay used for read discovery only.
+pub const DEFAULT_PUBLIC_RELAY: &str = "wss://radroots.org";
+
+/// Directional access authorized for one configured relay.
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+#[non_exhaustive]
+pub enum RelayAccess {
+ /// Queries and subscriptions are allowed; publication is never attempted.
+ ReadOnly,
+ /// Queries, subscriptions, and publication are allowed.
+ ReadWrite,
+}
+
+impl RelayAccess {
+ /// Returns whether the profile authorizes event reads.
+ #[must_use]
+ pub const fn can_read(self) -> bool {
+ true
+ }
+
+ /// Returns whether the profile authorizes event publication.
+ #[must_use]
+ pub const fn can_write(self) -> bool {
+ matches!(self, Self::ReadWrite)
+ }
+}
+
+/// Host environment whose network trust rules produced a relay profile.
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+#[non_exhaustive]
+pub enum RelayProfileKind {
+ /// Public-Internet profile with the bundled read-only relay.
+ Public,
+ /// Development-only profile restricted to exact loopback destinations.
+ Simulator,
+ /// Physical-device profile using explicit TLS endpoints on trusted networks.
+ Device,
+}
+
+/// One canonical relay endpoint with explicit network and access policy.
+#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct RelayEndpoint {
+ url: RelayUrl,
+ policy: RelayUrlPolicy,
+ access: RelayAccess,
+}
+
+impl RelayEndpoint {
+ fn new(
+ value: impl AsRef<str>,
+ policy: RelayUrlPolicy,
+ access: RelayAccess,
+ ) -> Result<Self, Error> {
+ Ok(Self {
+ url: RelayUrl::parse(value, policy)?,
+ policy,
+ access,
+ })
+ }
+
+ /// Returns the canonical relay URL.
+ #[must_use]
+ pub const fn url(&self) -> &RelayUrl {
+ &self.url
+ }
+
+ /// Returns the destination policy applied before and after DNS resolution.
+ #[must_use]
+ pub const fn policy(&self) -> RelayUrlPolicy {
+ self.policy
+ }
+
+ /// Returns the read/write authority declared by the profile.
+ #[must_use]
+ pub const fn access(&self) -> RelayAccess {
+ self.access
+ }
+}
+
+/// Complete validated relay selection for one host environment.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RelayProfile {
+ kind: RelayProfileKind,
+ endpoints: Vec<RelayEndpoint>,
+}
+
+impl RelayProfile {
+ /// Builds the ordinary public profile.
+ ///
+ /// `wss://radroots.org/` is always present as read-only. Every supplied
+ /// writable relay must be a TLS public-Internet destination. Supplying the
+ /// bundled relay as writable is rejected rather than silently broadening
+ /// its authority.
+ pub fn public<I, S>(writable_relays: I) -> Result<Self, Error>
+ where
+ I: IntoIterator<Item = S>,
+ S: AsRef<str>,
+ {
+ let mut endpoints = vec![RelayEndpoint::new(
+ DEFAULT_PUBLIC_RELAY,
+ RelayUrlPolicy::Public,
+ RelayAccess::ReadOnly,
+ )?];
+ endpoints.extend(parse_endpoints(
+ writable_relays,
+ RelayUrlPolicy::Public,
+ RelayAccess::ReadWrite,
+ )?);
+ Self::validated(RelayProfileKind::Public, endpoints)
+ }
+
+ /// Builds a development-only profile from exact loopback relays.
+ ///
+ /// Plaintext `ws://` is accepted only by this profile. At least one relay
+ /// is required and every endpoint is writable.
+ pub fn simulator<I, S>(loopback_relays: I) -> Result<Self, Error>
+ where
+ I: IntoIterator<Item = S>,
+ S: AsRef<str>,
+ {
+ let endpoints = parse_endpoints(
+ loopback_relays,
+ RelayUrlPolicy::Local,
+ RelayAccess::ReadWrite,
+ )?;
+ Self::validated(RelayProfileKind::Simulator, endpoints)
+ }
+
+ /// Builds a physical-device profile from explicit writable TLS endpoints.
+ ///
+ /// The bundled public relay remains read-only. Writable endpoints may
+ /// resolve to public or private addresses, but loopback, unspecified, and
+ /// multicast destinations remain forbidden before and after resolution.
+ pub fn device<I, S>(writable_relays: I) -> Result<Self, Error>
+ where
+ I: IntoIterator<Item = S>,
+ S: AsRef<str>,
+ {
+ let mut endpoints = vec![RelayEndpoint::new(
+ DEFAULT_PUBLIC_RELAY,
+ RelayUrlPolicy::Public,
+ RelayAccess::ReadOnly,
+ )?];
+ endpoints.extend(parse_endpoints(
+ writable_relays,
+ RelayUrlPolicy::PrivateNetwork,
+ RelayAccess::ReadWrite,
+ )?);
+ Self::validated(RelayProfileKind::Device, endpoints)
+ }
+
+ fn validated(kind: RelayProfileKind, endpoints: Vec<RelayEndpoint>) -> Result<Self, Error> {
+ if endpoints.is_empty() {
+ return Err(Error::EmptyRelaySet);
+ }
+ if endpoints.len() > crate::client::MAX_RELAYS {
+ return Err(Error::TooManyRelays {
+ max: crate::client::MAX_RELAYS,
+ actual: endpoints.len(),
+ });
+ }
+ let mut seen = BTreeSet::new();
+ for endpoint in &endpoints {
+ if !seen.insert(endpoint.url.clone()) {
+ return Err(Error::DuplicateRelayUrl {
+ url: endpoint.url.to_string(),
+ });
+ }
+ }
+ Ok(Self { kind, endpoints })
+ }
+
+ /// Returns the selected host-environment profile.
+ #[must_use]
+ pub const fn kind(&self) -> RelayProfileKind {
+ self.kind
+ }
+
+ /// Returns endpoints in deterministic profile order.
+ #[must_use]
+ pub fn endpoints(&self) -> &[RelayEndpoint] {
+ self.endpoints.as_slice()
+ }
+}
+
+fn parse_endpoints<I, S>(
+ values: I,
+ policy: RelayUrlPolicy,
+ access: RelayAccess,
+) -> Result<Vec<RelayEndpoint>, Error>
+where
+ I: IntoIterator<Item = S>,
+ S: AsRef<str>,
+{
+ values
+ .into_iter()
+ .map(|value| RelayEndpoint::new(value, policy, access))
+ .collect()
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[test]
+ fn public_profile_never_promotes_the_bundled_relay_to_writable() {
+ let profile = RelayProfile::public(["wss://write.example"]).expect("public profile");
+ assert_eq!(profile.kind(), RelayProfileKind::Public);
+ assert_eq!(profile.endpoints().len(), 2);
+ assert_eq!(profile.endpoints()[0].url().as_str(), DEFAULT_PUBLIC_RELAY);
+ assert_eq!(profile.endpoints()[0].access(), RelayAccess::ReadOnly);
+ assert_eq!(profile.endpoints()[1].access(), RelayAccess::ReadWrite);
+ assert!(RelayProfile::public([DEFAULT_PUBLIC_RELAY]).is_err());
+ assert!(RelayProfile::public(["ws://public.example"]).is_err());
+ }
+
+ #[test]
+ fn loopback_is_confined_to_the_simulator_profile() {
+ assert!(RelayProfile::simulator(["ws://127.0.0.1:7447"]).is_ok());
+ assert!(RelayProfile::simulator(["wss://localhost:7447"]).is_ok());
+ assert!(RelayProfile::simulator(Vec::<String>::new()).is_err());
+ assert!(RelayProfile::public(["ws://127.0.0.1:7447"]).is_err());
+ assert!(RelayProfile::device(["wss://127.0.0.1:7447"]).is_err());
+ }
+
+ #[test]
+ fn device_profile_requires_tls_and_rejects_duplicate_authority() {
+ let profile = RelayProfile::device(["wss://10.0.0.5:7447"]).expect("device profile");
+ assert_eq!(profile.kind(), RelayProfileKind::Device);
+ assert_eq!(
+ profile.endpoints()[1].policy(),
+ RelayUrlPolicy::PrivateNetwork
+ );
+ assert!(RelayProfile::device(["ws://10.0.0.5:7447"]).is_err());
+ assert!(RelayProfile::device([DEFAULT_PUBLIC_RELAY]).is_err());
+ }
+}
diff --git a/crates/transport_nostr/src/relay.rs b/crates/transport_nostr/src/relay.rs
@@ -1,6 +1,6 @@
//! Nostr relay identifiers and network policy.
-use crate::Error;
+use crate::{Error, RelayEndpoint};
use async_wsocket::futures_util::stream::SplitSink;
use async_wsocket::futures_util::{Sink, SinkExt, StreamExt, TryStreamExt};
use async_wsocket::{Message, WebSocket};
@@ -10,7 +10,9 @@ use nostr_relay_pool::ConnectionMode;
use nostr_relay_pool::transport::error::TransportError;
use nostr_relay_pool::transport::websocket::{WebSocketSink, WebSocketStream, WebSocketTransport};
use radroots_transport::{BoxFuture, Target, TransportId};
+use std::collections::BTreeMap;
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
+use std::sync::Arc;
use std::task::{Context, Poll};
use std::time::Duration;
use tokio::net::TcpStream;
@@ -97,14 +99,21 @@ impl fmt::Display for RelayUrl {
/// WebSocket connector that validates and pins DNS results before opening a
/// socket while retaining the original host name for TLS verification.
-#[derive(Clone, Copy, Debug)]
+#[derive(Clone, Debug)]
pub(crate) struct HardenedWebsocketTransport {
- policy: RelayUrlPolicy,
+ policies: Arc<BTreeMap<String, RelayUrlPolicy>>,
}
impl HardenedWebsocketTransport {
- pub(crate) const fn new(policy: RelayUrlPolicy) -> Self {
- Self { policy }
+ pub(crate) fn new(endpoints: &[RelayEndpoint]) -> Self {
+ Self {
+ policies: Arc::new(
+ endpoints
+ .iter()
+ .map(|endpoint| (endpoint.url().as_str().to_owned(), endpoint.policy()))
+ .collect(),
+ ),
+ }
}
}
@@ -128,7 +137,14 @@ impl WebSocketTransport for HardenedWebsocketTransport {
"proxy and Tor connection modes are not configured",
));
}
- let relay = RelayUrl::parse(url.as_str(), self.policy)
+ let target = Target::nostr_relay(url.as_str())
+ .map_err(|_| policy_error("relay URL is invalid"))?;
+ let policy = self
+ .policies
+ .get(target.uri().as_str())
+ .copied()
+ .ok_or_else(|| policy_error("relay URL is not configured"))?;
+ let relay = RelayUrl::parse(target.uri().as_str(), policy)
.map_err(|_| policy_error("relay URL is denied by network policy"))?;
let parsed =
Url::parse(relay.as_str()).map_err(|_| policy_error("relay URL is invalid"))?;
@@ -142,7 +158,7 @@ impl WebSocketTransport for HardenedWebsocketTransport {
let connect = async {
let addresses = resolve_bounded(host, port).await?;
relay
- .validate_resolved_addresses(self.policy, addresses.iter().map(SocketAddr::ip))
+ .validate_resolved_addresses(policy, addresses.iter().map(SocketAddr::ip))
.map_err(|_| policy_error("relay DNS result is denied by network policy"))?;
let tcp = connect_pinned(addresses.as_slice()).await?;
let (stream, _) = tokio_tungstenite::client_async_tls(relay.as_str(), tcp)
@@ -252,7 +268,7 @@ impl Sink<Message> for HardenedTransportSink {
}
/// Destination class authorized for relay connections.
-#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub enum RelayUrlPolicy {
/// TLS-only public Internet endpoints; resolved addresses must be global.
@@ -390,7 +406,8 @@ mod tests {
RelayUrl::from_target(&local, RelayUrlPolicy::Public),
Err(Error::UnexpectedTransport { .. })
));
- assert!(HardenedWebsocketTransport::new(RelayUrlPolicy::Public).support_ping());
+ let profile = crate::RelayProfile::public(["wss://relay.example.com"]).expect("profile");
+ assert!(HardenedWebsocketTransport::new(profile.endpoints()).support_ping());
assert!(!policy_error("denied").to_string().is_empty());
}
diff --git a/crates/transport_nostr/src/sink.rs b/crates/transport_nostr/src/sink.rs
@@ -2,6 +2,7 @@
use crate::{NostrTransport, RelayUrl, status};
use core::time::Duration;
+use futures::{StreamExt, stream};
use radroots_nostr::event::Event;
use radroots_transport::{
BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure,
@@ -21,6 +22,7 @@ pub(crate) trait RelayClient: Send + Sync {
&'a self,
relays: Vec<RelayUrl>,
event: Event,
+ max_connections: usize,
connect_timeout: Duration,
operation_timeout: Duration,
) -> BoxFuture<'a, Vec<RelayPublishResult>>;
@@ -52,47 +54,51 @@ impl RelayClient for LiveRelayClient {
&'a self,
relays: Vec<RelayUrl>,
event: Event,
+ max_connections: usize,
connect_timeout: Duration,
operation_timeout: Duration,
) -> BoxFuture<'a, Vec<RelayPublishResult>> {
Box::pin(async move {
- let mut results = Vec::with_capacity(relays.len());
- for relay in relays {
- let url = relay.as_str().to_owned();
- let attempt = async {
- self.client.add_relay(url.as_str()).await?;
- self.client
- .try_connect_relay(url.as_str(), connect_timeout)
- .await?;
- self.client.send_event_to([url.as_str()], &event).await
- };
- let outcome = match tokio::time::timeout(operation_timeout, 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()))
+ 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?;
+ self.client
+ .try_connect_relay(url.as_str(), connect_timeout)
+ .await?;
+ self.client.send_event_to([url.as_str()], &event).await
+ };
+ let outcome = match tokio::time::timeout(operation_timeout, 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")),
- };
- results.push(RelayPublishResult { relay, outcome });
- }
- results
+ .unwrap_or_else(|| status::delivery_failure("relay omitted result")),
+ };
+ RelayPublishResult { relay, outcome }
+ }
+ }))
+ .buffered(max_connections)
+ .collect()
+ .await
})
}
}
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(status::sink_status(&self.status, configured)) })
+ Box::pin(async move { Ok(status::sink_status(&self.status)) })
}
fn deliver(
@@ -100,14 +106,31 @@ impl EventSink for NostrTransport {
request: DeliveryRequest,
) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
Box::pin(async move {
+ let now_unix_ms = unix_time_ms();
let mut requested = Vec::new();
let mut skipped = Vec::new();
for target in request.target_set().targets() {
- match RelayUrl::from_target(target, self.config().relay_url_policy()) {
- Ok(relay) if self.config().relays().contains(&relay) => {
- requested.push((relay, target.clone()));
+ match self.config().endpoint_for_target(target) {
+ Some(endpoint) if endpoint.access().can_write() => {
+ let relay = endpoint.url().clone();
+ if self.status.may_write(&relay, now_unix_ms) {
+ requested.push((relay, target.clone()));
+ } else {
+ skipped.push(
+ DeliveryTargetReceipt::skipped(
+ target.clone(),
+ DeliveryOutcome::unavailable()
+ .with_detail(
+ "reconnect_backoff",
+ "relay reconnect backoff is active",
+ )
+ .map_err(|_| SinkFailure::invalid_contract(&request))?,
+ )
+ .map_err(|_| SinkFailure::invalid_contract(&request))?,
+ );
+ }
}
- _ => skipped.push(
+ None | Some(_) => skipped.push(
DeliveryTargetReceipt::skipped(
target.clone(),
DeliveryOutcome::rejected()
@@ -127,56 +150,61 @@ impl EventSink for NostrTransport {
Ok(event) => event,
Err(_) => return Err(SinkFailure::invalid_contract(&request)),
};
- let remaining_ms = request.deadline_unix_ms().saturating_sub(unix_time_ms());
+ let remaining_ms = request.deadline_unix_ms().saturating_sub(now_unix_ms);
let operation_timeout_ms = remaining_ms.min(self.config().request_timeout_ms());
if operation_timeout_ms == 0 {
+ for (relay, _) in &requested {
+ self.status.record_write(relay, false, true, now_unix_ms);
+ }
let timeout = status::delivery_failure("timeout");
skipped.extend(requested.into_iter().map(|(_, target)| {
DeliveryTargetReceipt::skipped(target, timeout.clone())
.expect("normalized timeout cannot satisfy delivery")
}));
- self.status.record_sink(0, skipped.len(), Some("timeout"));
return DeliveryReceipt::for_request(&request, skipped)
.map_err(|_| SinkFailure::invalid_contract(&request));
}
let expected: BTreeSet<_> = requested.iter().map(|(relay, _)| relay.clone()).collect();
+ for (relay, _) in &requested {
+ self.status.begin_write(relay, now_unix_ms);
+ }
let results = self
.client
.publish(
requested.iter().map(|(relay, _)| relay.clone()).collect(),
event,
+ self.config().max_connections(),
Duration::from_millis(self.config().connect_timeout_ms()),
Duration::from_millis(operation_timeout_ms),
)
.await;
let mut by_relay = BTreeMap::new();
+ let observed_at_unix_ms = unix_time_ms().max(now_unix_ms);
for result in results {
- if expected.contains(&result.relay)
- && by_relay.insert(result.relay, result.outcome).is_none()
- {
- continue;
+ if !expected.contains(&result.relay) || by_relay.contains_key(&result.relay) {
+ return Err(SinkFailure::invalid_contract(&request));
}
- return Err(SinkFailure::invalid_contract(&request));
+ let succeeded = status::delivery_succeeded(&result.outcome);
+ self.status.record_write(
+ &result.relay,
+ succeeded,
+ result.outcome.is_retryable(),
+ observed_at_unix_ms,
+ );
+ by_relay.insert(result.relay, result.outcome);
}
let mut receipts = skipped;
for (relay, target) in requested {
let outcome = by_relay.remove(&relay).unwrap_or_else(|| {
+ self.status
+ .record_write(&relay, false, true, observed_at_unix_ms);
DeliveryOutcome::unavailable()
.with_detail("missing_result", "relay returned no result")
.expect("static normalized outcome")
});
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)
.map_err(|_| SinkFailure::invalid_contract(&request))
})
@@ -194,7 +222,7 @@ fn unix_time_ms() -> u64 {
#[cfg(test)]
mod tests {
use super::*;
- use crate::{Config, RelayUrlPolicy};
+ use crate::{Config, RelayProfile, RelayUrlPolicy};
use radroots_transport::{
Target, TargetSet,
outcome::DeliveryOutcomeKind,
@@ -216,6 +244,7 @@ mod tests {
&'a self,
relays: Vec<RelayUrl>,
_event: Event,
+ _max_connections: usize,
_connect_timeout: Duration,
_operation_timeout: Duration,
) -> BoxFuture<'a, Vec<RelayPublishResult>> {
@@ -261,11 +290,9 @@ mod tests {
#[test]
fn sink_returns_normalized_per_relay_partial_success() {
- let config = Config::new(
- RelayUrlPolicy::Public,
- ["wss://one.example", "wss://two.example"],
- )
- .expect("config");
+ let config = Config::from_profile(
+ RelayProfile::public(["wss://one.example", "wss://two.example"]).expect("profile"),
+ );
let two = RelayUrl::parse("wss://two.example", RelayUrlPolicy::Public).expect("two");
let client = MockRelayClient {
outcomes: BTreeMap::from([(two, status::delivery_failure("rate limited"))]),
@@ -312,6 +339,7 @@ mod tests {
&'a self,
_relays: Vec<RelayUrl>,
_event: Event,
+ _max_connections: usize,
_connect_timeout: Duration,
_operation_timeout: Duration,
) -> BoxFuture<'a, Vec<RelayPublishResult>> {
@@ -321,11 +349,9 @@ mod tests {
}
let calls = Arc::new(AtomicUsize::new(0));
- let config = Config::new(
- RelayUrlPolicy::Public,
- ["wss://one.example", "wss://two.example"],
- )
- .expect("config");
+ let config = Config::from_profile(
+ RelayProfile::public(["wss://one.example", "wss://two.example"]).expect("profile"),
+ );
let transport =
NostrTransport::with_client(config, Arc::new(CountingRelayClient(Arc::clone(&calls))));
let delivery = transport.deliver(request());
@@ -343,6 +369,7 @@ mod tests {
&'a self,
_relays: Vec<RelayUrl>,
_event: Event,
+ _max_connections: usize,
_connect_timeout: Duration,
_operation_timeout: Duration,
) -> BoxFuture<'a, Vec<RelayPublishResult>> {
@@ -352,11 +379,9 @@ mod tests {
}
let calls = Arc::new(AtomicUsize::new(0));
- let config = Config::new(
- RelayUrlPolicy::Public,
- ["wss://one.example", "wss://two.example"],
- )
- .expect("config");
+ let config = Config::from_profile(
+ RelayProfile::public(["wss://one.example", "wss://two.example"]).expect("profile"),
+ );
let transport =
NostrTransport::with_client(config, Arc::new(CountingRelayClient(Arc::clone(&calls))));
let receipt = futures::executor::block_on(transport.deliver(request_with_deadline(1)))
@@ -378,6 +403,7 @@ mod tests {
&'a self,
_relays: Vec<RelayUrl>,
_event: Event,
+ _max_connections: usize,
_connect_timeout: Duration,
_operation_timeout: Duration,
) -> BoxFuture<'a, Vec<RelayPublishResult>> {
@@ -386,11 +412,9 @@ mod tests {
}
fn scripted(results: Vec<RelayPublishResult>) -> NostrTransport {
- let config = Config::new(
- RelayUrlPolicy::Public,
- ["wss://one.example", "wss://two.example"],
- )
- .expect("config");
+ let config = Config::from_profile(
+ RelayProfile::public(["wss://one.example", "wss://two.example"]).expect("profile"),
+ );
NostrTransport::with_client(config, Arc::new(ScriptedRelayClient(results)))
}
@@ -459,6 +483,7 @@ mod tests {
let results = futures::executor::block_on(client.publish(
vec![],
radroots_nostr::event::to_nostr(payload().event().envelope()).expect("nostr event"),
+ 1,
Duration::from_millis(1),
Duration::from_millis(1),
));
diff --git a/crates/transport_nostr/src/source.rs b/crates/transport_nostr/src/source.rs
@@ -1,19 +1,22 @@
//! Nostr implementation of the transport event source.
-use crate::{NostrTransport, RelayUrl, status};
+use crate::{NostrTransport, RelayCursor, RelayUrl, status};
use core::cmp::Ordering;
use core::time::Duration;
+use futures::{StreamExt, stream};
use nostr_sdk::prelude::{Filter, JsonUtil, Kind, Timestamp};
use radroots_transport::{
BoxFuture, EventSource, FetchPage, FetchRequest,
outcome::{FetchTargetOutcome, FetchTargetState},
source::{EventProvenance, FetchCursor, NextPage, ObservedEvent, SourceStatus},
};
+use sha2::{Digest, Sha256};
use std::collections::{BTreeMap, BTreeSet};
use std::time::{SystemTime, UNIX_EPOCH};
const UPSTREAM_FETCH_LIMIT: usize = 1_000;
-const CURSOR_PREFIX: &str = "nostr-v1";
+const CURSOR_PREFIX: &str = "nostr-v2";
+const CURSOR_SCOPE_DOMAIN: &[u8] = b"radroots.transport-nostr.fetch-cursor.v2\0";
#[derive(Clone, Debug)]
pub(crate) struct SourceQuery {
@@ -22,6 +25,7 @@ pub(crate) struct SourceQuery {
until_unix_seconds: Option<u64>,
connect_timeout: Duration,
timeout: Duration,
+ max_connections: usize,
}
#[derive(Clone, Debug)]
@@ -58,56 +62,68 @@ impl RelaySourceClient for LiveRelaySourceClient {
#[cfg_attr(coverage_nightly, coverage(off))]
fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> {
Box::pin(async move {
- let mut batches = Vec::with_capacity(query.relays.len());
- for relay in query.relays {
- let url = relay.as_str().to_owned();
- let result = async {
- let kinds = query
- .selector
- .kinds()
- .iter()
- .filter_map(|kind| u16::try_from(*kind).ok())
- .map(Kind::from)
- .collect::<Vec<_>>();
- if !query.selector.kinds().is_empty() && kinds.is_empty() {
- return Ok(Vec::new());
- }
- let authors = query
- .selector
- .authors()
- .iter()
- .filter_map(|author| radroots_nostr::key::public_key_to_nostr(*author).ok())
- .collect::<Vec<_>>();
- if authors.len() != query.selector.authors().len() {
- return Ok(Vec::new());
- }
- self.client.add_relay(url.as_str()).await?;
- self.client
- .try_connect_relay(url.as_str(), query.connect_timeout)
- .await?;
- let mut filter = Filter::new().limit(UPSTREAM_FETCH_LIMIT);
- if !kinds.is_empty() {
- filter = filter.kinds(kinds);
- }
- if !authors.is_empty() {
- filter = filter.authors(authors);
- }
- if let Some(since) = query.selector.since_unix_seconds() {
- filter = filter.since(Timestamp::from_secs(since));
- }
- if let Some(until) = query.until_unix_seconds {
- filter = filter.until(Timestamp::from_secs(until));
+ let SourceQuery {
+ relays,
+ selector,
+ until_unix_seconds,
+ connect_timeout,
+ timeout,
+ max_connections,
+ } = query;
+ stream::iter(relays.into_iter().map(|relay| {
+ let selector = selector.clone();
+ async move {
+ let url = relay.as_str().to_owned();
+ let result = async {
+ let kinds = selector
+ .kinds()
+ .iter()
+ .filter_map(|kind| u16::try_from(*kind).ok())
+ .map(Kind::from)
+ .collect::<Vec<_>>();
+ if !selector.kinds().is_empty() && kinds.is_empty() {
+ return Ok(Vec::new());
+ }
+ let authors = selector
+ .authors()
+ .iter()
+ .filter_map(|author| {
+ radroots_nostr::key::public_key_to_nostr(*author).ok()
+ })
+ .collect::<Vec<_>>();
+ if authors.len() != selector.authors().len() {
+ return Ok(Vec::new());
+ }
+ self.client.add_relay(url.as_str()).await?;
+ self.client
+ .try_connect_relay(url.as_str(), connect_timeout)
+ .await?;
+ let mut filter = Filter::new().limit(UPSTREAM_FETCH_LIMIT);
+ if !kinds.is_empty() {
+ filter = filter.kinds(kinds);
+ }
+ if !authors.is_empty() {
+ filter = filter.authors(authors);
+ }
+ if let Some(since) = selector.since_unix_seconds() {
+ filter = filter.since(Timestamp::from_secs(since));
+ }
+ if let Some(until) = until_unix_seconds {
+ filter = filter.until(Timestamp::from_secs(until));
+ }
+ self.client
+ .fetch_events_from([url.as_str()], filter, timeout)
+ .await
+ .map(|events| events.iter().map(JsonUtil::as_json).collect())
}
- self.client
- .fetch_events_from([url.as_str()], filter, query.timeout)
- .await
- .map(|events| events.iter().map(JsonUtil::as_json).collect())
+ .await
+ .map_err(|error: nostr_sdk::client::Error| error.to_string());
+ RelayFetchBatch { relay, result }
}
- .await
- .map_err(|error: nostr_sdk::client::Error| error.to_string());
- batches.push(RelayFetchBatch { relay, result });
- }
- batches
+ }))
+ .buffered(max_connections)
+ .collect()
+ .await
})
}
}
@@ -122,8 +138,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(status::source_status(&self.status, configured)) })
+ Box::pin(async move { Ok(status::source_status(&self.status)) })
}
fn fetch(
@@ -131,7 +146,11 @@ impl EventSource for NostrTransport {
request: FetchRequest,
) -> BoxFuture<'_, Result<FetchPage, radroots_transport::Error>> {
Box::pin(async move {
- let cursor = request.cursor().map(parse_cursor).transpose()?.flatten();
+ let cursor_scope = request_scope(&request);
+ let cursor = request
+ .cursor()
+ .map(|cursor| parse_cursor(cursor, cursor_scope.as_str()))
+ .transpose()?;
let selector_until = request.selector().until_unix_seconds();
let now_ms = unix_time_ms();
let remaining_ms = request.bounds().deadline_unix_ms().saturating_sub(now_ms);
@@ -140,11 +159,22 @@ impl EventSource for NostrTransport {
let mut targets = BTreeMap::new();
let mut outcomes = Vec::new();
for target in request.target_set().targets() {
- match RelayUrl::from_target(target, self.config().relay_url_policy()) {
- Ok(relay) if self.config().relays().contains(&relay) => {
- targets.insert(relay, target.clone());
+ match self.config().endpoint_for_target(target) {
+ Some(endpoint) if endpoint.access().can_read() => {
+ let relay = endpoint.url().clone();
+ if self.status.may_read(&relay, now_ms) {
+ targets.insert(relay, target.clone());
+ } else {
+ outcomes.push(
+ FetchTargetOutcome::new(
+ target.fingerprint().clone(),
+ FetchTargetState::FailedRetryable,
+ )
+ .with_message("relay reconnect backoff is active"),
+ );
+ }
}
- _ => outcomes.push(
+ None | Some(_) => outcomes.push(
FetchTargetOutcome::new(
target.fingerprint().clone(),
FetchTargetState::FailedTerminal,
@@ -155,6 +185,9 @@ impl EventSource for NostrTransport {
}
if timeout_ms == 0 {
+ for relay in targets.keys() {
+ self.status.record_read(relay, false, true, now_ms);
+ }
outcomes.extend(targets.values().map(|target| {
FetchTargetOutcome::new(
target.fingerprint().clone(),
@@ -162,33 +195,34 @@ 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);
}
+ for relay in targets.keys() {
+ self.status.begin_read(relay, now_ms);
+ }
let batches = self
.source_client
.fetch(SourceQuery {
relays: targets.keys().cloned().collect(),
selector: request.selector().clone(),
until_unix_seconds: match (selector_until, cursor.as_ref()) {
- (Some(until), Some(cursor)) => Some(until.min(cursor.created_at)),
+ (Some(until), Some(cursor)) => Some(until.min(cursor.created_at_unix_s())),
(Some(until), None) => Some(until),
- (None, Some(cursor)) => Some(cursor.created_at),
+ (None, Some(cursor)) => Some(cursor.created_at_unix_s()),
(None, None) => None,
},
connect_timeout: Duration::from_millis(
timeout_ms.min(self.config().connect_timeout_ms()),
),
timeout: Duration::from_millis(timeout_ms),
+ max_connections: self.config().max_connections(),
})
.await;
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;
+ let observed_at_unix_ms = unix_time_ms().max(now_ms);
for batch in batches {
let Some(target) = targets.get(&batch.relay) else {
return Err(radroots_transport::Error::UnexpectedFetchTargetOutcome);
@@ -198,7 +232,8 @@ impl EventSource for NostrTransport {
}
match batch.result {
Ok(raw_events) => {
- succeeded += 1;
+ self.status
+ .record_read(&batch.relay, true, false, observed_at_unix_ms);
for raw in raw_events {
match radroots_event_codec::decode::signed_event(raw.as_str()) {
Ok(event) if request.selector().matches(&event) => {
@@ -235,9 +270,13 @@ impl EventSource for NostrTransport {
outcomes.push(outcome);
}
Err(message) => {
- failed += 1;
let (state, safe) = status::fetch_failure(message.as_str());
- diagnostic.get_or_insert(safe);
+ self.status.record_read(
+ &batch.relay,
+ false,
+ state.is_retryable(),
+ observed_at_unix_ms,
+ );
outcomes.push(
FetchTargetOutcome::new(target.fingerprint().clone(), state)
.with_message(safe),
@@ -247,8 +286,8 @@ impl EventSource for NostrTransport {
}
for (relay, target) in &targets {
if !reported.contains(relay) {
- failed += 1;
- diagnostic.get_or_insert("relay returned no fetch result");
+ self.status
+ .record_read(relay, false, true, observed_at_unix_ms);
outcomes.push(
FetchTargetOutcome::new(
target.fingerprint().clone(),
@@ -258,8 +297,6 @@ impl EventSource for NostrTransport {
);
}
}
- self.status.record_source(succeeded, failed, diagnostic);
-
candidates.sort_by(compare_candidate);
if let Some(cursor) = &cursor {
candidates.retain(|candidate| candidate_is_after_cursor(candidate, cursor));
@@ -272,8 +309,8 @@ impl EventSource for NostrTransport {
let next_page = if has_more {
let last = candidates.last().expect("non-empty bounded page");
NextPage::Cursor(FetchCursor::parse(format!(
- "{CURSOR_PREFIX}:{}:{}",
- last.created_at, last.event_id
+ "{CURSOR_PREFIX}:{}:{}:{cursor_scope}",
+ last.created_at, last.event_id,
))?)
} else {
NextPage::Complete
@@ -301,30 +338,26 @@ impl EventSource for NostrTransport {
}
}
-#[derive(Clone, Debug)]
-struct CursorPosition {
- created_at: u64,
- event_id: String,
-}
-
-fn parse_cursor(cursor: &FetchCursor) -> Result<Option<CursorPosition>, radroots_transport::Error> {
+fn parse_cursor(
+ cursor: &FetchCursor,
+ expected_scope: &str,
+) -> Result<RelayCursor, radroots_transport::Error> {
let mut parts = cursor.as_str().split(':');
let valid = parts.next() == Some(CURSOR_PREFIX);
let created_at = parts.next().and_then(|value| value.parse::<u64>().ok());
let event_id = parts.next();
+ let scope = parts.next();
if !valid || parts.next().is_some() {
return Err(radroots_transport::Error::InvalidFetchCursor);
}
- let (Some(created_at), Some(event_id)) = (created_at, event_id) else {
+ let (Some(created_at), Some(event_id), Some(scope)) = (created_at, event_id, scope) else {
return Err(radroots_transport::Error::InvalidFetchCursor);
};
- if event_id.len() != 64 || !event_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
+ if scope != expected_scope {
return Err(radroots_transport::Error::InvalidFetchCursor);
}
- Ok(Some(CursorPosition {
- created_at,
- event_id: event_id.to_ascii_lowercase(),
- }))
+ RelayCursor::new(created_at, event_id)
+ .map_err(|_| radroots_transport::Error::InvalidFetchCursor)
}
fn compare_candidate(left: &Candidate, right: &Candidate) -> Ordering {
@@ -335,9 +368,48 @@ fn compare_candidate(left: &Candidate, right: &Candidate) -> Ordering {
.then_with(|| left.relay.cmp(&right.relay))
}
-fn candidate_is_after_cursor(candidate: &Candidate, cursor: &CursorPosition) -> bool {
- candidate.created_at < cursor.created_at
- || candidate.created_at == cursor.created_at && candidate.event_id < cursor.event_id
+fn candidate_is_after_cursor(candidate: &Candidate, cursor: &RelayCursor) -> bool {
+ cursor.page_precedes(candidate.created_at, candidate.event_id.as_str())
+}
+
+fn request_scope(request: &FetchRequest) -> String {
+ let mut hasher = Sha256::new();
+ hasher.update(CURSOR_SCOPE_DOMAIN);
+ for target in request.target_set().targets() {
+ hasher.update(target.fingerprint().as_str().as_bytes());
+ hasher.update([0]);
+ }
+ for kind in request.selector().kinds() {
+ hasher.update(kind.to_be_bytes());
+ }
+ hasher.update([0]);
+ for author in request.selector().authors() {
+ hasher.update(author.as_bytes());
+ }
+ hasher.update([0]);
+ hash_optional_u64(&mut hasher, request.selector().since_unix_seconds());
+ hash_optional_u64(&mut hasher, request.selector().until_unix_seconds());
+ hex_encode(&hasher.finalize())
+}
+
+fn hash_optional_u64(hasher: &mut Sha256, value: Option<u64>) {
+ match value {
+ Some(value) => {
+ hasher.update([1]);
+ hasher.update(value.to_be_bytes());
+ }
+ None => hasher.update([0]),
+ }
+}
+
+fn hex_encode(bytes: &[u8]) -> String {
+ const HEX: &[u8; 16] = b"0123456789abcdef";
+ let mut encoded = String::with_capacity(bytes.len() * 2);
+ for byte in bytes {
+ encoded.push(HEX[(byte >> 4) as usize] as char);
+ encoded.push(HEX[(byte & 0x0f) as usize] as char);
+ }
+ encoded
}
#[cfg_attr(coverage_nightly, coverage(off))]
@@ -351,7 +423,7 @@ fn unix_time_ms() -> u64 {
#[cfg(test)]
mod tests {
use super::*;
- use crate::{Config, RelayUrlPolicy};
+ use crate::{Config, RelayProfile, RelayUrlPolicy};
use radroots_transport::{
FetchRequest, Target, TargetSet,
source::{FetchBounds, FetchSelector, NextPage},
@@ -383,11 +455,9 @@ mod tests {
}
fn transport() -> NostrTransport {
- let config = Config::new(
- RelayUrlPolicy::Public,
- ["wss://one.example", "wss://two.example"],
- )
- .expect("config");
+ let config = Config::from_profile(
+ RelayProfile::public(["wss://one.example", "wss://two.example"]).expect("profile"),
+ );
NostrTransport::with_source_client(config, Arc::new(MockSourceClient))
}
@@ -477,7 +547,8 @@ mod tests {
}
let calls = Arc::new(AtomicUsize::new(0));
- let config = Config::new(RelayUrlPolicy::Public, ["wss://one.example"]).expect("config");
+ let config =
+ Config::from_profile(RelayProfile::public(["wss://one.example"]).expect("profile"));
let transport = NostrTransport::with_source_client(
config,
Arc::new(CountingSourceClient(Arc::clone(&calls))),
@@ -518,11 +589,9 @@ mod tests {
}
fn scripted(batches: Vec<RelayFetchBatch>) -> NostrTransport {
- let config = Config::new(
- RelayUrlPolicy::Public,
- ["wss://one.example", "wss://two.example"],
- )
- .expect("config");
+ let config = Config::from_profile(
+ RelayProfile::public(["wss://one.example", "wss://two.example"]).expect("profile"),
+ );
NostrTransport::with_source_client(config, Arc::new(ScriptedSourceClient(batches)))
}
@@ -599,22 +668,19 @@ mod tests {
fn cursor_and_candidate_ordering_cover_boundaries() {
for invalid in [
"nostr-v1",
- "nostr-v1:not-a-time:id",
- "nostr-v1:1",
- "nostr-v1:1:id:extra",
- "nostr-v1:1:abc",
- "nostr-v1:1:gggggggggggggggggggggggggggggggggggggggggggggggggggggggggggggggg",
+ "nostr-v2:not-a-time:id:scope",
+ "nostr-v2:1",
+ "nostr-v2:1:id:scope:extra",
+ "nostr-v2:1:abc:scope",
+ "nostr-v2:1:gggggggggggggggggggggggggggggggggggggggggggggggggggggggggggggggg:scope",
] {
let cursor = FetchCursor::parse(invalid).expect("opaque cursor");
assert!(matches!(
- parse_cursor(&cursor),
+ parse_cursor(&cursor, &"a".repeat(64)),
Err(radroots_transport::Error::InvalidFetchCursor)
));
}
- let cursor = CursorPosition {
- created_at: 10,
- event_id: "b".repeat(64),
- };
+ let cursor = RelayCursor::new(10, "b".repeat(64)).expect("cursor");
let relay = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).expect("relay");
let older = Candidate {
relay: relay.clone(),
@@ -641,6 +707,41 @@ mod tests {
}
#[test]
+ fn continuation_cursor_is_bound_to_exact_targets_and_selector() {
+ let first = futures::executor::block_on(transport().fetch(request(1))).expect("first page");
+ let NextPage::Cursor(cursor) = first.next_page() else {
+ panic!("cursor expected");
+ };
+ let different_selector = FetchSelector::all().with_kinds(vec![1]).expect("selector");
+ let error = futures::executor::block_on(
+ transport().fetch(
+ request(1)
+ .with_selector(different_selector)
+ .with_cursor(cursor.clone()),
+ ),
+ )
+ .expect_err("scope mismatch");
+ assert_eq!(error, radroots_transport::Error::InvalidFetchCursor);
+
+ let other_targets =
+ TargetSet::new(vec![Target::nostr_relay("wss://one.example").expect("one")])
+ .expect("targets");
+ let error = futures::executor::block_on(
+ transport().fetch(
+ FetchRequest::new(
+ "different-targets",
+ other_targets,
+ FetchBounds::new(1, u64::MAX).expect("bounds"),
+ )
+ .expect("request")
+ .with_cursor(cursor.clone()),
+ ),
+ )
+ .expect_err("target scope mismatch");
+ assert_eq!(error, radroots_transport::Error::InvalidFetchCursor);
+ }
+
+ #[test]
fn live_source_short_circuits_selectors_that_cannot_be_encoded() {
let client = LiveRelaySourceClient::isolated();
let relay = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).expect("relay");
@@ -653,6 +754,7 @@ mod tests {
until_unix_seconds: None,
connect_timeout: Duration::from_millis(1),
timeout: Duration::from_millis(1),
+ max_connections: 1,
}));
assert_eq!(batches.len(), 1);
assert_eq!(batches[0].result, Ok(vec![]));
diff --git a/crates/transport_nostr/src/status.rs b/crates/transport_nostr/src/status.rs
@@ -1,10 +1,12 @@
-//! Stable Nostr relay status and outcome normalization.
+//! Stable per-relay evidence, aggregate status, and outcome normalization.
+use crate::{Config, ReconnectBackoff, RelayEndpoint, RelayProfileKind, RelayUrl};
use radroots_transport::{
SinkStatus, SourceStatus,
capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities},
outcome::{DeliveryOutcome, DeliveryOutcomeKind, FetchTargetState, Retryability},
};
+use std::collections::BTreeMap;
use std::fmt;
use std::sync::Mutex;
@@ -63,94 +65,513 @@ impl fmt::Debug for RedactedDiagnostic {
}
}
+/// Evidence state for one relay capability direction.
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+#[non_exhaustive]
+pub enum RelayEvidenceState {
+ /// The profile does not authorize this capability direction.
+ Unsupported,
+ /// The capability is configured but no successful or failed attempt exists.
+ Unobserved,
+ /// An authorized operation is currently awaiting relay evidence.
+ Connecting,
+ /// The latest accepted observation proved the capability usable.
+ Available,
+ /// The latest accepted observation failed to prove the capability usable.
+ Unavailable,
+}
+
+/// Immutable evidence for one relay capability direction.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RelayCapabilityEvidence {
+ state: RelayEvidenceState,
+ last_attempt_unix_ms: Option<u64>,
+ last_success_unix_ms: Option<u64>,
+ consecutive_failures: u32,
+ next_attempt_unix_ms: Option<u64>,
+ last_failure_retryable: Option<bool>,
+}
+
+impl RelayCapabilityEvidence {
+ /// Returns whether this direction is unsupported, unobserved, or observed.
+ #[must_use]
+ pub const fn state(&self) -> RelayEvidenceState {
+ self.state
+ }
+
+ /// Returns the latest monotonic accepted attempt timestamp.
+ #[must_use]
+ pub const fn last_attempt_unix_ms(&self) -> Option<u64> {
+ self.last_attempt_unix_ms
+ }
+
+ /// Returns the latest successful observation timestamp.
+ #[must_use]
+ pub const fn last_success_unix_ms(&self) -> Option<u64> {
+ self.last_success_unix_ms
+ }
+
+ /// Returns consecutive failures since the latest success.
+ #[must_use]
+ pub const fn consecutive_failures(&self) -> u32 {
+ self.consecutive_failures
+ }
+
+ /// Returns the earliest adapter-permitted reconnect time after failure.
+ #[must_use]
+ pub const fn next_attempt_unix_ms(&self) -> Option<u64> {
+ self.next_attempt_unix_ms
+ }
+
+ /// Returns the retry class of the latest failure, when one exists.
+ #[must_use]
+ pub const fn last_failure_retryable(&self) -> Option<bool> {
+ self.last_failure_retryable
+ }
+}
+
+#[derive(Clone, Debug)]
+struct MutableEvidence {
+ public: RelayCapabilityEvidence,
+}
+
+impl MutableEvidence {
+ const fn unobserved() -> Self {
+ Self {
+ public: RelayCapabilityEvidence {
+ state: RelayEvidenceState::Unobserved,
+ last_attempt_unix_ms: None,
+ last_success_unix_ms: None,
+ consecutive_failures: 0,
+ next_attempt_unix_ms: None,
+ last_failure_retryable: None,
+ },
+ }
+ }
+
+ const fn unsupported() -> Self {
+ Self {
+ public: RelayCapabilityEvidence {
+ state: RelayEvidenceState::Unsupported,
+ last_attempt_unix_ms: None,
+ last_success_unix_ms: None,
+ consecutive_failures: 0,
+ next_attempt_unix_ms: None,
+ last_failure_retryable: None,
+ },
+ }
+ }
+
+ fn begin(&mut self, observed_at_unix_ms: u64) {
+ if matches!(self.public.state, RelayEvidenceState::Unsupported)
+ || self
+ .public
+ .last_attempt_unix_ms
+ .is_some_and(|current| observed_at_unix_ms < current)
+ {
+ return;
+ }
+ self.public.state = RelayEvidenceState::Connecting;
+ self.public.last_attempt_unix_ms = Some(observed_at_unix_ms);
+ }
+
+ fn record(
+ &mut self,
+ succeeded: bool,
+ retryable: bool,
+ observed_at_unix_ms: u64,
+ backoff: ReconnectBackoff,
+ ) {
+ if matches!(self.public.state, RelayEvidenceState::Unsupported)
+ || self
+ .public
+ .last_attempt_unix_ms
+ .is_some_and(|current| observed_at_unix_ms < current)
+ {
+ return;
+ }
+ self.public.last_attempt_unix_ms = Some(observed_at_unix_ms);
+ if succeeded {
+ self.public.state = RelayEvidenceState::Available;
+ self.public.last_success_unix_ms = Some(observed_at_unix_ms);
+ self.public.consecutive_failures = 0;
+ self.public.next_attempt_unix_ms = None;
+ self.public.last_failure_retryable = None;
+ } else {
+ self.public.state = RelayEvidenceState::Unavailable;
+ self.public.consecutive_failures = self.public.consecutive_failures.saturating_add(1);
+ self.public.last_failure_retryable = Some(retryable);
+ self.public.next_attempt_unix_ms = retryable.then(|| {
+ observed_at_unix_ms
+ .saturating_add(backoff.delay_ms(self.public.consecutive_failures))
+ });
+ }
+ }
+
+ fn may_attempt(&self, now_unix_ms: u64) -> bool {
+ !matches!(self.public.state, RelayEvidenceState::Unsupported)
+ && self.public.last_failure_retryable != Some(false)
+ && self
+ .public
+ .next_attempt_unix_ms
+ .is_none_or(|retry_at| now_unix_ms >= retry_at)
+ }
+}
+
+#[derive(Clone, Debug)]
+struct MutableRelayStatus {
+ endpoint: RelayEndpoint,
+ read: MutableEvidence,
+ write: MutableEvidence,
+}
+
+/// Passive typed status for one configured relay.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RelayStatus {
+ endpoint: RelayEndpoint,
+ read: RelayCapabilityEvidence,
+ write: RelayCapabilityEvidence,
+}
+
+impl RelayStatus {
+ /// Returns the canonical endpoint and its declared authority.
+ #[must_use]
+ pub const fn endpoint(&self) -> &RelayEndpoint {
+ &self.endpoint
+ }
+
+ /// Returns independent read evidence.
+ #[must_use]
+ pub const fn read(&self) -> &RelayCapabilityEvidence {
+ &self.read
+ }
+
+ /// Returns independent write evidence.
+ #[must_use]
+ pub const fn write(&self) -> &RelayCapabilityEvidence {
+ &self.write
+ }
+}
+
+/// Passive per-relay and aggregate status for one configured profile.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct RelayStatusReport {
+ profile_kind: RelayProfileKind,
+ state: RelayAggregateState,
+ relays: Vec<RelayStatus>,
+ read_availability: Availability,
+ write_availability: Availability,
+}
+
+/// Aggregate lifecycle derived from current per-relay evidence.
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+#[non_exhaustive]
+pub enum RelayAggregateState {
+ /// A profile is installed but no relay operation has started.
+ Configured,
+ /// At least one authorized relay operation is in flight.
+ Connecting,
+ /// Current evidence proves reads but not publication.
+ ReadOnly,
+ /// Current evidence proves both reads and publication.
+ Writable,
+ /// Some, but not all, authorized relay capabilities have current success.
+ Degraded,
+ /// Every attempted capability is temporarily unavailable and retry-bounded.
+ Offline,
+ /// Every attempted capability failed terminally for the unchanged request.
+ Failed,
+}
+
+impl RelayStatusReport {
+ /// Returns the host profile whose policies govern this report.
+ #[must_use]
+ pub const fn profile_kind(&self) -> RelayProfileKind {
+ self.profile_kind
+ }
+
+ /// Returns the aggregate lifecycle derived only from current evidence.
+ #[must_use]
+ pub const fn state(&self) -> RelayAggregateState {
+ self.state
+ }
+
+ /// Returns one status per configured relay in profile order.
+ #[must_use]
+ pub fn relays(&self) -> &[RelayStatus] {
+ self.relays.as_slice()
+ }
+
+ /// Returns aggregate read availability derived only from observations.
+ #[must_use]
+ pub const fn read_availability(&self) -> Availability {
+ self.read_availability
+ }
+
+ /// Returns aggregate write availability derived only from writable relays.
+ #[must_use]
+ pub const fn write_availability(&self) -> Availability {
+ self.write_availability
+ }
+}
+
#[derive(Clone, Debug)]
struct Snapshot {
- source: Availability,
- sink: Availability,
- source_diagnostic: Option<RedactedDiagnostic>,
- sink_diagnostic: Option<RedactedDiagnostic>,
+ order: Vec<RelayUrl>,
+ relays: BTreeMap<RelayUrl, MutableRelayStatus>,
}
-impl Default for Snapshot {
- fn default() -> Self {
+impl Snapshot {
+ fn new(config: &Config) -> Self {
Self {
- source: Availability::Available,
- sink: Availability::Available,
- source_diagnostic: None,
- sink_diagnostic: None,
+ order: config.relays().to_vec(),
+ relays: config
+ .endpoints()
+ .iter()
+ .map(|endpoint| {
+ (
+ endpoint.url().clone(),
+ MutableRelayStatus {
+ endpoint: endpoint.clone(),
+ read: MutableEvidence::unobserved(),
+ write: if endpoint.access().can_write() {
+ MutableEvidence::unobserved()
+ } else {
+ MutableEvidence::unsupported()
+ },
+ },
+ )
+ })
+ .collect(),
}
}
}
-#[derive(Debug, Default)]
+#[derive(Debug)]
pub(crate) struct StatusTracker {
+ initial: Snapshot,
snapshot: Mutex<Snapshot>,
+ profile_kind: RelayProfileKind,
+ backoff: ReconnectBackoff,
}
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 new(config: &Config) -> Self {
+ let initial = Snapshot::new(config);
+ Self {
+ snapshot: Mutex::new(initial.clone()),
+ initial,
+ profile_kind: config.profile_kind(),
+ backoff: config.reconnect_backoff(),
+ }
+ }
+
+ pub(crate) fn begin_read(&self, relay: &RelayUrl, observed_at_unix_ms: u64) {
+ if let Ok(mut snapshot) = self.snapshot.lock()
+ && let Some(status) = snapshot.relays.get_mut(relay)
+ {
+ status.read.begin(observed_at_unix_ms);
+ }
+ }
+
+ pub(crate) fn begin_write(&self, relay: &RelayUrl, observed_at_unix_ms: u64) {
+ if let Ok(mut snapshot) = self.snapshot.lock()
+ && let Some(status) = snapshot.relays.get_mut(relay)
+ {
+ status.write.begin(observed_at_unix_ms);
+ }
+ }
+
+ pub(crate) fn record_read(
+ &self,
+ relay: &RelayUrl,
+ succeeded: bool,
+ retryable: bool,
+ observed_at_unix_ms: u64,
+ ) {
+ if let Ok(mut snapshot) = self.snapshot.lock()
+ && let Some(status) = snapshot.relays.get_mut(relay)
+ {
+ status
+ .read
+ .record(succeeded, retryable, observed_at_unix_ms, self.backoff);
}
}
- 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);
+ pub(crate) fn record_write(
+ &self,
+ relay: &RelayUrl,
+ succeeded: bool,
+ retryable: bool,
+ observed_at_unix_ms: u64,
+ ) {
+ if let Ok(mut snapshot) = self.snapshot.lock()
+ && let Some(status) = snapshot.relays.get_mut(relay)
+ {
+ status
+ .write
+ .record(succeeded, retryable, observed_at_unix_ms, self.backoff);
}
}
- fn source_availability(&self) -> Availability {
+ pub(crate) fn may_read(&self, relay: &RelayUrl, now_unix_ms: u64) -> bool {
self.snapshot
.lock()
- .map(|snapshot| snapshot.source)
- .unwrap_or(Availability::Unavailable)
+ .ok()
+ .and_then(|snapshot| {
+ snapshot
+ .relays
+ .get(relay)
+ .map(|status| status.read.may_attempt(now_unix_ms))
+ })
+ .unwrap_or(false)
}
- fn sink_availability(&self) -> Availability {
+ pub(crate) fn may_write(&self, relay: &RelayUrl, now_unix_ms: u64) -> bool {
self.snapshot
.lock()
- .map(|snapshot| snapshot.sink)
- .unwrap_or(Availability::Unavailable)
+ .ok()
+ .and_then(|snapshot| {
+ snapshot
+ .relays
+ .get(relay)
+ .map(|status| status.write.may_attempt(now_unix_ms))
+ })
+ .unwrap_or(false)
+ }
+
+ pub(crate) fn report(&self) -> RelayStatusReport {
+ let snapshot = self
+ .snapshot
+ .lock()
+ .map(|snapshot| snapshot.clone())
+ .unwrap_or_else(|_| self.initial.clone());
+ let relays = self
+ .initial
+ .order
+ .iter()
+ .filter_map(|relay| snapshot.relays.get(relay))
+ .map(|status| RelayStatus {
+ endpoint: status.endpoint.clone(),
+ read: status.read.public.clone(),
+ write: status.write.public.clone(),
+ })
+ .collect::<Vec<_>>();
+ RelayStatusReport {
+ profile_kind: self.profile_kind,
+ state: aggregate_state(relays.as_slice()),
+ read_availability: aggregate(relays.iter().map(RelayStatus::read)),
+ write_availability: aggregate(relays.iter().map(RelayStatus::write)),
+ relays,
+ }
+ }
+}
+
+fn aggregate_state(relays: &[RelayStatus]) -> RelayAggregateState {
+ let evidence = relays
+ .iter()
+ .flat_map(|relay| [relay.read(), relay.write()])
+ .filter(|evidence| !matches!(evidence.state(), RelayEvidenceState::Unsupported))
+ .collect::<Vec<_>>();
+ if evidence
+ .iter()
+ .any(|evidence| matches!(evidence.state(), RelayEvidenceState::Connecting))
+ {
+ return RelayAggregateState::Connecting;
+ }
+ if evidence
+ .iter()
+ .all(|evidence| matches!(evidence.state(), RelayEvidenceState::Unobserved))
+ {
+ return RelayAggregateState::Configured;
+ }
+ let read = relays.iter().map(RelayStatus::read).collect::<Vec<_>>();
+ let write = relays
+ .iter()
+ .map(RelayStatus::write)
+ .filter(|evidence| !matches!(evidence.state(), RelayEvidenceState::Unsupported))
+ .collect::<Vec<_>>();
+ let read_available = read
+ .iter()
+ .filter(|evidence| matches!(evidence.state(), RelayEvidenceState::Available))
+ .count();
+ let write_available = write
+ .iter()
+ .filter(|evidence| matches!(evidence.state(), RelayEvidenceState::Available))
+ .count();
+ if read_available == read.len() && !write.is_empty() && write_available == write.len() {
+ RelayAggregateState::Writable
+ } else if read_available == read.len() && write_available == 0 {
+ RelayAggregateState::ReadOnly
+ } else if read_available + write_available > 0 {
+ RelayAggregateState::Degraded
+ } else if evidence
+ .iter()
+ .any(|evidence| evidence.last_failure_retryable() == Some(true))
+ {
+ RelayAggregateState::Offline
+ } else if evidence
+ .iter()
+ .any(|evidence| evidence.last_failure_retryable() == Some(false))
+ {
+ RelayAggregateState::Failed
+ } else {
+ RelayAggregateState::Configured
+ }
+}
+
+fn aggregate<'a>(evidence: impl Iterator<Item = &'a RelayCapabilityEvidence>) -> Availability {
+ let mut supported = 0usize;
+ let mut available = 0usize;
+ for evidence in evidence {
+ if !matches!(evidence.state, RelayEvidenceState::Unsupported) {
+ supported += 1;
+ if matches!(evidence.state, RelayEvidenceState::Available) {
+ available += 1;
+ }
+ }
+ }
+ match (supported, available) {
+ (0, _) | (_, 0) => Availability::Unavailable,
+ (supported, available) if supported == available => Availability::Available,
+ _ => Availability::Degraded,
}
}
-pub(crate) fn source_status(tracker: &StatusTracker, configured: bool) -> SourceStatus {
+pub(crate) fn source_status(tracker: &StatusTracker) -> SourceStatus {
+ let report = tracker.report();
SourceStatus::new(
radroots_transport::TransportId::NOSTR,
- configured,
+ !report.relays().is_empty(),
Maturity::Preview,
- if configured {
- tracker.source_availability()
- } else {
- Availability::Unavailable
- },
+ report.read_availability(),
SourceCapabilities::FETCH,
- if configured {
- "bounded Nostr event source configured"
- } else {
- "Nostr event source is not configured"
+ match report.read_availability() {
+ Availability::Available => "Nostr read capability has current successful evidence",
+ Availability::Degraded => "Nostr read capability has partial successful evidence",
+ Availability::Unavailable => "Nostr read capability has no current successful evidence",
},
)
}
-pub(crate) fn sink_status(tracker: &StatusTracker, configured: bool) -> SinkStatus {
+pub(crate) fn sink_status(tracker: &StatusTracker) -> SinkStatus {
+ let report = tracker.report();
+ let configured = report
+ .relays()
+ .iter()
+ .any(|status| status.endpoint().access().can_write());
SinkStatus::new(
radroots_transport::TransportId::NOSTR,
configured,
Maturity::Preview,
- if configured {
- tracker.sink_availability()
- } else {
- Availability::Unavailable
- },
+ report.write_availability(),
SinkCapabilities::DELIVER,
- if configured {
- "bounded Nostr event delivery configured"
- } else {
- "Nostr event sink is not configured"
+ match report.write_availability() {
+ Availability::Available => "Nostr write capability has current successful evidence",
+ Availability::Degraded => "Nostr write capability has partial successful evidence",
+ Availability::Unavailable => {
+ "Nostr write capability has no current successful evidence"
+ }
},
)
}
@@ -207,14 +628,6 @@ fn classify(upstream: &str) -> RedactedDiagnostic {
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(),
@@ -225,6 +638,16 @@ pub(crate) fn delivery_succeeded(outcome: &DeliveryOutcome) -> bool {
#[cfg(test)]
mod tests {
use super::*;
+ use crate::{ReconnectBackoff, RelayProfile};
+
+ fn tracker() -> (StatusTracker, RelayUrl, RelayUrl) {
+ let config =
+ Config::from_profile(RelayProfile::public(["wss://write.example"]).expect("profile"))
+ .with_reconnect_backoff(ReconnectBackoff::new(10, 40).expect("backoff"));
+ let read_only = config.relays()[0].clone();
+ let writable = config.relays()[1].clone();
+ (StatusTracker::new(&config), read_only, writable)
+ }
#[test]
fn every_upstream_class_maps_to_stable_secret_safe_output() {
@@ -248,45 +671,74 @@ mod tests {
}
#[test]
- fn status_tracks_available_degraded_and_unavailable_without_io() {
- let tracker = StatusTracker::default();
+ fn status_requires_directional_evidence_and_backoff_is_monotonic() {
+ let (tracker, read_only, writable) = tracker();
+ let initial = tracker.report();
+ assert_eq!(initial.read_availability(), Availability::Unavailable);
+ assert_eq!(initial.write_availability(), Availability::Unavailable);
+ assert_eq!(initial.state(), RelayAggregateState::Configured);
assert_eq!(
- sink_status(&tracker, true).availability(),
- Availability::Available
+ initial.relays()[0].write().state(),
+ RelayEvidenceState::Unsupported
);
- tracker.record_sink(1, 1, Some("token=secret timeout"));
+ assert!(!tracker.may_write(&read_only, 100));
+
+ tracker.begin_read(&read_only, 100);
+ assert_eq!(tracker.report().state(), RelayAggregateState::Connecting);
+ tracker.record_read(&read_only, true, false, 100);
+ tracker.record_read(&writable, false, true, 100);
+ tracker.record_write(&writable, false, true, 100);
+ let partial = tracker.report();
+ assert_eq!(partial.state(), RelayAggregateState::Degraded);
+ assert_eq!(partial.read_availability(), Availability::Degraded);
+ assert_eq!(partial.write_availability(), Availability::Unavailable);
+ assert!(!tracker.may_write(&writable, 109));
+ assert!(tracker.may_write(&writable, 110));
+
+ tracker.record_write(&writable, false, true, 110);
assert_eq!(
- sink_status(&tracker, true).availability(),
- Availability::Degraded
+ tracker.report().relays()[1].write().next_attempt_unix_ms(),
+ Some(130)
);
- tracker.record_source(0, 2, Some("token=secret offline"));
+ tracker.record_write(&writable, true, false, 109);
assert_eq!(
- source_status(&tracker, true).availability(),
- Availability::Unavailable
+ tracker.report().relays()[1].write().state(),
+ RelayEvidenceState::Unavailable
);
- assert!(!format!("{tracker:?}").contains("token=secret"));
+ tracker.record_write(&writable, true, false, 130);
+ let available = tracker.report();
+ assert_eq!(available.write_availability(), Availability::Available);
+ assert_eq!(available.relays()[1].write().consecutive_failures(), 0);
assert_eq!(
- sink_status(&tracker, false).availability(),
- Availability::Unavailable
+ available.relays()[1].write().last_success_unix_ms(),
+ Some(130)
);
+ assert!(tracker.may_write(&writable, 130));
assert_eq!(
- source_status(&tracker, false).availability(),
- Availability::Unavailable
+ source_status(&tracker).availability(),
+ Availability::Degraded
);
- tracker.record_sink(0, 0, None);
assert_eq!(
- sink_status(&tracker, true).availability(),
+ sink_status(&tracker).availability(),
Availability::Available
);
- tracker.record_source(2, 0, None);
+
+ tracker.record_write(&writable, false, false, 140);
+ assert!(!tracker.may_write(&writable, u64::MAX));
assert_eq!(
- source_status(&tracker, true).availability(),
- Availability::Available
+ tracker.report().relays()[1]
+ .write()
+ .last_failure_retryable(),
+ Some(false)
);
- assert!(delivery_succeeded(&DeliveryOutcome::accepted()));
- assert!(delivery_succeeded(&DeliveryOutcome::delivered()));
- assert!(!delivery_succeeded(&DeliveryOutcome::rejected()));
+ assert_eq!(
+ sink_status(&tracker).availability(),
+ Availability::Unavailable
+ );
+ }
+ #[test]
+ fn normalized_failures_and_success_helpers_cover_every_state() {
for message in [
"invalid event",
"restricted",
@@ -304,5 +756,74 @@ mod tests {
delivery_failure("malformed event").kind(),
DeliveryOutcomeKind::Rejected
);
+ assert!(delivery_succeeded(&DeliveryOutcome::accepted()));
+ assert!(delivery_succeeded(&DeliveryOutcome::delivered()));
+ assert!(!delivery_succeeded(&DeliveryOutcome::rejected()));
+ }
+
+ #[test]
+ fn all_success_is_available_and_no_writable_relay_is_honestly_unavailable() {
+ let config =
+ Config::from_profile(RelayProfile::public(Vec::<String>::new()).expect("profile"));
+ let tracker = StatusTracker::new(&config);
+ let relay = config.relays()[0].clone();
+ tracker.record_read(&relay, true, false, 1);
+ assert_eq!(
+ source_status(&tracker).availability(),
+ Availability::Available
+ );
+ let sink = sink_status(&tracker);
+ assert!(!sink.is_configured());
+ assert_eq!(sink.availability(), Availability::Unavailable);
+ assert_eq!(tracker.report().state(), RelayAggregateState::ReadOnly);
+ }
+
+ #[test]
+ fn aggregate_states_and_unconfigured_targets_cover_fail_closed_edges() {
+ let (live, read_only, writable) = tracker();
+ let unknown =
+ RelayUrl::parse("wss://unknown.example", crate::RelayUrlPolicy::Public).expect("relay");
+
+ live.begin_read(&unknown, 1);
+ live.begin_write(&unknown, 1);
+ live.record_read(&unknown, true, false, 1);
+ live.record_write(&unknown, true, false, 1);
+ live.begin_write(&read_only, 1);
+ live.record_write(&read_only, true, false, 1);
+ assert_eq!(
+ live.report().relays()[0].write().state(),
+ RelayEvidenceState::Unsupported
+ );
+
+ live.begin_read(&read_only, 2);
+ live.begin_read(&read_only, 1);
+ live.record_read(&read_only, true, false, 2);
+ live.record_read(&writable, true, false, 2);
+ live.record_write(&writable, true, false, 2);
+ let writable_report = live.report();
+ assert_eq!(writable_report.state(), RelayAggregateState::Writable);
+ assert_eq!(writable_report.read_availability(), Availability::Available);
+ assert_eq!(
+ writable_report.write_availability(),
+ Availability::Available
+ );
+
+ let (offline, read_only, writable) = tracker();
+ offline.record_read(&read_only, false, true, 10);
+ offline.record_read(&writable, false, true, 10);
+ offline.record_write(&writable, false, true, 10);
+ assert_eq!(offline.report().state(), RelayAggregateState::Offline);
+
+ let (failed, read_only, writable) = tracker();
+ failed.record_read(&read_only, false, false, 10);
+ failed.record_read(&writable, false, false, 10);
+ failed.record_write(&writable, false, false, 10);
+ assert_eq!(failed.report().state(), RelayAggregateState::Failed);
+
+ assert_eq!(delivery_failure("already have").code(), Some("duplicate"));
+ assert_eq!(
+ fetch_failure("timed out").0,
+ FetchTargetState::FailedRetryable
+ );
}
}
diff --git a/crates/transport_nostr/tests/local_io.rs b/crates/transport_nostr/tests/local_io.rs
@@ -0,0 +1,97 @@
+use futures::{SinkExt, StreamExt};
+use radroots_transport::{
+ EventSource, FetchRequest, TargetSet, capability::Availability, outcome::FetchTargetState,
+ source::FetchBounds,
+};
+use radroots_transport_nostr::{
+ Config, NostrTransport, RelayAggregateState, RelayEvidenceState, RelayProfile,
+};
+use serde_json::Value;
+use std::time::{Duration, SystemTime, UNIX_EPOCH};
+use tokio::net::TcpListener;
+use tokio_tungstenite::{accept_async, tungstenite::Message};
+
+#[tokio::test(flavor = "multi_thread")]
+async fn simulator_profile_proves_read_capability_against_a_real_loopback_socket() {
+ let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind relay");
+ let address = listener.local_addr().expect("relay address");
+ let relay_url = format!("ws://{address}");
+ let server = tokio::spawn(async move {
+ let (stream, _) = listener.accept().await.expect("accept relay client");
+ let mut websocket = accept_async(stream).await.expect("websocket handshake");
+ while let Some(message) = websocket.next().await {
+ let Message::Text(message) = message.expect("client message") else {
+ continue;
+ };
+ let parsed: Value = serde_json::from_str(message.as_str()).expect("Nostr message");
+ let Some(values) = parsed.as_array() else {
+ continue;
+ };
+ let [Value::String(kind), Value::String(subscription), ..] = values.as_slice() else {
+ continue;
+ };
+ if kind == "REQ" {
+ websocket
+ .send(Message::Text(
+ serde_json::to_string(&("EOSE", subscription))
+ .expect("EOSE message")
+ .into(),
+ ))
+ .await
+ .expect("send EOSE");
+ break;
+ }
+ }
+ });
+
+ let profile = RelayProfile::simulator([relay_url.as_str()]).expect("simulator profile");
+ let config = Config::from_profile(profile)
+ .with_timeouts(1_000, 2_000, 500)
+ .expect("timeouts");
+ let targets = TargetSet::new(
+ config
+ .read_relays()
+ .map(|relay| relay.to_target())
+ .collect::<Result<Vec<_>, _>>()
+ .expect("targets"),
+ )
+ .expect("target set");
+ let request = FetchRequest::new(
+ "local-real-io",
+ targets,
+ FetchBounds::new(1, unix_time_ms() + 5_000).expect("bounds"),
+ )
+ .expect("request");
+ let transport = NostrTransport::new(config);
+ let page = tokio::time::timeout(Duration::from_secs(5), transport.fetch(request))
+ .await
+ .expect("bounded fetch")
+ .expect("fetch page");
+
+ assert!(page.events().is_empty());
+ assert_eq!(page.target_outcomes().len(), 1);
+ assert_eq!(
+ page.target_outcomes()[0].state(),
+ FetchTargetState::Complete
+ );
+ let status = transport.relay_status();
+ assert_eq!(status.state(), RelayAggregateState::ReadOnly);
+ assert_eq!(status.read_availability(), Availability::Available);
+ assert_eq!(status.write_availability(), Availability::Unavailable);
+ assert_eq!(
+ status.relays()[0].read().state(),
+ RelayEvidenceState::Available
+ );
+ assert_eq!(
+ status.relays()[0].write().state(),
+ RelayEvidenceState::Unobserved
+ );
+ server.await.expect("relay task");
+}
+
+fn unix_time_ms() -> u64 {
+ SystemTime::now()
+ .duration_since(UNIX_EPOCH)
+ .map(|duration| u64::try_from(duration.as_millis()).unwrap_or(u64::MAX))
+ .unwrap_or_default()
+}
diff --git a/crates/transport_nostr/tests/network_hardening.rs b/crates/transport_nostr/tests/network_hardening.rs
@@ -1,5 +1,5 @@
use radroots_transport::{Target, TargetSet, source::FetchBounds};
-use radroots_transport_nostr::{Config, RelayUrl, RelayUrlPolicy};
+use radroots_transport_nostr::{RelayProfile, RelayUrl, RelayUrlPolicy};
const WORKSPACE_MANIFEST: &str = include_str!("../../../Cargo.toml");
const RELAY_SOURCE: &str = include_str!("../src/relay.rs");
@@ -43,7 +43,7 @@ fn page_and_target_limits_reject_oversized_requests() {
assert!(TargetSet::new(targets).is_err());
let relays = (0..=64).map(|index| format!("wss://r{index}.example.com"));
- assert!(Config::new(RelayUrlPolicy::Public, relays).is_err());
+ assert!(RelayProfile::public(relays).is_err());
}
#[test]
diff --git a/crates/transport_nostr/tests/package_boundary.rs b/crates/transport_nostr/tests/package_boundary.rs
@@ -35,12 +35,14 @@ fn manifest_and_root_match_the_governed_transport_boundary() {
assert_eq!(
private_modules(ROOT),
BTreeSet::from([
- "auth", "client", "error", "relay", "sink", "source", "status"
+ "auth", "client", "cursor", "error", "profile", "relay", "sink", "source", "status"
])
);
for export in [
- "pub use client::{Config, NostrTransport};",
+ "pub use client::{Config, NostrTransport, ReconnectBackoff};",
+ "pub use cursor::RelayCursor;",
"pub use error::Error;",
+ "pub use profile::{",
"pub use relay::{RelayUrl, RelayUrlPolicy};",
] {
assert!(ROOT.contains(export), "crate root is missing `{export}`");
@@ -65,8 +67,8 @@ fn documentation_example_and_reviewed_api_baseline_are_complete() {
assert!(README.contains(required), "README is missing `{required}`");
}
for required in [
- "Config::new(",
- "RelayUrlPolicy::Public",
+ "Config::from_profile(",
+ "RelayProfile::public(",
"NostrTransport::new(config)",
"let source: &dyn EventSource",
"let sink: &dyn EventSink",
@@ -81,6 +83,9 @@ fn documentation_example_and_reviewed_api_baseline_are_complete() {
for required in [
"pub struct radroots_transport_nostr::Config",
"pub struct radroots_transport_nostr::NostrTransport",
+ "pub enum radroots_transport_nostr::RelayAggregateState",
+ "pub struct radroots_transport_nostr::RelayProfile",
+ "pub struct radroots_transport_nostr::RelayStatusReport",
"pub struct radroots_transport_nostr::RelayUrl(_)",
"pub enum radroots_transport_nostr::RelayUrlPolicy",
"pub enum radroots_transport_nostr::Error",
@@ -171,8 +176,10 @@ fn adapter_owns_no_storage_outbox_or_orchestration_surface() {
BTreeSet::from([
"auth.rs".to_owned(),
"client.rs".to_owned(),
+ "cursor.rs".to_owned(),
"error.rs".to_owned(),
"lib.rs".to_owned(),
+ "profile.rs".to_owned(),
"relay.rs".to_owned(),
"sink.rs".to_owned(),
"source.rs".to_owned(),
diff --git a/docs/api/radroots_sdk.txt b/docs/api/radroots_sdk.txt
@@ -49,9 +49,11 @@ pub struct radroots_sdk::client::Client
impl radroots_sdk::client::Client
pub fn radroots_sdk::client::Client::capabilities(&self) -> radroots_sdk::capability::CapabilityReport
pub async fn radroots_sdk::client::Client::close(&self) -> radroots_sdk::error::Result<()>
+pub fn radroots_sdk::client::Client::configure_nostr(&self, radroots_transport_nostr::profile::RelayProfile) -> radroots_sdk::error::Result<()>
pub fn radroots_sdk::client::Client::farm(&self) -> radroots_sdk::error::Result<core::option::Option<radroots_sdk::farm::Operations<'_>>>
pub fn radroots_sdk::client::Client::is_closed(&self) -> bool
pub fn radroots_sdk::client::Client::listing(&self) -> radroots_sdk::error::Result<core::option::Option<radroots_sdk::listing::Operations<'_>>>
+pub fn radroots_sdk::client::Client::nostr_status(&self) -> radroots_sdk::error::Result<core::option::Option<radroots_transport_nostr::status::RelayStatusReport>>
pub fn radroots_sdk::client::Client::signer(&self) -> radroots_sdk::error::Result<core::option::Option<&dyn radroots_signing::signer::Signer>>
pub fn radroots_sdk::client::Client::signing(&self) -> radroots_sdk::error::Result<core::option::Option<radroots_sdk::signing::Operations<'_>>>
pub fn radroots_sdk::client::Client::sink(&self) -> radroots_sdk::error::Result<core::option::Option<&dyn radroots_transport::sink::EventSink>>
@@ -411,6 +413,19 @@ pub const fn radroots_sdk::trade::PrepareRequest::new(radroots_signing::actor::A
pub fn radroots_sdk::trade::prepare(radroots_sdk::trade::PrepareRequest) -> core::result::Result<radroots_sdk::trade::Plan, radroots_sdk::trade::PrepareError>
pub fn radroots_sdk::trade::project(radroots_trade::trade_contract_v1::RadrootsTradeReductionInputV1) -> radroots_trade::trade_contract_v1::RadrootsTradeProjectionV1
pub mod radroots_sdk::transport
+pub use radroots_sdk::transport::DEFAULT_PUBLIC_RELAY
+pub use radroots_sdk::transport::ReconnectBackoff
+pub use radroots_sdk::transport::RelayAccess
+pub use radroots_sdk::transport::RelayAggregateState
+pub use radroots_sdk::transport::RelayCapabilityEvidence
+pub use radroots_sdk::transport::RelayCursor
+pub use radroots_sdk::transport::RelayEndpoint
+pub use radroots_sdk::transport::RelayEvidenceState
+pub use radroots_sdk::transport::RelayProfile
+pub use radroots_sdk::transport::RelayProfileKind
+pub use radroots_sdk::transport::RelayStatus
+pub use radroots_sdk::transport::RelayStatusReport
+pub use radroots_sdk::transport::RelayUrl
pub use radroots_sdk::transport::RelayUrlPolicy
#[non_exhaustive] pub enum radroots_sdk::transport::DaemonAuth
pub radroots_sdk::transport::DaemonAuth::BearerToken(alloc::string::String)
@@ -444,9 +459,13 @@ pub fn radroots_sdk::transport::DaemonError::fmt(&self, &mut core::fmt::Formatte
pub struct radroots_sdk::transport::NostrSlot
impl radroots_sdk::transport::NostrSlot
pub fn radroots_sdk::transport::NostrSlot::clear(&self)
-pub fn radroots_sdk::transport::NostrSlot::configure<I, S>(&self, I) -> radroots_sdk::error::Result<()> where I: core::iter::traits::collect::IntoIterator<Item = S>, S: core::convert::AsRef<str>
-pub fn radroots_sdk::transport::NostrSlot::new(radroots_transport_nostr::relay::RelayUrlPolicy) -> Self
-pub fn radroots_sdk::transport::NostrSlot::targets(&self) -> core::option::Option<radroots_transport::target::TargetSet>
+pub fn radroots_sdk::transport::NostrSlot::configure(&self, radroots_transport_nostr::profile::RelayProfile) -> radroots_sdk::error::Result<()>
+pub fn radroots_sdk::transport::NostrSlot::new() -> Self
+pub fn radroots_sdk::transport::NostrSlot::read_targets(&self) -> core::option::Option<radroots_transport::target::TargetSet>
+pub fn radroots_sdk::transport::NostrSlot::relay_status(&self) -> core::option::Option<radroots_transport_nostr::status::RelayStatusReport>
+pub fn radroots_sdk::transport::NostrSlot::write_targets(&self) -> core::option::Option<radroots_transport::target::TargetSet>
+impl core::default::Default for radroots_sdk::transport::NostrSlot
+pub fn radroots_sdk::transport::NostrSlot::default() -> Self
impl core::fmt::Debug for radroots_sdk::transport::NostrSlot
pub fn radroots_sdk::transport::NostrSlot::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
impl radroots_transport::sink::EventSink for radroots_sdk::transport::NostrSlot
@@ -471,9 +490,11 @@ pub struct radroots_sdk::Client
impl radroots_sdk::client::Client
pub fn radroots_sdk::client::Client::capabilities(&self) -> radroots_sdk::capability::CapabilityReport
pub async fn radroots_sdk::client::Client::close(&self) -> radroots_sdk::error::Result<()>
+pub fn radroots_sdk::client::Client::configure_nostr(&self, radroots_transport_nostr::profile::RelayProfile) -> radroots_sdk::error::Result<()>
pub fn radroots_sdk::client::Client::farm(&self) -> radroots_sdk::error::Result<core::option::Option<radroots_sdk::farm::Operations<'_>>>
pub fn radroots_sdk::client::Client::is_closed(&self) -> bool
pub fn radroots_sdk::client::Client::listing(&self) -> radroots_sdk::error::Result<core::option::Option<radroots_sdk::listing::Operations<'_>>>
+pub fn radroots_sdk::client::Client::nostr_status(&self) -> radroots_sdk::error::Result<core::option::Option<radroots_transport_nostr::status::RelayStatusReport>>
pub fn radroots_sdk::client::Client::signer(&self) -> radroots_sdk::error::Result<core::option::Option<&dyn radroots_signing::signer::Signer>>
pub fn radroots_sdk::client::Client::signing(&self) -> radroots_sdk::error::Result<core::option::Option<radroots_sdk::signing::Operations<'_>>>
pub fn radroots_sdk::client::Client::sink(&self) -> radroots_sdk::error::Result<core::option::Option<&dyn radroots_transport::sink::EventSink>>
diff --git a/docs/api/radroots_transport_nostr.txt b/docs/api/radroots_transport_nostr.txt
@@ -17,6 +17,10 @@ pub radroots_transport_nostr::Error::EmptyResolution::url: alloc::string::String
pub radroots_transport_nostr::Error::InvalidAuthChallenge
pub radroots_transport_nostr::Error::InvalidConnectionLimit
pub radroots_transport_nostr::Error::InvalidConnectionLimit::value: usize
+pub radroots_transport_nostr::Error::InvalidReconnectBackoff
+pub radroots_transport_nostr::Error::InvalidReconnectBackoff::initial_delay_ms: u64
+pub radroots_transport_nostr::Error::InvalidReconnectBackoff::max_delay_ms: u64
+pub radroots_transport_nostr::Error::InvalidRelayCursor
pub radroots_transport_nostr::Error::InvalidRelayUrl
pub radroots_transport_nostr::Error::InvalidRelayUrl::reason: alloc::string::String
pub radroots_transport_nostr::Error::InvalidRelayUrl::url: alloc::string::String
@@ -40,6 +44,30 @@ pub radroots_transport_nostr::Error::UnexpectedTransport::actual: alloc::string:
impl core::error::Error for radroots_transport_nostr::Error
impl core::fmt::Display for radroots_transport_nostr::Error
pub fn radroots_transport_nostr::Error::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+#[non_exhaustive] pub enum radroots_transport_nostr::RelayAccess
+pub radroots_transport_nostr::RelayAccess::ReadOnly
+pub radroots_transport_nostr::RelayAccess::ReadWrite
+impl radroots_transport_nostr::RelayAccess
+pub const fn radroots_transport_nostr::RelayAccess::can_read(self) -> bool
+pub const fn radroots_transport_nostr::RelayAccess::can_write(self) -> bool
+#[non_exhaustive] pub enum radroots_transport_nostr::RelayAggregateState
+pub radroots_transport_nostr::RelayAggregateState::Configured
+pub radroots_transport_nostr::RelayAggregateState::Connecting
+pub radroots_transport_nostr::RelayAggregateState::Degraded
+pub radroots_transport_nostr::RelayAggregateState::Failed
+pub radroots_transport_nostr::RelayAggregateState::Offline
+pub radroots_transport_nostr::RelayAggregateState::ReadOnly
+pub radroots_transport_nostr::RelayAggregateState::Writable
+#[non_exhaustive] pub enum radroots_transport_nostr::RelayEvidenceState
+pub radroots_transport_nostr::RelayEvidenceState::Available
+pub radroots_transport_nostr::RelayEvidenceState::Connecting
+pub radroots_transport_nostr::RelayEvidenceState::Unavailable
+pub radroots_transport_nostr::RelayEvidenceState::Unobserved
+pub radroots_transport_nostr::RelayEvidenceState::Unsupported
+#[non_exhaustive] pub enum radroots_transport_nostr::RelayProfileKind
+pub radroots_transport_nostr::RelayProfileKind::Device
+pub radroots_transport_nostr::RelayProfileKind::Public
+pub radroots_transport_nostr::RelayProfileKind::Simulator
#[non_exhaustive] pub enum radroots_transport_nostr::RelayUrlPolicy
pub radroots_transport_nostr::RelayUrlPolicy::Local
pub radroots_transport_nostr::RelayUrlPolicy::PrivateNetwork
@@ -47,14 +75,19 @@ pub radroots_transport_nostr::RelayUrlPolicy::Public
pub struct radroots_transport_nostr::Config
impl radroots_transport_nostr::Config
pub const fn radroots_transport_nostr::Config::connect_timeout_ms(&self) -> u64
+pub fn radroots_transport_nostr::Config::endpoints(&self) -> &[radroots_transport_nostr::RelayEndpoint]
+pub fn radroots_transport_nostr::Config::from_profile(radroots_transport_nostr::RelayProfile) -> Self
pub const fn radroots_transport_nostr::Config::max_connections(&self) -> usize
-pub fn radroots_transport_nostr::Config::new<I, S>(radroots_transport_nostr::RelayUrlPolicy, I) -> core::result::Result<Self, radroots_transport_nostr::Error> where I: core::iter::traits::collect::IntoIterator<Item = S>, S: core::convert::AsRef<str>
-pub const fn radroots_transport_nostr::Config::relay_url_policy(&self) -> radroots_transport_nostr::RelayUrlPolicy
+pub const fn radroots_transport_nostr::Config::profile_kind(&self) -> radroots_transport_nostr::RelayProfileKind
+pub fn radroots_transport_nostr::Config::read_relays(&self) -> impl core::iter::traits::iterator::Iterator<Item = &radroots_transport_nostr::RelayUrl>
+pub const fn radroots_transport_nostr::Config::reconnect_backoff(&self) -> radroots_transport_nostr::ReconnectBackoff
pub fn radroots_transport_nostr::Config::relays(&self) -> &[radroots_transport_nostr::RelayUrl]
pub const fn radroots_transport_nostr::Config::request_timeout_ms(&self) -> u64
pub const fn radroots_transport_nostr::Config::status_timeout_ms(&self) -> u64
pub fn radroots_transport_nostr::Config::with_max_connections(self, usize) -> core::result::Result<Self, radroots_transport_nostr::Error>
+pub const fn radroots_transport_nostr::Config::with_reconnect_backoff(self, radroots_transport_nostr::ReconnectBackoff) -> Self
pub fn radroots_transport_nostr::Config::with_timeouts(self, u64, u64, u64) -> core::result::Result<Self, radroots_transport_nostr::Error>
+pub fn radroots_transport_nostr::Config::write_relays(&self) -> impl core::iter::traits::iterator::Iterator<Item = &radroots_transport_nostr::RelayUrl>
pub struct radroots_transport_nostr::NostrTransport
impl radroots_transport_nostr::NostrTransport
pub fn radroots_transport_nostr::NostrTransport::begin_authentication(&self, &radroots_transport_nostr::RelayUrl, impl core::convert::AsRef<str>, u64, u64) -> core::result::Result<alloc::string::String, radroots_transport_nostr::Error>
@@ -63,14 +96,60 @@ pub fn radroots_transport_nostr::NostrTransport::reject_authentication(&self, &r
impl radroots_transport_nostr::NostrTransport
pub const fn radroots_transport_nostr::NostrTransport::config(&self) -> &radroots_transport_nostr::Config
pub fn radroots_transport_nostr::NostrTransport::new(radroots_transport_nostr::Config) -> Self
+pub fn radroots_transport_nostr::NostrTransport::relay_status(&self) -> radroots_transport_nostr::RelayStatusReport
impl core::fmt::Debug for radroots_transport_nostr::NostrTransport
pub fn radroots_transport_nostr::NostrTransport::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
impl radroots_transport::sink::EventSink for radroots_transport_nostr::NostrTransport
-pub fn radroots_transport_nostr::NostrTransport::deliver(&self, radroots_transport::sink::DeliveryRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::error::Error>>
+pub fn radroots_transport_nostr::NostrTransport::deliver(&self, radroots_transport::sink::DeliveryRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::sink::DeliveryReceipt, radroots_transport::sink::SinkFailure>>
pub fn radroots_transport_nostr::NostrTransport::status(&self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::status::SinkStatus, radroots_transport::error::Error>>
impl radroots_transport::source::EventSource for radroots_transport_nostr::NostrTransport
pub fn radroots_transport_nostr::NostrTransport::fetch(&self, radroots_transport::source::FetchRequest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::source::FetchPage, radroots_transport::error::Error>>
pub fn radroots_transport_nostr::NostrTransport::status(&self) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_transport::status::SourceStatus, radroots_transport::error::Error>>
+pub struct radroots_transport_nostr::ReconnectBackoff
+impl radroots_transport_nostr::ReconnectBackoff
+pub const fn radroots_transport_nostr::ReconnectBackoff::initial_delay_ms(self) -> u64
+pub const fn radroots_transport_nostr::ReconnectBackoff::max_delay_ms(self) -> u64
+pub fn radroots_transport_nostr::ReconnectBackoff::new(u64, u64) -> core::result::Result<Self, radroots_transport_nostr::Error>
+impl core::default::Default for radroots_transport_nostr::ReconnectBackoff
+pub fn radroots_transport_nostr::ReconnectBackoff::default() -> Self
+pub struct radroots_transport_nostr::RelayCapabilityEvidence
+impl radroots_transport_nostr::RelayCapabilityEvidence
+pub const fn radroots_transport_nostr::RelayCapabilityEvidence::consecutive_failures(&self) -> u32
+pub const fn radroots_transport_nostr::RelayCapabilityEvidence::last_attempt_unix_ms(&self) -> core::option::Option<u64>
+pub const fn radroots_transport_nostr::RelayCapabilityEvidence::last_failure_retryable(&self) -> core::option::Option<bool>
+pub const fn radroots_transport_nostr::RelayCapabilityEvidence::last_success_unix_ms(&self) -> core::option::Option<u64>
+pub const fn radroots_transport_nostr::RelayCapabilityEvidence::next_attempt_unix_ms(&self) -> core::option::Option<u64>
+pub const fn radroots_transport_nostr::RelayCapabilityEvidence::state(&self) -> radroots_transport_nostr::RelayEvidenceState
+pub struct radroots_transport_nostr::RelayCursor
+impl radroots_transport_nostr::RelayCursor
+pub const fn radroots_transport_nostr::RelayCursor::created_at_unix_s(&self) -> u64
+pub fn radroots_transport_nostr::RelayCursor::event_id(&self) -> &str
+pub fn radroots_transport_nostr::RelayCursor::new(u64, impl core::convert::Into<alloc::string::String>) -> core::result::Result<Self, radroots_transport_nostr::Error>
+pub fn radroots_transport_nostr::RelayCursor::precedes(&self, u64, &str) -> bool
+pub struct radroots_transport_nostr::RelayEndpoint
+impl radroots_transport_nostr::RelayEndpoint
+pub const fn radroots_transport_nostr::RelayEndpoint::access(&self) -> radroots_transport_nostr::RelayAccess
+pub const fn radroots_transport_nostr::RelayEndpoint::policy(&self) -> radroots_transport_nostr::RelayUrlPolicy
+pub const fn radroots_transport_nostr::RelayEndpoint::url(&self) -> &radroots_transport_nostr::RelayUrl
+pub struct radroots_transport_nostr::RelayProfile
+impl radroots_transport_nostr::RelayProfile
+pub fn radroots_transport_nostr::RelayProfile::device<I, S>(I) -> core::result::Result<Self, radroots_transport_nostr::Error> where I: core::iter::traits::collect::IntoIterator<Item = S>, S: core::convert::AsRef<str>
+pub fn radroots_transport_nostr::RelayProfile::endpoints(&self) -> &[radroots_transport_nostr::RelayEndpoint]
+pub const fn radroots_transport_nostr::RelayProfile::kind(&self) -> radroots_transport_nostr::RelayProfileKind
+pub fn radroots_transport_nostr::RelayProfile::public<I, S>(I) -> core::result::Result<Self, radroots_transport_nostr::Error> where I: core::iter::traits::collect::IntoIterator<Item = S>, S: core::convert::AsRef<str>
+pub fn radroots_transport_nostr::RelayProfile::simulator<I, S>(I) -> core::result::Result<Self, radroots_transport_nostr::Error> where I: core::iter::traits::collect::IntoIterator<Item = S>, S: core::convert::AsRef<str>
+pub struct radroots_transport_nostr::RelayStatus
+impl radroots_transport_nostr::RelayStatus
+pub const fn radroots_transport_nostr::RelayStatus::endpoint(&self) -> &radroots_transport_nostr::RelayEndpoint
+pub const fn radroots_transport_nostr::RelayStatus::read(&self) -> &radroots_transport_nostr::RelayCapabilityEvidence
+pub const fn radroots_transport_nostr::RelayStatus::write(&self) -> &radroots_transport_nostr::RelayCapabilityEvidence
+pub struct radroots_transport_nostr::RelayStatusReport
+impl radroots_transport_nostr::RelayStatusReport
+pub const fn radroots_transport_nostr::RelayStatusReport::profile_kind(&self) -> radroots_transport_nostr::RelayProfileKind
+pub const fn radroots_transport_nostr::RelayStatusReport::read_availability(&self) -> radroots_transport::capability::Availability
+pub fn radroots_transport_nostr::RelayStatusReport::relays(&self) -> &[radroots_transport_nostr::RelayStatus]
+pub const fn radroots_transport_nostr::RelayStatusReport::state(&self) -> radroots_transport_nostr::RelayAggregateState
+pub const fn radroots_transport_nostr::RelayStatusReport::write_availability(&self) -> radroots_transport::capability::Availability
pub struct radroots_transport_nostr::RelayUrl(_)
impl radroots_transport_nostr::RelayUrl
pub fn radroots_transport_nostr::RelayUrl::as_str(&self) -> &str
@@ -80,3 +159,4 @@ pub fn radroots_transport_nostr::RelayUrl::to_target(&self) -> core::result::Res
pub fn radroots_transport_nostr::RelayUrl::validate_resolved_addresses(&self, radroots_transport_nostr::RelayUrlPolicy, impl core::iter::traits::collect::IntoIterator<Item = core::net::ip_addr::IpAddr>) -> core::result::Result<(), radroots_transport_nostr::Error>
impl core::fmt::Display for radroots_transport_nostr::RelayUrl
pub fn radroots_transport_nostr::RelayUrl::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result
+pub const radroots_transport_nostr::DEFAULT_PUBLIC_RELAY: &str