fix(connect): pace diagnostic network payload sends (#8096)

* fix(connect): pace diagnostic network payload sends

* test(connect): exercise native network pacing entry

---------

Co-authored-by: Hauser <housemecn@gmail.com>
This commit is contained in:
Chris
2026-09-27 05:11:35 +08:00
committed by GitHub
parent 82024907c5
commit cf096b3c22
3 changed files with 629 additions and 21 deletions
+350 -12
View File
@@ -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<Instant, NetworkPeerProbeError> {
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<NetworkPeerTarget>,
diagnostic_pacing: Option<Arc<DiagnosticPacing>>,
}
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<Self, NetworkPeerProbeError> {
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<NetworkPeerTarget> {
@@ -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<NetworkPeerProbeMeasurement, NetworkPeerProbeError> {
async fn probe_peer(
address: &str,
traffic_bytes: u64,
pacing: Option<&DiagnosticPacing>,
cancel: &CancellationToken,
) -> Result<NetworkPeerProbeMeasurement, NetworkPeerProbeError> {
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<NetworkPeerProb
validate_ping_response(&latency_response)?;
let latency = latency_started.elapsed();
let payload_len = usize::try_from(traffic_bytes).map_err(|_| NetworkPeerProbeError::LimitExceeded)?;
let payload = vec![0x5a; payload_len];
let mut bulk = client(address, ChannelClass::Bulk).await?;
let response = bulk
.ping(Request::new(ping_request(&payload)))
.await
.map_err(map_status)?
.into_inner();
validate_ping_response(&response)?;
send_payload(traffic_bytes, pacing, cancel, &mut bulk, |bulk, payload| {
Box::pin(async move {
let response = bulk
.ping(Request::new(ping_request(&payload)))
.await
.map_err(map_status)?
.into_inner();
validate_ping_response(&response)
})
})
.await?;
Ok(NetworkPeerProbeMeasurement {
transferred_bytes: traffic_bytes,
@@ -149,6 +227,59 @@ async fn probe_peer(address: &str, traffic_bytes: u64) -> Result<NetworkPeerProb
})
}
type PayloadSendFuture<'a> = Pin<Box<dyn Future<Output = Result<(), NetworkPeerProbeError>> + Send + 'a>>;
async fn send_payload<T, F>(
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<u8>) -> 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<DiagnosticPacing> {
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: 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());
+23 -4
View File
@@ -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<Vec<String>> {
@@ -433,14 +441,24 @@ pub async fn measure_network_with_harness(
request: &NetworkPerformanceRequest,
harness: &dyn NetworkPeerHarness,
cancel: &CancellationToken,
) -> Result<NetworkMeasurement, NetworkPerformanceError> {
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<NetworkMeasurement, NetworkPerformanceError> {
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));
}
+256 -5
View File
@@ -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<u16> = 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 {