diff --git a/tests/sync/live_sync.rs b/tests/sync/live_sync.rs index c626d25..dc1582d 100644 --- a/tests/sync/live_sync.rs +++ b/tests/sync/live_sync.rs @@ -147,7 +147,7 @@ async fn test_live_sync_batches_repo_filters_below_source_req_limit() { .await .expect("publish announcement to constrained source relay"); - wait_for_sync_connection(syncing.url(), 1, Duration::from_secs(5)) + wait_for_sync_connection(syncing.url(), 1, Duration::from_secs(30)) .await .expect("syncing relay should connect to constrained source"); wait_for_new_descendant_live_generation(&syncing, 0, Duration::from_secs(20)).await; @@ -169,7 +169,7 @@ async fn test_live_sync_batches_repo_filters_below_source_req_limit() { wait_for_event_on_relay( syncing.url(), Filter::new().id(readiness_issue.id), - Duration::from_secs(5), + Duration::from_secs(30), ) .await, "q-tagged readiness issue should sync before testing live delivery" @@ -189,7 +189,7 @@ async fn test_live_sync_batches_repo_filters_below_source_req_limit() { let synced = wait_for_event_on_relay( syncing.url(), Filter::new().id(issue.id), - Duration::from_secs(5), + Duration::from_secs(30), ) .await; @@ -221,7 +221,7 @@ async fn live_sync_regroups_after_filter_count_refusal() { .await .expect("publish announcement to strict source relay"); - let refusal_deadline = tokio::time::Instant::now() + Duration::from_secs(10); + let refusal_deadline = tokio::time::Instant::now() + Duration::from_secs(30); loop { let metrics = fetch_metrics(syncing.url()).await.unwrap_or_default(); if metrics.contains("ngit_sync_policy_refusals_total") && metrics.contains("filter_count") { @@ -248,7 +248,7 @@ async fn live_sync_regroups_after_filter_count_refusal() { wait_for_event_on_relay( syncing.url(), Filter::new().id(issue.id), - Duration::from_secs(10), + Duration::from_secs(30), ) .await, "exact live coverage should survive regrouping below the relay's limit" @@ -305,8 +305,12 @@ async fn test_live_sync_layer2_events() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b with git data (triggers discovery)"); - // 5. Wait for discovery to complete - tokio::time::sleep(Duration::from_secs(1)).await; + // 5. Wait for discovery: events published before the syncing relay has + // connected to its peer can only be picked up by later recovery passes, + // so wait on the observable connection rather than a fixed sleep. + wait_for_sync_connection(relay_b.url(), 1, Duration::from_secs(30)) + .await + .expect("relay_b should establish a sync connection after discovery"); // 6. Create and send a Layer 2 issue event (using helper) let repo_coordinate = repo_coord(&keys, repo_id); @@ -338,7 +342,7 @@ async fn test_live_sync_layer2_events() { .author(keys.public_key()) .id(issue_id); - let synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(5)).await; + let synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(30)).await; println!("Issue {} synced to relay_b: {}", issue_id, synced); @@ -399,8 +403,12 @@ async fn test_live_sync_layer3_events() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b with git data (triggers discovery)"); - // 3. Wait for discovery - tokio::time::sleep(Duration::from_secs(1)).await; + // 3. Wait for discovery: events published before the syncing relay has + // connected to its peer can only be picked up by later recovery passes, + // so wait on the observable connection rather than a fixed sleep. + wait_for_sync_connection(relay_b.url(), 1, Duration::from_secs(30)) + .await + .expect("relay_b should establish a sync connection after discovery"); // 4. Create and send Layer 2 issue let repo_coordinate = repo_coord(&keys, repo_id); @@ -446,12 +454,11 @@ async fn test_live_sync_layer3_events() { comment_id ); - // 6. Now wait for issue to sync to relay_b (this triggers Layer 3 filter creation) - tokio::time::sleep(Duration::from_secs(2)).await; - + // 6. Now wait for issue to sync to relay_b (this triggers Layer 3 filter + // creation); the bounded poll replaces a fixed sleep plus short deadline. let issue_filter = Filter::new().kind(Kind::GitIssue).id(issue_id); let issue_synced = - wait_for_event_on_relay(relay_b.url(), issue_filter, Duration::from_secs(3)).await; + wait_for_event_on_relay(relay_b.url(), issue_filter, Duration::from_secs(30)).await; println!("Issue synced to relay_b: {}", issue_synced); client_a.disconnect().await; @@ -571,8 +578,12 @@ async fn test_live_sync_event_ordering() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcements set up on both relays with git data"); - // 3. Wait for discovery - tokio::time::sleep(Duration::from_secs(1)).await; + // 3. Wait for discovery: events published before the syncing relay has + // connected to its peer can only be picked up by later recovery passes, + // so wait on the observable connection rather than a fixed sleep. + wait_for_sync_connection(relay_b.url(), 1, Duration::from_secs(30)) + .await + .expect("relay_b should establish a sync connection after discovery"); // 4. Create and send 3 issues with delays between them let repo_coordinate = repo_coord(&keys, repo_id); @@ -611,8 +622,19 @@ async fn test_live_sync_event_ordering() { client_a.disconnect().await; - // 5. Wait for all events to sync - tokio::time::sleep(Duration::from_secs(3)).await; + // 5. Wait for all events to sync (bounded poll per event rather than a + // fixed sleep, which is not a reliable proxy under CI load). + for issue_id in &issue_ids { + assert!( + wait_for_event_on_relay( + relay_b.url(), + Filter::new().id(*issue_id), + Duration::from_secs(30), + ) + .await, + "issue {issue_id} should sync to relay_b" + ); + } // 6. Fetch all events from relay_b let temp_keys = Keys::generate(); diff --git a/tests/sync/tag_variations.rs b/tests/sync/tag_variations.rs index 8272459..5fb764d 100644 --- a/tests/sync/tag_variations.rs +++ b/tests/sync/tag_variations.rs @@ -69,8 +69,12 @@ async fn test_layer2_sync_with_lowercase_a_tag() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b with git data (triggers discovery)"); - // 3. Wait for discovery - tokio::time::sleep(Duration::from_secs(1)).await; + // 3. Wait for discovery: events published before the syncing relay has + // connected to its peer can only be picked up by later recovery passes, + // so wait on the observable connection rather than a fixed sleep. + wait_for_sync_connection(relay_b.url(), 1, Duration::from_secs(30)) + .await + .expect("relay_b should establish a sync connection after discovery"); // 4. Create and send Layer 2 issue with lowercase 'a' tag let repo_coordinate = repo_coord(&keys, repo_id); @@ -103,7 +107,7 @@ async fn test_layer2_sync_with_lowercase_a_tag() { .author(keys.public_key()) .id(issue_id); - let synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(5)).await; + let synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(30)).await; println!("Issue {} synced to relay_b: {}", issue_id, synced); @@ -156,8 +160,12 @@ async fn test_layer2_sync_with_uppercase_a_tag() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b with git data (triggers discovery)"); - // 3. Wait for discovery - tokio::time::sleep(Duration::from_secs(1)).await; + // 3. Wait for discovery: events published before the syncing relay has + // connected to its peer can only be picked up by later recovery passes, + // so wait on the observable connection rather than a fixed sleep. + wait_for_sync_connection(relay_b.url(), 1, Duration::from_secs(30)) + .await + .expect("relay_b should establish a sync connection after discovery"); // 4. Create and send Layer 2 issue with uppercase 'A' tag let repo_coordinate = repo_coord(&keys, repo_id); @@ -193,7 +201,7 @@ async fn test_layer2_sync_with_uppercase_a_tag() { .author(keys.public_key()) .id(issue_id); - let synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(5)).await; + let synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(30)).await; println!("Issue {} synced to relay_b: {}", issue_id, synced); @@ -245,8 +253,12 @@ async fn test_layer2_sync_with_q_tag() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b with git data (triggers discovery)"); - // 3. Wait for discovery - tokio::time::sleep(Duration::from_secs(1)).await; + // 3. Wait for discovery: events published before the syncing relay has + // connected to its peer can only be picked up by later recovery passes, + // so wait on the observable connection rather than a fixed sleep. + wait_for_sync_connection(relay_b.url(), 1, Duration::from_secs(30)) + .await + .expect("relay_b should establish a sync connection after discovery"); // 4. Create and send Layer 2 issue with 'q' tag let repo_coordinate = repo_coord(&keys, repo_id); @@ -278,7 +290,7 @@ async fn test_layer2_sync_with_q_tag() { .author(keys.public_key()) .id(issue_id); - let synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(5)).await; + let synced = wait_for_event_on_relay(relay_b.url(), filter, Duration::from_secs(30)).await; println!("Issue {} synced to relay_b: {}", issue_id, synced); @@ -340,8 +352,12 @@ async fn test_layer3_sync_with_lowercase_e_tag() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b with git data (triggers discovery)"); - // 3. Wait for discovery - tokio::time::sleep(Duration::from_secs(1)).await; + // 3. Wait for discovery: events published before the syncing relay has + // connected to its peer can only be picked up by later recovery passes, + // so wait on the observable connection rather than a fixed sleep. + wait_for_sync_connection(relay_b.url(), 1, Duration::from_secs(30)) + .await + .expect("relay_b should establish a sync connection after discovery"); // 4. Create and send Layer 2 issue (parent event) let repo_coordinate = repo_coord(&keys, repo_id); @@ -363,7 +379,7 @@ async fn test_layer3_sync_with_lowercase_e_tag() { // 5. Wait for issue to sync to relay_b let issue_filter = Filter::new().kind(Kind::GitIssue).id(issue_id); 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!("Issue synced to relay_b: {}", issue_synced); assert!(issue_synced, "Layer 2 issue should sync first"); @@ -400,7 +416,7 @@ async fn test_layer3_sync_with_lowercase_e_tag() { .id(reply_id); let reply_synced = - wait_for_event_on_relay(relay_b.url(), reply_filter, Duration::from_secs(5)).await; + wait_for_event_on_relay(relay_b.url(), reply_filter, Duration::from_secs(30)).await; println!("Reply {} synced to relay_b: {}", reply_id, reply_synced); @@ -452,8 +468,12 @@ async fn test_layer3_sync_with_uppercase_e_tag() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b with git data (triggers discovery)"); - // 3. Wait for discovery - tokio::time::sleep(Duration::from_secs(1)).await; + // 3. Wait for discovery: events published before the syncing relay has + // connected to its peer can only be picked up by later recovery passes, + // so wait on the observable connection rather than a fixed sleep. + wait_for_sync_connection(relay_b.url(), 1, Duration::from_secs(30)) + .await + .expect("relay_b should establish a sync connection after discovery"); // 4. Create and send Layer 2 issue (parent event) let repo_coordinate = repo_coord(&keys, repo_id); @@ -475,7 +495,7 @@ async fn test_layer3_sync_with_uppercase_e_tag() { // 5. Wait for issue to sync to relay_b let issue_filter = Filter::new().kind(Kind::GitIssue).id(issue_id); 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!("Issue synced to relay_b: {}", issue_synced); assert!(issue_synced, "Layer 2 issue should sync first"); @@ -513,7 +533,7 @@ async fn test_layer3_sync_with_uppercase_e_tag() { .id(comment_id); let comment_synced = - wait_for_event_on_relay(relay_b.url(), comment_filter, Duration::from_secs(5)).await; + wait_for_event_on_relay(relay_b.url(), comment_filter, Duration::from_secs(30)).await; println!( "Comment {} synced to relay_b: {}", @@ -568,8 +588,12 @@ async fn test_layer3_sync_with_q_tag() { setup_announcement_on_relay(&relay_b, &keys, &domain_refs, repo_id).await; println!("Announcement set up on relay_b with git data (triggers discovery)"); - // 3. Wait for discovery - tokio::time::sleep(Duration::from_secs(1)).await; + // 3. Wait for discovery: events published before the syncing relay has + // connected to its peer can only be picked up by later recovery passes, + // so wait on the observable connection rather than a fixed sleep. + wait_for_sync_connection(relay_b.url(), 1, Duration::from_secs(30)) + .await + .expect("relay_b should establish a sync connection after discovery"); // 4. Create and send Layer 2 issue (parent event) let repo_coordinate = repo_coord(&keys, repo_id); @@ -591,7 +615,7 @@ async fn test_layer3_sync_with_q_tag() { // 5. Wait for issue to sync to relay_b let issue_filter = Filter::new().kind(Kind::GitIssue).id(issue_id); 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!("Issue synced to relay_b: {}", issue_synced); assert!(issue_synced, "Layer 2 issue should sync first"); @@ -628,7 +652,7 @@ async fn test_layer3_sync_with_q_tag() { .id(quote_id); let quote_synced = - wait_for_event_on_relay(relay_b.url(), quote_filter, Duration::from_secs(5)).await; + wait_for_event_on_relay(relay_b.url(), quote_filter, Duration::from_secs(30)).await; println!("Quote {} synced to relay_b: {}", quote_id, quote_synced);