lib

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

client.rs (15728B)


      1 //! Concrete Nostr transport composition.
      2 
      3 use crate::{Error, RelayEndpoint, RelayProfile, RelayProfileKind, RelayStatusReport, RelayUrl};
      4 use core::fmt;
      5 use std::sync::{
      6     Arc,
      7     atomic::{AtomicU64, Ordering},
      8 };
      9 
     10 /// Maximum relay targets accepted by one transport instance.
     11 pub(crate) const MAX_RELAYS: usize = 64;
     12 const MAX_TIMEOUT_MS: u64 = 120_000;
     13 const MAX_CONNECTIONS: usize = 64;
     14 const MAX_RECONNECT_DELAY_MS: u64 = 15 * 60 * 1_000;
     15 
     16 /// Deterministic exponential reconnect policy applied independently per relay
     17 /// and capability direction.
     18 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
     19 pub struct ReconnectBackoff {
     20     initial_delay_ms: u64,
     21     max_delay_ms: u64,
     22 }
     23 
     24 impl ReconnectBackoff {
     25     /// Creates a bounded reconnect policy.
     26     pub fn new(initial_delay_ms: u64, max_delay_ms: u64) -> Result<Self, Error> {
     27         if initial_delay_ms == 0
     28             || max_delay_ms < initial_delay_ms
     29             || max_delay_ms > MAX_RECONNECT_DELAY_MS
     30         {
     31             return Err(Error::InvalidReconnectBackoff {
     32                 initial_delay_ms,
     33                 max_delay_ms,
     34             });
     35         }
     36         Ok(Self {
     37             initial_delay_ms,
     38             max_delay_ms,
     39         })
     40     }
     41 
     42     /// Returns the first delay after a retryable failure.
     43     #[must_use]
     44     pub const fn initial_delay_ms(self) -> u64 {
     45         self.initial_delay_ms
     46     }
     47 
     48     /// Returns the upper bound for any computed reconnect delay.
     49     #[must_use]
     50     pub const fn max_delay_ms(self) -> u64 {
     51         self.max_delay_ms
     52     }
     53 
     54     pub(crate) fn delay_ms(self, consecutive_failures: u32) -> u64 {
     55         let exponent = consecutive_failures.saturating_sub(1).min(63);
     56         self.initial_delay_ms
     57             .saturating_mul(1_u64 << exponent)
     58             .min(self.max_delay_ms)
     59     }
     60 }
     61 
     62 impl Default for ReconnectBackoff {
     63     fn default() -> Self {
     64         Self {
     65             initial_delay_ms: 1_000,
     66             max_delay_ms: 60_000,
     67         }
     68     }
     69 }
     70 
     71 /// Validated configuration for a concrete Nostr transport.
     72 #[derive(Clone, Debug, Eq, PartialEq)]
     73 pub struct Config {
     74     profile_kind: RelayProfileKind,
     75     endpoints: Vec<RelayEndpoint>,
     76     relays: Vec<RelayUrl>,
     77     connect_timeout_ms: u64,
     78     request_timeout_ms: u64,
     79     status_timeout_ms: u64,
     80     max_connections: usize,
     81     reconnect_backoff: ReconnectBackoff,
     82 }
     83 
     84 impl Config {
     85     /// Builds inert transport configuration from one validated host profile.
     86     #[must_use]
     87     pub fn from_profile(profile: RelayProfile) -> Self {
     88         let relays: Vec<_> = profile
     89             .endpoints()
     90             .iter()
     91             .map(|endpoint| endpoint.url().clone())
     92             .collect();
     93         let max_connections = 8.min(relays.len());
     94         Self {
     95             profile_kind: profile.kind(),
     96             endpoints: profile.endpoints().to_vec(),
     97             relays,
     98             connect_timeout_ms: 10_000,
     99             request_timeout_ms: 30_000,
    100             status_timeout_ms: 5_000,
    101             max_connections,
    102             reconnect_backoff: ReconnectBackoff::default(),
    103         }
    104     }
    105 
    106     /// Sets explicit bounded connection, request, and status timeouts.
    107     pub fn with_timeouts(
    108         mut self,
    109         connect_timeout_ms: u64,
    110         request_timeout_ms: u64,
    111         status_timeout_ms: u64,
    112     ) -> Result<Self, Error> {
    113         validate_timeout("connect", connect_timeout_ms)?;
    114         validate_timeout("request", request_timeout_ms)?;
    115         validate_timeout("status", status_timeout_ms)?;
    116         self.connect_timeout_ms = connect_timeout_ms;
    117         self.request_timeout_ms = request_timeout_ms;
    118         self.status_timeout_ms = status_timeout_ms;
    119         Ok(self)
    120     }
    121 
    122     /// Sets the maximum simultaneous relay connections for one operation.
    123     pub fn with_max_connections(mut self, value: usize) -> Result<Self, Error> {
    124         if value == 0 || value > MAX_CONNECTIONS || value > self.relays.len() {
    125             return Err(Error::InvalidConnectionLimit { value });
    126         }
    127         self.max_connections = value;
    128         Ok(self)
    129     }
    130 
    131     /// Sets the deterministic per-relay reconnect policy.
    132     #[must_use]
    133     pub const fn with_reconnect_backoff(mut self, value: ReconnectBackoff) -> Self {
    134         self.reconnect_backoff = value;
    135         self
    136     }
    137 
    138     /// Returns the selected host profile kind.
    139     #[must_use]
    140     pub const fn profile_kind(&self) -> RelayProfileKind {
    141         self.profile_kind
    142     }
    143 
    144     /// Returns configured endpoints with directional access and network policy.
    145     #[must_use]
    146     pub fn endpoints(&self) -> &[RelayEndpoint] {
    147         self.endpoints.as_slice()
    148     }
    149 
    150     /// Returns relays in caller-specified order.
    151     pub fn relays(&self) -> &[RelayUrl] {
    152         self.relays.as_slice()
    153     }
    154 
    155     /// Returns relays authorized for reads in deterministic profile order.
    156     pub fn read_relays(&self) -> impl Iterator<Item = &RelayUrl> {
    157         self.endpoints
    158             .iter()
    159             .filter_map(|endpoint| endpoint.access().can_read().then_some(endpoint.url()))
    160     }
    161 
    162     /// Returns relays authorized for publication in deterministic profile order.
    163     pub fn write_relays(&self) -> impl Iterator<Item = &RelayUrl> {
    164         self.endpoints
    165             .iter()
    166             .filter_map(|endpoint| endpoint.access().can_write().then_some(endpoint.url()))
    167     }
    168 
    169     pub(crate) fn endpoint_for_target(
    170         &self,
    171         target: &radroots_transport::Target,
    172     ) -> Option<&RelayEndpoint> {
    173         (*target.kind() == radroots_transport::TransportId::NOSTR)
    174             .then(|| target.uri().as_str())
    175             .and_then(|url| {
    176                 self.endpoints
    177                     .iter()
    178                     .find(|endpoint| endpoint.url().as_str() == url)
    179             })
    180     }
    181 
    182     /// Returns the connection establishment deadline in milliseconds.
    183     pub const fn connect_timeout_ms(&self) -> u64 {
    184         self.connect_timeout_ms
    185     }
    186 
    187     /// Returns the bounded request deadline in milliseconds.
    188     pub const fn request_timeout_ms(&self) -> u64 {
    189         self.request_timeout_ms
    190     }
    191 
    192     /// Returns the passive status observation deadline in milliseconds.
    193     pub const fn status_timeout_ms(&self) -> u64 {
    194         self.status_timeout_ms
    195     }
    196 
    197     /// Returns the maximum simultaneous relay connections.
    198     pub const fn max_connections(&self) -> usize {
    199         self.max_connections
    200     }
    201 
    202     /// Returns the per-relay reconnect policy.
    203     #[must_use]
    204     pub const fn reconnect_backoff(&self) -> ReconnectBackoff {
    205         self.reconnect_backoff
    206     }
    207 }
    208 
    209 fn validate_timeout(field: &'static str, value_ms: u64) -> Result<(), Error> {
    210     if value_ms == 0 || value_ms > MAX_TIMEOUT_MS {
    211         return Err(Error::InvalidTimeout { field, value_ms });
    212     }
    213     Ok(())
    214 }
    215 
    216 /// Concrete Nostr implementation of the transport source and sink SPIs.
    217 #[derive(Clone)]
    218 pub struct NostrTransport {
    219     config: Config,
    220     pub(crate) client: Arc<dyn crate::sink::RelayClient>,
    221     pub(crate) source_client: Arc<dyn crate::source::RelaySourceClient>,
    222     pub(crate) subscription_client: Arc<dyn crate::subscription::RelaySubscriptionClient>,
    223     pub(crate) auth: Arc<crate::auth::AuthFlow>,
    224     pub(crate) status: Arc<crate::status::StatusTracker>,
    225     subscription_sequence: Arc<AtomicU64>,
    226 }
    227 
    228 impl NostrTransport {
    229     /// Creates an inert transport from validated explicit configuration.
    230     pub fn new(config: Config) -> Self {
    231         let connector = crate::relay::HardenedWebsocketTransport::new(config.endpoints());
    232         let writers = connector.writers.clone();
    233         let ingress = connector.ingress.clone();
    234         let client = nostr_sdk::Client::builder()
    235             .websocket_transport(connector)
    236             .build();
    237         client.automatic_authentication(false);
    238         let status = Arc::new(crate::status::StatusTracker::new(&config));
    239         Self {
    240             config,
    241             client: Arc::new(crate::sink::LiveRelayClient::new(client.clone(), writers)),
    242             source_client: Arc::new(crate::source::LiveRelaySourceClient::new(
    243                 client.clone(),
    244                 ingress,
    245             )),
    246             subscription_client: Arc::new(crate::subscription::LiveRelaySubscriptionClient::new(
    247                 client.clone(),
    248             )),
    249             auth: Arc::new(crate::auth::AuthFlow::new(Arc::new(
    250                 crate::auth::LiveAuthClient::new(client),
    251             ))),
    252             status,
    253             subscription_sequence: Arc::new(AtomicU64::new(0)),
    254         }
    255     }
    256 
    257     /// Returns the transport configuration.
    258     pub const fn config(&self) -> &Config {
    259         &self.config
    260     }
    261 
    262     /// Returns passive per-relay and aggregate evidence without network I/O.
    263     #[must_use]
    264     pub fn relay_status(&self) -> RelayStatusReport {
    265         self.status.report()
    266     }
    267 
    268     #[cfg(test)]
    269     pub(crate) fn with_client(config: Config, client: Arc<dyn crate::sink::RelayClient>) -> Self {
    270         let status = Arc::new(crate::status::StatusTracker::new(&config));
    271         Self {
    272             config,
    273             client,
    274             source_client: Arc::new(crate::source::LiveRelaySourceClient::isolated()),
    275             subscription_client: Arc::new(
    276                 crate::subscription::LiveRelaySubscriptionClient::isolated(),
    277             ),
    278             auth: Arc::new(crate::auth::AuthFlow::isolated()),
    279             status,
    280             subscription_sequence: Arc::new(AtomicU64::new(0)),
    281         }
    282     }
    283 
    284     #[cfg(test)]
    285     pub(crate) fn with_source_client(
    286         config: Config,
    287         source_client: Arc<dyn crate::source::RelaySourceClient>,
    288     ) -> Self {
    289         let status = Arc::new(crate::status::StatusTracker::new(&config));
    290         Self {
    291             config,
    292             client: Arc::new(crate::sink::LiveRelayClient::isolated()),
    293             source_client,
    294             subscription_client: Arc::new(
    295                 crate::subscription::LiveRelaySubscriptionClient::isolated(),
    296             ),
    297             auth: Arc::new(crate::auth::AuthFlow::isolated()),
    298             status,
    299             subscription_sequence: Arc::new(AtomicU64::new(0)),
    300         }
    301     }
    302 
    303     #[cfg(test)]
    304     pub(crate) fn with_subscription_client(
    305         config: Config,
    306         subscription_client: Arc<dyn crate::subscription::RelaySubscriptionClient>,
    307     ) -> Self {
    308         let status = Arc::new(crate::status::StatusTracker::new(&config));
    309         Self {
    310             config,
    311             client: Arc::new(crate::sink::LiveRelayClient::isolated()),
    312             source_client: Arc::new(crate::source::LiveRelaySourceClient::isolated()),
    313             subscription_client,
    314             auth: Arc::new(crate::auth::AuthFlow::isolated()),
    315             status,
    316             subscription_sequence: Arc::new(AtomicU64::new(0)),
    317         }
    318     }
    319 
    320     pub(crate) fn next_subscription_sequence(&self) -> Result<u64, radroots_transport::Error> {
    321         self.subscription_sequence
    322             .try_update(Ordering::SeqCst, Ordering::SeqCst, |current| {
    323                 current.checked_add(1)
    324             })
    325             .map(|previous| previous + 1)
    326             .map_err(|_| radroots_transport::Error::SubscriptionUnavailable)
    327     }
    328 }
    329 
    330 impl fmt::Debug for NostrTransport {
    331     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    332         formatter
    333             .debug_struct("NostrTransport")
    334             .field("config", &self.config)
    335             .finish_non_exhaustive()
    336     }
    337 }
    338 
    339 #[cfg(test)]
    340 mod tests {
    341     use super::*;
    342 
    343     #[test]
    344     fn config_rejects_empty_duplicate_and_excessive_relay_sets() {
    345         assert!(
    346             crate::profile::test_profile(
    347                 RelayProfileKind::Simulator,
    348                 crate::RelayUrlPolicy::Local,
    349                 Vec::<String>::new(),
    350             )
    351             .is_err()
    352         );
    353         assert!(
    354             crate::profile::test_profile(
    355                 RelayProfileKind::Public,
    356                 crate::RelayUrlPolicy::Public,
    357                 ["wss://relay.example.com", "wss://RELAY.EXAMPLE.COM:443/"],
    358             )
    359             .is_err()
    360         );
    361         let relays = (0..=MAX_RELAYS).map(|index| format!("wss://r{index}.example.com"));
    362         assert!(
    363             crate::profile::test_profile(
    364                 RelayProfileKind::Public,
    365                 crate::RelayUrlPolicy::Public,
    366                 relays,
    367             )
    368             .is_err()
    369         );
    370     }
    371 
    372     #[test]
    373     fn config_rejects_unbounded_limits() {
    374         let config = Config::from_profile(
    375             crate::profile::test_profile(
    376                 RelayProfileKind::Public,
    377                 crate::RelayUrlPolicy::Public,
    378                 ["wss://relay.example.com"],
    379             )
    380             .expect("profile"),
    381         );
    382         assert!(config.clone().with_timeouts(0, 1, 1).is_err());
    383         assert!(config.clone().with_timeouts(1, 120_001, 1).is_err());
    384         assert!(config.with_max_connections(3).is_err());
    385     }
    386 
    387     #[test]
    388     fn valid_configuration_accessors_and_transport_debug_are_complete() {
    389         let config = Config::from_profile(
    390             crate::profile::test_profile(
    391                 RelayProfileKind::Public,
    392                 crate::RelayUrlPolicy::Public,
    393                 ["wss://one.example", "wss://two.example"],
    394             )
    395             .expect("profile"),
    396         )
    397         .with_timeouts(1, 2, 3)
    398         .expect("timeouts")
    399         .with_max_connections(2)
    400         .expect("connections");
    401         assert_eq!(config.relays().len(), 2);
    402         assert_eq!(config.read_relays().count(), 2);
    403         assert_eq!(config.write_relays().count(), 2);
    404         assert_eq!(config.profile_kind(), RelayProfileKind::Public);
    405         assert_eq!(config.connect_timeout_ms(), 1);
    406         assert_eq!(config.request_timeout_ms(), 2);
    407         assert_eq!(config.status_timeout_ms(), 3);
    408         assert_eq!(config.max_connections(), 2);
    409         assert_eq!(config.reconnect_backoff(), ReconnectBackoff::default());
    410         assert!(ReconnectBackoff::new(0, 1).is_err());
    411         assert!(ReconnectBackoff::new(2, 1).is_err());
    412         assert!(ReconnectBackoff::new(1, MAX_RECONNECT_DELAY_MS + 1).is_err());
    413         let backoff = ReconnectBackoff::new(2, 5).expect("backoff");
    414         assert_eq!(backoff.delay_ms(0), 2);
    415         assert_eq!(backoff.delay_ms(1), 2);
    416         assert_eq!(backoff.delay_ms(2), 4);
    417         assert_eq!(backoff.delay_ms(3), 5);
    418         assert!(config.clone().with_timeouts(120_001, 1, 1).is_err());
    419         assert!(config.clone().with_timeouts(1, 1, 0).is_err());
    420         assert!(config.clone().with_max_connections(0).is_err());
    421         assert!(
    422             config
    423                 .clone()
    424                 .with_max_connections(MAX_CONNECTIONS + 1)
    425                 .is_err()
    426         );
    427 
    428         let transport = NostrTransport::new(config.clone());
    429         assert_eq!(transport.config(), &config);
    430         let debug = format!("{transport:?}");
    431         assert!(debug.contains("NostrTransport"));
    432         assert!(!debug.contains("client"));
    433     }
    434 
    435     #[test]
    436     fn subscription_sequence_is_monotonic_and_fails_closed_at_overflow() {
    437         let config = Config::from_profile(
    438             crate::profile::test_profile(
    439                 RelayProfileKind::Public,
    440                 crate::RelayUrlPolicy::Public,
    441                 ["wss://relay.example.com"],
    442             )
    443             .expect("profile"),
    444         );
    445         let transport = NostrTransport::new(config);
    446         assert_eq!(transport.next_subscription_sequence(), Ok(1));
    447         assert_eq!(transport.next_subscription_sequence(), Ok(2));
    448         transport
    449             .subscription_sequence
    450             .store(u64::MAX, Ordering::SeqCst);
    451         assert_eq!(
    452             transport.next_subscription_sequence(),
    453             Err(radroots_transport::Error::SubscriptionUnavailable)
    454         );
    455     }
    456 }