mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-06 03:28:24 +00:00
test(transport): wait for a writer parked on a deaf peer to stay parked
Three transport tests park a stream connection's writer behind a peer that does not read, then check that it is parked by reading the send queue's free capacity 300 ms after the fill. On FreeBSD the check failed intermittently with 1 to 15 free slots. The writer was parked; the queue was not full. FreeBSD delays the loopback ACK for the last data the peer received, by about 40 ms, or until a retransmission draws it. When that ACK lands after the fill has stopped offering frames, it frees send-buffer space, the writer moves a few more frames into the kernel and parks again, and nothing refills the queue. On Linux the proxied test's queue was full within 2 ms of the fill and stayed full. The fill and the check now live in one helper shared by the TCP and proxied tests. After the fill it waits 300 ms and reads the free capacity; if the queue is not full it tops it up and waits again, up to ten times. A pass still needs a queue that stayed full for 300 ms with nothing sending, which a writer still taking frames cannot produce, and a writer that never stays parked still fails the setup. On a FreeBSD 15.1 VM, each test run alone 50 times failed in none of them at the default ACK delay, where the old setup failed the proxied test in 10 of 30 runs. With the delay raised to 100 ms, each passed 50 of 50, where the old setup failed the two TCP tests in 30 and 29 of 30 runs and the proxied test in 18 of 30. A peer that reads still fails the setup.
This commit is contained in:
@@ -370,7 +370,7 @@ pub(crate) async fn proxied_receive_loop<S: ProxiedStats, M>(
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::transport::packet_channel;
|
||||
use crate::transport::stream::next_conn_id;
|
||||
use crate::transport::stream::{next_conn_id, park_writer};
|
||||
use portable_atomic::{AtomicU64, Ordering};
|
||||
use tokio::io::AsyncReadExt;
|
||||
use tokio::net::TcpListener;
|
||||
@@ -555,33 +555,19 @@ mod tests {
|
||||
},
|
||||
);
|
||||
|
||||
// Fill without stopping at the first refusal, yielding so the writer
|
||||
// runs, and never keep a sender past the fill.
|
||||
// Never keep a sender past the fill: the teardown below relies on the
|
||||
// pool entry holding the only one.
|
||||
let frame = vec![0xAB; 1400];
|
||||
let mut queued = 0usize;
|
||||
let mut refused = 0usize;
|
||||
for _ in 0..8000 {
|
||||
let sent = {
|
||||
let guard = pool.lock().await;
|
||||
guard
|
||||
let queued = park_writer(
|
||||
async || {
|
||||
pool.lock()
|
||||
.await
|
||||
.get(&remote)
|
||||
.map(|c| c.send_tx.try_send(frame.clone()).is_ok())
|
||||
};
|
||||
if sent == Some(true) {
|
||||
queued += 1;
|
||||
} else {
|
||||
refused += 1;
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(300)).await;
|
||||
let capacity = pool.lock().await.get(&remote).map(|c| c.send_tx.capacity());
|
||||
assert_eq!(
|
||||
capacity,
|
||||
Some(0),
|
||||
"setup did not park the writer: queued={queued} refused={refused}"
|
||||
);
|
||||
assert!(refused > 0, "setup never filled the queue: queued={queued}");
|
||||
.is_some_and(|c| c.send_tx.try_send(frame.clone()).is_ok())
|
||||
},
|
||||
async || pool.lock().await.get(&remote).map(|c| c.send_tx.capacity()),
|
||||
)
|
||||
.await;
|
||||
|
||||
peer.shutdown().await.unwrap();
|
||||
assert!(
|
||||
|
||||
@@ -75,6 +75,72 @@ pub(crate) fn drain_writer(send_task: JoinHandle<()>, bound: Duration) -> JoinHa
|
||||
})
|
||||
}
|
||||
|
||||
/// How long a filled send queue must stay full, with nothing sending, before a
|
||||
/// test treats its writer as parked.
|
||||
#[cfg(test)]
|
||||
pub(crate) const PARK_SETTLE: Duration = Duration::from_millis(300);
|
||||
|
||||
/// How many settle intervals a test waits for its writer to stay parked.
|
||||
#[cfg(test)]
|
||||
pub(crate) const PARK_ROUNDS: usize = 10;
|
||||
|
||||
/// Fill a connection's send queue behind a peer that does not read, until the
|
||||
/// writer is parked in `write_all`, and return how many frames were queued.
|
||||
///
|
||||
/// `offer` tries to queue one frame and says whether it was accepted;
|
||||
/// `capacity` reads the queue's free slots, `None` if the connection is gone.
|
||||
/// Every one of 8000 offers is made whatever the previous one returned, with a
|
||||
/// yield between them so the writer runs. Then the queue must read full after
|
||||
/// a settle interval with nothing sending. Free capacity only grows while
|
||||
/// nothing sends, so a full reading means the writer took no frame for the
|
||||
/// whole interval, which it can do only while blocked in `write_all`.
|
||||
///
|
||||
/// A queue that is not full is topped up and the interval repeated, because a
|
||||
/// late ACK can free send-buffer space after the fill ends and let the writer
|
||||
/// take a few frames before it parks again; FreeBSD delays that ACK on
|
||||
/// loopback. A writer that never stays parked for a whole interval panics.
|
||||
#[cfg(test)]
|
||||
pub(crate) async fn park_writer(
|
||||
mut offer: impl AsyncFnMut() -> bool,
|
||||
mut capacity: impl AsyncFnMut() -> Option<usize>,
|
||||
) -> usize {
|
||||
let mut queued = 0usize;
|
||||
let mut refused = 0usize;
|
||||
for _ in 0..8000 {
|
||||
if offer().await {
|
||||
queued += 1;
|
||||
} else {
|
||||
refused += 1;
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
let mut seen = Vec::with_capacity(PARK_ROUNDS);
|
||||
for _ in 0..PARK_ROUNDS {
|
||||
tokio::time::sleep(PARK_SETTLE).await;
|
||||
let free = capacity().await;
|
||||
seen.push(free);
|
||||
match free {
|
||||
Some(0) => {
|
||||
assert!(refused > 0, "setup never filled the queue: queued={queued}");
|
||||
return queued;
|
||||
}
|
||||
Some(n) => {
|
||||
for _ in 0..=n {
|
||||
if !offer().await {
|
||||
break;
|
||||
}
|
||||
queued += 1;
|
||||
}
|
||||
}
|
||||
None => break,
|
||||
}
|
||||
}
|
||||
panic!(
|
||||
"setup did not park the writer: queued={queued} refused={refused} \
|
||||
capacity after each settle={seen:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
@@ -116,4 +182,82 @@ mod tests {
|
||||
.expect("the drain timer kept running after the writer exited")
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// Depth of the modelled send queue in the `park_writer` tests.
|
||||
const MODEL_DEPTH: usize = 4;
|
||||
|
||||
/// A queue a late ACK drained after the fill is topped up, and the helper
|
||||
/// returns at the first settle that finds it still full.
|
||||
///
|
||||
/// The first capacity read drains 2 frames and later reads drain none, so
|
||||
/// the fill queues 4, the top-up queues 2 more, and the second read ends
|
||||
/// the wait.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn park_writer_tops_up_a_queue_a_late_ack_drained_and_returns_once_it_stays_full() {
|
||||
let len = std::cell::Cell::new(0usize);
|
||||
let drain = std::cell::Cell::new(2usize);
|
||||
let calls = std::cell::Cell::new(0usize);
|
||||
let queued = park_writer(
|
||||
async || {
|
||||
let accept = len.get() < MODEL_DEPTH;
|
||||
if accept {
|
||||
len.set(len.get() + 1);
|
||||
}
|
||||
accept
|
||||
},
|
||||
async || {
|
||||
calls.set(calls.get() + 1);
|
||||
len.set(len.get().saturating_sub(drain.take()));
|
||||
Some(MODEL_DEPTH - len.get())
|
||||
},
|
||||
)
|
||||
.await;
|
||||
assert_eq!(queued, 6, "fill of 4 plus a top-up of 2");
|
||||
assert_eq!(
|
||||
calls.get(),
|
||||
2,
|
||||
"the second settle should find the queue full"
|
||||
);
|
||||
}
|
||||
|
||||
/// A writer that takes a frame in every settle interval never counts as
|
||||
/// parked.
|
||||
#[tokio::test(start_paused = true)]
|
||||
#[should_panic(expected = "setup did not park the writer")]
|
||||
async fn park_writer_panics_when_the_writer_takes_a_frame_in_every_settle() {
|
||||
let len = std::cell::Cell::new(0usize);
|
||||
let drain = std::cell::Cell::new(1usize);
|
||||
park_writer(
|
||||
async || {
|
||||
let accept = len.get() < MODEL_DEPTH;
|
||||
if accept {
|
||||
len.set(len.get() + 1);
|
||||
}
|
||||
accept
|
||||
},
|
||||
async || {
|
||||
len.set(len.get().saturating_sub(drain.get()));
|
||||
Some(MODEL_DEPTH - len.get())
|
||||
},
|
||||
)
|
||||
.await;
|
||||
}
|
||||
|
||||
/// A connection that is gone by the settle check never counts as parked.
|
||||
#[tokio::test(start_paused = true)]
|
||||
#[should_panic(expected = "setup did not park the writer")]
|
||||
async fn park_writer_panics_when_the_connection_is_gone() {
|
||||
let len = std::cell::Cell::new(0usize);
|
||||
park_writer(
|
||||
async || {
|
||||
let accept = len.get() < MODEL_DEPTH;
|
||||
if accept {
|
||||
len.set(len.get() + 1);
|
||||
}
|
||||
accept
|
||||
},
|
||||
async || None,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
|
||||
+17
-31
@@ -1356,6 +1356,7 @@ mod tests {
|
||||
use super::*;
|
||||
use crate::transport::framing::build_msg1_frame;
|
||||
use crate::transport::packet_channel;
|
||||
use crate::transport::stream::park_writer;
|
||||
use tokio::time::{Duration, timeout};
|
||||
|
||||
/// Poll `f` every 10ms until it holds or `limit` elapses.
|
||||
@@ -2386,40 +2387,25 @@ mod tests {
|
||||
socket.listen(8).unwrap()
|
||||
}
|
||||
|
||||
/// Fill `remote`'s send queue behind a peer that does not read, until the
|
||||
/// Fill `remote`'s send queue behind a peer that does not read until the
|
||||
/// writer is parked in `write_all`, and return how many frames were
|
||||
/// queued.
|
||||
///
|
||||
/// Every send is attempted whatever the previous one returned, with a
|
||||
/// yield between sends so the writer runs. A queue still full 300 ms
|
||||
/// after the last send means the writer could not drain it, which is the
|
||||
/// state the caller's teardown needs; anything else panics rather than
|
||||
/// letting the caller pass without it.
|
||||
/// queued. The fill and the parked-writer check are `park_writer`'s.
|
||||
async fn park_tcp_writer(t: &TcpTransport, remote: &TransportAddr, frame: &[u8]) -> usize {
|
||||
let mut queued = 0usize;
|
||||
let mut refused = 0usize;
|
||||
for _ in 0..8000 {
|
||||
match timeout(Duration::from_secs(2), t.send_async(remote, frame)).await {
|
||||
Ok(Ok(_)) => queued += 1,
|
||||
Ok(Err(_)) => refused += 1,
|
||||
park_writer(
|
||||
async || match timeout(Duration::from_secs(2), t.send_async(remote, frame)).await {
|
||||
Ok(Ok(_)) => true,
|
||||
Ok(Err(_)) => false,
|
||||
Err(_) => panic!("send blocked on a peer that stopped reading"),
|
||||
}
|
||||
tokio::task::yield_now().await;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(300)).await;
|
||||
let capacity = t
|
||||
.pool
|
||||
.lock()
|
||||
.await
|
||||
.get(remote)
|
||||
.map(|c| c.send_tx.capacity());
|
||||
assert_eq!(
|
||||
capacity,
|
||||
Some(0),
|
||||
"setup did not park the writer: queued={queued} refused={refused}"
|
||||
);
|
||||
assert!(refused > 0, "setup never filled the queue: queued={queued}");
|
||||
queued
|
||||
},
|
||||
async || {
|
||||
t.pool
|
||||
.lock()
|
||||
.await
|
||||
.get(remote)
|
||||
.map(|c| c.send_tx.capacity())
|
||||
},
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Read `stream` to EOF within `limit`, returning the byte count, or
|
||||
|
||||
Reference in New Issue
Block a user