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
This commit is contained in:
DanConwayDev
2026-09-24 10:56:44 +00:00
parent e4e36862f3
commit 888cc70fee
3 changed files with 128 additions and 12 deletions
+3
View File
@@ -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.
@@ -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.
+114 -12
View File
@@ -1945,6 +1945,8 @@ struct Nip65DiscoveryState {
mailbox_repositories: HashMap<String, HashSet<String>>,
mailbox_probe_next_at: HashMap<String, Instant>,
mailbox_probe_next_filter: HashMap<String, usize>,
/// Failed mailbox probes back off independently of transport and live coverage.
mailbox_probe_failure_delay: HashMap<String, Duration>,
/// At most one ordinary fetch runs per mailbox relay. Different relays
/// remain independent and use their own connection capacity.
mailbox_probes_in_flight: HashSet<String>,
@@ -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<String> {
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<String> = 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<String> = {
@@ -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();