From 3c7d85420ce2533674ea73a8dabe56b4fe423b27 Mon Sep 17 00:00:00 2001 From: Chris Date: Tue, 29 Sep 2026 16:48:34 +0800 Subject: [PATCH] fix(test): control clocks in Connect network fixtures (#8228) --- rustfs/tests/connect_perf_network.rs | 52 +++++++++++++++++----------- 1 file changed, 32 insertions(+), 20 deletions(-) diff --git a/rustfs/tests/connect_perf_network.rs b/rustfs/tests/connect_perf_network.rs index 64844ed4a..4136a9ca8 100644 --- a/rustfs/tests/connect_perf_network.rs +++ b/rustfs/tests/connect_perf_network.rs @@ -19,7 +19,7 @@ use std::net::SocketAddr; #[cfg(unix)] use std::os::unix::fs::PermissionsExt as _; use std::sync::atomic::{AtomicUsize, Ordering}; -use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; +use std::time::{Duration, SystemTime, UNIX_EPOCH}; use base64_simd::URL_SAFE_NO_PAD; use p256::ecdsa::{Signature, VerifyingKey, signature::Verifier as _}; @@ -77,7 +77,7 @@ impl NetworkPeerHarness for TcpPeerHarness { fn probe<'a>(&'a self, peer_alias: &'a str, traffic_bytes: u64, cancel: &'a CancellationToken) -> PeerProbeFuture<'a> { Box::pin(async move { let address = *self.addresses.get(peer_alias).ok_or(PeerProbeError::ProtocolFailure)?; - let started = Instant::now(); + let started = tokio::time::Instant::now(); let mut stream = tokio::select! { () = cancel.cancelled() => return Err(PeerProbeError::Cancelled), result = TcpStream::connect(address) => result.map_err(|_| PeerProbeError::Unreachable)?, @@ -91,6 +91,7 @@ impl NetworkPeerHarness for TcpPeerHarness { if pong != [0x52] { return Err(PeerProbeError::ProtocolFailure); } + tokio::time::advance(Duration::from_millis(1)).await; let latency = started.elapsed(); let payload = vec![0x5a; usize::try_from(traffic_bytes).map_err(|_| PeerProbeError::ProtocolFailure)?]; stream @@ -98,6 +99,7 @@ impl NetworkPeerHarness for TcpPeerHarness { .await .map_err(|_| PeerProbeError::ProtocolFailure)?; stream.shutdown().await.map_err(|_| PeerProbeError::ProtocolFailure)?; + tokio::time::advance(Duration::from_millis(1)).await; Ok(PeerProbeMeasurement { transferred_bytes: traffic_bytes, duration: started.elapsed(), @@ -107,6 +109,20 @@ impl NetworkPeerHarness for TcpPeerHarness { } } +async fn with_frozen_clock(future: F) -> F::Output { + tokio::pin!(future); + std::future::poll_fn(|cx| { + let result = future.as_mut().poll(cx); + if result.is_pending() { + // Real socket readiness must not auto-advance paused time to the + // collection deadline. The fixture advances at its I/O milestones. + cx.waker().wake_by_ref(); + } + result + }) + .await +} + async fn echo_peer() -> (SocketAddr, tokio::task::JoinHandle) { let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind controlled peer"); let address = listener.local_addr().expect("peer address"); @@ -418,7 +434,7 @@ async fn native_network_source_is_explicitly_unsupported_without_runtime_topolog )); } -#[tokio::test] +#[tokio::test(start_paused = true)] async fn controlled_peer_reports_exact_bytes_duration_latency_and_attributed_failure() { let _guard = TEST_HARNESS_LOCK.lock().await; let (working_address, working_peer) = echo_peer().await; @@ -429,7 +445,7 @@ async fn controlled_peer_reports_exact_bytes_duration_latency_and_attributed_fai ("peer-2".to_owned(), unavailable_address), ]), }; - let measurement = measure_network_with_harness(&request(2), &harness, &CancellationToken::new()) + let measurement = with_frozen_clock(measure_network_with_harness(&request(2), &harness, &CancellationToken::new())) .await .expect("controlled network measurement"); assert_eq!(working_peer.await.expect("peer task"), 4_096); @@ -439,11 +455,12 @@ async fn controlled_peer_reports_exact_bytes_duration_latency_and_attributed_fai assert_eq!(data.transferred_bytes, 4_096); assert_eq!(data.error_count, 1); assert_eq!(data.peer_count, 2); - assert!(data.duration_millis >= 1); + assert_eq!(data.duration_millis, 2); assert_eq!(measurement.peers[0].peer_alias, "peer-1"); assert_eq!(measurement.peers[0].reason_code, PeerReasonCode::Complete); assert_eq!(measurement.peers[0].transferred_bytes, 4_096); - assert!(measurement.peers[0].latency_micros.is_some()); + assert_eq!(measurement.peers[0].duration_millis, 2); + assert_eq!(measurement.peers[0].latency_micros, Some(1_000)); assert_eq!(measurement.peers[1].peer_alias, "peer-2"); assert_eq!(measurement.peers[1].reason_code, PeerReasonCode::Unreachable); assert_eq!(measurement.peers[1].transferred_bytes, 0); @@ -613,30 +630,25 @@ impl NetworkPeerHarness for BlockingHarness { } } -#[tokio::test] +#[tokio::test(start_paused = true)] async fn cancellation_stops_collection_and_only_one_collector_runs() { let _guard = TEST_HARNESS_LOCK.lock().await; let first_cancel = CancellationToken::new(); let second_cancel = CancellationToken::new(); let request = request(1); let first = measure_network_with_harness(&request, &BlockingHarness, &first_cancel); - let second = async { - tokio::time::sleep(Duration::from_millis(5)).await; - measure_network_with_harness(&request, &BlockingHarness, &second_cancel).await - }; - let cancellation = async { - tokio::time::sleep(Duration::from_millis(20)).await; - first_cancel.cancel(); - }; - let (first, second, ()) = tokio::join!(first, second, cancellation); - let first = first.expect("typed cancelled result"); + tokio::pin!(first); + assert!(futures_util::poll!(first.as_mut()).is_pending(), "first collector must hold its lease"); + let second = measure_network_with_harness(&request, &BlockingHarness, &second_cancel).await; + assert!(matches!(second, Err(NetworkPerformanceError::Busy))); + first_cancel.cancel(); + let first = first.await.expect("typed cancelled result"); assert_eq!(first.result.outcome(), NetworkOutcome::Cancelled); assert_eq!(first.result.reason_code(), NetworkReasonCode::Cancelled); assert!(first.result.data().is_none()); - assert!(matches!(second, Err(NetworkPerformanceError::Busy))); } -#[tokio::test] +#[tokio::test(start_paused = true)] async fn successful_result_has_signed_bounded_private_offline_export() { let _guard = TEST_HARNESS_LOCK.lock().await; let request = request(1); @@ -644,7 +656,7 @@ async fn successful_result_has_signed_bounded_private_offline_export() { let harness = TcpPeerHarness { addresses: BTreeMap::from([("peer-1".to_owned(), address)]), }; - let measurement = measure_network_with_harness(&request, &harness, &CancellationToken::new()) + let measurement = with_frozen_clock(measure_network_with_harness(&request, &harness, &CancellationToken::new())) .await .expect("successful controlled result"); assert_eq!(peer.await.expect("peer task"), 4_096);