diff --git a/docs/explanation/grasp-02-proactive-sync.md b/docs/explanation/grasp-02-proactive-sync.md index 2249f72..0c2dad5 100644 --- a/docs/explanation/grasp-02-proactive-sync.md +++ b/docs/explanation/grasp-02-proactive-sync.md @@ -1048,8 +1048,8 @@ The [`RelayHealthTracker`](src/sync/health.rs:209) manages connection health wit Healthy <-> Disconnected: Normal connection/disconnection Disconnected -> Degraded: Connection failure Degraded -> Dead: 24h+ of continuous failures -Degraded -> Disconnected: Recovery (enters 5min stability period) -Disconnected -> Healthy: Stable for 5 minutes after recovery +Degraded -> Degraded: Handshake recovery starts a 5-minute stability period +Degraded -> Healthy: The recovered connection survives 5 minutes under normal sync load Any -> RateLimited: NOTICE message from relay indicating rate limiting RateLimited -> Probing: After 65-second cooldown expires Probing -> previous state: Recovery REQs succeed @@ -1067,6 +1067,8 @@ Probing -> RateLimited: Recovery REQ is rate limited again repeated notices during the same cooldown do not extend its deadline, while a rejection after that deadline starts a new cooldown - **Stability period**: 5 minutes after recovery before marking as Healthy + and clearing the failure streak. A short-lived successful handshake does not + reset exponential backoff; another disconnect continues the existing streak. ### Special Behaviors diff --git a/src/sync/health.rs b/src/sync/health.rs index 8f47ac6..7c47e40 100644 --- a/src/sync/health.rs +++ b/src/sync/health.rs @@ -274,7 +274,9 @@ impl RelayHealthTracker { /// Record a successful connection to a relay /// - /// Clears connection failure counters. Sets connected = true. + /// Sets connected = true. A relay recovering from failures keeps its + /// failure streak until it survives the stability period under normal + /// workload; a WebSocket handshake alone is not proof of recovery. /// /// A successful WebSocket connection does not prove that the relay will /// accept a new REQ, so it deliberately leaves any active rate-limit @@ -291,19 +293,35 @@ impl RelayHealthTracker { .flatten() .filter(|deadline| *deadline > now); - // Reset connection health. A live rate-limit cooldown is independent - // of whether the WebSocket handshake succeeded. + let was_connected = health.connected; + let recovering = health.consecutive_failures > 0 || health.last_failure_time.is_some(); + + // A live rate-limit cooldown is independent of whether the WebSocket + // handshake succeeded. health.connected = true; health.rate_limited = active_rate_limit.is_some(); - health.consecutive_failures = 0; - health.first_failure_time = None; - health.last_failure_time = None; - health.last_success_time = Some(now); + if !recovering { + health.consecutive_failures = 0; + health.first_failure_time = None; + health.last_failure_time = None; + } + // Keep one stability boundary per session if a duplicate connection + // notification reaches this idempotent accounting boundary. + if !was_connected { + health.last_success_time = Some(now); + } health.last_attempt_time = Some(now); health.next_retry_at = active_rate_limit; let new_state = health.state(); - if old_state != new_state { + if recovering && !was_connected { + tracing::info!( + relay = %relay_url, + consecutive_failures = health.consecutive_failures, + stability_secs = STABILITY_PERIOD_SECS, + "Relay connected; preserving failure streak until stable" + ); + } else if old_state != new_state { tracing::info!( "Relay {} connection recovered ({:?} -> {:?})", relay_url, @@ -313,6 +331,42 @@ impl RelayHealthTracker { } } + /// Promote recovered connections that have remained established through + /// the stability period under their normal workload. + pub fn promote_stable_connections(&self) -> Vec { + let now = Instant::now(); + let stability_period = Duration::from_secs(STABILITY_PERIOD_SECS); + let mut promoted = Vec::new(); + + for mut entry in self.health.iter_mut() { + let (relay_url, health) = entry.pair_mut(); + if !health.connected || health.consecutive_failures == 0 { + continue; + } + let Some(connected_at) = health.last_success_time else { + continue; + }; + if now.duration_since(connected_at) < stability_period { + continue; + } + + health.consecutive_failures = 0; + health.first_failure_time = None; + health.last_failure_time = None; + if !health.is_rate_limited_now() { + health.next_retry_at = None; + } + promoted.push(relay_url.clone()); + tracing::info!( + relay = %relay_url, + stability_secs = STABILITY_PERIOD_SECS, + "Relay connection proved stable; failure streak reset" + ); + } + + promoted + } + /// Record a connection failure for a relay /// /// Increments failure counter and calculates next retry time with exponential backoff. @@ -625,7 +679,7 @@ mod tests { } #[test] - fn test_record_success_resets_to_healthy() { + fn recovered_handshake_preserves_backoff_until_stable() { let tracker = RelayHealthTracker::with_defaults(); let url = "wss://test-relay.example.com"; @@ -637,9 +691,53 @@ mod tests { // Record success tracker.record_success(url); + assert_eq!(tracker.get_state(url), HealthState::Degraded); + assert_eq!(tracker.get_failure_count(url), 2); + assert!(tracker.should_attempt_connection(url)); + + // A further disconnect continues the existing failure streak rather + // than returning to the base backoff. + tracker.record_failure(url); + assert_eq!(tracker.get_failure_count(url), 3); + let remaining = tracker.get_remaining_backoff(url).unwrap(); + assert!(remaining > Duration::from_secs(10)); + } + + #[test] + fn stable_recovered_connection_resets_failure_history() { + let tracker = RelayHealthTracker::with_defaults(); + let url = "wss://stable-relay.example.com"; + + tracker.record_failure(url); + tracker.record_success(url); + { + let mut entry = tracker.health.get_mut(url).unwrap(); + entry.last_success_time = + Some(Instant::now() - Duration::from_secs(STABILITY_PERIOD_SECS + 1)); + } + + assert_eq!(tracker.promote_stable_connections(), vec![url.to_string()]); assert_eq!(tracker.get_state(url), HealthState::Healthy); assert_eq!(tracker.get_failure_count(url), 0); - assert!(tracker.should_attempt_connection(url)); + } + + #[test] + fn disconnected_relay_cannot_age_into_stable_recovery() { + let tracker = RelayHealthTracker::with_defaults(); + let url = "wss://disconnected-relay.example.com"; + + tracker.record_failure(url); + tracker.record_success(url); + tracker.record_failure(url); + { + let mut entry = tracker.health.get_mut(url).unwrap(); + entry.last_success_time = + Some(Instant::now() - Duration::from_secs(STABILITY_PERIOD_SECS + 1)); + } + + assert!(tracker.promote_stable_connections().is_empty()); + assert_eq!(tracker.get_state(url), HealthState::Degraded); + assert_eq!(tracker.get_failure_count(url), 2); } #[test] diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 00bbfe0..b56464e 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1120,14 +1120,18 @@ async fn run_health_and_metrics_checker( _ = tokio::time::sleep(interval) => { let mut manager = sync_manager.lock().await; - // 1. Check for disconnects and retry disconnected relays + // 1. Reset failure history only after a recovered connection + // has survived the stability period under normal sync load. + manager.health_tracker.promote_stable_connections(); + + // 2. Check for disconnects and retry disconnected relays manager.check_disconnects().await; manager.retry_disconnected_relays().await; - // 2. Check for rate limit recovery + // 3. Check for rate limit recovery manager.check_rate_limit_recovery().await; - // 3. Check for naughty list expiration + // 4. Check for naughty list expiration if let Some(naughty_list) = manager.health_tracker.naughty_list() { let recovered = naughty_list.expire_old_entries(); for url in recovered { @@ -1138,7 +1142,7 @@ async fn run_health_and_metrics_checker( } } - // 4. Update metrics with current health states and naughty list + // 5. Update metrics with current health states and naughty list if let Some(ref metrics) = manager.metrics { // Get all tracked relay URLs let relay_urls: Vec = { @@ -3351,8 +3355,8 @@ impl SyncManager { /// 3. Rebuilding L2+L3 from preserved RelaySyncIndex state /// 4. Computing actions for new items discovered during catchup /// - /// Basic connection state and metrics are managed by handle_connect_or_reconnect. - /// This method handles reconnect-specific concerns (health tracking, reconnect metrics). + /// Basic connection health is managed by the connection worker and its + /// stability timer. This method handles reconnect-specific sync and metrics. async fn quick_reconnect(&mut self, relay_url: &str, since: Timestamp) { self.cancel_deferred_consolidation(relay_url, "quick reconnect"); @@ -3363,9 +3367,6 @@ impl SyncManager { pending.remove(relay_url); } - // Record successful reconnection in health tracker - self.health_tracker.record_success(relay_url); - // Record reconnect-specific metrics (not basic connection metrics) if let Some(ref metrics) = self.metrics { metrics.record_health_state(relay_url, self.health_tracker.get_state(relay_url)); diff --git a/tests/common/flapping_relay.rs b/tests/common/flapping_relay.rs new file mode 100644 index 0000000..857fbdd --- /dev/null +++ b/tests/common/flapping_relay.rs @@ -0,0 +1,115 @@ +//! WebSocket relay fixture that accepts normal setup and then drops the session. + +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +use futures_util::StreamExt; +use tokio::net::TcpListener; +use tokio::sync::{oneshot, Notify}; +use tokio_tungstenite::tungstenite::Message; + +/// Relay-shaped endpoint for exercising reconnect health accounting. +pub struct FlappingRelay { + url: String, + connected_at: Arc>>, + connection_observed: Arc, + shutdown_tx: Option>, + handle: Option>, +} + +impl FlappingRelay { + /// Start an endpoint that drops each WebSocket after the first REQ frame. + pub async fn start() -> Self { + let listener = TcpListener::bind("127.0.0.1:0") + .await + .expect("FlappingRelay failed to bind"); + let port = listener.local_addr().expect("flapping local_addr").port(); + let connected_at = Arc::new(Mutex::new(Vec::new())); + let connection_observed = Arc::new(Notify::new()); + let (shutdown_tx, mut shutdown_rx) = oneshot::channel(); + let observed = connected_at.clone(); + let notify = connection_observed.clone(); + + let handle = tokio::spawn(async move { + loop { + tokio::select! { + accepted = listener.accept() => { + let Ok((stream, _)) = accepted else { break }; + let observed = observed.clone(); + let notify = notify.clone(); + tokio::spawn(async move { + let Ok(mut websocket) = tokio_tungstenite::accept_async(stream).await + else { + // NIP-11 HTTP probes are expected to fail this + // deliberately WebSocket-only fixture. + return; + }; + observed.lock().expect("flapping timestamps poisoned").push(Instant::now()); + notify.notify_waiters(); + + while let Some(message) = websocket.next().await { + let Ok(message) = message else { break }; + if matches!(message, Message::Text(ref text) if text.contains("\"REQ\"")) { + // Dropping the socket after observable normal + // setup models the production peer reset. + return; + } + } + }); + } + _ = &mut shutdown_rx => break, + } + } + }); + + Self { + url: format!("ws://127.0.0.1:{port}"), + connected_at, + connection_observed, + shutdown_tx: Some(shutdown_tx), + handle: Some(handle), + } + } + + pub fn url(&self) -> &str { + &self.url + } + + /// Wait for `count` successful WebSocket sessions using notifications, + /// bounded by `timeout`. + pub async fn wait_for_connections(&self, count: usize, timeout: Duration) -> Vec { + tokio::time::timeout(timeout, async { + loop { + let notified = self.connection_observed.notified(); + let timestamps = self + .connected_at + .lock() + .expect("flapping timestamps poisoned") + .clone(); + if timestamps.len() >= count { + return timestamps; + } + notified.await; + } + }) + .await + .expect("recovered relay did not make the expected connection attempts") + } + + pub async fn stop(mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + if let Some(handle) = self.handle.take() { + let _ = handle.await; + } + } +} + +impl Drop for FlappingRelay { + fn drop(&mut self) { + if let Some(tx) = self.shutdown_tx.take() { + let _ = tx.send(()); + } + } +} diff --git a/tests/common/mod.rs b/tests/common/mod.rs index 6af269a..aeaca7b 100644 --- a/tests/common/mod.rs +++ b/tests/common/mod.rs @@ -3,6 +3,7 @@ #![allow(unused_imports)] // Re-exports may not be used in all test configurations pub mod censoring_proxy; +pub mod flapping_relay; pub mod git_server; pub mod mock_relay; pub mod neg_limiting_proxy; diff --git a/tests/sync.rs b/tests/sync.rs index 6e3b20b..625f628 100644 --- a/tests/sync.rs +++ b/tests/sync.rs @@ -42,6 +42,7 @@ mod sync { pub mod naughty_list_scheduling; pub mod neg_concurrency; pub mod purgatory_fetch; + pub mod reconnect_backoff; pub mod req_concurrency; pub mod tag_variations; } diff --git a/tests/sync/reconnect_backoff.rs b/tests/sync/reconnect_backoff.rs new file mode 100644 index 0000000..6b0aacf --- /dev/null +++ b/tests/sync/reconnect_backoff.rs @@ -0,0 +1,78 @@ +//! Reconnect failure history must survive short-lived successful handshakes. + +use std::time::Duration; + +use nostr_sdk::prelude::*; + +use crate::common::flapping_relay::FlappingRelay; +use crate::common::{TestClient, TestRelay}; + +async fn wait_for_log(log_path: &std::path::Path, needle: &str, timeout: Duration) -> String { + tokio::time::timeout(timeout, async { + loop { + let contents = std::fs::read_to_string(log_path).unwrap_or_default(); + if contents.contains(needle) { + return contents; + } + tokio::task::yield_now().await; + } + }) + .await + .unwrap_or_else(|_| panic!("relay log never contained {needle:?}")) +} + +#[tokio::test] +async fn flapping_relay_handshakes_do_not_reset_exponential_backoff() { + let flapping = FlappingRelay::start().await; + let syncing = TestRelay::start_with_sync(None).await; + let keys = Keys::generate(); + let identifier = "flapping-reconnect-backoff"; + let npub = keys.public_key().to_bech32().expect("npub"); + let announcement = EventBuilder::new(Kind::GitRepoAnnouncement, "") + .tags(vec![ + Tag::identifier(identifier), + Tag::custom( + "clone", + vec![format!( + "http://{}/{npub}/{identifier}.git", + syncing.domain() + )], + ), + Tag::custom( + "relays", + vec![ + format!("ws://{}", syncing.domain()), + flapping.url().to_string(), + ], + ), + ]) + .finalize(&keys) + .expect("sign announcement"); + let client = TestClient::new(syncing.url(), keys) + .await + .expect("connect publishing client"); + client + .send_event(&announcement) + .await + .expect("publish flapping relay announcement"); + + // The recovered session must retain the first failure instead of treating + // its WebSocket handshake as proof of stability. Each session reaches a + // normal REQ before the fixture drops it. Unit coverage below the actor + // boundary verifies that the retained count drives the next backoff step. + let attempts = flapping + .wait_for_connections(2, Duration::from_secs(25)) + .await; + assert_eq!(attempts.len(), 2); + let logs = wait_for_log( + &syncing.log_path(), + "consecutive_failures=1", + Duration::from_secs(5), + ) + .await; + assert!(logs.contains("preserving failure streak until stable")); + + client.disconnect().await; + syncing.stop().await; + flapping.stop().await; +}