diff --git a/tests/common/nip09_helpers.rs b/tests/common/nip09_helpers.rs index 36d7e4b..51c5672 100644 --- a/tests/common/nip09_helpers.rs +++ b/tests/common/nip09_helpers.rs @@ -19,6 +19,24 @@ use super::purgatory_helpers::{ }; use super::sync_helpers::create_repo_announcement; +/// Wait for promotion rather than assuming the worker runs within a fixed delay. +async fn wait_until_served(client: &AuditClient, event_id: EventId) { + tokio::time::timeout(Duration::from_secs(10), async { + loop { + if client + .is_event_on_relay(event_id) + .await + .expect("query promoted event") + { + return; + } + tokio::time::sleep(Duration::from_millis(50)).await; + } + }) + .await + .unwrap_or_else(|_| panic!("event {event_id} was not served after git data arrived")); +} + /// Publish a repo announcement, submit a matching state event, and push the /// deterministic git data so the announcement is promoted out of purgatory and /// becomes queryable. @@ -130,16 +148,7 @@ pub async fn publish_served_repo(client: &AuditClient, test_name: &str) -> (Even Err(e) => panic!("git push error while promoting repo: {}", e), } - // Give the relay a moment to promote the events out of purgatory. - tokio::time::sleep(Duration::from_millis(300)).await; - - assert!( - client - .is_event_on_relay(announcement.id) - .await - .expect("query announcement"), - "announcement should be served after git data arrives" - ); + wait_until_served(client, announcement.id).await; (announcement, repo_id) } @@ -253,15 +262,7 @@ pub async fn publish_served_repo_with_maintainers( Err(e) => panic!("git push error while promoting repo: {}", e), } - tokio::time::sleep(Duration::from_millis(300)).await; - - assert!( - client - .is_event_on_relay(announcement.id) - .await - .expect("query announcement"), - "announcement should be served after git data arrives" - ); + wait_until_served(client, announcement.id).await; (announcement, repo_id) } @@ -320,22 +321,8 @@ pub async fn publish_served_repo_with_state_event( push_to_relay(temp_dir.path(), &relay_domain, &npub, &repo_id) .expect("git push should promote announcement + state event out of purgatory"); - tokio::time::sleep(Duration::from_millis(300)).await; - - assert!( - client - .is_event_on_relay(announcement.id) - .await - .expect("query announcement"), - "announcement should be served after git data arrives" - ); - assert!( - client - .is_event_on_relay(state_event.id) - .await - .expect("query state event"), - "state event should be served after git data arrives" - ); + wait_until_served(client, announcement.id).await; + wait_until_served(client, state_event.id).await; (announcement, repo_id, state_event) } @@ -406,22 +393,8 @@ pub async fn publish_served_repo_with_state_event_and_maintainers( push_to_relay(temp_dir.path(), &relay_domain, &npub, &repo_id) .expect("git push should promote announcement + state event out of purgatory"); - tokio::time::sleep(Duration::from_millis(300)).await; - - assert!( - client - .is_event_on_relay(announcement.id) - .await - .expect("query announcement"), - "announcement should be served after git data arrives" - ); - assert!( - client - .is_event_on_relay(state_event.id) - .await - .expect("query state event"), - "state event should be served after git data arrives" - ); + wait_until_served(client, announcement.id).await; + wait_until_served(client, state_event.id).await; (announcement, repo_id, state_event) } @@ -483,8 +456,6 @@ pub async fn publish_served_announcement_for_identifier( .await .expect("relay should accept state event"); - tokio::time::sleep(Duration::from_millis(300)).await; - if client .is_event_on_relay(announcement.id) .await @@ -536,13 +507,7 @@ pub async fn publish_served_announcement_for_identifier( Err(e) => panic!("git push error while promoting repo: {}", e), } - assert!( - client - .is_event_on_relay(announcement.id) - .await - .expect("query announcement"), - "announcement should be served after git data arrives" - ); + wait_until_served(client, announcement.id).await; announcement } @@ -599,8 +564,6 @@ pub async fn publish_served_announcement_with_state_for_identifier( .await .expect("relay should accept state event"); - tokio::time::sleep(Duration::from_millis(300)).await; - let announcement_served = client .is_event_on_relay(announcement.id) .await @@ -656,20 +619,8 @@ pub async fn publish_served_announcement_with_state_for_identifier( Err(e) => panic!("git push error while promoting repo: {}", e), } - assert!( - client - .is_event_on_relay(announcement.id) - .await - .expect("query announcement"), - "announcement should be served after git data arrives" - ); - assert!( - client - .is_event_on_relay(state_event.id) - .await - .expect("query state event"), - "state event should be served after git data arrives" - ); + wait_until_served(client, announcement.id).await; + wait_until_served(client, state_event.id).await; (announcement, state_event) } diff --git a/tests/common/sync_helpers.rs b/tests/common/sync_helpers.rs index e7dd955..75e9ae9 100644 --- a/tests/common/sync_helpers.rs +++ b/tests/common/sync_helpers.rs @@ -15,7 +15,7 @@ use std::time::Duration; use nostr_sdk::prelude::*; -use super::port::{self, PortReservation}; +use super::port::{self, PortReservation, UnavailableEndpoint}; use super::relay::TestRelay; const DESCENDANT_LIVE_LOG: &str = "Installed priority-bounded auxiliary live coverage"; @@ -476,80 +476,26 @@ pub async fn wait_for_sync_connection( expected_connections: usize, timeout: Duration, ) -> Result<(), String> { - // Convert ws:// URL to http:// for metrics endpoint - let http_url = syncing_relay_url - .replace("ws://", "http://") - .replace("/", "") - + "/metrics"; - - let start = std::time::Instant::now(); - let poll_interval = Duration::from_millis(100); - - while start.elapsed() < timeout { - // Fetch metrics - if let Ok(response) = reqwest::get(&http_url).await { - if let Ok(metrics) = response.text().await { - // Look for sync connection metrics - // The metric name pattern: ngit_sync_connections or similar - // We check for any indication of established connections + let wait = async { + loop { + if let Ok(metrics) = fetch_metrics(syncing_relay_url).await { if check_sync_connections_in_metrics(&metrics, expected_connections) { - return Ok(()); + return; } } + tokio::time::sleep(Duration::from_millis(100)).await; } - - tokio::time::sleep(poll_interval).await; - } - - Err(format!( - "Timeout waiting for {} sync connection(s) on {} after {:?}", - expected_connections, syncing_relay_url, timeout + }; + tokio::time::timeout(timeout, wait).await.map_err(|_| format!( + "Timeout waiting for {expected_connections} sync connection(s) on {syncing_relay_url} after {timeout:?}" )) } -/// Check metrics string for expected number of sync connections. -/// -/// Looks for various metric patterns that indicate sync connections: -/// - ngit_sync_connections (gauge) -/// - ngit_sync_relay_connections (gauge) -/// - Any metric containing "sync" and "connection" with count > 0 +/// Connection attempts and health states do not establish a live connection. fn check_sync_connections_in_metrics(metrics: &str, expected: usize) -> bool { - // Parse metrics line by line looking for connection counts - for line in metrics.lines() { - // Skip comments and empty lines - if line.starts_with('#') || line.is_empty() { - continue; - } - - // Look for sync connection metrics - // Format: metric_name{labels} value - // or: metric_name value - if line.contains("sync") && line.contains("connect") { - // Extract the value (last space-separated token) - if let Some(value_str) = line.split_whitespace().last() { - if let Ok(value) = value_str.parse::() { - if value as usize >= expected { - return true; - } - } - } - } - - // Also check for specific metric names that might indicate connections - // ngit_sync_health_state with value 1 or 2 (connecting/healthy) - if line.contains("ngit_sync_health") { - if let Some(value_str) = line.split_whitespace().last() { - if let Ok(value) = value_str.parse::() { - // Health state > 0 typically means connection attempt or established - if value > 0.0 && expected > 0 { - return true; - } - } - } - } - } - - false + ParsedMetrics::parse(metrics) + .relays_connected_total() + .is_some_and(|connected| connected >= expected as i64) } // ============================================================================ @@ -677,10 +623,14 @@ pub fn repo_coord(keys: &Keys, identifier: &str) -> String { /// assert!(metrics.contains("ngit_sync_")); /// ``` pub async fn fetch_metrics(relay_url: &str) -> Result { - // Convert ws:// URL to http:// for metrics endpoint - let http_url = relay_url.replace("ws://", "http://").replace("/", "") + "/metrics"; + reqwest::get(metrics_url(relay_url)).await?.text().await +} - reqwest::get(&http_url).await?.text().await +fn metrics_url(relay_url: &str) -> String { + let http_url = relay_url + .replacen("wss://", "https://", 1) + .replacen("ws://", "http://", 1); + format!("{}/metrics", http_url.trim_end_matches('/')) } // ============================================================================ @@ -837,10 +787,9 @@ impl ParsedMetrics { /// harness.stop_all().await; /// ``` pub struct MetricsTestHarness { - source_relays: Vec, + source_relays: Vec>, syncing_relay: Option, - #[allow(dead_code)] - nowhere_url: Option, + unavailable_endpoint: Option, } impl MetricsTestHarness { @@ -848,34 +797,36 @@ impl MetricsTestHarness { pub async fn with_sources(count: usize) -> Self { let mut source_relays = Vec::new(); for _ in 0..count { - source_relays.push(TestRelay::start().await); + source_relays.push(Some(TestRelay::start().await)); } Self { source_relays, syncing_relay: None, - nowhere_url: None, + unavailable_endpoint: None, } } /// Get source relay URL pub fn source_url(&self, idx: usize) -> &str { - self.source_relays[idx].url() + self.source_relay(idx).url() } /// Get source relay domain (for announcement tags) pub fn source_domain(&self, idx: usize) -> String { - self.source_relays[idx].domain() + self.source_relay(idx).domain() } /// Get a reference to a source relay (for advanced test operations) pub fn source_relay(&self, idx: usize) -> &TestRelay { - &self.source_relays[idx] + self.source_relays[idx] + .as_ref() + .expect("source relay has been stopped") } /// Submit events to a specific source relay pub async fn submit_events(&self, source_idx: usize, events: &[Event]) -> Result<(), String> { - let relay = &self.source_relays[source_idx]; + let relay = self.source_relay(source_idx); let keys = Keys::generate(); let client = TestClient::new(relay.url(), keys).await?; @@ -889,7 +840,7 @@ impl MetricsTestHarness { /// Start syncing relay pointing to source[idx] pub async fn start_syncing_relay(&mut self, source_idx: usize) { - let source_url = self.source_relays[source_idx].url().to_string(); + let source_url = self.source_relay(source_idx).url().to_string(); self.syncing_relay = Some(TestRelay::start_with_sync(Some(source_url)).await); } @@ -904,37 +855,26 @@ impl MetricsTestHarness { source_idx: usize, reservation: PortReservation, ) { - let source_url = self.source_relays[source_idx].url().to_string(); + let source_url = self.source_relay(source_idx).url().to_string(); self.syncing_relay = Some( TestRelay::start_on_reservation_with_options(reservation, Some(source_url), false) .await, ); } - /// Start syncing relay pointing to random unused port (for failure tests) + /// Start syncing against an owned endpoint that closes every connection. pub async fn start_syncing_relay_to_nowhere(&mut self) { - let port = random_unused_port(); - let nowhere_url = format!("ws://127.0.0.1:{}", port); - self.nowhere_url = Some(nowhere_url.clone()); + let endpoint = UnavailableEndpoint::new(); + let nowhere_url = format!("ws://127.0.0.1:{}", endpoint.port()); + self.unavailable_endpoint = Some(endpoint); self.syncing_relay = Some(TestRelay::start_with_sync(Some(nowhere_url)).await); } /// Stop a source relay pub async fn stop_source(&mut self, source_idx: usize) { - // We need to take ownership to stop, so we swap with a new relay - // that we immediately stop. This is a workaround since TestRelay::stop - // takes self by value. - let relay = std::mem::replace( - &mut self.source_relays[source_idx], - TestRelay::start().await, - ); - relay.stop().await; - // Stop the placeholder too - let placeholder = std::mem::replace( - &mut self.source_relays[source_idx], - TestRelay::start().await, - ); - placeholder.stop().await; + if let Some(relay) = self.source_relays[source_idx].take() { + relay.stop().await; + } } /// Fetch and parse metrics from syncing relay @@ -961,29 +901,41 @@ impl MetricsTestHarness { if let Some(relay) = self.syncing_relay.take() { relay.stop().await; } - for relay in self.source_relays.drain(..) { + for relay in self.source_relays.drain(..).flatten() { relay.stop().await; } } } -// ============================================================================ -// Port Helpers -// ============================================================================ - -/// Get a random unused port by binding to port 0 and letting the OS assign one -pub fn random_unused_port() -> u16 { - std::net::TcpListener::bind("127.0.0.1:0") - .expect("Failed to bind to random port") - .local_addr() - .expect("Failed to get local addr") - .port() -} - #[cfg(test)] mod tests { use super::*; + #[test] + fn metrics_url_preserves_scheme_and_base_path() { + assert_eq!( + metrics_url("ws://127.0.0.1:1234/"), + "http://127.0.0.1:1234/metrics" + ); + assert_eq!( + metrics_url("wss://example.test/relay/"), + "https://example.test/relay/metrics" + ); + } + + #[test] + fn connection_readiness_ignores_attempts_and_health() { + assert!(!check_sync_connections_in_metrics("ngit_sync_connection_attempts_total 5\nngit_sync_health 1\nngit_sync_relays_connected_total 0", 1)); + assert!(check_sync_connections_in_metrics( + "ngit_sync_relays_connected_total 2", + 2 + )); + assert!(!check_sync_connections_in_metrics( + "ngit_sync_relays_connected_total 1", + 2 + )); + } + #[test] fn test_repo_coord_format() { let keys = Keys::generate(); @@ -1268,8 +1220,15 @@ pub async fn push_git_data_to_relay( push_to_relay(git_temp_dir.path(), &relay.domain(), &npub, identifier) .expect("Failed to push git data to relay"); - // Brief wait for push processing - tokio::time::sleep(Duration::from_millis(500)).await; + assert!( + wait_for_event_on_relay( + relay.url(), + Filter::new().id(state_event.id), + Duration::from_secs(10) + ) + .await, + "pushed state event must leave purgatory" + ); git_temp_dir } @@ -1379,7 +1338,15 @@ pub async fn push_unique_git_data_to_relay( push_to_relay(path, &relay.domain(), &npub, identifier) .expect("Failed to push git data to relay"); - tokio::time::sleep(Duration::from_millis(500)).await; + assert!( + wait_for_event_on_relay( + relay.url(), + Filter::new().id(state_event.id), + Duration::from_secs(10) + ) + .await, + "pushed state event must leave purgatory" + ); git_temp_dir } @@ -1464,8 +1431,25 @@ pub async fn setup_announcement_on_relay( push_to_relay(git_temp_dir.path(), &relay.domain(), &npub, identifier) .expect("Failed to push git data to relay"); - // Brief wait for push processing - tokio::time::sleep(Duration::from_millis(500)).await; + assert!( + wait_for_event_on_relay( + relay.url(), + Filter::new().id(state_event.id), + Duration::from_secs(10) + ) + .await, + "pushed state event must leave purgatory" + ); + + assert!( + wait_for_event_on_relay( + relay.url(), + Filter::new().id(announcement.id), + Duration::from_secs(10) + ) + .await, + "pushed announcement must leave purgatory" + ); (announcement, git_temp_dir) } @@ -1587,7 +1571,15 @@ pub async fn run_sync_test(historic_events: &[Event], live_events: &[Event]) -> .expect("Failed to push git data to source relay"); // 8. Wait for source relay to process the push and release events from purgatory - tokio::time::sleep(Duration::from_secs(2)).await; + assert!( + wait_for_event_on_relay( + source.url(), + Filter::new().id(announcement.id), + Duration::from_secs(10) + ) + .await, + "source announcement must leave purgatory" + ); // 9. Send historic events to source BEFORE syncing relay starts for event in historic_events { @@ -1605,7 +1597,9 @@ pub async fn run_sync_test(historic_events: &[Event], live_events: &[Event]) -> .await; // 11. Wait for sync connection to establish - let _ = wait_for_sync_connection(syncing.url(), 1, Duration::from_secs(5)).await; + wait_for_sync_connection(syncing.url(), 1, Duration::from_secs(10)) + .await + .expect("sync connection must establish before publishing live events"); // 12. Send live events AFTER connection established for event in live_events { @@ -1614,12 +1608,16 @@ pub async fn run_sync_test(historic_events: &[Event], live_events: &[Event]) -> .expect("Failed to send live event"); } - // 13. Allow sync + purgatory promotion to complete on the syncing relay. - // The syncing relay receives the announcement (goes to purgatory) and state event. - // The purgatory sync loop (1s interval) fetches git data from source's clone URL - // (http://source-domain/npub/test-repo.git) and releases the announcement. - // We wait up to 8s to allow time for this. - tokio::time::sleep(Duration::from_secs(8)).await; + // 13. Observe announcement promotion rather than guessing a sync duration. + assert!( + wait_for_event_on_relay( + syncing.url(), + Filter::new().id(announcement.id), + Duration::from_secs(15) + ) + .await, + "synced announcement must leave purgatory" + ); // 14. Compute repo coordinate before moving keys let coordinate = repo_coord(&keys, "test-repo");