From 888cc70fee3c5c9e35adaf91e1fe646c551bd632 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Thu, 24 Sep 2026 10:51:50 +0000 Subject: [PATCH] fix(sync): back off repeatedly failed mailbox probes Persistently denied participant mailboxes retried every five minutes because probe failures did not enter subscription backoff and reconnects reset due times. Track a per-mailbox delay, doubling from five minutes to one hour, and reset it only after a successful fetch. Preserve active pauses across reconnects, inventory growth, and deferred admission. Retire delay state when its mailbox scope disappears. Keep cursor rotation and working live subscriptions unchanged; this assumes failed fetches merit slower probes without classifying every auth refusal as transport failure. No relay-wide bans, new settings, or wire-validation changes are included. Validation: 14 focused mailbox unit tests pass, including exponential cap, success reset, reconnect/inventory preservation, and scope-removal cleanup. Assisted-by: GPT-6 --- CHANGELOG.md | 3 + docs/explanation/sync-scaling-constraints.md | 11 ++ src/sync/mod.rs | 126 +++++++++++++++++-- 3 files changed, 128 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5da238c..193ec1f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Back off repeatedly failed participant mailbox probes, preserving their pauses + across reconnects and inventory changes. + - Honor policy and request-failure pauses when reconnecting, so refused peers cannot bypass transport backoff through their policy-limited health state. diff --git a/docs/explanation/sync-scaling-constraints.md b/docs/explanation/sync-scaling-constraints.md index 0a27344..621bfa0 100644 --- a/docs/explanation/sync-scaling-constraints.md +++ b/docs/explanation/sync-scaling-constraints.md @@ -694,3 +694,14 @@ query-rate refusal and then drains work smoothly at a learned rate. - [Defensive Measures & Rate Limiting](defensive-measures.md) — the serving-side counterpart. - [Monitoring Overview](monitoring.md) — metrics for observing sync health. + +### Failed participant mailbox probes + +Mailbox fetch failures, including rejected authentication, retry after five +minutes and double on each failed probe up to one hour. Only a successful fetch +resets this delay; reconnects and inventory growth preserve it. The cursor still +advances so one failed filter cannot monopolize the mailbox. Backoff is scoped +to that relay's mailbox work and does not close working live subscriptions. +Mailbox-only reconnects wait for the next probe deadline; independent discovery +or repository work can still connect. Removing a relay from the mailbox +inventory retires its backoff state. diff --git a/src/sync/mod.rs b/src/sync/mod.rs index a2b6ecc..85d67c7 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1945,6 +1945,8 @@ struct Nip65DiscoveryState { mailbox_repositories: HashMap>, mailbox_probe_next_at: HashMap, mailbox_probe_next_filter: HashMap, + /// Failed mailbox probes back off independently of transport and live coverage. + mailbox_probe_failure_delay: HashMap, /// At most one ordinary fetch runs per mailbox relay. Different relays /// remain independent and use their own connection capacity. mailbox_probes_in_flight: HashSet, @@ -1968,6 +1970,28 @@ impl Nip65DiscoveryState { .collect() } + /// Ownership survives a pause, but mailbox-only sockets should not be + /// redialed until their next probe. Other discovery work remains eligible. + fn due_connection_targets(&self, now: Instant) -> HashSet { + self.author_sources + .values() + .flatten() + .chain(self.in_flight.iter().map(|(relay, _)| relay)) + .chain(self.mailbox_probes_in_flight.iter()) + .chain( + self.mailbox_roots + .keys() + .chain(self.mailbox_repositories.keys()) + .filter(|relay| { + self.mailbox_probe_next_at + .get(*relay) + .is_none_or(|due| *due <= now) + }), + ) + .cloned() + .collect() + } + fn has_mailbox_scope(&self, relay: &str) -> bool { self.mailbox_roots.contains_key(relay) || self.mailbox_repositories.contains_key(relay) } @@ -2005,16 +2029,42 @@ impl Nip65DiscoveryState { .retain(|relay, _| current_relays.contains(relay)); self.mailbox_probe_next_filter .retain(|relay, _| current_relays.contains(relay)); + self.mailbox_probe_failure_delay + .retain(|relay, _| current_relays.contains(relay)); for relay in current_relays { if old_root_overlay.get(&relay) != self.mailbox_roots.get(&relay) || old_repository_overlay.get(&relay) != self.mailbox_repositories.get(&relay) { - self.mailbox_probe_next_at.insert(relay.clone(), now); + if !self.mailbox_probe_failure_delay.contains_key(&relay) { + self.mailbox_probe_next_at.insert(relay.clone(), now); + } self.mailbox_probe_next_filter.entry(relay).or_default(); } } } + fn mailbox_connected(&mut self, relay: &str, now: Instant) { + if self.has_mailbox_scope(relay) && !self.mailbox_probe_failure_delay.contains_key(relay) { + self.mailbox_probe_next_at.insert(relay.to_string(), now); + } + } + + fn mailbox_retry_delay(&mut self, relay: &str, succeeded: bool) -> Duration { + if succeeded { + self.mailbox_probe_failure_delay.remove(relay); + return Duration::ZERO; + } + if !self.has_mailbox_scope(relay) { + return mailbox_probe_retry_interval(); + } + let delay = self + .mailbox_probe_failure_delay + .entry(relay.to_string()) + .and_modify(|delay| *delay = delay.saturating_mul(2).min(Duration::from_secs(3600))) + .or_insert_with(mailbox_probe_retry_interval); + *delay + } + fn record_probe_completion( &mut self, relay: &str, @@ -5981,10 +6031,12 @@ impl SyncManager { fn defer_mailbox_relay(&mut self, relay: &str) { if self.nip65_discovery.has_mailbox_scope(relay) { - self.nip65_discovery.mailbox_probe_next_at.insert( - relay.to_string(), - Instant::now() + mailbox_probe_retry_interval(), - ); + let due = Instant::now() + mailbox_probe_retry_interval(); + self.nip65_discovery + .mailbox_probe_next_at + .entry(relay.to_string()) + .and_modify(|deadline| *deadline = (*deadline).max(due)) + .or_insert(due); } } @@ -6022,8 +6074,14 @@ impl SyncManager { .await; } - let (next_filter, completed_cycle, next_probe_in) = + let (next_filter, completed_cycle, mut next_probe_in) = mailbox_probe_completion(result.filter_index, result.filter_count, succeeded); + let failure_delay = self + .nip65_discovery + .mailbox_retry_delay(&result.source_relay, succeeded); + if !succeeded { + next_probe_in = failure_delay; + } self.nip65_discovery.record_probe_completion( &result.source_relay, next_filter, @@ -6382,11 +6440,8 @@ impl SyncManager { if let Some(ref metrics) = self.metrics { metrics.record_connection_attempt(&result.relay_url, true); } - if self.nip65_discovery.has_mailbox_scope(&result.relay_url) { - self.nip65_discovery - .mailbox_probe_next_at - .insert(result.relay_url.clone(), Instant::now()); - } + self.nip65_discovery + .mailbox_connected(&result.relay_url, Instant::now()); self.handle_connect_or_reconnect(&result.relay_url).await; } ConnectAttemptOutcome::PrivateService => { @@ -8413,7 +8468,7 @@ impl SyncManager { /// For each eligible relay, a reconnection is queued via schedule_connect_relay. async fn retry_disconnected_relays(&mut self) { let mut desired_relays: HashSet = self.derive_targets().await.into_keys().collect(); - desired_relays.extend(self.nip65_discovery.connection_targets()); + desired_relays.extend(self.nip65_discovery.due_connection_targets(Instant::now())); // Collect relays to reconnect let to_reconnect: Vec = { @@ -10217,6 +10272,53 @@ mod tests { assert!(!discovery.mailbox_probe_next_at.contains_key(&relay)); } + #[test] + fn failed_mailbox_probes_back_off_until_success_and_survive_inventory_changes() { + let relay = "wss://denied.example".to_string(); + let root = EventId::from_byte_array([34; 32]); + let now = Instant::now(); + let mut discovery = Nip65DiscoveryState::default(); + discovery.install_mailbox_overlay( + HashMap::from([(relay.clone(), HashSet::from([root]))]), + HashMap::new(), + now, + ); + let mut expected = mailbox_probe_retry_interval(); + for _ in 0..16 { + assert_eq!(discovery.mailbox_retry_delay(&relay, false), expected); + expected = expected.saturating_mul(2).min(Duration::from_secs(3600)); + } + let deadline = now + Duration::from_secs(3600); + discovery.record_probe_completion(&relay, 1, Duration::from_secs(3600), now); + discovery.install_mailbox_overlay( + HashMap::from([( + relay.clone(), + HashSet::from([root, EventId::from_byte_array([35; 32])]), + )]), + HashMap::new(), + now, + ); + discovery.mailbox_connected(&relay, now); + assert!(discovery.connection_targets().contains(&relay)); + assert!(!discovery.due_connection_targets(now).contains(&relay)); + assert!(discovery.due_connection_targets(deadline).contains(&relay)); + let author = Keys::generate().public_key(); + discovery + .author_sources + .insert(author, HashSet::from([relay.clone()])); + assert!(discovery.due_connection_targets(now).contains(&relay)); + discovery.author_sources.clear(); + assert_eq!(discovery.mailbox_probe_next_at[&relay], deadline); + assert_eq!(discovery.mailbox_retry_delay(&relay, true), Duration::ZERO); + assert_eq!( + discovery.mailbox_retry_delay(&relay, false), + mailbox_probe_retry_interval() + ); + discovery.install_mailbox_overlay(HashMap::new(), HashMap::new(), now); + discovery.mailbox_retry_delay(&relay, false); + assert!(discovery.mailbox_probe_failure_delay.is_empty()); + } + #[test] fn mailbox_completion_for_removed_relay_leaves_no_cursor_state() { let relay = "wss://mailbox.example".to_string();