test(sync): wait for observable fixture readiness

Fixed startup and promotion sleeps could expire before asynchronous work
completed. Wait with bounded deadlines for event visibility and an actual
connected gauge; connection attempts and health do not prove readiness.

Preserve URL schemes when constructing metrics endpoints, retain ownership
of unavailable endpoints, and remove stopped sources instead of replacing
them with a live placeholder. Keep immediate assertions where completion
already guarantees the result. Production sync policy is unchanged.

Validation: shared sync helper tests passed; relay_identity and the fixture
lifecycle target compiled and passed with the revised common helpers.
Full sync and deletion integration execution remains host validation.

Assisted-by: Codex (GPT-6)
This commit is contained in:
DanConwayDev
2026-09-12 14:48:26 +00:00
parent 75da9a8446
commit cf2ec36a68
2 changed files with 151 additions and 202 deletions
+27 -76
View File
@@ -19,6 +19,24 @@ use super::purgatory_helpers::{
};
use super::sync_helpers::create_repo_announcement;
/// Wait for promotion rather than assuming the worker runs within a fixed delay.
async fn wait_until_served(client: &AuditClient, event_id: EventId) {
tokio::time::timeout(Duration::from_secs(10), async {
loop {
if client
.is_event_on_relay(event_id)
.await
.expect("query promoted event")
{
return;
}
tokio::time::sleep(Duration::from_millis(50)).await;
}
})
.await
.unwrap_or_else(|_| panic!("event {event_id} was not served after git data arrived"));
}
/// Publish a repo announcement, submit a matching state event, and push the
/// deterministic git data so the announcement is promoted out of purgatory and
/// becomes queryable.
@@ -130,16 +148,7 @@ pub async fn publish_served_repo(client: &AuditClient, test_name: &str) -> (Even
Err(e) => panic!("git push error while promoting repo: {}", e),
}
// Give the relay a moment to promote the events out of purgatory.
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
client
.is_event_on_relay(announcement.id)
.await
.expect("query announcement"),
"announcement should be served after git data arrives"
);
wait_until_served(client, announcement.id).await;
(announcement, repo_id)
}
@@ -253,15 +262,7 @@ pub async fn publish_served_repo_with_maintainers(
Err(e) => panic!("git push error while promoting repo: {}", e),
}
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
client
.is_event_on_relay(announcement.id)
.await
.expect("query announcement"),
"announcement should be served after git data arrives"
);
wait_until_served(client, announcement.id).await;
(announcement, repo_id)
}
@@ -320,22 +321,8 @@ pub async fn publish_served_repo_with_state_event(
push_to_relay(temp_dir.path(), &relay_domain, &npub, &repo_id)
.expect("git push should promote announcement + state event out of purgatory");
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
client
.is_event_on_relay(announcement.id)
.await
.expect("query announcement"),
"announcement should be served after git data arrives"
);
assert!(
client
.is_event_on_relay(state_event.id)
.await
.expect("query state event"),
"state event should be served after git data arrives"
);
wait_until_served(client, announcement.id).await;
wait_until_served(client, state_event.id).await;
(announcement, repo_id, state_event)
}
@@ -406,22 +393,8 @@ pub async fn publish_served_repo_with_state_event_and_maintainers(
push_to_relay(temp_dir.path(), &relay_domain, &npub, &repo_id)
.expect("git push should promote announcement + state event out of purgatory");
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
client
.is_event_on_relay(announcement.id)
.await
.expect("query announcement"),
"announcement should be served after git data arrives"
);
assert!(
client
.is_event_on_relay(state_event.id)
.await
.expect("query state event"),
"state event should be served after git data arrives"
);
wait_until_served(client, announcement.id).await;
wait_until_served(client, state_event.id).await;
(announcement, repo_id, state_event)
}
@@ -483,8 +456,6 @@ pub async fn publish_served_announcement_for_identifier(
.await
.expect("relay should accept state event");
tokio::time::sleep(Duration::from_millis(300)).await;
if client
.is_event_on_relay(announcement.id)
.await
@@ -536,13 +507,7 @@ pub async fn publish_served_announcement_for_identifier(
Err(e) => panic!("git push error while promoting repo: {}", e),
}
assert!(
client
.is_event_on_relay(announcement.id)
.await
.expect("query announcement"),
"announcement should be served after git data arrives"
);
wait_until_served(client, announcement.id).await;
announcement
}
@@ -599,8 +564,6 @@ pub async fn publish_served_announcement_with_state_for_identifier(
.await
.expect("relay should accept state event");
tokio::time::sleep(Duration::from_millis(300)).await;
let announcement_served = client
.is_event_on_relay(announcement.id)
.await
@@ -656,20 +619,8 @@ pub async fn publish_served_announcement_with_state_for_identifier(
Err(e) => panic!("git push error while promoting repo: {}", e),
}
assert!(
client
.is_event_on_relay(announcement.id)
.await
.expect("query announcement"),
"announcement should be served after git data arrives"
);
assert!(
client
.is_event_on_relay(state_event.id)
.await
.expect("query state event"),
"state event should be served after git data arrives"
);
wait_until_served(client, announcement.id).await;
wait_until_served(client, state_event.id).await;
(announcement, state_event)
}
+124 -126
View File
@@ -15,7 +15,7 @@ use std::time::Duration;
use nostr_sdk::prelude::*;
use super::port::{self, PortReservation};
use super::port::{self, PortReservation, UnavailableEndpoint};
use super::relay::TestRelay;
const DESCENDANT_LIVE_LOG: &str = "Installed priority-bounded auxiliary live coverage";
@@ -476,80 +476,26 @@ pub async fn wait_for_sync_connection(
expected_connections: usize,
timeout: Duration,
) -> Result<(), String> {
// Convert ws:// URL to http:// for metrics endpoint
let http_url = syncing_relay_url
.replace("ws://", "http://")
.replace("/", "")
+ "/metrics";
let start = std::time::Instant::now();
let poll_interval = Duration::from_millis(100);
while start.elapsed() < timeout {
// Fetch metrics
if let Ok(response) = reqwest::get(&http_url).await {
if let Ok(metrics) = response.text().await {
// Look for sync connection metrics
// The metric name pattern: ngit_sync_connections or similar
// We check for any indication of established connections
let wait = async {
loop {
if let Ok(metrics) = fetch_metrics(syncing_relay_url).await {
if check_sync_connections_in_metrics(&metrics, expected_connections) {
return Ok(());
return;
}
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
tokio::time::sleep(poll_interval).await;
}
Err(format!(
"Timeout waiting for {} sync connection(s) on {} after {:?}",
expected_connections, syncing_relay_url, timeout
};
tokio::time::timeout(timeout, wait).await.map_err(|_| format!(
"Timeout waiting for {expected_connections} sync connection(s) on {syncing_relay_url} after {timeout:?}"
))
}
/// Check metrics string for expected number of sync connections.
///
/// Looks for various metric patterns that indicate sync connections:
/// - ngit_sync_connections (gauge)
/// - ngit_sync_relay_connections (gauge)
/// - Any metric containing "sync" and "connection" with count > 0
/// Connection attempts and health states do not establish a live connection.
fn check_sync_connections_in_metrics(metrics: &str, expected: usize) -> bool {
// Parse metrics line by line looking for connection counts
for line in metrics.lines() {
// Skip comments and empty lines
if line.starts_with('#') || line.is_empty() {
continue;
}
// Look for sync connection metrics
// Format: metric_name{labels} value
// or: metric_name value
if line.contains("sync") && line.contains("connect") {
// Extract the value (last space-separated token)
if let Some(value_str) = line.split_whitespace().last() {
if let Ok(value) = value_str.parse::<f64>() {
if value as usize >= expected {
return true;
}
}
}
}
// Also check for specific metric names that might indicate connections
// ngit_sync_health_state with value 1 or 2 (connecting/healthy)
if line.contains("ngit_sync_health") {
if let Some(value_str) = line.split_whitespace().last() {
if let Ok(value) = value_str.parse::<f64>() {
// Health state > 0 typically means connection attempt or established
if value > 0.0 && expected > 0 {
return true;
}
}
}
}
}
false
ParsedMetrics::parse(metrics)
.relays_connected_total()
.is_some_and(|connected| connected >= expected as i64)
}
// ============================================================================
@@ -677,10 +623,14 @@ pub fn repo_coord(keys: &Keys, identifier: &str) -> String {
/// assert!(metrics.contains("ngit_sync_"));
/// ```
pub async fn fetch_metrics(relay_url: &str) -> Result<String, reqwest::Error> {
// Convert ws:// URL to http:// for metrics endpoint
let http_url = relay_url.replace("ws://", "http://").replace("/", "") + "/metrics";
reqwest::get(metrics_url(relay_url)).await?.text().await
}
reqwest::get(&http_url).await?.text().await
fn metrics_url(relay_url: &str) -> String {
let http_url = relay_url
.replacen("wss://", "https://", 1)
.replacen("ws://", "http://", 1);
format!("{}/metrics", http_url.trim_end_matches('/'))
}
// ============================================================================
@@ -837,10 +787,9 @@ impl ParsedMetrics {
/// harness.stop_all().await;
/// ```
pub struct MetricsTestHarness {
source_relays: Vec<TestRelay>,
source_relays: Vec<Option<TestRelay>>,
syncing_relay: Option<TestRelay>,
#[allow(dead_code)]
nowhere_url: Option<String>,
unavailable_endpoint: Option<UnavailableEndpoint>,
}
impl MetricsTestHarness {
@@ -848,34 +797,36 @@ impl MetricsTestHarness {
pub async fn with_sources(count: usize) -> Self {
let mut source_relays = Vec::new();
for _ in 0..count {
source_relays.push(TestRelay::start().await);
source_relays.push(Some(TestRelay::start().await));
}
Self {
source_relays,
syncing_relay: None,
nowhere_url: None,
unavailable_endpoint: None,
}
}
/// Get source relay URL
pub fn source_url(&self, idx: usize) -> &str {
self.source_relays[idx].url()
self.source_relay(idx).url()
}
/// Get source relay domain (for announcement tags)
pub fn source_domain(&self, idx: usize) -> String {
self.source_relays[idx].domain()
self.source_relay(idx).domain()
}
/// Get a reference to a source relay (for advanced test operations)
pub fn source_relay(&self, idx: usize) -> &TestRelay {
&self.source_relays[idx]
self.source_relays[idx]
.as_ref()
.expect("source relay has been stopped")
}
/// Submit events to a specific source relay
pub async fn submit_events(&self, source_idx: usize, events: &[Event]) -> Result<(), String> {
let relay = &self.source_relays[source_idx];
let relay = self.source_relay(source_idx);
let keys = Keys::generate();
let client = TestClient::new(relay.url(), keys).await?;
@@ -889,7 +840,7 @@ impl MetricsTestHarness {
/// Start syncing relay pointing to source[idx]
pub async fn start_syncing_relay(&mut self, source_idx: usize) {
let source_url = self.source_relays[source_idx].url().to_string();
let source_url = self.source_relay(source_idx).url().to_string();
self.syncing_relay = Some(TestRelay::start_with_sync(Some(source_url)).await);
}
@@ -904,37 +855,26 @@ impl MetricsTestHarness {
source_idx: usize,
reservation: PortReservation,
) {
let source_url = self.source_relays[source_idx].url().to_string();
let source_url = self.source_relay(source_idx).url().to_string();
self.syncing_relay = Some(
TestRelay::start_on_reservation_with_options(reservation, Some(source_url), false)
.await,
);
}
/// Start syncing relay pointing to random unused port (for failure tests)
/// Start syncing against an owned endpoint that closes every connection.
pub async fn start_syncing_relay_to_nowhere(&mut self) {
let port = random_unused_port();
let nowhere_url = format!("ws://127.0.0.1:{}", port);
self.nowhere_url = Some(nowhere_url.clone());
let endpoint = UnavailableEndpoint::new();
let nowhere_url = format!("ws://127.0.0.1:{}", endpoint.port());
self.unavailable_endpoint = Some(endpoint);
self.syncing_relay = Some(TestRelay::start_with_sync(Some(nowhere_url)).await);
}
/// Stop a source relay
pub async fn stop_source(&mut self, source_idx: usize) {
// We need to take ownership to stop, so we swap with a new relay
// that we immediately stop. This is a workaround since TestRelay::stop
// takes self by value.
let relay = std::mem::replace(
&mut self.source_relays[source_idx],
TestRelay::start().await,
);
relay.stop().await;
// Stop the placeholder too
let placeholder = std::mem::replace(
&mut self.source_relays[source_idx],
TestRelay::start().await,
);
placeholder.stop().await;
if let Some(relay) = self.source_relays[source_idx].take() {
relay.stop().await;
}
}
/// Fetch and parse metrics from syncing relay
@@ -961,29 +901,41 @@ impl MetricsTestHarness {
if let Some(relay) = self.syncing_relay.take() {
relay.stop().await;
}
for relay in self.source_relays.drain(..) {
for relay in self.source_relays.drain(..).flatten() {
relay.stop().await;
}
}
}
// ============================================================================
// Port Helpers
// ============================================================================
/// Get a random unused port by binding to port 0 and letting the OS assign one
pub fn random_unused_port() -> u16 {
std::net::TcpListener::bind("127.0.0.1:0")
.expect("Failed to bind to random port")
.local_addr()
.expect("Failed to get local addr")
.port()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn metrics_url_preserves_scheme_and_base_path() {
assert_eq!(
metrics_url("ws://127.0.0.1:1234/"),
"http://127.0.0.1:1234/metrics"
);
assert_eq!(
metrics_url("wss://example.test/relay/"),
"https://example.test/relay/metrics"
);
}
#[test]
fn connection_readiness_ignores_attempts_and_health() {
assert!(!check_sync_connections_in_metrics("ngit_sync_connection_attempts_total 5\nngit_sync_health 1\nngit_sync_relays_connected_total 0", 1));
assert!(check_sync_connections_in_metrics(
"ngit_sync_relays_connected_total 2",
2
));
assert!(!check_sync_connections_in_metrics(
"ngit_sync_relays_connected_total 1",
2
));
}
#[test]
fn test_repo_coord_format() {
let keys = Keys::generate();
@@ -1268,8 +1220,15 @@ pub async fn push_git_data_to_relay(
push_to_relay(git_temp_dir.path(), &relay.domain(), &npub, identifier)
.expect("Failed to push git data to relay");
// Brief wait for push processing
tokio::time::sleep(Duration::from_millis(500)).await;
assert!(
wait_for_event_on_relay(
relay.url(),
Filter::new().id(state_event.id),
Duration::from_secs(10)
)
.await,
"pushed state event must leave purgatory"
);
git_temp_dir
}
@@ -1379,7 +1338,15 @@ pub async fn push_unique_git_data_to_relay(
push_to_relay(path, &relay.domain(), &npub, identifier)
.expect("Failed to push git data to relay");
tokio::time::sleep(Duration::from_millis(500)).await;
assert!(
wait_for_event_on_relay(
relay.url(),
Filter::new().id(state_event.id),
Duration::from_secs(10)
)
.await,
"pushed state event must leave purgatory"
);
git_temp_dir
}
@@ -1464,8 +1431,25 @@ pub async fn setup_announcement_on_relay(
push_to_relay(git_temp_dir.path(), &relay.domain(), &npub, identifier)
.expect("Failed to push git data to relay");
// Brief wait for push processing
tokio::time::sleep(Duration::from_millis(500)).await;
assert!(
wait_for_event_on_relay(
relay.url(),
Filter::new().id(state_event.id),
Duration::from_secs(10)
)
.await,
"pushed state event must leave purgatory"
);
assert!(
wait_for_event_on_relay(
relay.url(),
Filter::new().id(announcement.id),
Duration::from_secs(10)
)
.await,
"pushed announcement must leave purgatory"
);
(announcement, git_temp_dir)
}
@@ -1587,7 +1571,15 @@ pub async fn run_sync_test(historic_events: &[Event], live_events: &[Event]) ->
.expect("Failed to push git data to source relay");
// 8. Wait for source relay to process the push and release events from purgatory
tokio::time::sleep(Duration::from_secs(2)).await;
assert!(
wait_for_event_on_relay(
source.url(),
Filter::new().id(announcement.id),
Duration::from_secs(10)
)
.await,
"source announcement must leave purgatory"
);
// 9. Send historic events to source BEFORE syncing relay starts
for event in historic_events {
@@ -1605,7 +1597,9 @@ pub async fn run_sync_test(historic_events: &[Event], live_events: &[Event]) ->
.await;
// 11. Wait for sync connection to establish
let _ = wait_for_sync_connection(syncing.url(), 1, Duration::from_secs(5)).await;
wait_for_sync_connection(syncing.url(), 1, Duration::from_secs(10))
.await
.expect("sync connection must establish before publishing live events");
// 12. Send live events AFTER connection established
for event in live_events {
@@ -1614,12 +1608,16 @@ pub async fn run_sync_test(historic_events: &[Event], live_events: &[Event]) ->
.expect("Failed to send live event");
}
// 13. Allow sync + purgatory promotion to complete on the syncing relay.
// The syncing relay receives the announcement (goes to purgatory) and state event.
// The purgatory sync loop (1s interval) fetches git data from source's clone URL
// (http://source-domain/npub/test-repo.git) and releases the announcement.
// We wait up to 8s to allow time for this.
tokio::time::sleep(Duration::from_secs(8)).await;
// 13. Observe announcement promotion rather than guessing a sync duration.
assert!(
wait_for_event_on_relay(
syncing.url(),
Filter::new().id(announcement.id),
Duration::from_secs(15)
)
.await,
"synced announcement must leave purgatory"
);
// 14. Compute repo coordinate before moving keys
let coordinate = repo_coord(&keys, "test-repo");