status.rs (33558B)
1 //! Stable per-relay evidence, aggregate status, and outcome normalization. 2 3 use crate::{Config, ReconnectBackoff, RelayEndpoint, RelayProfileKind, RelayUrl}; 4 use radroots_transport::{ 5 SinkStatus, SourceStatus, 6 capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities}, 7 outcome::{DeliveryOutcome, DeliveryOutcomeKind, FetchTargetState, Retryability}, 8 }; 9 use std::collections::BTreeMap; 10 use std::fmt; 11 use std::sync::Mutex; 12 13 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 14 enum FailureClass { 15 Duplicate, 16 Rejected, 17 AuthRequired, 18 Quota, 19 RateLimited, 20 Timeout, 21 Connection, 22 Malformed, 23 Unknown, 24 } 25 26 impl FailureClass { 27 const fn code(self) -> &'static str { 28 match self { 29 Self::Duplicate => "duplicate", 30 Self::Rejected => "rejected", 31 Self::AuthRequired => "auth_required", 32 Self::Quota => "quota_exceeded", 33 Self::RateLimited => "rate_limited", 34 Self::Timeout => "timeout", 35 Self::Connection => "connection_failed", 36 Self::Malformed => "malformed_event", 37 Self::Unknown => "relay_failure", 38 } 39 } 40 41 const fn message(self) -> &'static str { 42 match self { 43 Self::Duplicate => "relay already has the event", 44 Self::Rejected => "relay rejected the event", 45 Self::AuthRequired => "relay authentication is required", 46 Self::Quota => "relay quota was exhausted", 47 Self::RateLimited => "relay rate limit was reached", 48 Self::Timeout => "relay operation timed out", 49 Self::Connection => "relay connection failed", 50 Self::Malformed => "relay returned a malformed event", 51 Self::Unknown => "relay operation failed", 52 } 53 } 54 } 55 56 #[derive(Clone)] 57 struct RedactedDiagnostic { 58 class: FailureClass, 59 } 60 61 impl fmt::Debug for RedactedDiagnostic { 62 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 63 formatter 64 .debug_struct("RedactedDiagnostic") 65 .field("class", &self.class) 66 .field("upstream", &"[redacted]") 67 .finish() 68 } 69 } 70 71 /// Evidence state for one relay capability direction. 72 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 73 #[non_exhaustive] 74 pub enum RelayEvidenceState { 75 /// The profile does not authorize this capability direction. 76 Unsupported, 77 /// The capability is configured but no successful or failed attempt exists. 78 Unobserved, 79 /// An authorized operation is currently awaiting relay evidence. 80 Connecting, 81 /// The latest accepted observation proved the capability usable. 82 Available, 83 /// The latest accepted observation failed to prove the capability usable. 84 Unavailable, 85 } 86 87 /// Immutable evidence for one relay capability direction. 88 #[derive(Clone, Debug, Eq, PartialEq)] 89 pub struct RelayCapabilityEvidence { 90 state: RelayEvidenceState, 91 last_attempt_unix_ms: Option<u64>, 92 last_success_unix_ms: Option<u64>, 93 consecutive_failures: u32, 94 next_attempt_unix_ms: Option<u64>, 95 last_failure_retryable: Option<bool>, 96 } 97 98 impl RelayCapabilityEvidence { 99 /// Returns whether this direction is unsupported, unobserved, or observed. 100 #[must_use] 101 pub const fn state(&self) -> RelayEvidenceState { 102 self.state 103 } 104 105 /// Returns the latest monotonic accepted attempt timestamp. 106 #[must_use] 107 pub const fn last_attempt_unix_ms(&self) -> Option<u64> { 108 self.last_attempt_unix_ms 109 } 110 111 /// Returns the latest successful observation timestamp. 112 #[must_use] 113 pub const fn last_success_unix_ms(&self) -> Option<u64> { 114 self.last_success_unix_ms 115 } 116 117 /// Returns consecutive failures since the latest success. 118 #[must_use] 119 pub const fn consecutive_failures(&self) -> u32 { 120 self.consecutive_failures 121 } 122 123 /// Returns the earliest adapter-permitted reconnect time after failure. 124 #[must_use] 125 pub const fn next_attempt_unix_ms(&self) -> Option<u64> { 126 self.next_attempt_unix_ms 127 } 128 129 /// Returns the retry class of the latest failure, when one exists. 130 #[must_use] 131 pub const fn last_failure_retryable(&self) -> Option<bool> { 132 self.last_failure_retryable 133 } 134 } 135 136 #[derive(Clone, Debug)] 137 struct MutableEvidence { 138 public: RelayCapabilityEvidence, 139 } 140 141 impl MutableEvidence { 142 const fn unobserved() -> Self { 143 Self { 144 public: RelayCapabilityEvidence { 145 state: RelayEvidenceState::Unobserved, 146 last_attempt_unix_ms: None, 147 last_success_unix_ms: None, 148 consecutive_failures: 0, 149 next_attempt_unix_ms: None, 150 last_failure_retryable: None, 151 }, 152 } 153 } 154 155 const fn unsupported() -> Self { 156 Self { 157 public: RelayCapabilityEvidence { 158 state: RelayEvidenceState::Unsupported, 159 last_attempt_unix_ms: None, 160 last_success_unix_ms: None, 161 consecutive_failures: 0, 162 next_attempt_unix_ms: None, 163 last_failure_retryable: None, 164 }, 165 } 166 } 167 168 fn begin(&mut self, observed_at_unix_ms: u64) { 169 if matches!(self.public.state, RelayEvidenceState::Unsupported) 170 || self 171 .public 172 .last_attempt_unix_ms 173 .is_some_and(|current| observed_at_unix_ms < current) 174 { 175 return; 176 } 177 self.public.state = RelayEvidenceState::Connecting; 178 self.public.last_attempt_unix_ms = Some(observed_at_unix_ms); 179 } 180 181 fn record( 182 &mut self, 183 succeeded: bool, 184 retryable: bool, 185 observed_at_unix_ms: u64, 186 backoff: ReconnectBackoff, 187 ) { 188 if matches!(self.public.state, RelayEvidenceState::Unsupported) 189 || self 190 .public 191 .last_attempt_unix_ms 192 .is_some_and(|current| observed_at_unix_ms < current) 193 { 194 return; 195 } 196 self.public.last_attempt_unix_ms = Some(observed_at_unix_ms); 197 if succeeded { 198 self.public.state = RelayEvidenceState::Available; 199 self.public.last_success_unix_ms = Some(observed_at_unix_ms); 200 self.public.consecutive_failures = 0; 201 self.public.next_attempt_unix_ms = None; 202 self.public.last_failure_retryable = None; 203 } else { 204 self.public.state = RelayEvidenceState::Unavailable; 205 self.public.consecutive_failures = self.public.consecutive_failures.saturating_add(1); 206 self.public.last_failure_retryable = Some(retryable); 207 self.public.next_attempt_unix_ms = retryable.then(|| { 208 observed_at_unix_ms 209 .saturating_add(backoff.delay_ms(self.public.consecutive_failures)) 210 }); 211 } 212 } 213 214 fn may_attempt(&self, now_unix_ms: u64) -> bool { 215 !matches!(self.public.state, RelayEvidenceState::Unsupported) 216 && self.public.last_failure_retryable != Some(false) 217 && self 218 .public 219 .next_attempt_unix_ms 220 .is_none_or(|retry_at| now_unix_ms >= retry_at) 221 } 222 } 223 224 #[derive(Clone, Debug)] 225 struct MutableRelayStatus { 226 endpoint: RelayEndpoint, 227 read: MutableEvidence, 228 write: MutableEvidence, 229 } 230 231 /// Passive typed status for one configured relay. 232 #[derive(Clone, Debug, Eq, PartialEq)] 233 pub struct RelayStatus { 234 endpoint: RelayEndpoint, 235 read: RelayCapabilityEvidence, 236 write: RelayCapabilityEvidence, 237 } 238 239 impl RelayStatus { 240 /// Returns the canonical endpoint and its declared authority. 241 #[must_use] 242 pub const fn endpoint(&self) -> &RelayEndpoint { 243 &self.endpoint 244 } 245 246 /// Returns independent read evidence. 247 #[must_use] 248 pub const fn read(&self) -> &RelayCapabilityEvidence { 249 &self.read 250 } 251 252 /// Returns independent write evidence. 253 #[must_use] 254 pub const fn write(&self) -> &RelayCapabilityEvidence { 255 &self.write 256 } 257 } 258 259 /// Passive per-relay and aggregate status for one configured profile. 260 #[derive(Clone, Debug, Eq, PartialEq)] 261 pub struct RelayStatusReport { 262 profile_kind: RelayProfileKind, 263 state: RelayAggregateState, 264 relays: Vec<RelayStatus>, 265 read_availability: Availability, 266 write_availability: Availability, 267 } 268 269 /// Aggregate lifecycle derived from current per-relay evidence. 270 #[derive(Clone, Copy, Debug, Eq, PartialEq)] 271 #[non_exhaustive] 272 pub enum RelayAggregateState { 273 /// A profile is installed but no relay operation has started. 274 Configured, 275 /// At least one authorized relay operation is in flight. 276 Connecting, 277 /// Current evidence proves reads but not publication. 278 ReadOnly, 279 /// Current evidence proves both reads and publication. 280 Writable, 281 /// Some, but not all, authorized relay capabilities have current success. 282 Degraded, 283 /// Every attempted capability is temporarily unavailable and retry-bounded. 284 Offline, 285 /// Every attempted capability failed terminally for the unchanged request. 286 Failed, 287 } 288 289 impl RelayStatusReport { 290 /// Returns the host profile whose policies govern this report. 291 #[must_use] 292 pub const fn profile_kind(&self) -> RelayProfileKind { 293 self.profile_kind 294 } 295 296 /// Returns the aggregate lifecycle derived only from current evidence. 297 #[must_use] 298 pub const fn state(&self) -> RelayAggregateState { 299 self.state 300 } 301 302 /// Returns one status per configured relay in profile order. 303 #[must_use] 304 pub fn relays(&self) -> &[RelayStatus] { 305 self.relays.as_slice() 306 } 307 308 /// Returns aggregate read availability derived only from observations. 309 #[must_use] 310 pub const fn read_availability(&self) -> Availability { 311 self.read_availability 312 } 313 314 /// Returns aggregate write availability derived only from writable relays. 315 #[must_use] 316 pub const fn write_availability(&self) -> Availability { 317 self.write_availability 318 } 319 } 320 321 #[derive(Clone, Debug)] 322 struct Snapshot { 323 order: Vec<RelayUrl>, 324 relays: BTreeMap<RelayUrl, MutableRelayStatus>, 325 } 326 327 impl Snapshot { 328 fn new(config: &Config) -> Self { 329 Self { 330 order: config.relays().to_vec(), 331 relays: config 332 .endpoints() 333 .iter() 334 .map(|endpoint| { 335 ( 336 endpoint.url().clone(), 337 MutableRelayStatus { 338 endpoint: endpoint.clone(), 339 read: MutableEvidence::unobserved(), 340 write: if endpoint.access().can_write() { 341 MutableEvidence::unobserved() 342 } else { 343 MutableEvidence::unsupported() 344 }, 345 }, 346 ) 347 }) 348 .collect(), 349 } 350 } 351 } 352 353 #[derive(Debug)] 354 pub(crate) struct StatusTracker { 355 initial: Snapshot, 356 snapshot: Mutex<Snapshot>, 357 profile_kind: RelayProfileKind, 358 backoff: ReconnectBackoff, 359 } 360 361 impl StatusTracker { 362 pub(crate) fn new(config: &Config) -> Self { 363 let initial = Snapshot::new(config); 364 Self { 365 snapshot: Mutex::new(initial.clone()), 366 initial, 367 profile_kind: config.profile_kind(), 368 backoff: config.reconnect_backoff(), 369 } 370 } 371 372 pub(crate) fn begin_read(&self, relay: &RelayUrl, observed_at_unix_ms: u64) { 373 if let Ok(mut snapshot) = self.snapshot.lock() 374 && let Some(status) = snapshot.relays.get_mut(relay) 375 { 376 status.read.begin(observed_at_unix_ms); 377 } 378 } 379 380 pub(crate) fn begin_write(&self, relay: &RelayUrl, observed_at_unix_ms: u64) { 381 if let Ok(mut snapshot) = self.snapshot.lock() 382 && let Some(status) = snapshot.relays.get_mut(relay) 383 { 384 status.write.begin(observed_at_unix_ms); 385 } 386 } 387 388 pub(crate) fn record_read( 389 &self, 390 relay: &RelayUrl, 391 succeeded: bool, 392 retryable: bool, 393 observed_at_unix_ms: u64, 394 ) { 395 if let Ok(mut snapshot) = self.snapshot.lock() 396 && let Some(status) = snapshot.relays.get_mut(relay) 397 { 398 status 399 .read 400 .record(succeeded, retryable, observed_at_unix_ms, self.backoff); 401 } 402 } 403 404 pub(crate) fn record_write( 405 &self, 406 relay: &RelayUrl, 407 succeeded: bool, 408 retryable: bool, 409 observed_at_unix_ms: u64, 410 ) { 411 if let Ok(mut snapshot) = self.snapshot.lock() 412 && let Some(status) = snapshot.relays.get_mut(relay) 413 { 414 status 415 .write 416 .record(succeeded, retryable, observed_at_unix_ms, self.backoff); 417 } 418 } 419 420 pub(crate) fn may_read(&self, relay: &RelayUrl, now_unix_ms: u64) -> bool { 421 self.snapshot 422 .lock() 423 .ok() 424 .and_then(|snapshot| { 425 snapshot 426 .relays 427 .get(relay) 428 .map(|status| status.read.may_attempt(now_unix_ms)) 429 }) 430 .unwrap_or(false) 431 } 432 433 pub(crate) fn may_write(&self, relay: &RelayUrl, now_unix_ms: u64) -> bool { 434 self.snapshot 435 .lock() 436 .ok() 437 .and_then(|snapshot| { 438 snapshot 439 .relays 440 .get(relay) 441 .map(|status| status.write.may_attempt(now_unix_ms)) 442 }) 443 .unwrap_or(false) 444 } 445 446 pub(crate) fn report(&self) -> RelayStatusReport { 447 let snapshot = self 448 .snapshot 449 .lock() 450 .map(|snapshot| snapshot.clone()) 451 .unwrap_or_else(|_| self.initial.clone()); 452 let relays = self 453 .initial 454 .order 455 .iter() 456 .filter_map(|relay| snapshot.relays.get(relay)) 457 .map(|status| RelayStatus { 458 endpoint: status.endpoint.clone(), 459 read: status.read.public.clone(), 460 write: status.write.public.clone(), 461 }) 462 .collect::<Vec<_>>(); 463 RelayStatusReport { 464 profile_kind: self.profile_kind, 465 state: aggregate_state(relays.as_slice()), 466 read_availability: aggregate(relays.iter().map(RelayStatus::read)), 467 write_availability: aggregate(relays.iter().map(RelayStatus::write)), 468 relays, 469 } 470 } 471 } 472 473 fn aggregate_state(relays: &[RelayStatus]) -> RelayAggregateState { 474 let evidence = relays 475 .iter() 476 .flat_map(|relay| [relay.read(), relay.write()]) 477 .filter(|evidence| !matches!(evidence.state(), RelayEvidenceState::Unsupported)) 478 .collect::<Vec<_>>(); 479 if evidence 480 .iter() 481 .any(|evidence| matches!(evidence.state(), RelayEvidenceState::Connecting)) 482 { 483 return RelayAggregateState::Connecting; 484 } 485 if evidence 486 .iter() 487 .all(|evidence| matches!(evidence.state(), RelayEvidenceState::Unobserved)) 488 { 489 return RelayAggregateState::Configured; 490 } 491 let read = relays.iter().map(RelayStatus::read).collect::<Vec<_>>(); 492 let write = relays 493 .iter() 494 .map(RelayStatus::write) 495 .filter(|evidence| !matches!(evidence.state(), RelayEvidenceState::Unsupported)) 496 .collect::<Vec<_>>(); 497 let read_available = read 498 .iter() 499 .filter(|evidence| matches!(evidence.state(), RelayEvidenceState::Available)) 500 .count(); 501 let write_available = write 502 .iter() 503 .filter(|evidence| matches!(evidence.state(), RelayEvidenceState::Available)) 504 .count(); 505 if read_available == read.len() && !write.is_empty() && write_available == write.len() { 506 RelayAggregateState::Writable 507 } else if read_available == read.len() && write_available == 0 { 508 RelayAggregateState::ReadOnly 509 } else if read_available + write_available > 0 { 510 RelayAggregateState::Degraded 511 } else if evidence 512 .iter() 513 .any(|evidence| evidence.last_failure_retryable() == Some(true)) 514 { 515 RelayAggregateState::Offline 516 } else if evidence 517 .iter() 518 .any(|evidence| evidence.last_failure_retryable() == Some(false)) 519 { 520 RelayAggregateState::Failed 521 } else { 522 RelayAggregateState::Configured 523 } 524 } 525 526 fn aggregate<'a>(evidence: impl Iterator<Item = &'a RelayCapabilityEvidence>) -> Availability { 527 let mut supported = 0usize; 528 let mut available = 0usize; 529 for evidence in evidence { 530 if !matches!(evidence.state, RelayEvidenceState::Unsupported) { 531 supported += 1; 532 if matches!(evidence.state, RelayEvidenceState::Available) { 533 available += 1; 534 } 535 } 536 } 537 match (supported, available) { 538 (0, _) | (_, 0) => Availability::Unavailable, 539 (supported, available) if supported == available => Availability::Available, 540 _ => Availability::Degraded, 541 } 542 } 543 544 pub(crate) fn source_status(tracker: &StatusTracker) -> SourceStatus { 545 let report = tracker.report(); 546 SourceStatus::new( 547 radroots_transport::TransportId::NOSTR, 548 !report.relays().is_empty(), 549 Maturity::Preview, 550 report.read_availability(), 551 SourceCapabilities::FETCH, 552 match report.read_availability() { 553 Availability::Available => "Nostr read capability has current successful evidence", 554 Availability::Degraded => "Nostr read capability has partial successful evidence", 555 Availability::Unavailable => "Nostr read capability has no current successful evidence", 556 }, 557 ) 558 } 559 560 pub(crate) fn sink_status(tracker: &StatusTracker) -> SinkStatus { 561 let report = tracker.report(); 562 let configured = report 563 .relays() 564 .iter() 565 .any(|status| status.endpoint().access().can_write()); 566 SinkStatus::new( 567 radroots_transport::TransportId::NOSTR, 568 configured, 569 Maturity::Preview, 570 report.write_availability(), 571 SinkCapabilities::DELIVER, 572 match report.write_availability() { 573 Availability::Available => "Nostr write capability has current successful evidence", 574 Availability::Degraded => "Nostr write capability has partial successful evidence", 575 Availability::Unavailable => { 576 "Nostr write capability has no current successful evidence" 577 } 578 }, 579 ) 580 } 581 582 pub(crate) fn delivery_failure(upstream: &str) -> DeliveryOutcome { 583 let class = classify(upstream).class; 584 let outcome = match class { 585 FailureClass::Duplicate => DeliveryOutcome::accepted(), 586 FailureClass::Rejected | FailureClass::Malformed | FailureClass::Quota => { 587 DeliveryOutcome::rejected() 588 } 589 FailureClass::AuthRequired => DeliveryOutcome::failed(Retryability::Retryable) 590 .expect("retryable authentication outcome"), 591 FailureClass::RateLimited 592 | FailureClass::Timeout 593 | FailureClass::Connection 594 | FailureClass::Unknown => DeliveryOutcome::unavailable(), 595 }; 596 outcome 597 .with_detail(class.code(), class.message()) 598 .expect("static normalized relay outcome") 599 } 600 601 pub(crate) fn fetch_failure(upstream: &str) -> (FetchTargetState, &'static str) { 602 let class = classify(upstream).class; 603 let state = match class { 604 FailureClass::Rejected | FailureClass::Malformed | FailureClass::Quota => { 605 FetchTargetState::FailedTerminal 606 } 607 _ => FetchTargetState::FailedRetryable, 608 }; 609 (state, class.message()) 610 } 611 612 fn classify(upstream: &str) -> RedactedDiagnostic { 613 let message = upstream.to_ascii_lowercase(); 614 let class = if message.contains("duplicate") || message.contains("already have") { 615 FailureClass::Duplicate 616 } else if message.contains("auth") { 617 FailureClass::AuthRequired 618 } else if message.contains("quota") { 619 FailureClass::Quota 620 } else if message.contains("blocked") 621 || message.contains("restricted") 622 || message.contains("invalid") 623 || message.contains("reject") 624 { 625 FailureClass::Rejected 626 } else if message.contains("rate") { 627 FailureClass::RateLimited 628 } else if message.contains("timeout") || message.contains("timed out") { 629 FailureClass::Timeout 630 } else if message.contains("connect") || message.contains("offline") { 631 FailureClass::Connection 632 } else if message.contains("malformed") || message.contains("decode") { 633 FailureClass::Malformed 634 } else { 635 FailureClass::Unknown 636 }; 637 RedactedDiagnostic { class } 638 } 639 640 pub(crate) fn delivery_succeeded(outcome: &DeliveryOutcome) -> bool { 641 matches!( 642 outcome.kind(), 643 DeliveryOutcomeKind::Accepted | DeliveryOutcomeKind::Delivered 644 ) 645 } 646 647 #[cfg(test)] 648 mod tests { 649 use super::*; 650 use crate::ReconnectBackoff; 651 652 fn tracker() -> (StatusTracker, RelayUrl, RelayUrl) { 653 let config = Config::from_profile( 654 crate::profile::test_profile( 655 crate::RelayProfileKind::Public, 656 crate::RelayUrlPolicy::Public, 657 ["wss://read.example", "wss://write.example"], 658 ) 659 .expect("profile"), 660 ) 661 .with_reconnect_backoff(ReconnectBackoff::new(10, 40).expect("backoff")); 662 let readable = config.relays()[0].clone(); 663 let writable = config.relays()[1].clone(); 664 (StatusTracker::new(&config), readable, writable) 665 } 666 667 #[test] 668 fn every_upstream_class_maps_to_stable_secret_safe_output() { 669 let secret = "token=very-secret-value"; 670 let cases = [ 671 ("duplicate: already have", "duplicate"), 672 ("blocked by policy", "rejected"), 673 ("auth required", "auth_required"), 674 ("quota exceeded", "quota_exceeded"), 675 ("blocked: account QUOTA exhausted", "quota_exceeded"), 676 ("rate limited", "rate_limited"), 677 ("connection timeout", "timeout"), 678 ("connection offline", "connection_failed"), 679 ("unknown failure", "relay_failure"), 680 ]; 681 for (message, code) in cases { 682 let outcome = delivery_failure(format!("{message} {secret}").as_str()); 683 assert_eq!(outcome.code(), Some(code)); 684 assert!(!outcome.message().expect("message").contains(secret)); 685 let diagnostic = classify(format!("{message} {secret}").as_str()); 686 assert!(!format!("{diagnostic:?}").contains(secret)); 687 } 688 } 689 690 #[test] 691 fn quota_refusal_requires_action_without_reclassifying_other_failures() { 692 for text in ["quota exceeded", "blocked: account QUOTA exhausted"] { 693 let outcome = delivery_failure(text); 694 assert_eq!(outcome.code(), Some("quota_exceeded")); 695 assert_eq!(outcome.kind(), DeliveryOutcomeKind::Rejected); 696 assert_eq!(outcome.retryability(), Retryability::Terminal); 697 assert_eq!( 698 fetch_failure(text), 699 ( 700 FetchTargetState::FailedTerminal, 701 "relay quota was exhausted" 702 ) 703 ); 704 } 705 assert_eq!( 706 delivery_failure("rate limited").code(), 707 Some("rate_limited") 708 ); 709 assert!(delivery_failure("rate limited").is_retryable()); 710 assert_eq!( 711 delivery_failure("auth required").code(), 712 Some("auth_required") 713 ); 714 assert_eq!( 715 delivery_failure("malformed event").code(), 716 Some("malformed_event") 717 ); 718 assert_eq!( 719 delivery_failure("unknown failure").code(), 720 Some("relay_failure") 721 ); 722 assert!(delivery_failure("unknown failure").is_retryable()); 723 } 724 725 #[test] 726 fn status_requires_directional_evidence_and_backoff_is_monotonic() { 727 let (tracker, canonical, writable) = tracker(); 728 let initial = tracker.report(); 729 assert_eq!(initial.read_availability(), Availability::Unavailable); 730 assert_eq!(initial.write_availability(), Availability::Unavailable); 731 assert_eq!(initial.state(), RelayAggregateState::Configured); 732 assert_eq!( 733 initial.relays()[0].write().state(), 734 RelayEvidenceState::Unobserved 735 ); 736 assert!(tracker.may_write(&canonical, 100)); 737 738 tracker.begin_read(&canonical, 100); 739 assert_eq!(tracker.report().state(), RelayAggregateState::Connecting); 740 tracker.record_read(&canonical, true, false, 100); 741 tracker.record_read(&writable, false, true, 100); 742 tracker.record_write(&writable, false, true, 100); 743 let partial = tracker.report(); 744 assert_eq!(partial.state(), RelayAggregateState::Degraded); 745 assert_eq!(partial.read_availability(), Availability::Degraded); 746 assert_eq!(partial.write_availability(), Availability::Unavailable); 747 assert!(!tracker.may_write(&writable, 109)); 748 assert!(tracker.may_write(&writable, 110)); 749 750 tracker.record_write(&writable, false, true, 110); 751 assert_eq!( 752 tracker.report().relays()[1].write().next_attempt_unix_ms(), 753 Some(130) 754 ); 755 tracker.record_write(&writable, true, false, 109); 756 assert_eq!( 757 tracker.report().relays()[1].write().state(), 758 RelayEvidenceState::Unavailable 759 ); 760 tracker.record_write(&writable, true, false, 130); 761 let partially_available = tracker.report(); 762 assert_eq!( 763 partially_available.write_availability(), 764 Availability::Degraded 765 ); 766 assert_eq!( 767 partially_available.relays()[1] 768 .write() 769 .consecutive_failures(), 770 0 771 ); 772 assert_eq!( 773 partially_available.relays()[1] 774 .write() 775 .last_success_unix_ms(), 776 Some(130) 777 ); 778 assert!(tracker.may_write(&writable, 130)); 779 tracker.record_write(&canonical, true, false, 130); 780 assert_eq!( 781 tracker.report().write_availability(), 782 Availability::Available 783 ); 784 assert_eq!( 785 source_status(&tracker).availability(), 786 Availability::Degraded 787 ); 788 assert_eq!( 789 sink_status(&tracker).availability(), 790 Availability::Available 791 ); 792 793 tracker.record_write(&writable, false, false, 140); 794 tracker.record_write(&canonical, false, false, 140); 795 assert!(!tracker.may_write(&writable, u64::MAX)); 796 assert!(!tracker.may_write(&canonical, u64::MAX)); 797 assert_eq!( 798 tracker.report().relays()[1] 799 .write() 800 .last_failure_retryable(), 801 Some(false) 802 ); 803 assert_eq!( 804 sink_status(&tracker).availability(), 805 Availability::Unavailable 806 ); 807 } 808 809 #[test] 810 fn normalized_failures_and_success_helpers_cover_every_state() { 811 for message in [ 812 "invalid event", 813 "restricted", 814 "rejected", 815 "malformed event", 816 "decode failed", 817 ] { 818 assert_eq!(fetch_failure(message).0, FetchTargetState::FailedTerminal); 819 } 820 assert_eq!( 821 fetch_failure("offline").0, 822 FetchTargetState::FailedRetryable 823 ); 824 assert_eq!( 825 delivery_failure("malformed event").kind(), 826 DeliveryOutcomeKind::Rejected 827 ); 828 assert!(delivery_succeeded(&DeliveryOutcome::accepted())); 829 assert!(delivery_succeeded(&DeliveryOutcome::delivered())); 830 assert!(!delivery_succeeded(&DeliveryOutcome::rejected())); 831 } 832 833 #[test] 834 fn canonical_relay_read_and_write_success_is_available() { 835 let config = Config::from_profile( 836 crate::profile::test_profile( 837 crate::RelayProfileKind::Public, 838 crate::RelayUrlPolicy::Public, 839 ["wss://relay.example"], 840 ) 841 .expect("profile"), 842 ); 843 let tracker = StatusTracker::new(&config); 844 let relay = config.relays()[0].clone(); 845 tracker.record_read(&relay, true, false, 1); 846 tracker.record_write(&relay, true, false, 1); 847 assert_eq!( 848 source_status(&tracker).availability(), 849 Availability::Available 850 ); 851 let sink = sink_status(&tracker); 852 assert!(sink.is_configured()); 853 assert_eq!(sink.availability(), Availability::Available); 854 assert_eq!(tracker.report().state(), RelayAggregateState::Writable); 855 } 856 857 #[test] 858 fn aggregate_states_and_unconfigured_targets_cover_fail_closed_edges() { 859 let (live, canonical, writable) = tracker(); 860 let unknown = 861 RelayUrl::parse("wss://unknown.example", crate::RelayUrlPolicy::Public).expect("relay"); 862 863 live.begin_read(&unknown, 1); 864 live.begin_write(&unknown, 1); 865 live.record_read(&unknown, true, false, 1); 866 live.record_write(&unknown, true, false, 1); 867 live.begin_write(&canonical, 1); 868 live.record_write(&canonical, true, false, 1); 869 assert_eq!( 870 live.report().relays()[0].write().state(), 871 RelayEvidenceState::Available 872 ); 873 874 live.begin_read(&canonical, 2); 875 live.begin_read(&canonical, 1); 876 live.record_read(&canonical, true, false, 2); 877 live.record_read(&writable, true, false, 2); 878 live.record_write(&writable, true, false, 2); 879 let writable_report = live.report(); 880 assert_eq!(writable_report.state(), RelayAggregateState::Writable); 881 assert_eq!(writable_report.read_availability(), Availability::Available); 882 assert_eq!( 883 writable_report.write_availability(), 884 Availability::Available 885 ); 886 887 let (offline, canonical, writable) = tracker(); 888 offline.record_read(&canonical, false, true, 10); 889 offline.record_read(&writable, false, true, 10); 890 offline.record_write(&canonical, false, true, 10); 891 offline.record_write(&writable, false, true, 10); 892 assert_eq!(offline.report().state(), RelayAggregateState::Offline); 893 894 let (failed, canonical, writable) = tracker(); 895 failed.record_read(&canonical, false, false, 10); 896 failed.record_read(&writable, false, false, 10); 897 failed.record_write(&canonical, false, false, 10); 898 failed.record_write(&writable, false, false, 10); 899 assert_eq!(failed.report().state(), RelayAggregateState::Failed); 900 901 assert_eq!(delivery_failure("already have").code(), Some("duplicate")); 902 assert_eq!( 903 fetch_failure("timed out").0, 904 FetchTargetState::FailedRetryable 905 ); 906 } 907 908 #[test] 909 fn capability_evidence_rejects_unsupported_and_stale_observations() { 910 let backoff = ReconnectBackoff::new(10, 40).expect("backoff"); 911 912 let mut unsupported = MutableEvidence::unsupported(); 913 unsupported.begin(10); 914 unsupported.record(true, false, 10, backoff); 915 assert_eq!(unsupported.public.state(), RelayEvidenceState::Unsupported); 916 assert!(!unsupported.may_attempt(u64::MAX)); 917 918 let mut evidence = MutableEvidence::unobserved(); 919 evidence.begin(10); 920 evidence.begin(9); 921 assert_eq!(evidence.public.last_attempt_unix_ms(), Some(10)); 922 923 evidence.record(false, true, 10, backoff); 924 evidence.record(true, false, 9, backoff); 925 assert_eq!(evidence.public.state(), RelayEvidenceState::Unavailable); 926 assert!(!evidence.may_attempt(19)); 927 assert!(evidence.may_attempt(20)); 928 929 let (read_only, canonical, writable) = tracker(); 930 read_only.record_read(&canonical, true, false, 1); 931 read_only.record_read(&writable, true, false, 1); 932 let report = read_only.report(); 933 assert_eq!(report.state(), RelayAggregateState::ReadOnly); 934 assert_eq!(report.read_availability(), Availability::Available); 935 assert_eq!(report.write_availability(), Availability::Unavailable); 936 937 let read_only_profile = crate::RelayProfile::explicit( 938 crate::RelayProfileKind::Public, 939 [crate::RelayEndpoint::new( 940 "wss://read-only.example", 941 crate::RelayUrlPolicy::Public, 942 crate::RelayAccess::ReadOnly, 943 ) 944 .expect("read-only endpoint")], 945 ) 946 .expect("read-only profile"); 947 let read_only_config = Config::from_profile(read_only_profile); 948 let read_only_relay = read_only_config.relays()[0].clone(); 949 let read_only_tracker = StatusTracker::new(&read_only_config); 950 read_only_tracker.record_read(&read_only_relay, true, false, 1); 951 let report = read_only_tracker.report(); 952 assert_eq!(report.state(), RelayAggregateState::ReadOnly); 953 assert_eq!(report.read_availability(), Availability::Available); 954 assert_eq!(report.write_availability(), Availability::Unavailable); 955 assert_eq!( 956 report.relays()[0].write().state(), 957 RelayEvidenceState::Unsupported 958 ); 959 } 960 }