use std::{
collections::HashMap,
fmt::{Debug, Display},
future::Future,
sync::Arc,
time::Duration,
};
use crate::nq_core::{
ConnectionTiming, ConnectionType, Network, ScopedHeaders, Time, Timestamp,
client::{Direction, ThroughputClient, wait_for_finish},
};
use crate::nq_load_generator::{LoadConfig, LoadGenerator, LoadedConnection};
use crate::nq_stats::{TimeSeries, instant_minus_intervals};
use humansize::{DECIMAL, format_size};
use tokio::{select, sync::mpsc};
use tokio_util::sync::CancellationToken;
use tracing::{Instrument, debug, error, info, warn};
use url::Url;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ConnectionErrorPolicy {
#[default]
Retire,
Abort,
}
#[derive(Debug, Clone)]
pub struct ResponsivenessConfig {
pub large_download_url: Url,
pub small_download_url: Url,
pub upload_url: Url,
pub moving_average_distance: usize,
pub interval_duration: Duration,
pub test_duration: Duration,
pub trimmed_mean_percent: f64,
pub std_tolerance: f64,
pub max_loaded_connections: usize,
pub conn_type: ConnectionType,
pub determine_load_only: bool,
pub upload_bytes_per_request: usize,
pub on_connection_error: ConnectionErrorPolicy,
pub scoped_headers: Option<ScopedHeaders>,
}
impl ResponsivenessConfig {
pub fn load_config(&self) -> LoadConfig {
LoadConfig {
headers: HashMap::default(),
scoped_headers: self.scoped_headers.clone(),
download_url: self.large_download_url.clone(),
upload_url: self.upload_url.clone(),
}
}
}
pub const DEFAULT_UPLOAD_BYTES_PER_REQUEST: usize = 100_000_000;
impl Default for ResponsivenessConfig {
fn default() -> Self {
Self {
large_download_url: "https://h3.speed.cloudflare.com/__down?bytes=10000000000"
.parse()
.unwrap(),
small_download_url: "https://h3.speed.cloudflare.com/__down?bytes=10"
.parse()
.unwrap(),
upload_url: "https://h3.speed.cloudflare.com/__up".parse().unwrap(),
moving_average_distance: 4,
interval_duration: Duration::from_millis(1000),
test_duration: Duration::from_secs(20),
trimmed_mean_percent: 0.95,
std_tolerance: 0.05,
max_loaded_connections: 16,
conn_type: ConnectionType::H2,
determine_load_only: false,
upload_bytes_per_request: DEFAULT_UPLOAD_BYTES_PER_REQUEST,
on_connection_error: ConnectionErrorPolicy::default(),
scoped_headers: None,
}
}
}
pub struct Responsiveness {
start: Timestamp,
config: ResponsivenessConfig,
load_generator: LoadGenerator,
foreign_probe_results: ForeignProbeResults,
self_probe_results: SelfProbeResults,
average_goodput_series: TimeSeries,
rpm_series: TimeSeries,
goodput_saturated: bool,
rpm_saturated: bool,
direction: Direction,
rpm: Option<f64>,
last_rpm: Option<f64>,
capacity: f64,
failed_connections: usize,
starved_intervals: usize,
}
impl Responsiveness {
pub fn new(config: ResponsivenessConfig, download: bool) -> anyhow::Result<Self> {
let load_generator = LoadGenerator::new(config.load_config())?;
let upload_bytes_per_request = config.upload_bytes_per_request;
Ok(Self {
start: Timestamp::now(),
config,
load_generator,
foreign_probe_results: Default::default(),
self_probe_results: Default::default(),
average_goodput_series: TimeSeries::new(),
rpm_series: TimeSeries::new(),
failed_connections: 0,
starved_intervals: 0,
goodput_saturated: false,
rpm_saturated: false,
direction: if download {
Direction::Down
} else {
Direction::Up(upload_bytes_per_request)
},
rpm: None,
last_rpm: None,
capacity: 0.0,
})
}
}
impl Responsiveness {
pub async fn run_test(
mut self,
network: Arc<dyn Network>,
time: Arc<dyn Time>,
shutdown: CancellationToken,
) -> anyhow::Result<ResponsivenessResult> {
let env = Env { time, network };
self.start = env.time.now();
info!("running responsiveness test: {:?}", self.config);
let mut interval = None;
let mut interval_timer = tokio::time::interval(self.config.interval_duration);
let (event_tx, mut event_rx) = mpsc::channel(1024);
self.new_load_generating_connection(event_tx.clone(), &env, shutdown.clone())?;
if !self.config.determine_load_only {
self.send_foreign_probe(event_tx.clone(), &env, shutdown.clone())?;
}
loop {
select! {
Some(event) = event_rx.recv() => {
match event {
Event::NewLoadedConnection(connection) => {
self.load_generator.push(connection);
}
Event::ForeignProbe(f) => {
self.foreign_probe_results.add(f);
if !self.send_self_probe(event_tx.clone(), &env, shutdown.clone())? {
self.send_foreign_probe(event_tx.clone(), &env, shutdown.clone())?;
}
}
Event::SelfProbe(s) => {
self.self_probe_results.add(s);
self.send_foreign_probe(event_tx.clone(), &env, shutdown.clone())?;
}
Event::Error(e) => {
error!("error: {e}");
}
}
}
_ = interval_timer.tick() => {
self.load_generator.update();
if let Some(interval) = interval.as_mut() {
if self.on_interval(*interval, event_tx.clone(), &env, shutdown.clone()).await? {
break;
}
*interval += 1;
} else {
interval = Some(0);
}
}
_ = shutdown.cancelled() => {
debug!("shutdown requested");
break;
}
};
if env.time.now().duration_since(self.start) > self.config.test_duration {
break;
}
}
let now = env.time.now();
self.rpm = select_reported_rpm(self.rpm, self.last_rpm);
let mut loads = self.load_generator.into_connections();
loads.iter_mut().for_each(|load| load.stop());
Ok(ResponsivenessResult {
capacity: self.capacity,
rpm: self.rpm,
self_probe_latencies: self.self_probe_results.http,
loaded_connections: loads,
failed_connections: self.failed_connections,
duration: now.duration_since(self.start),
average_goodput_series: self.average_goodput_series,
})
}
async fn on_interval(
&mut self,
interval: usize,
event_tx: mpsc::Sender<Event>,
env: &Env,
shutdown: CancellationToken,
) -> anyhow::Result<bool> {
let end_data_interval = self.start + self.config.interval_duration * interval as u32;
let start_data_interval = instant_minus_intervals(
end_data_interval,
self.config.moving_average_distance,
self.config.interval_duration,
);
self.enforce_connection_error_policy()?;
if self.load_generator.count_loads() < self.config.max_loaded_connections
&& interval % 2 == 0
{
self.new_load_generating_connection(event_tx, env, shutdown)?;
}
let current_goodput = self.current_average_throughput(end_data_interval);
self.average_goodput_series
.add(end_data_interval, current_goodput);
let std_goodput = self
.average_goodput_series
.interval_std(start_data_interval, end_data_interval)
.unwrap_or(f64::MAX);
let goodput_saturated = std_goodput < current_goodput * self.config.std_tolerance;
if goodput_saturated {
self.capacity = current_goodput;
self.goodput_saturated = true;
}
let current_rpm = compute_responsiveness(
&self.foreign_probe_results,
&self.self_probe_results,
start_data_interval,
end_data_interval,
self.config.trimmed_mean_percent,
);
let current_rpm_or_zero = current_rpm.unwrap_or(0.0);
self.rpm_series.add(end_data_interval, current_rpm_or_zero);
if let Some(current_rpm) = current_rpm {
self.last_rpm = Some(current_rpm);
}
let std_rpm = self
.rpm_series
.interval_std(start_data_interval, end_data_interval);
let is_rpm_saturated = if let Some(std_rpm) = std_rpm {
if std_rpm < current_rpm_or_zero * self.config.std_tolerance {
self.rpm = Some(current_rpm_or_zero);
self.rpm_saturated = true;
true
} else {
false
}
} else {
false
};
self.log_interval(
interval,
current_goodput,
std_goodput,
goodput_saturated,
current_rpm_or_zero,
current_rpm.is_some(),
std_rpm,
is_rpm_saturated,
);
Ok(self.goodput_saturated && self.rpm_saturated)
}
#[allow(clippy::too_many_arguments)]
fn log_interval(
&mut self,
interval: usize,
current_goodput: f64,
std_goodput: f64,
goodput_saturated: bool,
current_rpm: f64,
rpm_measured: bool,
std_rpm: Option<f64>,
is_rpm_saturated: bool,
) {
let custom_options = humansize::FormatSizeOptions::from(DECIMAL)
.base_unit(humansize::BaseUnit::Bit)
.long_units(false)
.decimal_places(2);
info!(
interval,
loads = self.load_generator.count_loads(),
throughput = format_size(current_goodput as usize, custom_options),
rpm = current_rpm,
throughput_saturated = goodput_saturated,
rpm_saturated = is_rpm_saturated,
"interval finished"
);
if !rpm_measured {
warn!(
interval,
"no probe measurements in this interval's window; recorded 0 RPM, \
which inflates the stability std for the next MAD intervals"
);
}
info!(
interval,
throughput_std = format_size(std_goodput as usize, custom_options),
throughput_target_std = format_size(
(current_goodput * self.config.std_tolerance) as usize,
custom_options
),
rpm_std = std_rpm.unwrap_or(f64::NAN),
rpm_target_std = current_rpm * self.config.std_tolerance,
"interval stats"
);
}
fn current_average_throughput(&self, end_data_interval: Timestamp) -> f64 {
let start_data_interval =
instant_minus_intervals(end_data_interval, 4, self.config.interval_duration);
let mut bytes_seen = 0.0;
for connection in self.load_generator.connections() {
bytes_seen += connection
.total_bytes_series()
.interval_sum(start_data_interval, end_data_interval);
}
let total_time = end_data_interval
.duration_since(start_data_interval)
.as_secs_f64();
8.0 * bytes_seen / total_time
}
fn enforce_connection_error_policy(&mut self) -> anyhow::Result<()> {
let failed = self.load_generator.count_failed_loads();
let newly_failed = failed.saturating_sub(self.failed_connections);
self.failed_connections = failed;
if newly_failed > 0 {
let reason = self
.load_generator
.connections()
.filter_map(|c| c.failure_reason())
.last()
.unwrap_or("connection terminated early")
.to_owned();
warn!(
newly_failed,
total_failed = failed,
reason = %reason,
"load-generating connection(s) terminated with an error"
);
if self.config.on_connection_error == ConnectionErrorPolicy::Abort {
anyhow::bail!(
"aborting test: {failed} load-generating connection(s) terminated with an \
error (most recent: {reason})"
);
}
}
if failed > 0 && self.load_generator.count_loads() == 0 {
self.starved_intervals += 1;
if self.starved_intervals >= 2 {
anyhow::bail!(
"aborting test: no load-generating connections could be sustained \
({failed} terminated with an error); the link was never saturated so a \
responsiveness result would be meaningless"
);
}
} else {
self.starved_intervals = 0;
}
Ok(())
}
#[tracing::instrument(skip_all)]
fn new_load_generating_connection(
&self,
event_tx: mpsc::Sender<Event>,
env: &Env,
shutdown: CancellationToken,
) -> anyhow::Result<()> {
let oneshot_res = self.load_generator.new_loaded_connection(
self.direction,
self.config.conn_type,
Arc::clone(&env.network),
Arc::clone(&env.time),
shutdown,
)?;
tokio::spawn(
async move {
let _ = match oneshot_res.await {
Ok(conn) => event_tx.send(Event::NewLoadedConnection(conn)),
Err(e) => event_tx.send(Event::Error(e)),
}
.await;
}
.in_current_span(),
);
Ok(())
}
fn send_foreign_probe(
&mut self,
event_tx: mpsc::Sender<Event>,
env: &Env,
shutdown: CancellationToken,
) -> anyhow::Result<()> {
let client = ThroughputClient::download()
.new_connection(ConnectionType::H2)
.scoped_headers(self.config.scoped_headers.clone());
let inflight_body_fut = client.send(
self.config.small_download_url.as_str().parse()?,
Arc::clone(&env.network),
Arc::clone(&env.time),
shutdown,
)?;
tokio::spawn(report_err(
event_tx.clone(),
async move {
let inflight_body = inflight_body_fut.await?;
let finished_result = wait_for_finish(inflight_body.events).await?;
let Some(connection_timing) = inflight_body.timing else {
anyhow::bail!("a new connection with timing should have been created");
};
let (tcp, tls, http) =
foreign_probe_phases(&connection_timing, finished_result.finished_at);
if event_tx
.send(Event::ForeignProbe(ForeignProbeResult {
start: connection_timing.start(),
tcp,
tls,
http,
}))
.await
.is_err()
{
anyhow::bail!("unable to send foreign probe result");
}
Ok(())
}
.in_current_span(),
));
Ok(())
}
fn send_self_probe(
&mut self,
event_tx: mpsc::Sender<Event>,
env: &Env,
shutdown: CancellationToken,
) -> anyhow::Result<bool> {
let Some(connection) = self.load_generator.random_connection() else {
return Ok(false);
};
let client = ThroughputClient::download()
.with_connection(connection)
.scoped_headers(self.config.scoped_headers.clone());
let inflight_body_fut = client.send(
self.config.small_download_url.as_str().parse()?,
Arc::clone(&env.network),
Arc::clone(&env.time),
shutdown,
)?;
tokio::spawn(report_err(
event_tx.clone(),
async move {
let inflight_body = inflight_body_fut.await?;
let finish_result = wait_for_finish(inflight_body.events).await?;
debug!("self_probe_finished: {finish_result:?}");
if event_tx
.send(Event::SelfProbe(SelfProbeResult {
start: inflight_body.start,
time_body: finish_result
.finished_at
.duration_since(inflight_body.start),
}))
.await
.is_err()
{
anyhow::bail!("unable to send self probe result");
}
Ok(())
}
.in_current_span(),
));
Ok(true)
}
}
async fn report_err(event_tx: mpsc::Sender<Event>, f: impl Future<Output = anyhow::Result<()>>) {
if let Err(e) = f.await {
let _ = event_tx.send(Event::Error(e)).await;
}
}
#[derive(Default)]
pub struct ForeignProbeResults {
connect: TimeSeries,
secure: TimeSeries,
http: TimeSeries,
}
impl ForeignProbeResults {
pub fn add(&mut self, result: ForeignProbeResult) {
self.connect
.add(result.start, result.tcp.as_secs_f64() * 1000.0);
self.secure
.add(result.start, result.tls.as_secs_f64() * 1000.0);
self.http
.add(result.start, result.http.as_secs_f64() * 1000.0);
}
pub fn connect(&self) -> &TimeSeries {
&self.connect
}
pub fn secure(&self) -> &TimeSeries {
&self.secure
}
pub fn http(&self) -> &TimeSeries {
&self.http
}
}
#[derive(Default)]
pub struct SelfProbeResults {
http: TimeSeries,
}
impl SelfProbeResults {
pub fn add(&mut self, result: SelfProbeResult) {
self.http
.add(result.start, result.time_body.as_secs_f64() * 1000.0);
}
pub fn http(&self) -> &TimeSeries {
&self.http
}
}
fn select_reported_rpm(saturated: Option<f64>, last_interval: Option<f64>) -> Option<f64> {
saturated.or(last_interval)
}
fn compute_responsiveness(
foreign_results: &ForeignProbeResults,
self_results: &SelfProbeResults,
from: Timestamp,
to: Timestamp,
percentile: f64,
) -> Option<f64> {
let tm = |ts: &TimeSeries| ts.interval_trimmed_mean(from, to, percentile);
let tcp_f = tm(foreign_results.connect())?;
let tls_f = tm(foreign_results.secure())?;
let http_f = tm(foreign_results.http())?;
let http_l = tm(self_results.http())?;
let foreign_rtt = (tcp_f + tls_f + http_f) / 3.0;
let loaded_rtt = http_l;
if foreign_rtt <= 0.0 || loaded_rtt <= 0.0 {
return None;
}
let foreign_rpm = 60_000.0 / foreign_rtt;
let loaded_rpm = 60_000.0 / loaded_rtt;
let responsiveness = (foreign_rpm + loaded_rpm) / 2.0;
responsiveness.is_finite().then_some(responsiveness)
}
#[derive(Debug)]
pub struct ForeignProbeResult {
start: Timestamp,
tcp: Duration,
tls: Duration,
http: Duration,
}
fn foreign_probe_phases(
timing: &ConnectionTiming,
finished_at: Timestamp,
) -> (Duration, Duration, Duration) {
let tcp_f = timing.tcp_handshake();
let tls_f = timing.tls_handshake() / timing.tls_round_trips();
let request_issued = timing.start() + timing.time_application();
let http_f = finished_at.duration_since(request_issued);
(tcp_f, tls_f, http_f)
}
#[derive(Debug)]
pub struct SelfProbeResult {
start: Timestamp,
time_body: Duration,
}
enum Event {
ForeignProbe(ForeignProbeResult),
SelfProbe(SelfProbeResult),
NewLoadedConnection(LoadedConnection),
Error(anyhow::Error),
}
impl Debug for Event {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::ForeignProbe(_) => f.debug_tuple("ForeignProbe").finish(),
Self::SelfProbe(_) => f.debug_tuple("SelfProbe").finish(),
Self::NewLoadedConnection(_) => f.debug_tuple("NewLoadedConnection").finish(),
Self::Error(_) => f.debug_tuple("Error").finish(),
}
}
}
#[derive(Clone)]
struct Env {
time: Arc<dyn Time>,
network: Arc<dyn Network>,
}
#[derive(Default, Debug)]
pub struct ResponsivenessResult {
pub duration: Duration,
pub capacity: f64,
pub rpm: Option<f64>,
pub self_probe_latencies: TimeSeries,
pub loaded_connections: Vec<LoadedConnection>,
pub average_goodput_series: TimeSeries,
pub failed_connections: usize,
}
impl ResponsivenessResult {
pub fn throughput(&self) -> Option<usize> {
self.average_goodput_series
.quantile(0.90)
.map(|t| t as usize)
}
}
impl Display for ResponsivenessResult {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let custom_options = humansize::FormatSizeOptions::from(DECIMAL)
.base_unit(humansize::BaseUnit::Bit)
.long_units(false)
.decimal_places(2);
writeln!(
f,
"{:8}: {}/s",
"capacity",
format_size(self.capacity as usize, custom_options)
)?;
match self.rpm {
Some(rpm) => write!(f, "{:>8}: {}", "rpm", rpm.round() as usize),
None => write!(f, "{:>8}: unavailable", "rpm"),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::time::Duration;
fn ms(v: f64) -> Duration {
Duration::from_secs_f64(v / 1000.0)
}
fn series(
tcp_ms: f64,
tls_ms: f64,
http_f_ms: f64,
http_l_ms: f64,
) -> (ForeignProbeResults, SelfProbeResults, Timestamp, Timestamp) {
let start = Timestamp::now();
let mut foreign = ForeignProbeResults::default();
let mut selfp = SelfProbeResults::default();
for i in 0..10u64 {
let at = start + Duration::from_millis(i);
foreign.add(ForeignProbeResult {
start: at,
tcp: ms(tcp_ms),
tls: ms(tls_ms),
http: ms(http_f_ms),
});
selfp.add(SelfProbeResult {
start: at,
time_body: ms(http_l_ms),
});
}
(foreign, selfp, start, start + Duration::from_millis(100))
}
fn draft03(tcp: f64, tls: f64, http_f: f64, http_l: f64) -> f64 {
let foreign_sum = tcp + tls + http_f;
60_000.0 / (foreign_sum / 6.0 + http_l / 2.0)
}
#[test]
fn arithmetic_mean_of_the_two_rpms() {
let (f, s, from, to) = series(30.0, 30.0, 30.0, 30.0);
let rpm = compute_responsiveness(&f, &s, from, to, 0.95).unwrap();
assert!((rpm - 2000.0).abs() < 1e-6, "got {rpm}");
}
#[test]
fn equals_draft03_only_when_foreign_equals_loaded() {
let (f, s, from, to) = series(30.0, 30.0, 30.0, 30.0);
let rpm = compute_responsiveness(&f, &s, from, to, 0.95).unwrap();
assert!((rpm - draft03(30.0, 30.0, 30.0, 30.0)).abs() < 1e-6);
}
#[test]
fn reports_higher_than_draft03_when_rtts_diverge() {
let (f, s, from, to) = series(60.0, 60.0, 60.0, 20.0);
let rpm = compute_responsiveness(&f, &s, from, to, 0.95).unwrap();
let old = draft03(60.0, 60.0, 60.0, 20.0);
assert!((rpm - 2000.0).abs() < 1e-6, "got {rpm}");
assert!(rpm > old, "new {rpm} should exceed draft-03 {old}");
}
#[test]
fn returns_none_without_samples() {
let f = ForeignProbeResults::default();
let s = SelfProbeResults::default();
let start = Timestamp::now();
let to = start + Duration::from_millis(100);
assert!(compute_responsiveness(&f, &s, start, to, 0.95).is_none());
}
#[test]
fn returns_none_on_zero_rtt() {
let (f, s, from, to) = series(0.0, 0.0, 0.0, 0.0);
assert!(compute_responsiveness(&f, &s, from, to, 0.95).is_none());
}
fn conn_timing(
connect_ms: u64,
secure_ms: u64,
application_ms: u64,
tls_round_trips: u32,
) -> (ConnectionTiming, Timestamp) {
let start = Timestamp::now();
let mut t = ConnectionTiming::new(start);
t.set_connect(start + Duration::from_millis(connect_ms));
t.set_secure(start + Duration::from_millis(secure_ms));
t.set_application(start + Duration::from_millis(application_ms));
t.set_tls_round_trips(tls_round_trips);
(t, start)
}
#[test]
fn foreign_phases_are_independent_single_rtt_each() {
let (t, start) = conn_timing(30, 60, 62, 1);
let finished_at = start + Duration::from_millis(92);
let (tcp, tls, http) = foreign_probe_phases(&t, finished_at);
assert_eq!(tcp, Duration::from_millis(30));
assert_eq!(tls, Duration::from_millis(30));
assert_eq!(http, Duration::from_millis(30));
}
#[test]
fn foreign_tls_phase_normalized_by_round_trips() {
let (t, start) = conn_timing(30, 90, 92, 2);
let finished_at = start + Duration::from_millis(122);
let (tcp, tls, http) = foreign_probe_phases(&t, finished_at);
assert_eq!(tcp, Duration::from_millis(30));
assert_eq!(tls, Duration::from_millis(30)); assert_eq!(http, Duration::from_millis(30));
}
#[test]
fn foreign_phases_differ_from_cumulative_measurement() {
let (t, start) = conn_timing(30, 60, 62, 1);
let finished_at = start + Duration::from_millis(92);
let (tcp, tls, http) = foreign_probe_phases(&t, finished_at);
let new_sum = (tcp + tls + http).as_secs_f64() * 1000.0;
let old_tcp = t.time_connect().as_secs_f64() * 1000.0; let old_tls = t.time_secure().as_secs_f64() * 1000.0; let old_http = finished_at.duration_since(t.start()).as_secs_f64() * 1000.0; let old_sum = old_tcp + old_tls + old_http;
assert!(new_sum < old_sum, "new {new_sum} should be < old {old_sum}");
assert!((new_sum - 90.0).abs() < 1e-6);
assert!((old_sum - 182.0).abs() < 1e-6);
}
#[test]
fn reports_the_saturated_value_when_responsiveness_converged() {
assert_eq!(select_reported_rpm(Some(340.0), Some(999.0)), Some(340.0));
}
#[test]
fn reports_the_last_interval_when_the_time_limit_is_reached() {
assert_eq!(select_reported_rpm(None, Some(347.9)), Some(347.9));
}
#[test]
fn reports_nothing_when_no_interval_ever_measured() {
assert_eq!(select_reported_rpm(None, None), None);
}
}