diff --git a/tests/common/sync_helpers.rs b/tests/common/sync_helpers.rs index 668b5fc..4084a05 100644 --- a/tests/common/sync_helpers.rs +++ b/tests/common/sync_helpers.rs @@ -135,52 +135,58 @@ impl TestClient { Err("Connection loop exited unexpectedly".to_string()) } - /// Send an event with retry logic. + /// Send an event with bounded retry for transient transport failures. /// - /// Attempts to send up to 3 times with exponential backoff: - /// - Attempt 1: immediate - /// - Attempt 2: after 200ms - /// - Attempt 3: after 400ms + /// Attempts to send up to 3 times with short backoff (200ms, then 400ms), + /// reconnecting between attempts. Only errors that leave delivery in doubt + /// (disconnection, transport failure, OK timeout) are retried; a relay + /// `OK false` rejection is deterministic and fails immediately so tests + /// cannot mask write-policy regressions. /// /// # Arguments /// * `event` - The signed event to send /// /// # Returns /// * `Ok(EventId)` on successful send - /// * `Err(String)` if all attempts fail + /// * `Err(String)` on relay rejection or after all attempts fail pub async fn send_event(&self, event: &Event) -> Result { - let delays = [0, 200, 400]; // Exponential backoff in ms + let delays = [0, 200, 400]; // Backoff in ms + let mut last_error = String::new(); for (attempt, delay_ms) in delays.iter().enumerate() { if *delay_ms > 0 { tokio::time::sleep(Duration::from_millis(*delay_ms)).await; } - match self.client.send_event(event).await { - Ok(output) => { - if !output.success.is_empty() { - return Ok(output.value); - } - // Log failures for debugging - if !output.failed.is_empty() { - eprintln!( - " Send attempt {} - failures: {:?}", - attempt + 1, - output.failed - ); - // Try reconnecting if relay disconnected - self.client.connect().await; - } + // Send through the single tracked relay rather than the pool so + // the SDK error kind survives and OK-false can be told apart from + // transport failures. + let relay = self + .client + .relay(self.relay_url.as_str()) + .await + .map_err(|e| format!("Failed to look up relay {}: {}", self.relay_url, e))? + .ok_or_else(|| format!("Relay {} missing from client pool", self.relay_url))?; + + match relay.send_event(event).await { + Ok(output) => return Ok(*output.id()), + Err(e) if e.kind() == ErrorKind::Rejected => { + return Err(format!("Relay rejected event {}: {}", event.id, e)); } Err(e) => { - eprintln!(" Send attempt {} - error: {}", attempt + 1, e); + eprintln!(" Send attempt {} - transient error: {}", attempt + 1, e); + last_error = e.to_string(); + // Re-establish the connection before the next attempt. + let _ = self.connect().await; } } } Err(format!( - "Failed to send event {} after 3 attempts", - event.id + "Failed to send event {} after {} attempts: {}", + event.id, + delays.len(), + last_error )) } diff --git a/tests/sync/discovery.rs b/tests/sync/discovery.rs index bfcc2cc..84f549a 100644 --- a/tests/sync/discovery.rs +++ b/tests/sync/discovery.rs @@ -106,14 +106,13 @@ async fn test_discovers_layer3_via_layer2() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b (should trigger discovery of relay_a)"); - // 9. Wait for relay_b to discover relay_a and sync the patch - println!("Waiting 3s for relay_b to discover relay_a and sync patch..."); - tokio::time::sleep(Duration::from_secs(3)).await; - - // 10. Verify patch was synced to relay_b + // 9/10. Verify the patch syncs to relay_b. The bounded poll on the synced + // event replaces a fixed discovery sleep, which is not a reliable proxy + // under CI load. let filter = Filter::new().kind(Kind::GitPatch).author(keys.public_key()); - let patch_synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(5)).await; + let patch_synced = + wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(30)).await; if patch_synced { println!( @@ -212,14 +211,12 @@ async fn test_relay_discovery_via_announcements_with_historic_sync() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b (should trigger discovery of relay_a)"); - // 7. Wait for sync - println!("Waiting 3s for Layer 2 sync..."); - tokio::time::sleep(Duration::from_secs(3)).await; - - // 8. Verify Layer 2 event synced to relay_b + // 7/8. Verify the Layer 2 event syncs to relay_b. The bounded poll on the + // synced event replaces a fixed discovery sleep, which is not a reliable + // proxy under CI load. let issue_filter = Filter::new().kind(Kind::GitIssue).author(keys.public_key()); let issue_synced = - wait_for_event_on_relay(relay_b.url(), issue_filter, Duration::from_secs(5)).await; + wait_for_event_on_relay(relay_b.url(), issue_filter, Duration::from_secs(30)).await; println!("Sync result:"); println!(" Issue {} synced: {}", issue_id, issue_synced); diff --git a/tests/sync/historic_sync.rs b/tests/sync/historic_sync.rs index a4ca7df..8db7a43 100644 --- a/tests/sync/historic_sync.rs +++ b/tests/sync/historic_sync.rs @@ -33,7 +33,7 @@ async fn test_bootstrap_syncs_existing_layer2_events() { .author(result.maintainer_keys.public_key()); let synced = - wait_for_event_on_relay(result.syncing_relay.url(), filter, Duration::from_secs(5)).await; + wait_for_event_on_relay(result.syncing_relay.url(), filter, Duration::from_secs(30)).await; // Cleanup result.syncing_relay.stop().await; @@ -70,15 +70,15 @@ async fn test_relay_replays_events_after_restart() { let synced_first = wait_for_event_on_relay( result.syncing_relay.url(), filter.clone(), - Duration::from_secs(5), + Duration::from_secs(30), ) .await; println!("First sync check: {}", synced_first); - // Stop syncing relay (simulates restart) + // Stop syncing relay (simulates restart). `stop` waits on the process, so + // the replacement instance can start immediately. result.syncing_relay.stop().await; - tokio::time::sleep(Duration::from_millis(500)).await; // Restart syncing relay (new instance with same bootstrap config) // Note: The new syncing relay will have a different domain, so it may not @@ -90,10 +90,10 @@ async fn test_relay_replays_events_after_restart() { syncing_new.domain() ); - // Wait for re-sync - tokio::time::sleep(Duration::from_secs(2)).await; - - // Verify announcement is available on restarted syncing relay + // Check whether the announcement re-syncs to the restarted relay. The + // bounded poll is kept short because this outcome is informational only: + // the new instance has a different domain, so the announcement may + // legitimately never be accepted (see the note below). let synced_after_restart = wait_for_event_on_relay(syncing_new.url(), filter, Duration::from_secs(5)).await; @@ -135,7 +135,7 @@ async fn test_announcement_not_listing_relay_is_not_synced() { let keys = Keys::generate(); // Wait for sync connection to establish - match wait_for_sync_connection(syncing.url(), 1, Duration::from_secs(5)).await { + match wait_for_sync_connection(syncing.url(), 1, Duration::from_secs(30)).await { Ok(()) => println!("Sync connection established (verified via metrics)"), Err(e) => println!("Sync connection check: {} (continuing with test)", e), } @@ -168,15 +168,14 @@ async fn test_announcement_not_listing_relay_is_not_synced() { client.disconnect().await; - // Wait for potential sync attempt - tokio::time::sleep(Duration::from_secs(3)).await; - - // Verify announcement did NOT sync to syncing relay + // Verify announcement did NOT sync to syncing relay. The bounded wait is + // the observation window for a wrongful sync; the sync connection was + // already confirmed above, so any policy failure has this long to appear. let filter = Filter::new() .kind(Kind::GitRepoAnnouncement) .author(keys.public_key()); - let synced = wait_for_event_on_relay(syncing.url(), filter, Duration::from_secs(2)).await; + let synced = wait_for_event_on_relay(syncing.url(), filter, Duration::from_secs(5)).await; // Cleanup syncing.stop().await; @@ -246,8 +245,17 @@ async fn test_history_sync_without_negentropy() { announcement_id ); - // Wait to ensure event is stored - tokio::time::sleep(Duration::from_millis(500)).await; + // Confirm the announcement is stored and served by the source before the + // syncing relay ever connects, so this genuinely exercises history sync. + assert!( + wait_for_event_on_relay( + source.url(), + Filter::new().id(announcement.id), + Duration::from_secs(30), + ) + .await, + "announcement should be served by the source before the syncing relay starts" + ); // NOW start syncing relay on the reserved port, with negentropy DISABLED // This syncing relay has never connected before - it needs to do HISTORY sync @@ -263,15 +271,13 @@ async fn test_history_sync_without_negentropy() { syncing.domain() ); - // Wait for history sync to complete (using REQ+EOSE, not negentropy) - tokio::time::sleep(Duration::from_secs(3)).await; - - // Verify announcement synced to syncing relay via HISTORY sync + // Verify announcement syncs to the syncing relay via HISTORY sync. The + // bounded poll replaces a fixed sleep plus short deadline. let filter = Filter::new() .kind(Kind::GitRepoAnnouncement) .author(keys.public_key()); - let synced = wait_for_event_on_relay(syncing.url(), filter, Duration::from_secs(5)).await; + let synced = wait_for_event_on_relay(syncing.url(), filter, Duration::from_secs(30)).await; // Cleanup syncing.stop().await; @@ -368,8 +374,17 @@ async fn test_pagination_for_large_historic_sync() { .expect("Failed to send announcement to source"); println!("Announcement sent to source"); - // Wait for announcement to be stored - tokio::time::sleep(Duration::from_millis(200)).await; + // Confirm the announcement is stored before sending the issues that + // reference it. + assert!( + wait_for_event_on_relay( + source.url(), + Filter::new().id(announcement.id), + Duration::from_secs(30), + ) + .await, + "announcement should be stored on the source relay" + ); // Send all 40 issue events to source (before syncing relay starts) println!("Sending {} issues to source relay...", issue_events.len()); @@ -391,8 +406,18 @@ async fn test_pagination_for_large_historic_sync() { client.disconnect().await; - // Wait to ensure all events are stored - tokio::time::sleep(Duration::from_millis(500)).await; + // Confirm the last-sent issue is stored so every event exists before the + // syncing relay connects. + let last_issue = issue_events.last().expect("at least one issue"); + assert!( + wait_for_event_on_relay( + source.url(), + Filter::new().id(last_issue.id), + Duration::from_secs(30), + ) + .await, + "all issues should be stored on the source relay" + ); // NOW start syncing relay on the reserved port, with negentropy DISABLED // This forces it to use REQ+EOSE historic sync with pagination @@ -408,17 +433,14 @@ async fn test_pagination_for_large_historic_sync() { syncing.domain() ); - // Wait for historic sync with pagination to complete - println!("Waiting for historic sync with pagination to complete..."); - tokio::time::sleep(Duration::from_secs(8)).await; - - // Verify announcement synced + // Verify announcement synced. The bounded poll replaces a fixed + // "wait for pagination" sleep plus short deadline. let announcement_filter = Filter::new() .kind(Kind::GitRepoAnnouncement) .author(keys.public_key()); let announcement_synced = - wait_for_event_on_relay(syncing.url(), announcement_filter, Duration::from_secs(3)).await; + wait_for_event_on_relay(syncing.url(), announcement_filter, Duration::from_secs(30)).await; // Verify ALL 40 issues synced let issues_filter = Filter::new().kind(Kind::GitIssue).author(keys.public_key()); @@ -432,18 +454,21 @@ async fn test_pagination_for_large_historic_sync() { .add_relay(syncing.url()) .await .expect("Failed to add syncing relay to client"); - client.connect().await; + client.connect().and_wait(Duration::from_secs(30)).await; - // Wait for connection - tokio::time::sleep(Duration::from_millis(500)).await; - - let synced_issues = client - .fetch_events(issues_filter) - .timeout(Duration::from_secs(5)) - .await - .expect("Failed to fetch issues from syncing relay"); - - let synced_count = synced_issues.len(); + // Poll until every issue has paginated across, up to a bounded deadline. + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + let synced_count = loop { + let synced_issues = client + .fetch_events(issues_filter.clone()) + .timeout(Duration::from_secs(5)) + .await + .expect("Failed to fetch issues from syncing relay"); + if synced_issues.len() >= 40 || tokio::time::Instant::now() >= deadline { + break synced_issues.len(); + } + tokio::time::sleep(Duration::from_millis(200)).await; + }; println!("Synced {} out of 40 expected issues", synced_count); client.disconnect().await; diff --git a/tests/sync/live_sync.rs b/tests/sync/live_sync.rs index dc1582d..36dd2a3 100644 --- a/tests/sync/live_sync.rs +++ b/tests/sync/live_sync.rs @@ -486,8 +486,7 @@ async fn test_live_sync_layer3_events() { .authenticator(SignerAuthenticator::new(temp_keys)) .build(); if client.add_relay(relay_b.url()).await.is_ok() { - client.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + client.connect().and_wait(Duration::from_secs(30)).await; let fetch_filter = Filter::new().kind(Kind::Comment).id(comment_id); @@ -644,8 +643,7 @@ async fn test_live_sync_event_ordering() { let events_found: Vec; if client.add_relay(relay_b.url()).await.is_ok() { - client.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + client.connect().and_wait(Duration::from_secs(30)).await; let filter = Filter::new().kind(Kind::GitIssue).author(keys.public_key()); diff --git a/tests/sync/maintainer_reprocessing.rs b/tests/sync/maintainer_reprocessing.rs index 75438b3..4abbc05 100644 --- a/tests/sync/maintainer_reprocessing.rs +++ b/tests/sync/maintainer_reprocessing.rs @@ -508,6 +508,8 @@ async fn test_invitee_only_acceptance_recovers_cold_owner_invitation() { .await, "Invitee-only server should reject and index the owner invitation" ); + // Passage of the configured one-second hot-cache TTL is required so the + // acceptance below exercises cold exact-ID recovery, not the hot copy. tokio::time::sleep(Duration::from_secs(2)).await; let pushes_before_acceptance = [ @@ -767,6 +769,8 @@ async fn test_invitation_applies_newer_invitee_state_to_owner_before_acceptance( .await; } + // Deliberate timestamp spacing: the invitee events below are dated one + // second after the owner state, so let wall-clock time pass that mark. tokio::time::sleep(Duration::from_secs(1)).await; let invitee_git = tempfile::tempdir().expect("Failed to create invitee repository directory"); @@ -982,6 +986,8 @@ async fn test_acceptance_replaces_existing_invitee_repository_with_newer_owner_s .await; } + // Deliberate timestamp spacing: the owner events below are dated one + // second after the invitee state, so let wall-clock time pass that mark. tokio::time::sleep(Duration::from_secs(1)).await; let owner_git = tempfile::tempdir().expect("Failed to create owner repository directory"); @@ -1062,6 +1068,8 @@ async fn test_acceptance_replaces_existing_invitee_repository_with_newer_owner_s .await, "Invitee-only server should process the unilateral invitation before acceptance" ); + // Deliberate observation window: give a wrongful state replacement time + // to manifest before asserting the invitee state is still intact. tokio::time::sleep(Duration::from_millis(500)).await; for relay in invitee_servers { assert_exact_event_served(relay, &invitee_state, "Invitee state before acceptance").await; @@ -1395,9 +1403,22 @@ async fn test_maintainer_announcement_reprocessed_immediately() { let relay_b = TestRelay::start_with_sync(Some(relay_a.url().to_string())).await; println!("relay_b started at {}", relay_b.url()); - // Give relay_b's SyncManager time to complete the initial negentropy sync with relay_a. - tokio::time::sleep(Duration::from_secs(3)).await; - println!("✓ relay_b synced from relay_a (maintainer announcement should be in hot cache)"); + // Wait for relay_b's initial negentropy sync on its observable outcome: + // the rejected maintainer announcement is logged by its note ID. + let maintainer_note = maintainer_announcement + .id + .to_bech32() + .expect("Failed to encode maintainer announcement ID"); + assert!( + wait_for_log( + &relay_b.log_path(), + &maintainer_note, + Duration::from_secs(30), + ) + .await, + "relay_b should sync, reject and index the maintainer announcement" + ); + println!("✓ relay_b synced from relay_a (maintainer announcement in hot cache)"); let start = std::time::Instant::now(); @@ -1441,19 +1462,15 @@ async fn test_maintainer_announcement_reprocessed_immediately() { "✓ Owner git data pushed to relay_b (owner announcement promoted, hot cache re-processed)" ); - // Step 5: Wait briefly for async processing to complete. - tokio::time::sleep(Duration::from_secs(1)).await; - - let elapsed = start.elapsed(); - - // Step 6: Verify both announcements are in relay_b's database. + // Step 5/6: Verify both announcements are in relay_b's database. Bounded + // polls on the served events replace a fixed post-push sleep. let owner_filter = Filter::new() .kind(Kind::GitRepoAnnouncement) .author(owner_keys.public_key()) .identifier(identifier); let owner_found = - wait_for_event_on_relay(relay_b.url(), owner_filter, Duration::from_secs(2)).await; + wait_for_event_on_relay(relay_b.url(), owner_filter, Duration::from_secs(30)).await; assert!(owner_found, "Owner announcement should be in relay_b"); let maintainer_filter = Filter::new() @@ -1462,12 +1479,14 @@ async fn test_maintainer_announcement_reprocessed_immediately() { .identifier(identifier); let maintainer_found = - wait_for_event_on_relay(relay_b.url(), maintainer_filter, Duration::from_secs(2)).await; + wait_for_event_on_relay(relay_b.url(), maintainer_filter, Duration::from_secs(30)).await; assert!( maintainer_found, "Maintainer announcement should be re-processed and accepted in relay_b" ); + // Immediate hot-cache re-processing must beat the scheduled retry path. + let elapsed = start.elapsed(); assert!( elapsed.as_secs() < 15, "Re-processing should happen in <15 seconds, took {:?}", @@ -1514,6 +1533,7 @@ async fn test_multiple_maintainers_all_reprocessed() { // land in relay_a's DB. Each announcement lists relay_a only, so relay_b will reject // them when syncing (no owner announcement in relay_b's DB yet). let mut git_dirs = Vec::new(); + let mut maintainer_notes = Vec::new(); for (idx, maintainer_keys) in [&maintainer1_keys, &maintainer2_keys, &maintainer3_keys] .iter() .enumerate() @@ -1541,6 +1561,12 @@ async fn test_multiple_maintainers_all_reprocessed() { ]) .finalize(*maintainer_keys) .unwrap(); + maintainer_notes.push( + announcement + .id + .to_bech32() + .expect("Failed to encode maintainer announcement ID"), + ); send_to_relay(&relay_a, &announcement).await.unwrap(); // Use push_unique_git_data_to_relay so each maintainer gets a distinct commit // hash. Identical hashes cause git to skip pack transfer when the object @@ -1584,11 +1610,15 @@ async fn test_multiple_maintainers_all_reprocessed() { let relay_b = TestRelay::start_with_sync(Some(relay_a.url().to_string())).await; println!("relay_b started at {}", relay_b.url()); - // Give relay_b's SyncManager time to complete the initial negentropy sync with relay_a. - // The negentropy sync completes within ~200ms (NGIT_TEST=1 sets batch window to 200ms), but we - // allow extra time for slow CI environments. - tokio::time::sleep(Duration::from_secs(3)).await; - println!("✓ relay_b synced from relay_a (maintainer announcements should be in hot cache)"); + // Wait for relay_b's initial negentropy sync on its observable outcome: + // every rejected maintainer announcement is logged by its note ID. + for note in &maintainer_notes { + assert!( + wait_for_log(&relay_b.log_path(), note, Duration::from_secs(30)).await, + "relay_b should sync, reject and index maintainer announcement {note}" + ); + } + println!("✓ relay_b synced from relay_a (maintainer announcements in hot cache)"); // Step 3: Send owner announcement to relay_b → goes to purgatory. let owner_npub = owner_keys @@ -1634,10 +1664,8 @@ async fn test_multiple_maintainers_all_reprocessed() { push_git_data_to_relay(&relay_b, &owner_keys, identifier, &[&relay_b.domain()]).await; println!("✓ Owner git data pushed to relay_b (hot-cache re-processing should fire)"); - // Step 5: Wait briefly for async processing to complete. - tokio::time::sleep(Duration::from_secs(1)).await; - - // Step 6: Verify all four announcements are in relay_b's database. + // Step 5/6: Verify all four announcements are in relay_b's database. + // Bounded polls on the served events replace a fixed post-push sleep. for (name, keys) in [ ("owner", &owner_keys), ("maintainer1", &maintainer1_keys), @@ -1649,7 +1677,7 @@ async fn test_multiple_maintainers_all_reprocessed() { .author(keys.public_key()) .identifier(identifier); - let found = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(2)).await; + let found = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(30)).await; assert!(found, "{} announcement should be in relay_b", name); } @@ -1695,9 +1723,10 @@ async fn test_invalid_maintainer_pubkey_handled_gracefully() { .finalize(&maintainer_keys) .unwrap(); - // Send maintainer announcement - expect it to be rejected + // Send maintainer announcement - expect it to be rejected. `send_event` + // waits for the relay's OK response, so the rejection has already been + // processed when it returns. let _ = client.send_event(&maintainer_announcement).await; - tokio::time::sleep(Duration::from_millis(200)).await; // Step 2: Send owner announcement with INVALID maintainer hex, then push git data. // The announcement goes to purgatory first; the git push promotes it. @@ -1728,16 +1757,16 @@ async fn test_invalid_maintainer_pubkey_handled_gracefully() { send_to_relay(&relay, &owner_announcement).await.unwrap(); let _git_dir = push_git_data_to_relay(&relay, &owner_keys, identifier, &[&relay.domain()]).await; - tokio::time::sleep(Duration::from_millis(500)).await; - // Step 3: Verify owner announcement accepted, maintainer not re-processed + // Step 3: Verify owner announcement accepted, maintainer not re-processed. + // The bounded poll on the served announcement replaces a fixed sleep. let owner_filter = Filter::new() .kind(Kind::GitRepoAnnouncement) .author(owner_keys.public_key()) .identifier(identifier); let owner_found = - wait_for_event_on_relay(relay.url(), owner_filter, Duration::from_secs(2)).await; + wait_for_event_on_relay(relay.url(), owner_filter, Duration::from_secs(30)).await; assert!( owner_found, "Owner announcement should be accepted despite invalid maintainer" diff --git a/tests/sync/metrics.rs b/tests/sync/metrics.rs index 9c60bff..cbfb9be 100644 --- a/tests/sync/metrics.rs +++ b/tests/sync/metrics.rs @@ -28,15 +28,31 @@ use crate::common::{ // Format and Availability Tests (Keepers) // ============================================================================ +/// Poll the metrics endpoint until it responds, returning the first scrape. +/// +/// The HTTP endpoint can come up slightly after the relay's WebSocket +/// listener, so a bounded poll replaces the fixed startup sleeps these tests +/// used to rely on. +async fn wait_for_metrics_ready(relay_url: &str, timeout: Duration) -> String { + let deadline = tokio::time::Instant::now() + timeout; + loop { + if let Ok(metrics) = fetch_metrics(relay_url).await { + return metrics; + } + assert!( + tokio::time::Instant::now() < deadline, + "metrics endpoint at {relay_url} did not become ready before the deadline" + ); + tokio::time::sleep(Duration::from_millis(100)).await; + } +} + /// Test that Prometheus text format is valid #[tokio::test] async fn test_prometheus_format_valid() { let relay = TestRelay::start().await; - tokio::time::sleep(Duration::from_millis(500)).await; - let metrics = fetch_metrics(relay.url()) - .await - .expect("Failed to fetch metrics"); + let metrics = wait_for_metrics_ready(relay.url(), Duration::from_secs(30)).await; relay.stop().await; @@ -65,9 +81,10 @@ async fn test_metrics_availability_during_sync() { let source_relay = TestRelay::start().await; let sync_relay = TestRelay::start_with_sync(Some(source_relay.url().into())).await; - tokio::time::sleep(Duration::from_millis(500)).await; + wait_for_metrics_ready(sync_relay.url(), Duration::from_secs(30)).await; - // Make multiple metrics requests while sync is active + // Make multiple metrics requests while sync is active. The 200ms gap is + // deliberate scrape spacing, not a wait for a condition. for i in 0..3 { let metrics = fetch_metrics(sync_relay.url()).await; assert!( @@ -91,7 +108,7 @@ async fn test_concurrent_metrics_requests() { let source_relay = TestRelay::start().await; let sync_relay = TestRelay::start_with_sync(Some(source_relay.url().into())).await; - tokio::time::sleep(Duration::from_secs(1)).await; + wait_for_metrics_ready(sync_relay.url(), Duration::from_secs(30)).await; // Clone the URL string so we have an owned value for spawned tasks let sync_url: String = sync_relay.url().to_string(); @@ -135,11 +152,8 @@ async fn test_concurrent_metrics_requests() { #[tokio::test] async fn test_metric_values_are_numeric() { let relay = TestRelay::start().await; - tokio::time::sleep(Duration::from_millis(500)).await; - let metrics = fetch_metrics(relay.url()) - .await - .expect("Should fetch metrics"); + let metrics = wait_for_metrics_ready(relay.url(), Duration::from_secs(30)).await; relay.stop().await; @@ -276,9 +290,18 @@ async fn test_startup_sync_event_count() { setup_announcement_on_relay(&syncing_relay, &keys, &domain_refs, repo_id).await; println!("Announcement set up on syncing relay (triggers discovery of source)"); - // 9. Wait for discovery + sync to complete - println!("Waiting 5s for discovery and sync..."); - tokio::time::sleep(Duration::from_secs(5)).await; + // 9. Wait for discovery + sync on the observable condition the test + // cares about: the patches arriving on the syncing relay. + let patches_filter = Filter::new() + .kind(Kind::Custom(Kind::GitPatch.as_u16())) + .author(keys.public_key()); + let patches_synced = crate::common::sync_helpers::wait_for_event_on_relay( + syncing_relay.url(), + patches_filter, + Duration::from_secs(30), + ) + .await; + println!("Patches synced to syncing relay: {}", patches_synced); // 10. Fetch and parse metrics let raw_metrics = fetch_metrics(syncing_relay.url()) @@ -305,19 +328,6 @@ async fn test_startup_sync_event_count() { println!("Relays connected: {:?}", connected); println!("Events synced total: {:?}", events_synced); - // 12. Verify patches actually synced (functional check) - let filter = Filter::new() - .kind(Kind::Custom(Kind::GitPatch.as_u16())) - .author(keys.public_key()); - - let patches_synced = crate::common::sync_helpers::wait_for_event_on_relay( - syncing_relay.url(), - filter, - Duration::from_secs(2), - ) - .await; - println!("Patches synced to syncing relay: {}", patches_synced); - // Cleanup syncing_relay.stop().await; source_relay.stop().await; @@ -364,27 +374,30 @@ 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 attempt to the unreachable bootstrap relay - tokio::time::sleep(Duration::from_secs(2)).await; - - let metrics = harness.get_metrics().await.unwrap(); - - // Failure counter should be recorded when connecting to unreachable relay - let failures = metrics - .counter( - "ngit_sync_connection_attempts_total", - &[("result", "failure")], - ) - .unwrap_or(0); + // Poll for the failure counter recorded by the connection attempt to the + // unreachable bootstrap relay, bounded rather than sleeping a fixed time. + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + let failures = loop { + let metrics = harness.get_metrics().await.unwrap(); + let failures = metrics + .counter( + "ngit_sync_connection_attempts_total", + &[("result", "failure")], + ) + .unwrap_or(0); + if failures >= 1 { + break failures; + } + assert!( + tokio::time::Instant::now() < deadline, + "Expected at least 1 connection failure to be recorded, got {}", + failures + ); + tokio::time::sleep(Duration::from_millis(200)).await; + }; println!("Connection failures recorded: {}", failures); - assert!( - failures >= 1, - "Expected at least 1 connection failure to be recorded, got {}", - failures - ); - harness.stop_all().await; } @@ -473,16 +486,19 @@ async fn test_live_sync_event_count() { client.disconnect().await; println!("Two patches sent to source relay (live mode)"); - // Wait for live events to be processed and metrics updated - tokio::time::sleep(Duration::from_secs(4)).await; - - // Fetch metrics from syncing relay - let raw_metrics = fetch_metrics(&sync_url) - .await - .expect("Failed to fetch metrics"); - let metrics = ParsedMetrics::parse(&raw_metrics); - - let synced_count = metrics.events_synced_total(); + // Poll for the live events being processed and counted in metrics, + // bounded rather than sleeping a fixed time. + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + let synced_count = loop { + let raw_metrics = fetch_metrics(&sync_url) + .await + .expect("Failed to fetch metrics"); + let synced_count = ParsedMetrics::parse(&raw_metrics).events_synced_total(); + if synced_count.is_some_and(|count| count >= 2) || tokio::time::Instant::now() >= deadline { + break synced_count; + } + tokio::time::sleep(Duration::from_millis(200)).await; + }; println!("Events synced total: {:?}", synced_count); // Cleanup @@ -521,7 +537,7 @@ async fn test_relay_connected_status() { // Stop the source harness.stop_source(0).await; - let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); let metrics = loop { let metrics = harness.get_metrics().await.unwrap(); if metrics.relay_connected(&source_url).is_none() { @@ -559,28 +575,24 @@ async fn test_health_state_degrades_on_failure() { let mut harness = MetricsTestHarness::with_sources(0).await; harness.start_syncing_relay_to_nowhere().await; - // Initially might be trying to connect - tokio::time::sleep(Duration::from_secs(1)).await; - let initial = harness.get_metrics().await.unwrap(); - - // After several failures, should degrade (status = 2 or 3) - tokio::time::sleep(Duration::from_secs(5)).await; - let later = harness.get_metrics().await.unwrap(); - - // Get the relay status (1=healthy, 2=degraded, 3=dead) - let status = later.gauge("ngit_sync_relay_status", &[]).unwrap_or(0); - - println!( - "Initial metrics: {:?}", - initial.gauge("ngit_sync_relay_status", &[]) - ); - println!("Later status: {}", status); - - assert!( - status >= 2, - "Health should degrade to 2 (degraded) or 3 (dead), got {}", - status - ); + // Poll the relay status gauge (1=healthy, 2=degraded, 3=dead) until it + // reflects the repeated connection failures, bounded rather than sleeping + // a fixed time. + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); + let status = loop { + let metrics = harness.get_metrics().await.unwrap(); + let status = metrics.gauge("ngit_sync_relay_status", &[]).unwrap_or(0); + if status >= 2 { + break status; + } + assert!( + tokio::time::Instant::now() < deadline, + "Health should degrade to 2 (degraded) or 3 (dead), got {}", + status + ); + tokio::time::sleep(Duration::from_millis(200)).await; + }; + println!("Degraded status: {}", status); harness.stop_all().await; } @@ -617,13 +629,14 @@ async fn test_multi_source_aggregate_counts() { ); harness.submit_events(0, &[announcement]).await.unwrap(); - // Now start syncing relay - it should sync the existing announcement + // Now start syncing relay - it should sync the existing announcement. + // The loop below already polls for the connected state, so no fixed + // settling sleep is needed first. harness .start_syncing_relay_on_reservation(0, sync_reservation) .await; - tokio::time::sleep(Duration::from_secs(2)).await; - let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + let deadline = tokio::time::Instant::now() + Duration::from_secs(30); let metrics = loop { let metrics = harness.get_metrics().await.unwrap(); if metrics.relays_tracked_total() == Some(1) && metrics.relays_connected_total() == Some(1) diff --git a/tests/sync/purgatory_fetch.rs b/tests/sync/purgatory_fetch.rs index 0f09fcb..76d5788 100644 --- a/tests/sync/purgatory_fetch.rs +++ b/tests/sync/purgatory_fetch.rs @@ -158,8 +158,10 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() { .add_relay(source.url()) .await .expect("add source relay"); - source_client.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + source_client + .connect() + .and_wait(Duration::from_secs(30)) + .await; source_client .send_event(&source_announcement) .await @@ -236,8 +238,10 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() { .add_relay(mock.url()) .await .expect("add mock relay"); - mock_client.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + mock_client + .connect() + .and_wait(Duration::from_secs(30)) + .await; mock_client .send_event(&syncing_announcement) .await