fix(sync): require stable recovery before resetting backoff

Production showed nostr-pub.wellorder.net repeatedly completing a WebSocket handshake, resetting under normal REQ load 30-42 seconds later, and returning to the five-second base reconnect delay. The health tracker cleared its failure streak at handshake time even though its state model and documentation already required five minutes of stable recovery.

Preserve connection-failure history across recovered handshakes, promote a connected relay only after the existing five-minute stability period, and remove the redundant quick-reconnect success record. A relay continues normal live and historic work while degraded; only its next reconnect delay escalates. Active rate-limit cooldown accounting remains independent.

A WebSocket/REQ fixture reproduces the short-lived-success lifecycle and verifies the recovered session retains its failure count. Unit tests cover escalated backoff, stable promotion, and the rule that disconnected relays cannot age into recovery.

Correctness assumes an uninterrupted five-minute session under normal workload is sufficient evidence to forgive the streak. Intermittent sessions remain one instability episode, including for the existing 24-hour dead-relay threshold. Changing workload pace, stability duration, backoff configuration, or relay subscription policy is deliberately excluded.

Validated with cargo test --lib (650 passed), cargo test --test sync (88 passed, 1 ignored), the focused reconnect scenario, and nix build .#ngit-grasp.
This commit is contained in:
DanConwayDev
2026-08-07 16:25:21 +00:00
parent 50033e0157
commit bccbb54b62
7 changed files with 317 additions and 21 deletions
+4 -2
View File
@@ -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
+108 -10
View File
@@ -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<String> {
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]
+10 -9
View File
@@ -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<String> = {
@@ -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));
+115
View File
@@ -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<Mutex<Vec<Instant>>>,
connection_observed: Arc<Notify>,
shutdown_tx: Option<oneshot::Sender<()>>,
handle: Option<tokio::task::JoinHandle<()>>,
}
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<Instant> {
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(());
}
}
}
+1
View File
@@ -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;
+1
View File
@@ -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;
}
+78
View File
@@ -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;
}