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 }