lib

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

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 }