Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 40 additions & 3 deletions src/defs/bytes_counter.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
//! Counts bytes transferred during the download and upload tests.

use std::sync::Mutex;
use std::sync::{Arc, Mutex};
use std::time::Instant;

use crate::http::WriteMeter;

/// Tracks total bytes transferred and derives the average transfer rate.
///
/// The total is a mutex rather than an `AtomicU64` because 32-bit targets such
Expand All @@ -15,6 +17,16 @@ pub struct BytesCounter {
start: Mutex<Option<Instant>>,
mebi: bool,
upload_size: usize,
/// The ceiling on the total, for the upload phase; see `set_wire`.
wire: Option<Wire>,
}

/// What a client's connections had written when a phase started, and the meter
/// to ask again.
#[derive(Debug)]
struct Wire {
meter: Arc<WriteMeter>,
start: u64,
}

impl BytesCounter {
Expand All @@ -24,6 +36,7 @@ impl BytesCounter {
start: Mutex::new(None),
mebi: false,
upload_size: 0,
wire: None,
}
}

Expand All @@ -41,6 +54,26 @@ impl BytesCounter {
self.upload_size
}

/// Caps the total at what `meter`'s connections write from now on.
///
/// For the upload phase. Its total is body bytes counted as hyper takes
/// each frame, and hyper takes frames ahead of writing them, so without a
/// ceiling the total includes whatever is still queued -- hundreds of
/// kilobytes per connection, which is a quarter to a half of what a 1 Mbit
/// uplink carries in a whole test.
///
/// The meter counts request heads and chunk framing along with body bytes,
/// so as a measure of body bytes written it is high by the overhead: a
/// head is 119 bytes for the request this client sends, and chunk framing
/// is 8 bytes per 16 KiB frame. That is the whole error this leaves -- the
/// total lands between what the connections wrote of the body and that
/// plus the overhead -- and a request the peer took whole is counted
/// exactly, its body bytes alone being fewer than head and body together.
pub fn set_wire(&mut self, meter: Arc<WriteMeter>) {
let start = meter.written();
self.wire = Some(Wire { meter, start });
}

