mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
test: replace fixed sleeps with observable waits in purgatory sync and archive tests
tests/purgatory_sync.rs and tests/archive_grasp_services.rs still carried 34 fixed sleeps: 500 ms after connecting a client, 200-500 ms after sending an event before checking its side effect, and a 100 ms grace pad after a raw TCP connect to a subprocess relay. The project's fixture policy in docs/how-to/test-fixtures.md reserves fixed sleeps for cases where elapsed time is itself under test. Post-send sleeps are removed outright. The relay runs the GRASP write policy, including bare-repository creation and purgatory parking, inside event admission before it replies, and the client awaits that reply. The accepted or rejected send is therefore the barrier, including for the "repository must not be created" case. Effects that run in spawned tasks (promotion after a push, sync to a peer) already used bounded wait_for_event_served waits, which stay. Post-connect sleeps become connect_client, a shared helper built on try_connect, which awaits the connection attempts themselves and fails the test on refusal. connect().and_wait() was rejected: it checks relay status and then subscribes to status notifications, so a loopback connection that completes in that gap is missed and the caller stalls for the full timeout. The subprocess readiness pad becomes wait_for_http_ready, lifted from TestRelay::probe_http_ready without behaviour change, so readiness means a served HTTP 200 rather than a bound socket. Validation on an 8-core box: 20 unloaded and 20 loaded runs of both suites passed, as did 30 samples with three instances running concurrently under CPU spinners. Per-binary wall time under that contention fell from 14.4 s to 12.2 s (purgatory_sync) and 8.2 s to 6.3 s (archive). Assisted-by: Claude Fable 5.1 Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Fable 5.1
parent
2dbb8dfdaf
commit
86047d78d6
@@ -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)
|
||||
|
||||
+64
-19
@@ -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,15 +1078,35 @@ 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("/");
|
||||
probe_http_ready(self.port, base_path).await
|
||||
}
|
||||
|
||||
/// Stop the relay
|
||||
pub async fn stop(mut self) {
|
||||
// kill() sends SIGKILL; reap directly instead of guessing a grace period.
|
||||
let _ = self.process.kill();
|
||||
let _ = self.process.wait();
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for TestRelay {
|
||||
fn drop(&mut self) {
|
||||
// Ensure process is killed when TestRelay is dropped
|
||||
let _ = self.process.kill();
|
||||
let _ = self.process.wait();
|
||||
}
|
||||
}
|
||||
|
||||
/// 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:{}\r\nConnection: close\r\n\r\n",
|
||||
self.port
|
||||
"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?;
|
||||
|
||||
@@ -1122,20 +1142,45 @@ impl TestRelay {
|
||||
})
|
||||
}
|
||||
|
||||
/// Stop the relay
|
||||
pub async fn stop(mut self) {
|
||||
// kill() sends SIGKILL; reap directly instead of guessing a grace period.
|
||||
let _ = self.process.kill();
|
||||
let _ = self.process.wait();
|
||||
/// 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,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for TestRelay {
|
||||
fn drop(&mut self) {
|
||||
// Ensure process is killed when TestRelay is dropped
|
||||
let _ = self.process.kill();
|
||||
let _ = self.process.wait();
|
||||
}
|
||||
/// 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.
|
||||
|
||||
+8
-41
@@ -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/<event-id> 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/<event-id> 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)
|
||||
|
||||
Reference in New Issue
Block a user