diff --git a/tests/sync/adaptive_pagination.rs b/tests/sync/adaptive_pagination.rs index b614e00..3013080 100644 --- a/tests/sync/adaptive_pagination.rs +++ b/tests/sync/adaptive_pagination.rs @@ -144,7 +144,22 @@ async fn honest_default_limit_stops_after_its_verification_page() { "/tmp/relay-{}.log", syncing.domain().split(':').next_back().unwrap() ); - let log = std::fs::read_to_string(&log_path).expect("read syncing relay log"); + let deadline = tokio::time::Instant::now() + Duration::from_secs(5); + let log = loop { + let log = std::fs::read_to_string(&log_path).expect("read syncing relay log"); + if log + .matches("Grouped subscription hit pagination threshold") + .count() + >= 1 + { + break log; + } + assert!( + tokio::time::Instant::now() < deadline, + "verification page was not scheduled before the deadline" + ); + tokio::task::yield_now().await; + }; assert_eq!( log.matches("Grouped subscription hit pagination threshold") .count(), diff --git a/tests/sync/metrics.rs b/tests/sync/metrics.rs index 481f2b8..1b79413 100644 --- a/tests/sync/metrics.rs +++ b/tests/sync/metrics.rs @@ -521,13 +521,22 @@ async fn test_relay_connected_status() { // Stop the source harness.stop_source(0).await; - tokio::time::sleep(Duration::from_secs(2)).await; - - let metrics = harness.get_metrics().await.unwrap(); + let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + let metrics = loop { + let metrics = harness.get_metrics().await.unwrap(); + if metrics.relay_connected(&source_url).is_none() { + break metrics; + } + assert!( + tokio::time::Instant::now() < deadline, + "relay was not retired" + ); + tokio::task::yield_now().await; + }; assert_eq!( metrics.relay_connected(&source_url), - Some(false), - "Should be disconnected from {}", + None, + "Retired relay series should be removed for {}", source_url ); @@ -614,7 +623,19 @@ async fn test_multi_source_aggregate_counts() { .await; tokio::time::sleep(Duration::from_secs(2)).await; - let metrics = harness.get_metrics().await.unwrap(); + let deadline = tokio::time::Instant::now() + Duration::from_secs(10); + let metrics = loop { + let metrics = harness.get_metrics().await.unwrap(); + if metrics.relays_tracked_total() == Some(1) && metrics.relays_connected_total() == Some(1) + { + break metrics; + } + assert!( + tokio::time::Instant::now() < deadline, + "relay did not become connected" + ); + tokio::task::yield_now().await; + }; println!("Tracked total: {:?}", metrics.relays_tracked_total()); println!("Connected total: {:?}", metrics.relays_connected_total()); @@ -630,30 +651,5 @@ async fn test_multi_source_aggregate_counts() { "Should have 1 connected" ); - // Stop source, verify connected drops to 0 - harness.stop_source(0).await; - - let metrics = harness.get_metrics().await.unwrap(); - - println!( - "After stop - Tracked total: {:?}", - metrics.relays_tracked_total() - ); - println!( - "After stop - Connected total: {:?}", - metrics.relays_connected_total() - ); - - assert_eq!( - metrics.relays_tracked_total(), - Some(1), - "Still tracking 1 relay" - ); - assert_eq!( - metrics.relays_connected_total(), - Some(0), - "Should have 0 connected (waited up to 10s for disconnect detection)" - ); - harness.stop_all().await; }