From c82684092c7b4f81e49833b0888500fcb9851218 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Thu, 11 Dec 2025 11:19:38 +0000 Subject: fix(sync): improve metrics recording and connection failure detection Changes: - Fix connection attempt metrics: record success/failure based on actual connection result instead of pre-emptively recording failure - Add health tracker integration on connection failure: call record_failure() and record_health_state() in error path - Add connection verification in relay_connection.rs: wait 500ms after connect() then verify is_connected() to detect silent failures - Add configurable disconnect check interval via NGIT_SYNC_DISCONNECT_CHECK_INTERVAL_SECS env var - Update TestRelay with fast test settings: startup_delay=0, jitter=0, disconnect_check_interval=1s - Add debug output to metrics tests for investigation Note: Tests may still fail due to 5-second base backoff in health tracker. A follow-up task will add NGIT_SYNC_BASE_BACKOFF_SECS config parameter to allow faster test cycles. Related: metrics-wiring-plan.md Tasks 1 & 2 --- src/config.rs | 6 ++++ src/sync/mod.rs | 76 ++++++++++++++++++++++++++++---------------- src/sync/relay_connection.rs | 18 ++++++++++- tests/common/relay.rs | 3 ++ tests/sync/metrics.rs | 34 ++++++++++++++++++-- 5 files changed, 107 insertions(+), 30 deletions(-) diff --git a/src/config.rs b/src/config.rs index 69a160a..5e74471 100644 --- a/src/config.rs +++ b/src/config.rs @@ -109,6 +109,11 @@ pub struct Config { /// Set to 0 to disable jitter (useful for testing) #[arg(long, env = "NGIT_SYNC_STARTUP_JITTER_MS", default_value_t = 10_000)] pub sync_startup_jitter_ms: u64, + + /// Interval in seconds for checking disconnected relays and attempting reconnection (default: 60) + /// Set to lower value for faster reconnection testing + #[arg(long, env = "NGIT_SYNC_DISCONNECT_CHECK_INTERVAL_SECS", default_value_t = 60)] + pub sync_disconnect_check_interval_secs: u64, } impl Config { @@ -170,6 +175,7 @@ impl Config { sync_reconnect_delay_secs: 10, sync_reconnect_lookback_days: 3, sync_startup_jitter_ms: 10_000, + sync_disconnect_check_interval_secs: 60, } } } diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 5039c04..16ad833 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -264,23 +264,30 @@ async fn run_daily_timer( // Disconnect Checker (Phase 8) // ============================================================================= -/// Check interval for empty relay cleanup in seconds -const DISCONNECT_CHECK_INTERVAL_SECS: u64 = 60; - /// Run the disconnect checker for periodic cleanup of empty relays /// -/// This function runs in a loop, checking every 60 seconds for relays +/// This function runs in a loop, checking at the configured interval for relays /// that have no repos or root events to sync. Non-bootstrap relays /// that are empty will be disconnected to free up resources. /// /// Bootstrap relays are never disconnected, even if empty. +/// +/// The check interval is configurable via `NGIT_SYNC_DISCONNECT_CHECK_INTERVAL_SECS` +/// (default: 60 seconds). Set to a lower value for faster reconnection testing. async fn run_disconnect_checker( sync_manager: Arc>, mut shutdown_rx: broadcast::Receiver<()>, + check_interval_secs: u64, ) { + let interval = Duration::from_secs(check_interval_secs); + tracing::info!( + interval_secs = check_interval_secs, + "Disconnect checker started with configured interval" + ); + loop { tokio::select! { - _ = tokio::time::sleep(Duration::from_secs(DISCONNECT_CHECK_INTERVAL_SECS)) => { + _ = tokio::time::sleep(interval) => { tracing::debug!("Disconnect checker running"); let mut manager = sync_manager.lock().await; @@ -609,21 +616,24 @@ impl SyncManager { self.spawn_relay_connection(bootstrap_url.clone()).await; } - // 7. Wrap self in Arc for sharing with timer task + // 7. Capture config values before moving self into Arc + let disconnect_check_interval_secs = self.config.sync_disconnect_check_interval_secs; + + // 8. Wrap self in Arc for sharing with timer task let sync_manager = Arc::new(Mutex::new(self)); - // 8. Spawn daily timer task with shutdown receiver + // 9. Spawn daily timer task with shutdown receiver let timer_manager = Arc::clone(&sync_manager); let timer_shutdown = shutdown_tx.subscribe(); tokio::spawn(async move { run_daily_timer(timer_manager, timer_shutdown).await; }); - // 9. Spawn disconnect checker task with shutdown receiver + // 10. Spawn disconnect checker task with shutdown receiver let checker_manager = Arc::clone(&sync_manager); let checker_shutdown = shutdown_tx.subscribe(); tokio::spawn(async move { - run_disconnect_checker(checker_manager, checker_shutdown).await; + run_disconnect_checker(checker_manager, checker_shutdown, disconnect_check_interval_secs).await; }); // 10. Main loop - handle actions from self-subscriber, disconnect, EOSE, and connect notifications @@ -1174,27 +1184,39 @@ impl SyncManager { // Create relay connection let connection = RelayConnection::new(relay_url.clone()); - // Record connection attempt - if let Some(ref metrics) = self.metrics { - metrics.record_connection_attempt(&relay_url, false); - } - // Connect and subscribe to Layer 1 - if let Err(e) = connection.connect_and_subscribe(None).await { - tracing::error!(relay = %relay_url, error = %e, "Failed to connect to relay"); - // Update state to disconnected on failure - { - let mut index = relay_sync_index.write().await; - if let Some(state) = index.get_mut(&relay_url) { - state.connection_status = ConnectionStatus::Disconnected; + match connection.connect_and_subscribe(None).await { + Ok(_) => { + // Record successful connection attempt + if let Some(ref metrics) = self.metrics { + metrics.record_connection_attempt(&relay_url, true); } } - return; - } - - // If successful, update connection attempt metric to success - if let Some(ref metrics) = self.metrics { - metrics.record_connection_attempt(&relay_url, true); + Err(e) => { + tracing::error!(relay = %relay_url, error = %e, "Failed to connect to relay"); + + // Record failed connection attempt + if let Some(ref metrics) = self.metrics { + metrics.record_connection_attempt(&relay_url, false); + } + + // Record failure in health tracker + self.health_tracker.record_failure(&relay_url); + + // Record health state in metrics + if let Some(ref metrics) = self.metrics { + metrics.record_health_state(&relay_url, self.health_tracker.get_state(&relay_url)); + } + + // Update state to disconnected on failure + { + let mut index = relay_sync_index.write().await; + if let Some(state) = index.get_mut(&relay_url) { + state.connection_status = ConnectionStatus::Disconnected; + } + } + return; + } } // Mark as connected in relay sync index diff --git a/src/sync/relay_connection.rs b/src/sync/relay_connection.rs index 09c9887..d69e112 100644 --- a/src/sync/relay_connection.rs +++ b/src/sync/relay_connection.rs @@ -55,7 +55,8 @@ impl RelayConnection { /// This method: /// 1. Adds the relay to the client /// 2. Establishes the WebSocket connection - /// 3. Subscribes to Layer 1 filter (kinds 30617 + 30618) + /// 3. Verifies connection was established + /// 4. Subscribes to Layer 1 filter (kinds 30617 + 30618) /// /// # Arguments /// * `since` - Optional timestamp for incremental sync on reconnect @@ -76,6 +77,21 @@ impl RelayConnection { // Establish connection self.client.connect().await; + // Wait briefly for connection to establish and check status + // nostr-sdk's connect() is async and may not immediately reflect failure + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + + // Check if relay is actually connected + let relay = self.client.relay(&self.url).await + .map_err(|e| format!("Failed to get relay handle for {}: {}", self.url, e))?; + + if !relay.is_connected() { + return Err(format!( + "Failed to connect to relay {}: connection not established after timeout", + self.url + )); + } + // Subscribe to Layer 1 (announcements) let filter = build_announcement_filter(since); let output = self diff --git a/tests/common/relay.rs b/tests/common/relay.rs index d954e34..b11dde3 100644 --- a/tests/common/relay.rs +++ b/tests/common/relay.rs @@ -94,6 +94,9 @@ impl TestRelay { .env("NGIT_DATABASE_BACKEND", "memory") // Force in-memory database for isolation .env("NGIT_OWNER_NPUB", &test_npub) .env("NGIT_SYNC_BATCH_WINDOW_MS", "200") // Fast batch window for tests (200ms instead of 5s default) + .env("NGIT_SYNC_STARTUP_DELAY_SECS", "0") // No startup delay for faster tests + .env("NGIT_SYNC_STARTUP_JITTER_MS", "0") // No jitter for tests + .env("NGIT_SYNC_DISCONNECT_CHECK_INTERVAL_SECS", "1") // Fast reconnect attempts for tests .env("RUST_LOG", "info") // Enable INFO logging for diagnostics .stdout(Stdio::null()) // Disable stderr for cleaner test output // .stdout(Stdio::inherit()) // Show stdout for diagnostics diff --git a/tests/sync/metrics.rs b/tests/sync/metrics.rs index 82d681e..e11fe58 100644 --- a/tests/sync/metrics.rs +++ b/tests/sync/metrics.rs @@ -368,17 +368,48 @@ async fn test_startup_sync_event_count() { /// NOTE: This test may fail until sync metrics recording is fully wired up. /// The test documents the expected behavior. #[tokio::test] -#[ignore] // Enable when metrics recording is implemented async fn test_connection_failure_increments_counter() { let mut harness = MetricsTestHarness::with_sources(0).await; // No sources harness.start_syncing_relay_to_nowhere().await; // Wait for initial connection attempts tokio::time::sleep(Duration::from_secs(2)).await; + + // Fetch raw metrics to debug + let syncing_url = harness.syncing_relay_url().expect("Syncing relay should be started"); + let raw_1 = fetch_metrics(syncing_url) + .await + .expect("Failed to fetch metrics"); + + // Print all sync-related metrics + println!("\n=== RAW METRICS (t1) ==="); + for line in raw_1.lines() { + if line.contains("sync") || line.contains("connection") { + println!("{}", line); + } + } + println!("========================\n"); + let metrics_1 = harness.get_metrics().await.unwrap(); // Wait for more attempts tokio::time::sleep(Duration::from_secs(2)).await; + + // Fetch raw metrics again + let syncing_url = harness.syncing_relay_url().expect("Syncing relay should be started"); + let raw_2 = fetch_metrics(syncing_url) + .await + .expect("Failed to fetch metrics"); + + // Print all sync-related metrics + println!("\n=== RAW METRICS (t2) ==="); + for line in raw_2.lines() { + if line.contains("sync") || line.contains("connection") { + println!("{}", line); + } + } + println!("========================\n"); + let metrics_2 = harness.get_metrics().await.unwrap(); // Failure counter should have increased @@ -496,7 +527,6 @@ async fn test_relay_connected_status() { /// NOTE: This test may fail until sync metrics recording is fully wired up. /// The test documents the expected behavior. #[tokio::test] -#[ignore] // Ignored until sync metrics are fully wired up async fn test_health_state_degrades_on_failure() { use crate::common::sync_helpers::MetricsTestHarness; -- cgit v1.2.3