/// Starts the clock used for the average.
pub fn start(&self) {
*self.start.lock().unwrap() = Some(Instant::now());
Expand All @@ -51,9 +84,13 @@ impl BytesCounter {
*self.total.lock().unwrap() += n;
}

/// Total bytes read or written.
/// Total bytes read or written, at most what the wire ceiling allows.
pub fn total(&self) -> u64 {
*self.total.lock().unwrap()
let counted = *self.total.lock().unwrap();
match &self.wire {
Some(wire) => counted.min(wire.meter.written().saturating_sub(wire.start)),
None => counted,
}
}

fn elapsed_secs(&self) -> f64 {
Expand Down
152 changes: 86 additions & 66 deletions src/defs/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,9 @@ use crate::{write_debug, write_ui};
/// The stagger between starting concurrent transfer streams.
const RAMP_UP_DELAY: Duration = Duration::from_millis(200);
/// The chunk size the upload body is fed to the connection in. hyper takes a
/// whole chunk whenever its write buffer has room, so this is also how much of
/// a connection's counted upload can sit past that buffer, unsent.
/// whole chunk whenever its write buffer has room and writes up to sixteen of
/// them at once, so this is the granularity of the upload's write syscalls, not
/// their size.
const UPLOAD_CHUNK: usize = 16 * 1024;

/// A speed test server, as described by the server list JSON.
Expand Down Expand Up @@ -449,6 +450,10 @@ impl Server {
let mut counter = BytesCounter::new();
counter.set_mebi(opts.use_mebi);
counter.set_upload_size(opts.upload_size);
// The body counts frames as hyper takes them, which runs ahead of what
// the connections have written; this is what keeps the total from
// reporting the difference as sent.
counter.set_wire(client.wire());
let counter = Arc::new(counter);

// Pre-allocating one random blob and reusing it keeps the CPU out of the
Expand All @@ -464,9 +469,6 @@ impl Server {
};

let url = url_join_path(&self.get_url()?, &self.upload_url);
// Upload over the pool whose buffers are capped, so little is counted
// but still unsent when the window closes.
let client = client.for_uploads();

counter.start();
let spinner = self.start_transfer_spinner("Uploading... ", opts, &counter);
Expand Down Expand Up @@ -809,10 +811,9 @@ mod upload_body {
/// Each frame is added to the upload total as hyper takes it, so the total
/// is request body bytes, which is what the Go client's TeeReader around
/// its request body counts: no request heads, chunk framing or TLS
/// overhead. A frame taken is not yet sent. hyper takes one whenever its
/// write buffer has room, and that buffer is capped (see
/// `http::client_builder`), so a connection holds at most the buffer and
/// one frame counted but unsent when the window closes.
/// overhead. A frame taken is not yet sent -- hyper takes one whenever its
/// write buffer has room -- so the phase caps the total at what the
/// connections have written; see `BytesCounter::set_wire`.
pub(super) struct UploadBody {
payload: Option<Bytes>,
pos: usize,
Expand Down Expand Up @@ -1278,19 +1279,21 @@ mod transfer_tests {
);
}

/// hyper takes upload frames while its write buffer has room, and each
/// frame is counted as it is taken, so the buffer decides how much a
/// connection has counted but not sent when the window closes. hyper's
/// default of ~400 KB overstated a slow uplink by a quarter to a half.
/// hyper takes upload frames while its write buffer has room, which runs
/// ahead of writing them, and each frame is counted as it is taken. Without
/// the ceiling that queue is reported as sent: hyper writes sixteen frames
/// at a time, and on a slow uplink filling the queue is most of what a
/// window does.
#[tokio::test]
async fn a_connection_holds_at_most_a_buffer_and_a_frame_of_counted_upload() {
async fn an_upload_counts_no_more_than_the_connection_wrote() {
use crate::http::MeteredConnector;
use hyper_util::client::legacy::connect::{Connected, Connection};
use std::pin::Pin;
use std::task::{Context, Poll};

/// What the peer takes before it stops reading. Not a whole number of
/// 16-frame batches: hyper also stops at 16 queued buffers, and a
/// budget ending on that boundary would hide an uncapped buffer.
/// budget ending on that boundary would leave nothing queued.
const TAKES: usize = 300_000;

/// A connection whose peer takes `TAKES` bytes and then nothing more,
Expand Down Expand Up @@ -1361,82 +1364,99 @@ mod transfer_tests {
}
}

let taken = Arc::new(AtomicUsize::new(0));
let client = crate::http::client_builder(1, false, true)
.build::<_, crate::http::ReqBody>(Connect(taken.clone()));
let counter = Arc::new(BytesCounter::new());
let body = BodyExt::boxed(upload_body::UploadBody::new(None, counter.clone()));
let request = http::Request::post("http://stalled.invalid/")
.body(body)
.unwrap();
// Nothing ever answers, so this only ends at the timeout.
let _ = tokio::time::timeout(Duration::from_millis(500), client.request(request)).await;
// The same upload twice: without the ceiling to show the queue is
// counted, and with it.
let mut counted = Vec::new();
for ceiling in [false, true] {
let taken = Arc::new(AtomicUsize::new(0));
let meter = Arc::new(crate::http::WriteMeter::new());
let client = crate::http::client_builder(1, false).build::<_, crate::http::ReqBody>(
MeteredConnector::new(Connect(taken.clone()), meter.clone()),
);
let mut counter = BytesCounter::new();
if ceiling {
counter.set_wire(meter.clone());
}
let counter = Arc::new(counter);
let body = BodyExt::boxed(upload_body::UploadBody::new(None, counter.clone()));
let request = http::Request::post("http://stalled.invalid/")
.body(body)
.unwrap();
// Nothing ever answers, so this only ends at the timeout.
let _ = tokio::time::timeout(Duration::from_millis(500), client.request(request)).await;

let taken = taken.load(Ordering::SeqCst) as u64;
assert_eq!(
taken, TAKES as u64,
"the connection took less than it could"
);
assert_eq!(
meter.written(),
taken,
"the meter counted something other than what the peer took"
);
counted.push(counter.total());
}

let taken = taken.load(Ordering::SeqCst) as u64;
assert_eq!(
taken, TAKES as u64,
"the connection took less than it could"
);
// What the peer took includes the request head and chunk framing, so
// this slightly understates what is held; the bound has room to spare.
let held = counter.total().saturating_sub(taken);
let bound = (crate::http::H1_MAX_BUF + UPLOAD_CHUNK) as u64;
assert!(
held <= bound,
"{held} bytes were counted but not sent, more than the {bound} a buffer and a frame hold"
counted[0] > TAKES as u64,
"nothing was queued unsent, so this test cannot show a ceiling working"
);
// The peer took the request head and the chunk framing along with the
// body, so the ceiling is a little above the body bytes it received;
// what matters is that nothing still in hyper is in the total.
assert_eq!(
counted[1], TAKES as u64,
"{} bytes were counted, more than the {TAKES} the connection wrote",
counted[1]
);
}

/// The upload phase sends over the capped pool: a response head bigger
/// than the cap fails its request there, and a failed request is not
/// started again. Over the default pool it would be read and repeated.
/// The phase asks for the ceiling, which is the half of the accounting the
/// body cannot do for itself: with the peer no longer reading, the total is
/// what the connection wrote and not what the body handed over.
#[tokio::test]
async fn the_upload_phase_sends_over_the_capped_pool() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
async fn the_upload_phase_counts_no_more_than_its_connections_wrote() {
use tokio::io::AsyncReadExt;

let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let requests = Arc::new(AtomicUsize::new(0));
let seen = requests.clone();
tokio::spawn(async move {
while let Ok((mut stream, _)) = listener.accept().await {
let seen = seen.clone();
tokio::spawn(async move {
// Answer only once the head and the 1 KiB body are in.
let mut request = Vec::new();
// Read the head, then stop reading and hold the connection
// open, so the client fills the socket and hyper queues the
// rest.
let mut seen = Vec::new();
let mut buf = [0u8; 8192];
while !request
.windows(4)
.position(|w| w == b"\r\n\r\n")
.is_some_and(|head| request.len() >= head + 4 + 1024)
{
while !seen.windows(4).any(|w| w == b"\r\n\r\n") {
match stream.read(&mut buf).await {
Ok(n) if n > 0 => request.extend_from_slice(&buf[..n]),
Ok(n) if n > 0 => seen.extend_from_slice(&buf[..n]),
_ => return,
}
}
seen.fetch_add(1, Ordering::SeqCst);
let reply = format!(
"HTTP/1.1 200 OK\r\nX-Pad: {}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n",
"a".repeat(2 * crate::http::H1_MAX_BUF)
);
let _ = stream.write_all(reply.as_bytes()).await;
while matches!(stream.read(&mut buf).await, Ok(n) if n > 0) {}
std::future::pending::<()>().await;
});
}
});
let mut opts = opts(Duration::from_millis(500));

let mut opts = opts(Duration::from_secs(1));
opts.requests = 1;
// An endless body, so the window closes with the stream mid-transfer
// and whatever hyper holds is held for good.
opts.no_prealloc = true;

server_at(addr)
.upload(&client(), &TelemetryLog::new(), &opts)
let client = client_with_timeout(Duration::ZERO);
let (_, total) = server_at(addr)
.upload(&client, &TelemetryLog::new(), &opts)
.await
.unwrap();

assert_eq!(
requests.load(Ordering::SeqCst),
1,
"the upload read a response head bigger than the cap, so its pool is not capped"
let written = client.wire().written();
assert!(total > 0, "nothing was counted at all");
assert!(
total <= written,
"the phase counted {total} bytes, more than the {written} its connections wrote"
);
}
}
Loading
Loading