diff --git a/tests/archive_grasp_services.rs b/tests/archive_grasp_services.rs index 8e36324..892463c 100644 --- a/tests/archive_grasp_services.rs +++ b/tests/archive_grasp_services.rs @@ -103,28 +103,16 @@ async fn start_relay_with_grasp_services(services: &str) -> (Child, String, Path (process, url, git_data_path) } -/// Wait for the relay to be ready to accept connections +/// Wait for the relay to be ready to accept connections. +/// +/// A successful TCP connect only proves the socket is bound: the accept loop +/// and HTTP service may still be coming up. Poll the relay's HTTP handler +/// instead, bounded by a deadline, so the immediately-following WebSocket +/// client is racing a relay that has already served a request. async fn wait_for_relay_ready(port: u16) { - let max_attempts = 50; // 5 seconds total - let delay = Duration::from_millis(100); - - for attempt in 0..max_attempts { - // Try to connect to the relay - match tokio::net::TcpStream::connect(format!("127.0.0.1:{}", port)).await { - Ok(_) => { - // Connection successful, relay is ready - // Give it a tiny bit more time to fully initialize - tokio::time::sleep(Duration::from_millis(100)).await; - return; - } - Err(_) => { - if attempt == max_attempts - 1 { - panic!("Relay failed to start after {} attempts", max_attempts); - } - tokio::time::sleep(delay).await; - } - } - } + common::relay::wait_for_http_ready(port, Duration::from_secs(10)) + .await + .expect("Relay should become ready"); } /// Test that announcements with matching GRASP service domains are accepted. @@ -161,18 +149,16 @@ async fn test_archive_accepts_matching_grasp_service() { .authenticator(SignerAuthenticator::new(keys.clone())) .build(); client.add_relay(&url).await.expect("Failed to add relay"); - client.connect().await; - - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&client).await; client .send_event(&announcement) .await .expect("Failed to send announcement"); - tokio::time::sleep(Duration::from_millis(500)).await; - - // Verify repository was created (announcement was accepted) + // The relay runs the archive write policy and creates the bare repository + // inside event admission, before it replies. The completed send is + // therefore the barrier: the directory already exists, or never will. let repo_path = git_data_path.join(format!("{}/{}.git", npub, identifier)); assert!( @@ -220,18 +206,16 @@ async fn test_archive_rejects_non_matching_grasp_service() { .authenticator(SignerAuthenticator::new(keys.clone())) .build(); client.add_relay(&url).await.expect("Failed to add relay"); - client.connect().await; - - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&client).await; client .send_event(&announcement) .await .expect("Failed to send announcement"); - tokio::time::sleep(Duration::from_millis(500)).await; - - // Verify repository was NOT created (announcement was rejected) + // Admission completes before the relay replies, and a rejected + // announcement never reaches repository creation, so the absence check + // needs no grace period after the send. let repo_path = git_data_path.join(format!("{}/{}.git", npub, identifier)); assert!( @@ -281,14 +265,12 @@ async fn test_archive_multiple_grasp_services() { .authenticator(SignerAuthenticator::new(keys1.clone())) .build(); client1.add_relay(&url).await.expect("Failed to add relay"); - client1.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&client1).await; client1 .send_event(&announcement1) .await .expect("Failed to send announcement"); - tokio::time::sleep(Duration::from_millis(500)).await; // Test second service (gitlab.example.org) let keys2 = Keys::generate(); @@ -316,14 +298,12 @@ async fn test_archive_multiple_grasp_services() { .authenticator(SignerAuthenticator::new(keys2.clone())) .build(); client2.add_relay(&url).await.expect("Failed to add relay"); - client2.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&client2).await; client2 .send_event(&announcement2) .await .expect("Failed to send announcement"); - tokio::time::sleep(Duration::from_millis(500)).await; // Test non-listed service (github.com) let keys3 = Keys::generate(); @@ -348,16 +328,16 @@ async fn test_archive_multiple_grasp_services() { .authenticator(SignerAuthenticator::new(keys3.clone())) .build(); client3.add_relay(&url).await.expect("Failed to add relay"); - client3.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&client3).await; client3 .send_event(&announcement3) .await .expect("Failed to send announcement"); - tokio::time::sleep(Duration::from_millis(500)).await; - // Verify first service announcement was accepted + // Each send above completed only after the relay replied, and the bare + // repository is created inside admission, so all three outcomes are + // already settled. let repo_path1 = git_data_path.join(format!("{}/{}.git", npub1, identifier1)); assert!( repo_path1.exists(), @@ -431,10 +411,7 @@ async fn test_archive_read_only_creates_bare_repo() { .add_relay(source_relay.url()) .await .expect("Failed to add source relay"); - source_client.connect().await; - - // Wait for connection - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&source_client).await; // Send announcement to source relay source_client @@ -442,8 +419,6 @@ async fn test_archive_read_only_creates_bare_repo() { .await .expect("Failed to send announcement to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // 4. Create and send state event let clone_urls = [ format!( @@ -477,8 +452,6 @@ async fn test_archive_read_only_creates_bare_repo() { .await .expect("Failed to send state event to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // 5. Push git data to source relay // The state event in purgatory authorizes this push push_to_relay(temp_dir.path(), &source_relay.domain(), &npub, identifier) diff --git a/tests/common/purgatory_helpers.rs b/tests/common/purgatory_helpers.rs index b9dd77f..8e37e17 100644 --- a/tests/common/purgatory_helpers.rs +++ b/tests/common/purgatory_helpers.rs @@ -367,78 +367,55 @@ pub fn create_announcement_event( .map_err(|e| format!("Failed to sign announcement event: {}", e)) } -/// Wait for an event to be served by a relay (not in purgatory). +/// Wait until `relay_url` serves `event_id`, polling at a fixed short +/// interval with one connection, bounded by `timeout`. /// -/// Polls the relay until the event is queryable, indicating it has -/// been released from purgatory. Uses exponential backoff for polling. -/// -/// # Arguments -/// * `relay_url` - WebSocket URL of the relay -/// * `event_id` - Event ID to wait for -/// * `timeout` - Maximum time to wait -/// -/// # Returns -/// * `Ok(Event)` - The event was found -/// * `Err(String)` - Timeout or error +/// A healthy relay answers the first poll, so the wait costs one round +/// trip. When the event is still being promoted, a fixed interval keeps +/// the detection delay near the interval instead of letting an exponential +/// backoff stretch it under load. pub async fn wait_for_event_served( relay_url: &str, event_id: &EventId, timeout: Duration, ) -> Result { + let deadline = std::time::Instant::now() + timeout; + let poll_interval = Duration::from_millis(100); + let fetch_timeout = Duration::from_secs(1); + let temp_keys = Keys::generate(); let client = Client::builder() .authenticator(SignerAuthenticator::new(temp_keys)) .build(); - client .add_relay(relay_url) .await .map_err(|e| format!("Failed to add relay: {}", e))?; - - client.connect().await; - - // Wait for connection - let mut connected = false; - for _ in 0..20 { - tokio::time::sleep(Duration::from_millis(100)).await; - let relays = client.relays().await; - if relays.values().any(|r| r.status().is_connected()) { - connected = true; - break; - } - } - - if !connected { + let output = client.try_connect().timeout(timeout).await; + if output.success.is_empty() { client.disconnect().await; - return Err("Failed to connect to relay".to_string()); + return Err(format!( + "Failed to connect to relay {}: {:?}", + relay_url, output.failed + )); } - // Poll for the event with exponential backoff - let start = std::time::Instant::now(); - let mut poll_interval = Duration::from_millis(100); - let max_interval = Duration::from_secs(2); - - while start.elapsed() < timeout { - let filter = Filter::new().id(*event_id); - - match client - .fetch_events(filter) - .timeout(Duration::from_secs(2)) + let filter = Filter::new().id(*event_id); + loop { + if let Ok(events) = client + .fetch_events(filter.clone()) + .timeout(fetch_timeout) .await { - Ok(events) => { - if let Some(event) = events.into_iter().next() { - client.disconnect().await; - return Ok(event); - } - } - Err(_) => { - // Ignore fetch errors, will retry + if let Some(event) = events.into_iter().next() { + client.disconnect().await; + return Ok(event); } } - + if std::time::Instant::now() >= deadline { + break; + } tokio::time::sleep(poll_interval).await; - poll_interval = std::cmp::min(poll_interval * 2, max_interval); } client.disconnect().await; diff --git a/tests/common/relay.rs b/tests/common/relay.rs index cd3f909..dbe4d12 100644 --- a/tests/common/relay.rs +++ b/tests/common/relay.rs @@ -12,7 +12,7 @@ //! `start_on_reservation_*` constructors — the listener stays bound for //! the duration of the reservation, eliminating the same-process race. -use nostr_sdk::prelude::{Keys, ToBech32}; +use nostr_sdk::prelude::{Client, Keys, ToBech32}; use std::path::PathBuf; use std::process::{Child, Command, Stdio}; use std::time::{Duration, Instant}; @@ -1078,48 +1078,8 @@ impl TestRelay { } async fn probe_http_ready(&self) -> std::io::Result<()> { - let probe = async { - let mut stream = tokio::net::TcpStream::connect(("127.0.0.1", self.port)).await?; - // Use the relay's HTTP handler rather than a bare TCP connect: - // this verifies the accept loop and Hyper service are both live, - // which is what the immediately-following WebSocket tests need. - let base_path = self.options.base_path.as_deref().unwrap_or("/"); - let request = format!( - "GET {base_path} HTTP/1.1\r\nHost: 127.0.0.1:{}\r\nConnection: close\r\n\r\n", - self.port - ); - stream.write_all(request.as_bytes()).await?; - - let mut response = Vec::new(); - let mut reader = tokio::io::BufReader::new(stream); - let read = - tokio::io::AsyncBufReadExt::read_until(&mut reader, b'\n', &mut response).await?; - if read == 0 { - return Err(std::io::Error::new( - std::io::ErrorKind::UnexpectedEof, - "readiness probe got empty HTTP response", - )); - } - - let status = String::from_utf8_lossy(&response[..read]); - if status.starts_with("HTTP/1.1 200") || status.starts_with("HTTP/1.0 200") { - Ok(()) - } else { - Err(std::io::Error::new( - std::io::ErrorKind::InvalidData, - format!("readiness probe got non-200 response: {status:?}"), - )) - } - }; - - tokio::time::timeout(READY_PROBE_TIMEOUT, probe) - .await - .unwrap_or_else(|_| { - Err(std::io::Error::new( - std::io::ErrorKind::TimedOut, - "HTTP readiness probe timed out", - )) - }) + let base_path = self.options.base_path.as_deref().unwrap_or("/"); + probe_http_ready(self.port, base_path).await } /// Stop the relay @@ -1138,6 +1098,91 @@ impl Drop for TestRelay { } } +/// Probe a relay's HTTP handler once, returning `Ok(())` only on a 200 +/// response. A raw TCP connect can succeed before Hyper has installed the +/// service that accepts test clients, so only a handled request proves the +/// accept loop and relay wiring are live. +pub async fn probe_http_ready(port: u16, base_path: &str) -> std::io::Result<()> { + let probe = async { + let mut stream = tokio::net::TcpStream::connect(("127.0.0.1", port)).await?; + let request = format!( + "GET {base_path} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\nConnection: close\r\n\r\n" + ); + stream.write_all(request.as_bytes()).await?; + + let mut response = Vec::new(); + let mut reader = tokio::io::BufReader::new(stream); + let read = + tokio::io::AsyncBufReadExt::read_until(&mut reader, b'\n', &mut response).await?; + if read == 0 { + return Err(std::io::Error::new( + std::io::ErrorKind::UnexpectedEof, + "readiness probe got empty HTTP response", + )); + } + + let status = String::from_utf8_lossy(&response[..read]); + if status.starts_with("HTTP/1.1 200") || status.starts_with("HTTP/1.0 200") { + Ok(()) + } else { + Err(std::io::Error::new( + std::io::ErrorKind::InvalidData, + format!("readiness probe got non-200 response: {status:?}"), + )) + } + }; + + tokio::time::timeout(READY_PROBE_TIMEOUT, probe) + .await + .unwrap_or_else(|_| { + Err(std::io::Error::new( + std::io::ErrorKind::TimedOut, + "HTTP readiness probe timed out", + )) + }) +} + +/// Poll [`probe_http_ready`] until the relay answers, bounded by `timeout`. +/// +/// Used by test-local relay fixtures that spawn the binary themselves and so +/// cannot reuse [`TestRelay`]'s own readiness wait. +pub async fn wait_for_http_ready(port: u16, timeout: Duration) -> Result<(), String> { + let deadline = Instant::now() + timeout; + loop { + match probe_http_ready(port, "/").await { + Ok(()) => return Ok(()), + Err(e) if Instant::now() >= deadline => { + return Err(format!( + "relay at 127.0.0.1:{port} did not answer HTTP within {timeout:?}: {e}" + )) + } + Err(_) => sleep(READY_POLL).await, + } + } +} + +/// Connect `client` to its configured relays and return only once every +/// relay reports a completed connection attempt, failing the test if any +/// attempt fails. +/// +/// `Client::connect().and_wait(..)` observes connection state through status +/// notifications after a separate status check; a loopback relay can connect +/// in that gap under scheduler pressure, and the missed notification then +/// stalls the caller for the whole timeout. `try_connect` awaits the attempt +/// itself, so healthy runs continue immediately and failures are explicit. +pub async fn connect_client(client: &Client) { + let output = client.try_connect().timeout(Duration::from_secs(30)).await; + assert!( + output.failed.is_empty(), + "relay connection attempts failed: {:?}", + output.failed + ); + assert!( + !output.success.is_empty(), + "no relay connection was attempted; add relays before connecting" + ); +} + /// Internal outcome of a single [`TestRelay::try_start_once`] attempt. /// /// `EarlyExit` is the retry-eligible case (subprocess died before diff --git a/tests/purgatory_sync.rs b/tests/purgatory_sync.rs index ca70605..7da9d93 100644 --- a/tests/purgatory_sync.rs +++ b/tests/purgatory_sync.rs @@ -273,10 +273,7 @@ async fn test_state_event_syncs_from_remote() { .add_relay(source_relay.url()) .await .expect("Failed to add source relay"); - source_client.connect().await; - - // Wait for connection - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&source_client).await; // Send announcement to source relay source_client @@ -284,8 +281,6 @@ async fn test_state_event_syncs_from_remote() { .await .expect("Failed to send announcement to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // 4. Create and send state event BEFORE pushing // The state event goes to purgatory on source relay, which authorizes the push let clone_urls = [ @@ -320,8 +315,6 @@ async fn test_state_event_syncs_from_remote() { .await .expect("Failed to send state event to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // 5. Push git data to source relay // The state event in purgatory authorizes this push push_to_relay(temp_dir.path(), &source_relay.domain(), &npub, identifier) @@ -439,10 +432,7 @@ async fn test_pr_event_syncs_from_remote() { .add_relay(source_relay.url()) .await .expect("Failed to add source relay"); - source_client.connect().await; - - // Wait for connection - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&source_client).await; // Step 1: Send announcement to source relay → purgatory (StateOnly) source_client @@ -450,8 +440,6 @@ async fn test_pr_event_syncs_from_remote() { .await .expect("Failed to send announcement to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // Step 2: Create and send state event → purgatory (no git data yet) let clone_urls = [ format!( @@ -484,8 +472,6 @@ async fn test_pr_event_syncs_from_remote() { .await .expect("Failed to send state event to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // Step 3: Push git data to source relay // This promotes the announcement from StateOnly to Full AND releases state event push_to_relay(temp_dir.path(), &source_relay.domain(), &npub, identifier) @@ -518,16 +504,13 @@ async fn test_pr_event_syncs_from_remote() { .add_relay(source_relay.url()) .await .expect("Failed to add source relay for PR"); - pr_client.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&pr_client).await; pr_client .send_event(&pr_event) .await .expect("Failed to send PR event to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // Step 5: Push PR commit to refs/nostr/ on source relay // This releases the PR event from purgatory let ref_name = format!("refs/nostr/{}", pr_event_id.to_hex()); @@ -656,10 +639,7 @@ async fn test_concurrent_state_and_pr_sync() { .add_relay(source_relay.url()) .await .expect("Failed to add source relay"); - source_client.connect().await; - - // Wait for connection - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&source_client).await; // Step 1: Send announcement to source relay → purgatory (StateOnly) source_client @@ -667,8 +647,6 @@ async fn test_concurrent_state_and_pr_sync() { .await .expect("Failed to send announcement to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // Step 2: Create and send state event → purgatory (no git data yet) let clone_urls = [ format!( @@ -704,8 +682,6 @@ async fn test_concurrent_state_and_pr_sync() { .await .expect("Failed to send state event to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // Step 3: Push git data to source relay // This promotes the announcement from StateOnly to Full AND releases state event push_to_relay(temp_dir.path(), &source_relay.domain(), &npub, identifier) @@ -738,16 +714,13 @@ async fn test_concurrent_state_and_pr_sync() { .add_relay(source_relay.url()) .await .expect("Failed to add source relay for PR"); - pr_client.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&pr_client).await; pr_client .send_event(&pr_event) .await .expect("Failed to send PR event to source"); - tokio::time::sleep(Duration::from_millis(200)).await; - // Step 5: Push PR commit to refs/nostr/ on source relay // This releases the PR event from purgatory let pr_ref_name = format!("refs/nostr/{}", pr_event_id.to_hex()); @@ -995,14 +968,12 @@ async fn test_pr_event_clone_tag_sync_with_partial_oid_aggregation_from_multiple .add_relay(source_grasp.url()) .await .expect("Failed to add source_grasp relay"); - source_client.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&source_client).await; source_client .send_event(&announcement) .await .expect("Failed to send announcement to source_grasp"); - tokio::time::sleep(Duration::from_millis(200)).await; // Create state event referencing commit_a let state_event = create_state_event( @@ -1026,7 +997,6 @@ async fn test_pr_event_clone_tag_sync_with_partial_oid_aggregation_from_multiple .send_event(&state_event) .await .expect("Failed to send state event to source_grasp"); - tokio::time::sleep(Duration::from_millis(200)).await; // Push main branch (commit_a) to source_grasp - releases state event push_to_relay(repo_a.path(), &source_grasp.domain(), &npub, identifier) @@ -1051,14 +1021,12 @@ async fn test_pr_event_clone_tag_sync_with_partial_oid_aggregation_from_multiple .add_relay(mock_relay.url()) .await .expect("Failed to add mock_relay for announcement"); - mock_client.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&mock_client).await; mock_client .send_event(&announcement) .await .expect("Failed to send announcement to mock_relay"); - tokio::time::sleep(Duration::from_millis(200)).await; let repo_coord = build_repo_coord(&owner_keys, identifier); @@ -1084,8 +1052,7 @@ async fn test_pr_event_clone_tag_sync_with_partial_oid_aggregation_from_multiple .add_relay(mock_relay.url()) .await .expect("Failed to add mock_relay"); - pr_client.connect().await; - tokio::time::sleep(Duration::from_millis(500)).await; + common::relay::connect_client(&pr_client).await; pr_client .send_event(&pr_event) diff --git a/tests/sync/historic_sync.rs b/tests/sync/historic_sync.rs index 8db7a43..efd2e24 100644 --- a/tests/sync/historic_sync.rs +++ b/tests/sync/historic_sync.rs @@ -454,7 +454,7 @@ async fn test_pagination_for_large_historic_sync() { .add_relay(syncing.url()) .await .expect("Failed to add syncing relay to client"); - client.connect().and_wait(Duration::from_secs(30)).await; + crate::common::relay::connect_client(&client).await; // Poll until every issue has paginated across, up to a bounded deadline. let deadline = tokio::time::Instant::now() + Duration::from_secs(30); diff --git a/tests/sync/live_sync.rs b/tests/sync/live_sync.rs index b815355..a12966c 100644 --- a/tests/sync/live_sync.rs +++ b/tests/sync/live_sync.rs @@ -486,7 +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().and_wait(Duration::from_secs(30)).await; + crate::common::relay::connect_client(&client).await; let fetch_filter = Filter::new().kind(Kind::Comment).id(comment_id); @@ -642,7 +642,7 @@ async fn test_live_sync_event_ordering() { let events_found: Vec; if client.add_relay(relay_b.url()).await.is_ok() { - client.connect().and_wait(Duration::from_secs(30)).await; + crate::common::relay::connect_client(&client).await; let filter = Filter::new().kind(Kind::GitIssue).author(keys.public_key()); diff --git a/tests/sync/purgatory_fetch.rs b/tests/sync/purgatory_fetch.rs index 3fc4f90..92bee40 100644 --- a/tests/sync/purgatory_fetch.rs +++ b/tests/sync/purgatory_fetch.rs @@ -159,10 +159,7 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() { .add_relay(source.url()) .await .expect("add source relay"); - source_client - .connect() - .and_wait(Duration::from_secs(30)) - .await; + crate::common::relay::connect_client(&source_client).await; source_client .send_event(&source_announcement) .await @@ -239,10 +236,7 @@ async fn purgatory_fetch_batches_available_tips_and_isolates_missing_oids() { .add_relay(mock.url()) .await .expect("add mock relay"); - mock_client - .connect() - .and_wait(Duration::from_secs(30)) - .await; + crate::common::relay::connect_client(&mock_client).await; mock_client .send_event(&syncing_announcement) .await