0.3.5-alpha.10
This commit is contained in:
20
crates/apps/game-realtime-transport-measure/Cargo.toml
Normal file
20
crates/apps/game-realtime-transport-measure/Cargo.toml
Normal file
@@ -0,0 +1,20 @@
|
||||
# file: crates/apps/game-realtime-transport-measure/Cargo.toml
|
||||
# version: 1
|
||||
|
||||
[package]
|
||||
name = "game-realtime-transport-measure"
|
||||
version.workspace = true
|
||||
edition.workspace = true
|
||||
license.workspace = true
|
||||
repository.workspace = true
|
||||
authors.workspace = true
|
||||
publish.workspace = true
|
||||
|
||||
[dependencies]
|
||||
game-realtime-transport-lib = { path = "../../common/game-realtime-transport-lib" }
|
||||
game-realtime-websocket-lib = { path = "../../common/game-realtime-websocket-lib" }
|
||||
game-realtime-webtransport-lib = { path = "../../common/game-realtime-webtransport-lib" }
|
||||
tokio = { workspace = true, features = ["macros", "rt", "time"] }
|
||||
|
||||
[lints]
|
||||
workspace = true
|
||||
514
crates/apps/game-realtime-transport-measure/src/main.rs
Normal file
514
crates/apps/game-realtime-transport-measure/src/main.rs
Normal file
@@ -0,0 +1,514 @@
|
||||
// file: crates/apps/game-realtime-transport-measure/src/main.rs
|
||||
// version: 1
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
#![forbid(unsafe_code)]
|
||||
|
||||
//! Bounded localhost characterization tool for WebSocket and WebTransport realtime backends.
|
||||
|
||||
use game_realtime_transport_lib::RealtimeConnection; // rust-rules: trait-import
|
||||
use game_realtime_transport_lib::RealtimeReceiver; // rust-rules: trait-import
|
||||
use game_realtime_transport_lib::RealtimeSender; // rust-rules: trait-import
|
||||
|
||||
const ACK_PAYLOAD: &[u8] = b"ack";
|
||||
const DATAGRAM_PAYLOAD_BYTES: usize = 256;
|
||||
const DATAGRAM_RECEIVE_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(250);
|
||||
const DATAGRAM_SAMPLES: usize = 64;
|
||||
const ESTABLISHMENT_SAMPLES: usize = 8;
|
||||
const IN_FLIGHT_MESSAGES: usize = 64;
|
||||
const IN_FLIGHT_PAYLOAD_BYTES: usize = 1024;
|
||||
const MEASUREMENT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
|
||||
const RTT_PAYLOAD_BYTES: usize = 32;
|
||||
const RTT_SAMPLES: usize = 128;
|
||||
const RTT_WARMUP: usize = 16;
|
||||
const THROUGHPUT_MESSAGES: usize = 128;
|
||||
const THROUGHPUT_PAYLOAD_BYTES: usize = 64 * 1024;
|
||||
|
||||
#[derive(Clone, Copy)]
|
||||
struct LatencySummary {
|
||||
min_us: u128,
|
||||
median_us: u128,
|
||||
p95_us: u128,
|
||||
max_us: u128,
|
||||
}
|
||||
|
||||
struct TransferSummary {
|
||||
elapsed: std::time::Duration,
|
||||
messages: usize,
|
||||
payload_bytes: usize,
|
||||
}
|
||||
|
||||
struct DatagramSummary {
|
||||
elapsed: std::time::Duration,
|
||||
attempted: usize,
|
||||
received: usize,
|
||||
payload_bytes: usize,
|
||||
client_max: usize,
|
||||
server_max: usize,
|
||||
}
|
||||
|
||||
#[tokio::main(flavor = "current_thread")]
|
||||
async fn main() -> std::process::ExitCode {
|
||||
let result = tokio::time::timeout(MEASUREMENT_TIMEOUT, run_measurements()).await;
|
||||
return match result {
|
||||
std::result::Result::Ok(std::result::Result::Ok(())) => {
|
||||
println!("CONCLUSION webtransport=retain scope=second-backend reason=reliable-browser-fallback-datagram-capabilities");
|
||||
println!("game-realtime-transport-measure: PASS");
|
||||
std::process::ExitCode::SUCCESS
|
||||
},
|
||||
std::result::Result::Ok(std::result::Result::Err(error)) => {
|
||||
eprintln!("game-realtime-transport-measure: FAIL: {error}");
|
||||
std::process::ExitCode::FAILURE
|
||||
},
|
||||
std::result::Result::Err(_) => {
|
||||
eprintln!("game-realtime-transport-measure: FAIL: measurement timeout after {} ms", MEASUREMENT_TIMEOUT.as_millis());
|
||||
std::process::ExitCode::FAILURE
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
async fn run_measurements() -> std::result::Result<(), String> {
|
||||
let websocket_listener = match game_realtime_websocket_lib::WebSocketListener::bind(std::net::SocketAddr::from(([127, 0, 0, 1], 0))).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebSocket measurement listener bind failed: {error}")),
|
||||
};
|
||||
let websocket_endpoint = format!("ws://{}/measure", websocket_listener.local_addr());
|
||||
let identity = match game_realtime_webtransport_lib::WebTransportServerIdentity::generate_loopback() {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebTransport measurement identity generation failed: {error}")),
|
||||
};
|
||||
let certificate_hash = identity.certificate_hash().clone();
|
||||
let server_config = game_realtime_webtransport_lib::WebTransportServerConfig::new(std::net::SocketAddr::from(([127, 0, 0, 1], 0)), identity);
|
||||
let mut webtransport_listener = match game_realtime_webtransport_lib::WebTransportListener::bind(server_config) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebTransport measurement listener bind failed: {error}")),
|
||||
};
|
||||
let webtransport_endpoint = format!("https://{}/measure", webtransport_listener.local_addr());
|
||||
let webtransport_client_config = match game_realtime_webtransport_lib::WebTransportClientConfig::new(webtransport_endpoint.as_str(), certificate_hash) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebTransport measurement client configuration failed: {error}")),
|
||||
};
|
||||
let websocket_establishment = match measure_websocket_establishment(&websocket_listener, websocket_endpoint.as_str()).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
print_latency("websocket", "establishment", ESTABLISHMENT_SAMPLES, 0, websocket_establishment);
|
||||
let webtransport_establishment = match measure_webtransport_establishment(&mut webtransport_listener, &webtransport_client_config).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
print_latency("webtransport", "establishment", ESTABLISHMENT_SAMPLES, 0, webtransport_establishment);
|
||||
let (websocket_client, websocket_server) = match establish_websocket_pair(&websocket_listener, websocket_endpoint.as_str()).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let websocket_rtt = match measure_rtt(websocket_client, websocket_server).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
print_latency("websocket", "rtt", RTT_SAMPLES, RTT_PAYLOAD_BYTES, websocket_rtt);
|
||||
let (webtransport_client, webtransport_server) = match establish_webtransport_pair(&mut webtransport_listener, &webtransport_client_config).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let webtransport_rtt = match measure_rtt(webtransport_client, webtransport_server).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
print_latency("webtransport", "rtt", RTT_SAMPLES, RTT_PAYLOAD_BYTES, webtransport_rtt);
|
||||
let (websocket_client, websocket_server) = match establish_websocket_pair(&websocket_listener, websocket_endpoint.as_str()).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let websocket_throughput = match measure_transfer(websocket_client, websocket_server, THROUGHPUT_MESSAGES, THROUGHPUT_PAYLOAD_BYTES).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
print_transfer("websocket", "throughput", &websocket_throughput);
|
||||
let (webtransport_client, webtransport_server) = match establish_webtransport_pair(&mut webtransport_listener, &webtransport_client_config).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let webtransport_throughput = match measure_transfer(webtransport_client, webtransport_server, THROUGHPUT_MESSAGES, THROUGHPUT_PAYLOAD_BYTES).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
print_transfer("webtransport", "throughput", &webtransport_throughput);
|
||||
let (websocket_client, websocket_server) = match establish_websocket_pair(&websocket_listener, websocket_endpoint.as_str()).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let websocket_in_flight = match measure_transfer(websocket_client, websocket_server, IN_FLIGHT_MESSAGES, IN_FLIGHT_PAYLOAD_BYTES).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
print_transfer("websocket", "in_flight_window", &websocket_in_flight);
|
||||
let (webtransport_client, webtransport_server) = match establish_webtransport_pair(&mut webtransport_listener, &webtransport_client_config).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let webtransport_in_flight = match measure_transfer(webtransport_client, webtransport_server, IN_FLIGHT_MESSAGES, IN_FLIGHT_PAYLOAD_BYTES).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
print_transfer("webtransport", "in_flight_window", &webtransport_in_flight);
|
||||
let datagram = match measure_webtransport_datagrams(&mut webtransport_listener, &webtransport_client_config).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
print_datagram(&datagram);
|
||||
return std::result::Result::Ok(());
|
||||
}
|
||||
|
||||
async fn measure_websocket_establishment(
|
||||
listener: &game_realtime_websocket_lib::WebSocketListener,
|
||||
endpoint: &str,
|
||||
) -> std::result::Result<LatencySummary, String> {
|
||||
let mut samples = Vec::with_capacity(ESTABLISHMENT_SAMPLES);
|
||||
for _ in 0..ESTABLISHMENT_SAMPLES {
|
||||
let started = std::time::Instant::now();
|
||||
let (server, client) = tokio::join!(listener.accept(), game_realtime_websocket_lib::connect(endpoint));
|
||||
let _server = match server {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebSocket establishment server failed: {error}")),
|
||||
};
|
||||
let _client = match client {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebSocket establishment client failed: {error}")),
|
||||
};
|
||||
samples.push(started.elapsed());
|
||||
}
|
||||
return summarize_latency(samples.as_slice());
|
||||
}
|
||||
|
||||
async fn measure_webtransport_establishment(
|
||||
listener: &mut game_realtime_webtransport_lib::WebTransportListener,
|
||||
client_config: &game_realtime_webtransport_lib::WebTransportClientConfig,
|
||||
) -> std::result::Result<LatencySummary, String> {
|
||||
let mut samples = Vec::with_capacity(ESTABLISHMENT_SAMPLES);
|
||||
for _ in 0..ESTABLISHMENT_SAMPLES {
|
||||
let started = std::time::Instant::now();
|
||||
let _pair = match establish_webtransport_pair(listener, client_config).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
samples.push(started.elapsed());
|
||||
}
|
||||
return summarize_latency(samples.as_slice());
|
||||
}
|
||||
|
||||
async fn establish_websocket_pair(
|
||||
listener: &game_realtime_websocket_lib::WebSocketListener,
|
||||
endpoint: &str,
|
||||
) -> std::result::Result<(game_realtime_websocket_lib::WebSocketConnection, game_realtime_websocket_lib::WebSocketConnection), String> {
|
||||
let (server, client) = tokio::join!(listener.accept(), game_realtime_websocket_lib::connect(endpoint));
|
||||
let server = match server {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebSocket server establishment failed: {error}")),
|
||||
};
|
||||
let client = match client {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebSocket client establishment failed: {error}")),
|
||||
};
|
||||
return std::result::Result::Ok((client, server));
|
||||
}
|
||||
|
||||
async fn establish_webtransport_pair(
|
||||
listener: &mut game_realtime_webtransport_lib::WebTransportListener,
|
||||
client_config: &game_realtime_webtransport_lib::WebTransportClientConfig,
|
||||
) -> std::result::Result<(game_realtime_webtransport_lib::WebTransportConnection, game_realtime_webtransport_lib::WebTransportConnection), String> {
|
||||
let server = async {
|
||||
let session = match listener.accept().await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebTransport server session establishment failed: {error}")),
|
||||
};
|
||||
return match session.accept_primary_connection().await {
|
||||
std::result::Result::Ok(value) => std::result::Result::Ok(value),
|
||||
std::result::Result::Err(error) => std::result::Result::Err(format!("WebTransport server primary stream establishment failed: {error}")),
|
||||
};
|
||||
};
|
||||
let client = async {
|
||||
let session = match game_realtime_webtransport_lib::connect(client_config).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("WebTransport client session establishment failed: {error}")),
|
||||
};
|
||||
return match session.open_primary_connection().await {
|
||||
std::result::Result::Ok(value) => std::result::Result::Ok(value),
|
||||
std::result::Result::Err(error) => std::result::Result::Err(format!("WebTransport client primary stream establishment failed: {error}")),
|
||||
};
|
||||
};
|
||||
let (server, client) = tokio::join!(server, client);
|
||||
return match (client, server) {
|
||||
(std::result::Result::Ok(client), std::result::Result::Ok(server)) => std::result::Result::Ok((client, server)),
|
||||
(std::result::Result::Err(error), _) => std::result::Result::Err(error),
|
||||
(_, std::result::Result::Err(error)) => std::result::Result::Err(error),
|
||||
};
|
||||
}
|
||||
|
||||
async fn measure_rtt<ClientConnection, ServerConnection>(client: ClientConnection, server: ServerConnection) -> std::result::Result<LatencySummary, String>
|
||||
where
|
||||
ClientConnection: game_realtime_transport_lib::RealtimeConnection,
|
||||
ServerConnection: game_realtime_transport_lib::RealtimeConnection,
|
||||
{
|
||||
let total_messages = RTT_WARMUP + RTT_SAMPLES;
|
||||
let payload = vec![0x52_u8; RTT_PAYLOAD_BYTES];
|
||||
let (mut client_sender, mut client_receiver) = client.split();
|
||||
let (mut server_sender, mut server_receiver) = server.split();
|
||||
let server = async move {
|
||||
for _ in 0..total_messages {
|
||||
let message = match server_receiver.receive().await {
|
||||
std::result::Result::Ok(game_realtime_transport_lib::TransportReceive::Message(value)) => value,
|
||||
std::result::Result::Ok(game_realtime_transport_lib::TransportReceive::Closed) => {
|
||||
return std::result::Result::Err(String::from("RTT server observed an early clean close"));
|
||||
},
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("RTT server receive failed: {error}")),
|
||||
};
|
||||
if let std::result::Result::Err(error) = server_sender.send(message).await {
|
||||
return std::result::Result::Err(format!("RTT server echo failed: {error}"));
|
||||
}
|
||||
}
|
||||
let client_close = match server_receiver.receive().await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("RTT server close observation failed: {error}")),
|
||||
};
|
||||
if client_close != game_realtime_transport_lib::TransportReceive::Closed {
|
||||
return std::result::Result::Err(String::from("RTT server did not observe the client close"));
|
||||
}
|
||||
if let std::result::Result::Err(error) = server_sender.close().await {
|
||||
return std::result::Result::Err(format!("RTT server close failed: {error}"));
|
||||
}
|
||||
return std::result::Result::Ok(());
|
||||
};
|
||||
let client = async move {
|
||||
let mut samples = Vec::with_capacity(RTT_SAMPLES);
|
||||
for index in 0..total_messages {
|
||||
let started = std::time::Instant::now();
|
||||
if let std::result::Result::Err(error) = client_sender.send(game_realtime_transport_lib::TransportMessage::new(payload.clone())).await {
|
||||
return std::result::Result::Err(format!("RTT client send failed: {error}"));
|
||||
}
|
||||
let response = match client_receiver.receive().await {
|
||||
std::result::Result::Ok(game_realtime_transport_lib::TransportReceive::Message(value)) => value,
|
||||
std::result::Result::Ok(game_realtime_transport_lib::TransportReceive::Closed) => {
|
||||
return std::result::Result::Err(String::from("RTT client observed an early clean close"));
|
||||
},
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("RTT client receive failed: {error}")),
|
||||
};
|
||||
if response.as_bytes() != payload.as_slice() {
|
||||
return std::result::Result::Err(String::from("RTT echo payload mismatch"));
|
||||
}
|
||||
if index >= RTT_WARMUP {
|
||||
samples.push(started.elapsed());
|
||||
}
|
||||
}
|
||||
if let std::result::Result::Err(error) = client_sender.close().await {
|
||||
return std::result::Result::Err(format!("RTT client close failed: {error}"));
|
||||
}
|
||||
let server_close = match client_receiver.receive().await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("RTT client close observation failed: {error}")),
|
||||
};
|
||||
if server_close != game_realtime_transport_lib::TransportReceive::Closed {
|
||||
return std::result::Result::Err(String::from("RTT client did not observe the server close"));
|
||||
}
|
||||
return summarize_latency(samples.as_slice());
|
||||
};
|
||||
let (server, client) = tokio::join!(server, client);
|
||||
match server {
|
||||
std::result::Result::Ok(()) => {},
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
async fn measure_transfer<ClientConnection, ServerConnection>(
|
||||
client: ClientConnection,
|
||||
server: ServerConnection,
|
||||
messages: usize,
|
||||
payload_bytes: usize,
|
||||
) -> std::result::Result<TransferSummary, String>
|
||||
where
|
||||
ClientConnection: game_realtime_transport_lib::RealtimeConnection,
|
||||
ServerConnection: game_realtime_transport_lib::RealtimeConnection,
|
||||
{
|
||||
let payload = vec![0x54_u8; payload_bytes];
|
||||
let (mut client_sender, mut client_receiver) = client.split();
|
||||
let (mut server_sender, mut server_receiver) = server.split();
|
||||
let server = async move {
|
||||
for _ in 0..messages {
|
||||
let received = match server_receiver.receive().await {
|
||||
std::result::Result::Ok(game_realtime_transport_lib::TransportReceive::Message(value)) => value,
|
||||
std::result::Result::Ok(game_realtime_transport_lib::TransportReceive::Closed) => {
|
||||
return std::result::Result::Err(String::from("transfer server observed an early clean close"));
|
||||
},
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("transfer server receive failed: {error}")),
|
||||
};
|
||||
if received.len() != payload_bytes {
|
||||
return std::result::Result::Err(format!("transfer server payload size mismatch: expected {payload_bytes}, got {}", received.len()));
|
||||
}
|
||||
}
|
||||
if let std::result::Result::Err(error) = server_sender.send(game_realtime_transport_lib::TransportMessage::new(ACK_PAYLOAD.to_vec())).await {
|
||||
return std::result::Result::Err(format!("transfer server ACK failed: {error}"));
|
||||
}
|
||||
let client_close = match server_receiver.receive().await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("transfer server close observation failed: {error}")),
|
||||
};
|
||||
if client_close != game_realtime_transport_lib::TransportReceive::Closed {
|
||||
return std::result::Result::Err(String::from("transfer server did not observe the client close"));
|
||||
}
|
||||
if let std::result::Result::Err(error) = server_sender.close().await {
|
||||
return std::result::Result::Err(format!("transfer server close failed: {error}"));
|
||||
}
|
||||
return std::result::Result::Ok(());
|
||||
};
|
||||
let client = async move {
|
||||
let started = std::time::Instant::now();
|
||||
for _ in 0..messages {
|
||||
if let std::result::Result::Err(error) = client_sender.send(game_realtime_transport_lib::TransportMessage::new(payload.clone())).await {
|
||||
return std::result::Result::Err(format!("transfer client send failed: {error}"));
|
||||
}
|
||||
}
|
||||
let ack = match client_receiver.receive().await {
|
||||
std::result::Result::Ok(game_realtime_transport_lib::TransportReceive::Message(value)) => value,
|
||||
std::result::Result::Ok(game_realtime_transport_lib::TransportReceive::Closed) => {
|
||||
return std::result::Result::Err(String::from("transfer client observed close before ACK"));
|
||||
},
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("transfer client ACK receive failed: {error}")),
|
||||
};
|
||||
if ack.as_bytes() != ACK_PAYLOAD {
|
||||
return std::result::Result::Err(String::from("transfer ACK payload mismatch"));
|
||||
}
|
||||
let elapsed = started.elapsed();
|
||||
if let std::result::Result::Err(error) = client_sender.close().await {
|
||||
return std::result::Result::Err(format!("transfer client close failed: {error}"));
|
||||
}
|
||||
let server_close = match client_receiver.receive().await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("transfer client close observation failed: {error}")),
|
||||
};
|
||||
if server_close != game_realtime_transport_lib::TransportReceive::Closed {
|
||||
return std::result::Result::Err(String::from("transfer client did not observe the server close"));
|
||||
}
|
||||
return std::result::Result::Ok(TransferSummary { elapsed, messages, payload_bytes });
|
||||
};
|
||||
let (server, client) = tokio::join!(server, client);
|
||||
match server {
|
||||
std::result::Result::Ok(()) => {},
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
async fn measure_webtransport_datagrams(
|
||||
listener: &mut game_realtime_webtransport_lib::WebTransportListener,
|
||||
client_config: &game_realtime_webtransport_lib::WebTransportClientConfig,
|
||||
) -> std::result::Result<DatagramSummary, String> {
|
||||
let (server, client) = tokio::join!(listener.accept(), game_realtime_webtransport_lib::connect(client_config));
|
||||
let server = match server {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("datagram measurement server session failed: {error}")),
|
||||
};
|
||||
let client = match client {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(format!("datagram measurement client session failed: {error}")),
|
||||
};
|
||||
let client_max = client.max_datagram_size();
|
||||
let server_max = server.max_datagram_size();
|
||||
let negotiated_max = core::cmp::min(client_max, server_max);
|
||||
if negotiated_max == 0 {
|
||||
return std::result::Result::Err(String::from("datagram measurement negotiated a zero-byte payload bound"));
|
||||
}
|
||||
let payload_bytes = core::cmp::min(DATAGRAM_PAYLOAD_BYTES, negotiated_max);
|
||||
let payload = vec![0x44_u8; payload_bytes];
|
||||
let started = std::time::Instant::now();
|
||||
for _ in 0..DATAGRAM_SAMPLES {
|
||||
if let std::result::Result::Err(error) = client.send_datagram(payload.as_slice()) {
|
||||
return std::result::Result::Err(format!("datagram measurement send failed: {error}"));
|
||||
}
|
||||
}
|
||||
let mut received = 0_usize;
|
||||
while received < DATAGRAM_SAMPLES {
|
||||
let datagram = tokio::time::timeout(DATAGRAM_RECEIVE_TIMEOUT, server.receive_datagram()).await;
|
||||
match datagram {
|
||||
std::result::Result::Ok(std::result::Result::Ok(value)) => {
|
||||
if value.len() != payload_bytes {
|
||||
return std::result::Result::Err(format!("datagram measurement payload size mismatch: expected {payload_bytes}, got {}", value.len()));
|
||||
}
|
||||
received += 1;
|
||||
},
|
||||
std::result::Result::Ok(std::result::Result::Err(error)) => {
|
||||
return std::result::Result::Err(format!("datagram measurement receive failed: {error}"));
|
||||
},
|
||||
std::result::Result::Err(_) => break,
|
||||
}
|
||||
}
|
||||
return std::result::Result::Ok(DatagramSummary {
|
||||
elapsed: started.elapsed(),
|
||||
attempted: DATAGRAM_SAMPLES,
|
||||
received,
|
||||
payload_bytes,
|
||||
client_max,
|
||||
server_max,
|
||||
});
|
||||
}
|
||||
|
||||
fn summarize_latency(samples: &[std::time::Duration]) -> std::result::Result<LatencySummary, String> {
|
||||
if samples.is_empty() {
|
||||
return std::result::Result::Err(String::from("latency summary requires at least one sample"));
|
||||
}
|
||||
let mut values: Vec<u128> = samples.iter().map(std::time::Duration::as_micros).collect();
|
||||
values.sort_unstable();
|
||||
let min_us = match values.first() {
|
||||
Some(value) => *value,
|
||||
None => return std::result::Result::Err(String::from("latency summary unexpectedly lost all samples")),
|
||||
};
|
||||
let max_us = match values.last() {
|
||||
Some(value) => *value,
|
||||
None => return std::result::Result::Err(String::from("latency summary unexpectedly lost all samples")),
|
||||
};
|
||||
let median_us = match values.get(values.len() / 2) {
|
||||
Some(value) => *value,
|
||||
None => return std::result::Result::Err(String::from("latency median index is invalid")),
|
||||
};
|
||||
let p95_rank = values.len().saturating_mul(95).div_ceil(100);
|
||||
let p95_index = p95_rank.saturating_sub(1);
|
||||
let p95_us = match values.get(p95_index) {
|
||||
Some(value) => *value,
|
||||
None => return std::result::Result::Err(String::from("latency p95 index is invalid")),
|
||||
};
|
||||
return std::result::Result::Ok(LatencySummary { min_us, median_us, p95_us, max_us });
|
||||
}
|
||||
|
||||
fn print_latency(transport: &str, metric: &str, samples: usize, payload_bytes: usize, summary: LatencySummary) {
|
||||
println!(
|
||||
"MEASURE transport={transport} metric={metric} samples={samples} payload_bytes={payload_bytes} min_us={} median_us={} p95_us={} max_us={}",
|
||||
summary.min_us, summary.median_us, summary.p95_us, summary.max_us
|
||||
);
|
||||
}
|
||||
|
||||
fn print_transfer(transport: &str, metric: &str, summary: &TransferSummary) {
|
||||
let elapsed_seconds = summary.elapsed.as_secs_f64();
|
||||
let total_bytes = summary.messages.saturating_mul(summary.payload_bytes);
|
||||
let mib_per_second = if elapsed_seconds > 0.0 { (total_bytes as f64 / (1024.0 * 1024.0)) / elapsed_seconds } else { 0.0 };
|
||||
let messages_per_second = if elapsed_seconds > 0.0 { summary.messages as f64 / elapsed_seconds } else { 0.0 };
|
||||
println!(
|
||||
"MEASURE transport={transport} metric={metric} messages={} payload_bytes={} total_bytes={total_bytes} elapsed_us={} mib_per_s={mib_per_second:.3} messages_per_s={messages_per_second:.3}",
|
||||
summary.messages,
|
||||
summary.payload_bytes,
|
||||
summary.elapsed.as_micros()
|
||||
);
|
||||
}
|
||||
|
||||
fn print_datagram(summary: &DatagramSummary) {
|
||||
let ratio = if summary.attempted == 0 { 0.0 } else { summary.received as f64 / summary.attempted as f64 };
|
||||
println!(
|
||||
"MEASURE transport=webtransport-datagram metric=bounded_burst attempted={} received={} receive_ratio={ratio:.6} payload_bytes={} client_max={} server_max={} elapsed_us={}",
|
||||
summary.attempted,
|
||||
summary.received,
|
||||
summary.payload_bytes,
|
||||
summary.client_max,
|
||||
summary.server_max,
|
||||
summary.elapsed.as_micros()
|
||||
);
|
||||
}
|
||||
Reference in New Issue
Block a user