mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 15:08:24 +00:00
Merge #fc7ab4b2: test: replace fixed sleeps with bounded observable wai…
test: replace fixed sleeps with bounded observable waits nostr:nevent1qgsx2lyl2e4zvfadwcvkd9fkrcwczj7mf858hy85mwqclwgut8wpg2spz3mhxue69uhhyetvv9ujumn8d96zuer9wcq3yamnwvaz7tm8d96xummnw3ezucm0d5q3kamnwvaz7tmwva5hgtnyv9hxxmmwwashjer9wchxxmmdqqs0c745kfq20lx5z8ypdh77mz2szaw6epc5l9dd0cgdpy4kfy87r9sg2l33m PR-Author: DanConwayDev's Agent nostr:npub1v47f74n2ycn66asev62nv8sas99akj0g0wg0fkup37u3ckwuzs4q7cwtp0 PR description: Replaces the 34 remaining fixed sleeps in tests/purgatory_sync.rs and tests/archive_grasp_services.rs with either nothing, where the relay's awaited reply already guarantees the effect, or a deadline-bounded wait on an observable. It also fixes two latent problems that the sleeps had been masking: an SDK connection wait that can stall for its full timeout, and a helper whose exponential backoff makes detection slow under load. Why the reply is a barrier: the relay runs the GRASP write policy inside event admission, including `ensure_bare_repository` and purgatory parking, before it sends OK, and the client awaits that OK. A rejected announcement therefore never creates a directory, so the negative "must not be created" case needs no grace period either. Effects that run in spawned tasks (promotion after a push, sync to a peer relay) keep their existing bounded `wait_for_event_served` waits. Commits: 1. Remove the sleeps in the two files. Post-connect sleeps become `connect_client`, built on `try_connect`, which awaits the connection attempts and fails on refusal. The subprocess readiness pad becomes `wait_for_http_ready`, lifted from `TestRelay::probe_http_ready` unchanged, so readiness means a served HTTP 200. 2. Tighten `wait_for_event_served`: one connection, fixed 100 ms poll, 1 s fetch timeout, same deadline semantics. Its exponential backoff to 2 s only stayed invisible because the sleeps landed the first poll after the relay's work. 3. Move the five existing `connect().and_wait(..)` sites in the sync tests onto `connect_client`. `Relay::wait_for_connection` checks status and then subscribes to notifications, so a loopback connection that completes in between is missed and the caller sleeps for the whole timeout. Validation on an 8-core box, comparing master against this branch: - 20 unloaded and 20 loaded runs of both suites: 20/20 in every variant. - 30 samples with three instances running concurrently under CPU spinners, invoking the test binaries directly: 30/30 for both. Per-binary wall time fell from 14.4 s to 12.5 s (purgatory_sync) and 8.2 s to 6.3 s (archive_grasp_services). - Per-test sequential timing under spinners: total 70 s on master, 53 s here, no test slower. - The sync test binary (114 tests) passes; formatting and all-targets Clippy with warnings denied pass. Two failures were seen during the investigation in cargo-driven concurrent runs of intermediate versions: one 30 s `wait_for_event_served` timeout (before commit 2) and one connection refused where a port was taken between reservation and bind. Neither reproduced in 90 direct-binary samples across three variants. The port race is pre-existing and independent of this change. Not changed: fixed sleeps in other test files, including the 1100 ms wait in replaceable_history that deliberately crosses a second boundary. 🤖 Generated with [Claude Code](https://claude.com/claude-code)
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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<Event, String> {
|
||||
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;
|
||||
|
||||
+88
-43
@@ -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
|
||||
|
||||
+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)
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<Event>;
|
||||
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());
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user