diff --git a/crates/ecstore/src/cluster/rpc/network_probe.rs b/crates/ecstore/src/cluster/rpc/network_probe.rs index a3c8e9882..b4e58baad 100644 --- a/crates/ecstore/src/cluster/rpc/network_probe.rs +++ b/crates/ecstore/src/cluster/rpc/network_probe.rs @@ -14,13 +14,18 @@ //! Bounded, authenticated inter-node network probes. -use std::time::{Duration, Instant}; +use std::future::Future; +use std::pin::Pin; +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; use bytes::Bytes; use rustfs_protos::ChannelClass; use rustfs_protos::models::{PingBody, PingBodyBuilder}; use rustfs_protos::proto_gen::node_service::{PingRequest, PingResponse}; use thiserror::Error; +use tokio::time::Instant; use tokio_util::sync::CancellationToken; use tonic::{Code, Request}; @@ -33,6 +38,40 @@ pub const MAX_NETWORK_PROBE_DURATION: Duration = Duration::from_secs(30); const PING_PROTOCOL_VERSION: u64 = 1; const LATENCY_PAYLOAD: &[u8] = b"network-probe-latency-v1"; const EXPECTED_RESPONSE_PAYLOAD: &[u8] = b"hello, caller"; +const DIAGNOSTIC_BYTES_PER_SECOND: u64 = 1_048_576; +const DIAGNOSTIC_CHUNK_BYTES: u64 = 16_384; + +#[derive(Debug)] +struct DiagnosticPacing { + started: Instant, + deadline: Instant, + charged_bytes: AtomicU64, +} + +impl DiagnosticPacing { + fn reserve(&self, bytes: u64) -> Result { + if bytes == 0 { + return Err(NetworkPeerProbeError::LimitExceeded); + } + // Reservation is never refunded: a failed or dropped RPC may have sent some payload. + let previous = self + .charged_bytes + .fetch_update(Ordering::AcqRel, Ordering::Acquire, |charged| { + charged.checked_add(bytes).filter(|total| *total <= MAX_NETWORK_PROBE_BYTES) + }) + .map_err(|_| NetworkPeerProbeError::LimitExceeded)?; + let charged = previous.checked_add(bytes).ok_or(NetworkPeerProbeError::LimitExceeded)?; + let millis = u128::from(charged) + .checked_mul(1_000) + .ok_or(NetworkPeerProbeError::LimitExceeded)? + .div_ceil(u128::from(DIAGNOSTIC_BYTES_PER_SECOND)); + self.started + .checked_add(Duration::from_millis( + u64::try_from(millis).map_err(|_| NetworkPeerProbeError::LimitExceeded)?, + )) + .ok_or(NetworkPeerProbeError::LimitExceeded) + } +} #[derive(Clone, Debug, PartialEq, Eq)] pub struct NetworkPeerTarget { @@ -66,6 +105,7 @@ pub enum NetworkPeerProbeError { #[derive(Clone, Debug)] pub struct NetworkPeerProbeClient { targets: Vec, + diagnostic_pacing: Option>, } impl NetworkPeerProbeClient { @@ -80,7 +120,25 @@ impl NetworkPeerProbeClient { address, }) .collect(); - Self { targets } + Self { + targets, + diagnostic_pacing: None, + } + } + + /// Opt in one diagnostic run; clones share its charged payload and absolute deadline. + /// Default/admin probes retain their existing unpaced behavior. + pub fn with_diagnostic_pacing(mut self, started: Instant, duration: Duration) -> Result { + if self.diagnostic_pacing.is_some() || duration.is_zero() || duration > MAX_NETWORK_PROBE_DURATION { + return Err(NetworkPeerProbeError::LimitExceeded); + } + let deadline = started.checked_add(duration).ok_or(NetworkPeerProbeError::LimitExceeded)?; + self.diagnostic_pacing = Some(Arc::new(DiagnosticPacing { + started, + deadline, + charged_bytes: AtomicU64::new(0), + })); + Ok(self) } pub fn targets(&self) -> Vec { @@ -110,7 +168,19 @@ impl NetworkPeerProbeClient { .find(|target| target.alias == peer_alias) .map(|target| target.address.as_str()) .ok_or(NetworkPeerProbeError::UnknownPeer)?; - let probe = probe_peer(address, traffic_bytes); + let probe = probe_peer(address, traffic_bytes, self.diagnostic_pacing.as_deref(), cancel); + if let Some(pacing) = &self.diagnostic_pacing { + let deadline = Instant::now() + .checked_add(max_duration) + .ok_or(NetworkPeerProbeError::LimitExceeded)? + .min(pacing.deadline); + return tokio::select! { + biased; + () = cancel.cancelled() => Err(NetworkPeerProbeError::Cancelled), + () = tokio::time::sleep_until(deadline) => Err(NetworkPeerProbeError::TimedOut), + result = probe => result, + }; + } tokio::select! { () = cancel.cancelled() => Err(NetworkPeerProbeError::Cancelled), result = tokio::time::timeout(max_duration, probe) => { @@ -120,7 +190,12 @@ impl NetworkPeerProbeClient { } } -async fn probe_peer(address: &str, traffic_bytes: u64) -> Result { +async fn probe_peer( + address: &str, + traffic_bytes: u64, + pacing: Option<&DiagnosticPacing>, + cancel: &CancellationToken, +) -> Result { let started = Instant::now(); let mut control = client(address, ChannelClass::Control).await?; let latency_started = Instant::now(); @@ -132,15 +207,18 @@ async fn probe_peer(address: &str, traffic_bytes: u64) -> Result Result = Pin> + Send + 'a>>; + +async fn send_payload( + traffic_bytes: u64, + pacing: Option<&DiagnosticPacing>, + cancel: &CancellationToken, + sender: &mut T, + mut send: F, +) -> Result<(), NetworkPeerProbeError> +where + T: Send, + F: for<'a> FnMut(&'a mut T, Vec) -> PayloadSendFuture<'a> + Send, +{ + let mut remaining = traffic_bytes; + while remaining > 0 { + let bytes = if pacing.is_some() { + remaining.min(DIAGNOSTIC_CHUNK_BYTES) + } else { + remaining + }; + if let Some(pacing) = pacing { + if cancel.is_cancelled() { + return Err(NetworkPeerProbeError::Cancelled); + } + if Instant::now() >= pacing.deadline { + return Err(NetworkPeerProbeError::TimedOut); + } + let due = pacing.reserve(bytes)?; + tokio::select! { + biased; + () = cancel.cancelled() => return Err(NetworkPeerProbeError::Cancelled), + () = tokio::time::sleep_until(pacing.deadline) => return Err(NetworkPeerProbeError::TimedOut), + () = tokio::time::sleep_until(due) => {}, + } + } + let payload_len = usize::try_from(bytes).map_err(|_| NetworkPeerProbeError::LimitExceeded)?; + let payload = vec![0x5a; payload_len]; + if let Some(pacing) = pacing { + // Check before polling the send future, including a cancellation/deadline tie. + tokio::select! { + biased; + () = cancel.cancelled() => return Err(NetworkPeerProbeError::Cancelled), + () = tokio::time::sleep_until(pacing.deadline) => return Err(NetworkPeerProbeError::TimedOut), + result = send(sender, payload) => result?, + } + } else { + send(sender, payload).await?; + } + remaining -= bytes; + } + Ok(()) +} + async fn client( address: &str, class: ChannelClass, @@ -208,6 +339,205 @@ mod tests { use crate::layout::endpoint::Endpoint; use crate::layout::endpoints::{Endpoints, PoolEndpoints}; + fn pacing(duration: Duration) -> Arc { + NetworkPeerProbeClient::from_endpoint_pools(&topology()) + .with_diagnostic_pacing(Instant::now(), duration) + .expect("bounded diagnostic pacing window") + .diagnostic_pacing + .expect("opt-in pacing is present") + } + + #[test] + fn pacing_uses_exact_millisecond_ceiling_at_n_and_n_plus_one() { + for duration in [Duration::ZERO, MAX_NETWORK_PROBE_DURATION + Duration::from_nanos(1)] { + assert!(matches!( + NetworkPeerProbeClient::from_endpoint_pools(&topology()).with_diagnostic_pacing(Instant::now(), duration), + Err(NetworkPeerProbeError::LimitExceeded) + )); + } + for (bytes, millis) in [(1, 1), (1_047_527, 999), (1_047_528, 1_000), (1_048_576, 1_000)] { + let pacing = pacing(Duration::from_secs(2)); + assert_eq!( + pacing.reserve(bytes).expect("bounded reservation") - pacing.started, + Duration::from_millis(millis) + ); + } + for bytes in [0, MAX_NETWORK_PROBE_BYTES + 1, u64::MAX] { + let pacing = pacing(Duration::from_secs(2)); + assert_eq!(pacing.reserve(bytes), Err(NetworkPeerProbeError::LimitExceeded)); + assert_eq!(pacing.charged_bytes.load(Ordering::Acquire), 0); + } + } + + #[test] + fn cloned_diagnostic_clients_share_one_atomic_budget_without_reset() { + let client = NetworkPeerProbeClient::from_endpoint_pools(&topology()) + .with_diagnostic_pacing(Instant::now(), Duration::from_secs(2)) + .expect("pacing"); + let cloned = client.clone(); + let first = client.diagnostic_pacing.as_ref().expect("first budget"); + let second = cloned.diagnostic_pacing.as_ref().expect("cloned budget"); + assert!(Arc::ptr_eq(first, second)); + let mut due = std::thread::scope(|scope| { + let first = scope.spawn(|| first.reserve(MAX_NETWORK_PROBE_BYTES / 2)); + let second = scope.spawn(|| second.reserve(MAX_NETWORK_PROBE_BYTES / 2)); + [ + first.join().expect("first reservation thread").expect("first half"), + second.join().expect("second reservation thread").expect("second half"), + ] + }); + due.sort(); + assert_eq!( + due, + [ + first.started + Duration::from_millis(500), + first.started + Duration::from_secs(1) + ] + ); + assert_eq!(first.reserve(1), Err(NetworkPeerProbeError::LimitExceeded)); + assert_eq!(first.charged_bytes.load(Ordering::Acquire), MAX_NETWORK_PROBE_BYTES); + assert!(matches!( + cloned.with_diagnostic_pacing(Instant::now(), Duration::from_secs(2)), + Err(NetworkPeerProbeError::LimitExceeded) + )); + } + + #[tokio::test(start_paused = true)] + async fn failed_peer_payload_stays_charged_before_the_next_peer_send() { + let pacing = pacing(Duration::from_secs(2)); + let cancel = CancellationToken::new(); + let mut observed_partial_bytes = 0; + let failed = send_payload(32_768, Some(&pacing), &cancel, &mut observed_partial_bytes, |partial, payload| { + Box::pin(async move { + assert_eq!(payload.len(), 16_384); + // Model failure after a partial write: the sender cannot safely refund the remainder. + *partial += 8_192; + Err(NetworkPeerProbeError::Unreachable) + }) + }) + .await; + assert_eq!(failed, Err(NetworkPeerProbeError::Unreachable)); + assert_eq!(observed_partial_bytes, 8_192); + assert_eq!(pacing.charged_bytes.load(Ordering::Acquire), 16_384); + let started = pacing.started; + let mut second_bytes = 0; + send_payload(16_384, Some(&pacing), &cancel, &mut second_bytes, move |bytes, payload| { + Box::pin(async move { + assert!(started.elapsed() >= Duration::from_millis(32)); + *bytes += payload.len(); + Ok(()) + }) + }) + .await + .expect("second peer uses remaining shared budget"); + assert_eq!(second_bytes, 16_384); + assert_eq!(pacing.charged_bytes.load(Ordering::Acquire), 32_768); + } + + #[tokio::test(start_paused = true)] + async fn cancellation_and_deadline_take_priority_over_an_earned_send() { + for cancelled in [false, true] { + let pacing = pacing(Duration::from_millis(16)); + let cancel = CancellationToken::new(); + let mut sends = 0; + { + let mut sending = + std::pin::pin!(send_payload(8_192, Some(&pacing), &cancel, &mut sends, |sends, _| Box::pin(async move { + *sends += 1; + Ok(()) + }))); + assert!(futures::poll!(&mut sending).is_pending()); + tokio::time::advance(Duration::from_millis(16)).await; + if cancelled { + cancel.cancel(); + } + let expected = if cancelled { + NetworkPeerProbeError::Cancelled + } else { + NetworkPeerProbeError::TimedOut + }; + assert_eq!(sending.await, Err(expected)); + } + assert_eq!(sends, 0); + assert_eq!(pacing.charged_bytes.load(Ordering::Acquire), 8_192); + } + } + + #[tokio::test(start_paused = true)] + async fn cancellation_during_a_payload_rpc_retains_its_attempted_charge() { + let pacing = pacing(Duration::from_secs(2)); + let cancel = CancellationToken::new(); + let mut sends = 0; + { + let mut sending = std::pin::pin!(send_payload(16_384, Some(&pacing), &cancel, &mut sends, |sends, _| Box::pin( + async move { + *sends += 1; + std::future::pending().await + } + ))); + assert!(futures::poll!(&mut sending).is_pending()); + tokio::time::advance(Duration::from_millis(16)).await; + assert!(futures::poll!(&mut sending).is_pending()); + cancel.cancel(); + assert_eq!(sending.await, Err(NetworkPeerProbeError::Cancelled)); + } + assert_eq!(sends, 1); + assert_eq!(pacing.charged_bytes.load(Ordering::Acquire), 16_384); + } + + #[tokio::test(start_paused = true)] + async fn default_payload_send_remains_one_unpaced_rpc() { + let started = Instant::now(); + let mut sends = 0; + send_payload(32_768, None, &CancellationToken::new(), &mut sends, move |sends, payload| { + Box::pin(async move { + assert_eq!(started.elapsed(), Duration::ZERO); + assert_eq!(payload.len(), 32_768); + *sends += 1; + Ok(()) + }) + }) + .await + .expect("legacy admin send"); + assert_eq!(sends, 1); + assert!( + NetworkPeerProbeClient::from_endpoint_pools(&topology()) + .diagnostic_pacing + .is_none() + ); + } + + #[tokio::test(start_paused = true)] + async fn diagnostic_payload_is_charged_before_each_actual_send_poll() { + let started = tokio::time::Instant::now(); + let pacing = DiagnosticPacing { + started, + deadline: started + Duration::from_secs(2), + charged_bytes: AtomicU64::new(0), + }; + let mut attempted = 0_u64; + send_payload( + 32_768, + Some(&pacing), + &CancellationToken::new(), + &mut attempted, + move |attempted, payload| { + Box::pin(async move { + *attempted += u64::try_from(payload.len()).expect("bounded payload length"); + let available = u128::from(MAX_NETWORK_PROBE_BYTES) * started.elapsed().as_millis() / 1_000; + assert!( + u128::from(*attempted) <= available, + "send boundary emitted {attempted} bytes with only {available} bytes earned" + ); + Ok(()) + }) + }, + ) + .await + .expect("bounded diagnostic payload"); + assert_eq!(attempted, 32_768); + } + fn endpoint(value: &str, is_local: bool) -> Endpoint { let mut endpoint = Endpoint::try_from(value).expect("test endpoint should parse"); endpoint.is_local = is_local; @@ -230,6 +560,14 @@ mod tests { }]) } + #[test] + fn probe_future_remains_send() { + fn require_send(_: T) {} + let client = NetworkPeerProbeClient::from_endpoint_pools(&topology()); + let cancel = CancellationToken::new(); + require_send(client.probe("peer-1", 1, Duration::from_secs(1), &cancel)); + } + #[test] fn targets_are_unique_sorted_remote_nodes_with_stable_aliases() { let client = NetworkPeerProbeClient::from_endpoint_pools(&topology()); diff --git a/rustfs/src/connect/diagnostics/perf_network.rs b/rustfs/src/connect/diagnostics/perf_network.rs index 04e184379..f158d4809 100644 --- a/rustfs/src/connect/diagnostics/perf_network.rs +++ b/rustfs/src/connect/diagnostics/perf_network.rs @@ -20,7 +20,7 @@ use std::io::{Cursor, Write as _}; use std::path::Path; use std::pin::Pin; use std::sync::atomic::{AtomicBool, 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, SigningKey, signature::Signer as _}; @@ -29,6 +29,7 @@ use serde::Serialize; use sha2::{Digest as _, Sha256}; use thiserror::Error; use time::{OffsetDateTime, format_description::well_known::Rfc3339}; +use tokio::time::Instant; use tokio_util::sync::CancellationToken; use uuid::{Uuid, Variant, Version}; use zip::{CompressionMethod, ZipWriter, write::SimpleFileOptions}; @@ -416,7 +417,14 @@ pub async fn measure_network( if aliases != request.peer_aliases { return Err(NetworkPerformanceError::InvalidRequest); } - measure_network_with_harness(request, &harness, cancel).await + let started = Instant::now(); + let harness = RuntimeNetworkPeerHarness { + client: harness + .client + .with_diagnostic_pacing(started, request.duration) + .map_err(|_| NetworkPerformanceError::LimitExceeded)?, + }; + measure_network_with_harness_at(request, &harness, cancel, started).await } pub(crate) fn runtime_network_peer_aliases() -> Option> { @@ -433,14 +441,24 @@ pub async fn measure_network_with_harness( request: &NetworkPerformanceRequest, harness: &dyn NetworkPeerHarness, cancel: &CancellationToken, +) -> Result { + measure_network_with_harness_at(request, harness, cancel, Instant::now()).await +} + +async fn measure_network_with_harness_at( + request: &NetworkPerformanceRequest, + harness: &dyn NetworkPeerHarness, + cancel: &CancellationToken, + started: Instant, ) -> Result { request.validate(unix_now()?)?; if cancel.is_cancelled() { return Ok(terminal_measurement(request, NetworkOutcome::Cancelled, NetworkReasonCode::Cancelled)); } let _lease = CollectorLease::acquire()?; - let started = Instant::now(); - let deadline = tokio::time::Instant::now() + request.duration; + let deadline = started + .checked_add(request.duration) + .ok_or(NetworkPerformanceError::LimitExceeded)?; let mut peers = Vec::with_capacity(request.peer_aliases.len()); let mut transferred_bytes = 0_u64; let mut completed_units = 0_u32; @@ -448,6 +466,7 @@ pub async fn measure_network_with_harness( for alias in &request.peer_aliases { let peer_started = Instant::now(); let outcome = tokio::select! { + biased; () = cancel.cancelled() => { return Ok(cancelled_measurement(request, started.elapsed(), completed_units, peers)); } diff --git a/rustfs/tests/connect_perf_network.rs b/rustfs/tests/connect_perf_network.rs index 53606369b..64844ed4a 100644 --- a/rustfs/tests/connect_perf_network.rs +++ b/rustfs/tests/connect_perf_network.rs @@ -136,6 +136,200 @@ async fn unused_address() -> SocketAddr { address } +// The native collector is Unix-only. Fresh processes exercise normal runtime +// publication without replacing its immutable first-writer context. +#[cfg(unix)] +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn native_network_entry_uses_distributed_runtime() { + use std::path::PathBuf; + use std::process::Stdio; + use tokio::process::Command; + + const CHILD: &str = "RUSTFS_NETWORK_ENTRY_CHILD"; + const ROOT: &str = "RUSTFS_NETWORK_ENTRY_ROOT"; + const PORTS: &str = "RUSTFS_NETWORK_ENTRY_PORTS"; + const TEST: &str = "native_network_entry_uses_distributed_runtime"; + if let Ok(node) = std::env::var(CHILD) { + let root = PathBuf::from(std::env::var_os(ROOT).expect("owned root")); + let ports: Vec = std::env::var(PORTS) + .expect("owned ports") + .split(',') + .map(|port| port.parse().expect("port")) + .collect(); + let index: usize = node.parse().expect("node index"); + let volumes = ports + .iter() + .enumerate() + .flat_map(|(node, port)| { + let root = &root; + (0..2).map(move |disk| format!("http://127.0.0.1:{port}{}/node-{node}/disk-{disk}", root.display())) + }) + .collect(); + let server = tokio::time::timeout( + Duration::from_secs(90), + rustfs::embedded::RustFSServerBuilder::new() + .address(format!("127.0.0.1:{}", ports[index])) + .access_key("network-entry-access") + .secret_key("network-entry-secret") + .volumes(volumes) + .build(), + ) + .await + .expect("bounded distributed startup") + .expect("distributed startup"); + fs::write(root.join(format!("ready-{index}")), b"ready").expect("ready marker"); + tokio::time::timeout(Duration::from_secs(120), async { + loop { + if root.join("stop").exists() { + break; + } + if index == 0 && root.join("measure").exists() { + use rustfs_ecstore::api::disk::Endpoint; + use rustfs_ecstore::api::layout::{EndpointServerPools, Endpoints, PoolEndpoints}; + use rustfs_ecstore::api::rpc::{NetworkPeerProbeClient, NetworkPeerProbeError}; + + // This control uses the same live peer and authentication, but + // leaves ordinary/admin transport behavior unpaced. Its full + // payload must hit the real receiver's configured codec limit. + let mut local = Endpoint::try_from(format!("http://127.0.0.1:{}/local", ports[0]).as_str()) + .expect("local control endpoint"); + local.is_local = true; + let mut remote = Endpoint::try_from(format!("http://127.0.0.1:{}/remote", ports[1]).as_str()) + .expect("remote control endpoint"); + remote.is_local = false; + let topology = EndpointServerPools::from(vec![PoolEndpoints { + legacy: false, + set_count: 1, + drives_per_set: 2, + endpoints: Endpoints::from(vec![local, remote]), + cmd_line: String::new(), + platform: String::new(), + }]); + assert_eq!( + NetworkPeerProbeClient::from_endpoint_pools(&topology) + .probe("peer-1", 65_536, Duration::from_secs(2), &CancellationToken::new()) + .await, + Err(NetworkPeerProbeError::ProtocolFailure), + "unpaced control must exceed the real peer's message limit" + ); + let mut input = request(1); + input.duration = Duration::from_secs(2); + input.traffic_bytes_per_peer = 65_536; + let measured = measure_network(&input, &CancellationToken::new()) + .await + .expect("actual service entry"); + assert_eq!(measured.result.outcome(), NetworkOutcome::Succeeded, "{measured:?}"); + assert_eq!(measured.result.data().expect("real network data").transferred_bytes, 65_536); + fs::write(root.join("measured"), b"success").expect("result marker"); + while !root.join("stop").exists() { + tokio::time::sleep(Duration::from_millis(25)).await; + } + break; + } + tokio::time::sleep(Duration::from_millis(25)).await; + } + }) + .await + .expect("bounded child lifetime"); + server.shutdown().await; + return; + } + + let root = tempfile::Builder::new() + .prefix("rustfs-network-entry-") + .tempdir() + .expect("owned root"); + // Reserve both ports together so the OS cannot hand back the same port. + let reservations = [ + std::net::TcpListener::bind("127.0.0.1:0").expect("reserve node port"), + std::net::TcpListener::bind("127.0.0.1:0").expect("reserve peer port"), + ]; + let ports = reservations + .each_ref() + .map(|listener| listener.local_addr().expect("reserved address").port()); + drop(reservations); + let mut children = Vec::new(); + let result = tokio::time::timeout(Duration::from_secs(110), async { + let setup: Result<(), std::io::Error> = async { + for node in 0..2 { + for disk in 0..2 { + fs::create_dir_all(root.path().join(format!("node-{node}/disk-{disk}")))?; + } + let log = fs::File::create(root.path().join(format!("node-{node}.log")))?; + children.push( + Command::new(std::env::current_exe()?) + .args(["--exact", TEST, "--nocapture"]) + .env(CHILD, node.to_string()) + .env(ROOT, root.path()) + .env(PORTS, format!("{},{}", ports[0], ports[1])) + .env("RUSTFS_UNSAFE_BYPASS_DISK_CHECK", "true") + .env("RUSTFS_CONSOLE_ENABLE", "false") + // A full 64 KiB probe cannot fit, but production diagnostic chunks can. + .env("RUSTFS_INTERNODE_RPC_MAX_MESSAGE_SIZE", "32768") + .env("NO_PROXY", "localhost,127.0.0.1,::1") + .env("no_proxy", "localhost,127.0.0.1,::1") + .env_remove("HTTP_PROXY") + .env_remove("HTTPS_PROXY") + .env_remove("ALL_PROXY") + .env_remove("http_proxy") + .env_remove("https_proxy") + .env_remove("all_proxy") + .current_dir(root.path()) + .stdin(Stdio::null()) + .stdout(log.try_clone()?) + .stderr(log) + .kill_on_drop(true) + .spawn()?, + ); + } + Ok(()) + } + .await; + setup.map_err(|error| format!("child setup failed: {error}"))?; + loop { + for child in &mut children { + if let Some(status) = child.try_wait().map_err(|error| format!("child status failed: {error}"))? { + return Err(format!("child exited: {status}")); + } + } + if (0..2).all(|node| root.path().join(format!("ready-{node}")).exists()) { + fs::write(root.path().join("measure"), b"go").map_err(|error| format!("measurement marker failed: {error}"))?; + } + if root.path().join("measured").exists() { + return Ok(()); + } + tokio::time::sleep(Duration::from_millis(25)).await; + } + }) + .await; + let mut cleanup_errors = Vec::new(); + if let Err(error) = fs::write(root.path().join("stop"), b"stop") { + cleanup_errors.push(format!("stop marker failed: {error}")); + } + for child in &mut children { + match tokio::time::timeout(Duration::from_secs(10), child.wait()).await { + Ok(Ok(status)) if status.success() => {} + Ok(Ok(status)) => cleanup_errors.push(format!("child exited unsuccessfully: {status}")), + waited => { + cleanup_errors.push(format!("child failed to stop cleanly: {waited:?}")); + if let Err(error) = child.kill().await { + cleanup_errors.push(format!("kill child failed: {error}")); + } + if let Err(error) = child.wait().await { + cleanup_errors.push(format!("reap child failed: {error}")); + } + } + } + } + if !matches!(result, Ok(Ok(()))) || !cleanup_errors.is_empty() { + let path = root.keep(); + panic!( + "distributed entry attempt failed: {result:?}; cleanup: {cleanup_errors:?}; owned logs at {}", + path.display() + ); + } +} + #[tokio::test] async fn real_rustfs_peer_accepts_authenticated_payload_and_attributes_disconnect() { use rustfs::embedded::{RustFSServerBuilder, find_available_port}; @@ -144,11 +338,7 @@ async fn real_rustfs_peer_accepts_authenticated_payload_and_attributes_disconnec use rustfs_ecstore::api::rpc::{NetworkPeerProbeClient, NetworkPeerProbeError}; let _guard = TEST_HARNESS_LOCK.lock().await; - let port = match find_available_port() { - Ok(port) => port, - Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => return, - Err(error) => panic!("find free port: {error}"), - }; + let port = find_available_port().expect("real peer test requires an available loopback port"); let server = RustFSServerBuilder::new() .address(format!("127.0.0.1:{port}")) .access_key("network-probe-access") @@ -179,6 +369,19 @@ async fn real_rustfs_peer_accepts_authenticated_payload_and_attributes_disconnec assert!(!measurement.duration.is_zero()); assert!(!measurement.latency.is_zero()); + let pacing_started = tokio::time::Instant::now(); + let paced = client + .clone() + .with_diagnostic_pacing(pacing_started, Duration::from_secs(2)) + .expect("opt-in diagnostic pacing"); + let paced_measurement = paced + .probe("peer-1", 65_536, Duration::from_secs(2), &CancellationToken::new()) + .await + .expect("real authenticated peer accepts paced chunks"); + assert_eq!(paced_measurement.transferred_bytes, 65_536); + assert!(pacing_started.elapsed() >= Duration::from_millis(63)); + assert!(paced_measurement.latency <= paced_measurement.duration); + server.shutdown().await; assert_eq!( client @@ -345,6 +548,54 @@ async fn invalid_and_over_budget_requests_fail_before_peer_io() { assert_eq!(harness.0.load(Ordering::Relaxed), 0); } +#[tokio::test] +async fn network_admission_preserves_the_999_millisecond_n_and_n_plus_one_boundary() { + let _guard = TEST_HARNESS_LOCK.lock().await; + let harness = CountingHarness(AtomicUsize::new(0)); + let mut bounded = request(1); + bounded.duration = Duration::from_millis(999); + bounded.traffic_bytes_per_peer = 1_047_527; + let measurement = measure_network_with_harness(&bounded, &harness, &CancellationToken::new()) + .await + .expect("N is admitted"); + assert_eq!(measurement.result.outcome(), NetworkOutcome::Succeeded); + bounded.traffic_bytes_per_peer += 1; + assert!(matches!( + measure_network_with_harness(&bounded, &harness, &CancellationToken::new()).await, + Err(NetworkPerformanceError::LimitExceeded) + )); + assert_eq!(harness.0.load(Ordering::Relaxed), 1); +} + +#[tokio::test(start_paused = true)] +async fn run_and_peer_durations_share_the_tokio_clock_with_submillisecond_remainders() { + struct ClockHarness; + impl NetworkPeerHarness for ClockHarness { + fn probe<'a>(&'a self, _peer_alias: &'a str, traffic_bytes: u64, _cancel: &'a CancellationToken) -> PeerProbeFuture<'a> { + Box::pin(async move { + let started = tokio::time::Instant::now(); + tokio::time::advance(Duration::from_micros(1_500)).await; + Ok(PeerProbeMeasurement { + transferred_bytes: traffic_bytes, + duration: started.elapsed(), + latency: Duration::from_micros(500), + }) + }) + } + } + let _guard = TEST_HARNESS_LOCK.lock().await; + let mut request = request(1); + request.traffic_bytes_per_peer = 1; + let measurement = measure_network_with_harness(&request, &ClockHarness, &CancellationToken::new()) + .await + .expect("clock-bound result"); + let value = serde_json::to_value(&measurement.result).expect("result JSON"); + assert_eq!(value["durationMillis"], 1); + assert_eq!(value["data"]["durationMillis"], 1); + assert_eq!(measurement.peers[0].duration_millis, 1); + assert_eq!(measurement.peers[0].latency_micros, Some(500)); +} + struct BlockingHarness; impl NetworkPeerHarness for BlockingHarness {