diff --git a/CHANGELOG.md b/CHANGELOG.md index 311a2a2..4ac0d60 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 +- Preserve background discovery relay backoff across idle-connection cleanup + so unavailable mailbox and profile sources do not restart their retry history. + - Render metrics on a blocking worker and serialize scrapes so repository counting does not occupy async workers serving relay and Git requests. diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 3186132..d49fd77 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -744,6 +744,14 @@ is deliberately historical-only: owning or maintaining a repository never adds that author's inbox relays to ordinary persistent live repository targets. Private instances omit repository-coordinate mailbox expansion. +Discovery source and mailbox scope also own connection retry state. The empty +relay checker and reconnect scheduler include that scope, including in-flight +fetches and future probe deadlines, without adding persistent subscriptions. +Healthy idle discovery sessions can retire after their fetch; failed or +policy-limited sessions retain their backoff until recovery or scope removal. +This prevents cleanup from erasing a failure every two seconds and allowing +discovery to dial the same unavailable endpoint from a fresh retry history. + ### Rejected Events Index The rejected events index solves two critical problems during sync: diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 41cc3d7..fb4972f 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1960,6 +1960,20 @@ struct Nip65DiscoveryState { } impl Nip65DiscoveryState { + /// Connection ownership includes historical discovery work, even while + /// its next query is deferred. It is not persistent subscription scope. + fn connection_targets(&self) -> HashSet { + self.author_sources + .values() + .flatten() + .chain(self.mailbox_roots.keys()) + .chain(self.mailbox_repositories.keys()) + .chain(self.mailbox_probes_in_flight.iter()) + .chain(self.in_flight.iter().map(|(relay, _)| relay)) + .cloned() + .collect() + } + fn has_mailbox_scope(&self, relay: &str) -> bool { self.mailbox_roots.contains_key(relay) || self.mailbox_repositories.contains_key(relay) } @@ -6422,6 +6436,22 @@ impl SyncManager { if !self.nip65_discovery_only_relays.contains(source) { return; } + // Retirement forgets connection health. A failed or policy-limited + // source must retain its session state until recovery or removal from + // discovery ownership, otherwise the next probe bypasses its backoff. + if self.health_tracker.get_failure_count(source) > 0 + || self.health_tracker.is_subscription_paused(source) + { + return; + } + if let Some(connection) = self.connections.get(source) { + if !connection.is_connected().await { + // The disconnect notification may still be queued behind a + // fetch result. Let it record the unexpected failure instead + // of relabeling the session as an intentional retirement. + return; + } + } let now = Instant::now(); let has_author_work = self.nip65_discovery @@ -7221,7 +7251,12 @@ impl SyncManager { // Once the connection ends naturally, however, reconcile its // confirmed state with the latest index before deciding whether it // should ever reconnect. - let desired = self.derive_targets().await.remove(relay_url); + let desired = self.derive_targets().await.remove(relay_url).or_else(|| { + self.nip65_discovery + .connection_targets() + .contains(relay_url) + .then(RelaySyncNeeds::default) + }); let Some(desired) = desired else { tracing::info!( relay = %relay_url, @@ -8420,6 +8455,7 @@ impl SyncManager { let mut desired_relays: HashSet = self.derive_targets().await.into_keys().collect(); desired_relays.extend(self.dependency_relay_deadlines.keys().cloned()); + desired_relays.extend(self.nip65_discovery.connection_targets()); // Collect relays to disconnect let to_disconnect: Vec = { @@ -8531,7 +8567,8 @@ impl SyncManager { /// /// For each eligible relay, a reconnection is queued via schedule_connect_relay. async fn retry_disconnected_relays(&mut self) { - let desired_relays: HashSet = self.derive_targets().await.into_keys().collect(); + let mut desired_relays: HashSet = self.derive_targets().await.into_keys().collect(); + desired_relays.extend(self.nip65_discovery.connection_targets()); // Collect relays to reconnect let to_reconnect: Vec = { @@ -10971,6 +11008,37 @@ mod tests { ); } + #[test] + fn discovery_connection_ownership_survives_retry_deadlines_and_in_flight_removal() { + let author = Keys::generate().public_key(); + let profile = "wss://profile.example".to_string(); + let mailbox = "wss://mailbox.example".to_string(); + let retired = "wss://retired.example".to_string(); + let mut discovery = Nip65DiscoveryState::default(); + discovery + .author_sources + .insert(author, HashSet::from([profile.clone()])); + discovery + .mailbox_repositories + .insert(mailbox.clone(), HashSet::from(["repo".into()])); + discovery + .mailbox_probe_next_at + .insert(mailbox.clone(), Instant::now() + Duration::from_secs(3600)); + discovery.in_flight.insert((retired.clone(), author)); + + let disconnected = RelayState::default(); + for relay in [&profile, &mailbox, &retired] { + assert!(!disconnected + .is_disconnect_candidate(false, discovery.connection_targets().contains(relay))); + } + discovery.in_flight.clear(); + assert!(!discovery.connection_targets().contains(&retired)); + discovery.author_sources.clear(); + discovery.install_mailbox_overlay(HashMap::new(), HashMap::new(), Instant::now()); + assert!(discovery.connection_targets().is_empty()); + assert!(disconnected.is_disconnect_candidate(false, false)); + } + #[test] fn disconnected_empty_relay_can_still_be_cleaned_up() { let mut source = RelayState::default(); diff --git a/tests/sync/reconnect_backoff.rs b/tests/sync/reconnect_backoff.rs index ed496fc..3730563 100644 --- a/tests/sync/reconnect_backoff.rs +++ b/tests/sync/reconnect_backoff.rs @@ -7,6 +7,47 @@ use nostr_sdk::prelude::*; use crate::common::flapping_relay::FlappingRelay; use crate::common::{TestClient, TestRelay}; +#[tokio::test] +async fn discovery_only_mailbox_keeps_failure_history_through_cleanup() { + use crate::common::{send_to_relay_url, setup_announcement_on_relay, MockRelay}; + + let index = MockRelay::start().await; + let flapping = FlappingRelay::start().await; + let owner = Keys::generate(); + let relay_list = EventBuilder::new(Kind::RelayList, "") + .tags([Tag::custom("r", vec![flapping.url(), "read"])]) + .finalize(&owner) + .unwrap(); + send_to_relay_url(index.url(), &relay_list).await.unwrap(); + let syncing = TestRelay::start_with_sync(Some(index.url().to_string())).await; + let domain = syncing.domain(); + let (_announcement, _git) = + setup_announcement_on_relay(&syncing, &owner, &[&domain], "discovery-mailbox-backoff") + .await; + + // Owner inbox discovery has no ordinary repository/root live target. + // The two-second cleanup pass previously erased its first failure before + // the next handshake, so it could never reach this recovery state. + flapping + .wait_for_connections(2, Duration::from_secs(60)) + .await; + let logs = wait_for_log( + &syncing.log_path(), + "consecutive_failures=1", + Duration::from_secs(15), + ) + .await; + assert!(logs.lines().any(|line| { + line.contains(flapping.url()) + && line.contains("consecutive_failures=1") + && line.contains("preserving failure streak until stable") + })); + + syncing.stop().await; + flapping.stop().await; + index.stop().await; +} + async fn wait_for_log(log_path: &std::path::Path, needle: &str, timeout: Duration) -> String { tokio::time::timeout(timeout, async { loop {