From 3ffa250b358faa6f9c341891ef4d885e4b95fc6d Mon Sep 17 00:00:00 2001 From: Benjamin Saunders Date: Tue, 19 Nov 2019 21:04:40 -0800 Subject: [PATCH] Add minimal self-contained benchmark for easy profiling --- Cargo.toml | 2 +- bench/Cargo.toml | 24 ++++++++++ bench/src/bulk.rs | 115 ++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 140 insertions(+), 1 deletion(-) create mode 100644 bench/Cargo.toml create mode 100644 bench/src/bulk.rs diff --git a/Cargo.toml b/Cargo.toml index 0213fbff4..f7404b35f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = ["quinn", "quinn-proto", "quinn-h3", "interop"] +members = ["quinn", "quinn-proto", "quinn-h3", "interop", "bench"] [profile.bench] debug = true diff --git a/bench/Cargo.toml b/bench/Cargo.toml new file mode 100644 index 000000000..a290a1c5e --- /dev/null +++ b/bench/Cargo.toml @@ -0,0 +1,24 @@ +[package] +name = "bench" +version = "0.1.0" +authors = ["Benjamin Saunders "] +edition = "2018" +license = "MIT/Apache-2.0" +publish = false + +[dependencies] +quinn = { path = "../quinn" } +tokio = { version = "0.2.0-alpha.6", default-features = false, features = ["rt-current-thread"] } +tracing = "0.1.10" +tracing-subscriber = "0.1.5" +anyhow = "1.0.22" +rcgen = "0.7" +futures = { package = "futures-preview", version = "0.3.0-alpha.18" } +rustls = "0.16" + +[[bin]] +name = "bulk" +path = "src/bulk.rs" + +[profile.release] +debug = true diff --git a/bench/src/bulk.rs b/bench/src/bulk.rs new file mode 100644 index 000000000..ac4615270 --- /dev/null +++ b/bench/src/bulk.rs @@ -0,0 +1,115 @@ +use std::net::{IpAddr, Ipv6Addr, SocketAddr}; +use std::time::Instant; + +use anyhow::{anyhow, Context, Result}; +use futures::StreamExt; +use tokio::runtime::current_thread::Runtime; +use tracing::trace; + +fn main() { + tracing::subscriber::set_global_default( + tracing_subscriber::FmtSubscriber::builder() + .with_env_filter(tracing_subscriber::EnvFilter::from_default_env()) + .finish(), + ) + .unwrap(); + + let cert = rcgen::generate_simple_self_signed(vec!["localhost".into()]).unwrap(); + let key = quinn::PrivateKey::from_der(&cert.serialize_private_key_der()).unwrap(); + let cert = quinn::Certificate::from_der(&cert.serialize_der().unwrap()).unwrap(); + + let mut server_config = quinn::ServerConfigBuilder::default(); + server_config + .certificate(quinn::CertificateChain::from_certs(vec![cert.clone()]), key) + .unwrap(); + let mut endpoint = quinn::EndpointBuilder::default(); + endpoint.listen(server_config.build()); + let (driver, endpoint, incoming) = endpoint + .bind(&SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 0)) + .unwrap(); + let server_addr = endpoint.local_addr().unwrap(); + drop(endpoint); // Ensure server shuts down when finished + let thread = std::thread::spawn(|| { + let mut runtime = Runtime::new().unwrap(); + runtime.spawn(async { + driver.await.expect("server endpoint driver"); + }); + if let Err(e) = runtime.block_on(server(incoming)) { + eprintln!("server failed: {:#}", e); + } + runtime.run().expect("server run"); + }); + + let mut runtime = Runtime::new().unwrap(); + if let Err(e) = runtime.block_on(client(server_addr, cert)) { + eprintln!("client failed: {:#}", e); + } + runtime.run().expect("client run"); + + thread.join().expect("server thread"); +} + +async fn server(mut incoming: quinn::Incoming) -> Result<()> { + let handshake = incoming.next().await.unwrap(); + let quinn::NewConnection { + driver, + mut uni_streams, + .. + } = handshake.await.context("handshake failed")?; + tokio::spawn(async { + driver.await.expect("server conn driver"); + }); + let mut stream = uni_streams + .next() + .await + .ok_or(anyhow!("accepting stream failed"))??; + trace!("stream established"); + let start = Instant::now(); + let mut n = 0; + while let Some((data, offset)) = stream.read_unordered().await? { + n = n.max(offset + data.len() as u64); + } + println!("recvd {} bytes in {:?}", n, start.elapsed()); + Ok(()) +} + +async fn client(server_addr: SocketAddr, server_cert: quinn::Certificate) -> Result<()> { + let (driver, endpoint, _) = quinn::EndpointBuilder::default() + .bind(&SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 0)) + .unwrap(); + tokio::spawn(async { + driver.await.expect("client endpoint driver"); + }); + + let mut client_config = quinn::ClientConfigBuilder::default(); + client_config + .add_certificate_authority(server_cert) + .unwrap(); + let quinn::NewConnection { + driver, connection, .. + } = endpoint + .connect_with(client_config.build(), &server_addr, "localhost") + .unwrap() + .await + .context("unable to connect")?; + tokio::spawn(async { + let _ = driver.await; + }); + trace!("connected"); + + let mut stream = connection + .open_uni() + .await + .context("failed to open stream")?; + const DATA: &[u8] = &[0xAB; 1024 * 1024]; + let start = Instant::now(); + for _ in 0..1024 { + stream + .write_all(DATA) + .await + .context("failed sending data")?; + } + stream.finish().await.context("failed finishing stream")?; + println!("sent {} bytes in {:?}", 1024 * DATA.len(), start.elapsed()); + Ok(()) +}