test(sync): replace remaining fixed sleeps with observable waits

Sweep the fixed sleeps left in the sync suite after the tag_variations /
live_sync hardening, following the same pattern: wait on the condition
the test actually cares about with a bounded deadline, so healthy runs
are no slower and only the failure bound widens to 30 seconds.

- discovery, historic_sync, metrics, maintainer_reprocessing: replace
  "wait for discovery/sync" sleeps with bounded polls on the synced
  event, a metrics counter/gauge, or the syncing relay's log (rejected
  announcements are observable by their note ID), and widen the
  short verification deadlines that followed them.
- purgatory_fetch, live_sync, historic_sync: replace fixed post-connect
  sleeps on raw clients with `connect().and_wait(timeout)`.
- metrics: poll the metrics endpoint for readiness instead of sleeping
  through relay startup.

TestClient::send_event is reworked to send through its single tracked
relay so the SDK error kind survives: transient transport failures
(disconnect, timeout, not-connected) are retried with the existing
bounded backoff plus a reconnect, while a relay OK-false rejection now
fails immediately instead of being retried, so tests cannot mask
write-policy regressions. A transient client error under load
previously killed adaptive_pagination's seeding despite the retry loop,
because rejections and transport failures were indistinguishable at the
pool level.

Deliberate time-as-behavior sleeps stay and are now commented: the
one-second rejected-hot-cache TTL expiries, created_at spacing before
+1s-dated events, and fixed observation windows for events that must
NOT appear. req_concurrency is deliberately untouched: its proper fix
needs a relay observable for descendant-sweep completion (issue
b052fa70).

Validated under CI's hostile global git config (GIT_CONFIG_GLOBAL
mirror of the workflow step): three consecutive green runs of
`cargo test --locked --test sync`, plus `cargo fmt --check`,
`cargo clippy --workspace --all-targets -- -D warnings`, and a full
`cargo test --locked` (its sync pass green a fourth time; one unrelated
load flake in live_sync_regroups_after_filter_count_refusal passed on
isolated rerun and in the repeat full run).
This commit is contained in:
DanConwayDev
2026-08-15 14:21:43 +00:00
parent 22abda22a8
commit 6ef105d95f
7 changed files with 266 additions and 194 deletions
+31 -25
View File
@@ -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<EventId, String> {
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
))
}
+9 -12
View File
@@ -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);
+67 -42
View File
@@ -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;
+2 -4
View File
@@ -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<Event>;
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());
+55 -26
View File
@@ -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"
+94 -81
View File
@@ -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)
+8 -4
View File
@@ -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