mirror of
https://github.com/rustfs/rustfs.git
synced 2026-10-04 12:31:36 +00:00
fix(test): control clocks in Connect network fixtures (#8228)
This commit is contained in:
@@ -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<F: std::future::Future>(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<usize>) {
|
||||
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);
|
||||
|
||||
Reference in New Issue
Block a user