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/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)