diff --git a/interop/Cargo.toml b/interop/Cargo.toml index 9add786d1..8e0eaef7c 100644 --- a/interop/Cargo.toml +++ b/interop/Cargo.toml @@ -20,6 +20,7 @@ quinn-proto = { path = "../quinn-proto" } rustls = { version = "0.17", features = ["dangerous_configuration"] } structopt = "0.3.0" tokio = { version = "0.2.2", features = ["macros", "rt-core", "io-util"] } +tokio-rustls = "0.13" tracing = "0.1.10" tracing-subscriber = { version = "0.2.3", default-features = false, features = ["env-filter", "fmt", "ansi", "chrono"]} tracing-futures = { version = "0.2.0", default-features = false, features = ["std-future"] } diff --git a/interop/src/main.rs b/interop/src/main.rs index 553fb0488..1ffdef505 100644 --- a/interop/src/main.rs +++ b/interop/src/main.rs @@ -527,7 +527,7 @@ impl State { "hq time {:+.2}% of h2's", percentage ); - if percentage > 30.0 { + if percentage > 10.0 { return Err(anyhow!("Throughput {} is {:+.2}% slower", size, percentage)); } } @@ -600,7 +600,7 @@ impl State { "h3 time {:+.2}% of h2's", percentage ); - if percentage > 30.0 { + if percentage > 10.0 { return Err(anyhow!("Throughput {} is {:+.2}% slower", size, percentage)); } } diff --git a/interop/src/server.rs b/interop/src/server.rs index 3b3673477..854b613cd 100644 --- a/interop/src/server.rs +++ b/interop/src/server.rs @@ -4,14 +4,19 @@ use std::{ fs, net::{IpAddr, Ipv4Addr, SocketAddr}, path::PathBuf, - str, + pin::Pin, + str, sync, + task::{Context, Poll}, }; -use anyhow::{anyhow, bail, Context, Result}; +use anyhow::{anyhow, bail, Context as _, Result}; use bytes::Bytes; -use futures::{AsyncReadExt, AsyncWriteExt, StreamExt, TryFutureExt}; +use futures::{ready, AsyncReadExt, Future, StreamExt, TryFutureExt}; use http::{Response, StatusCode}; +use hyper::service::{make_service_fn, service_fn}; use structopt::{self, StructOpt}; +use tokio::net::{TcpListener, TcpStream}; +use tokio_rustls::{server::TlsStream, TlsAcceptor}; use tracing::{error, info, info_span}; use tracing_futures::Instrument as _; @@ -77,7 +82,7 @@ async fn main() -> Result<()> { server_config.use_stateless_retry(true); let retry = server(server_config.clone(), 4434); - tokio::try_join!(main, default, retry)?; + tokio::try_join!(main, default, retry, h2_server(server_config.clone()))?; Ok(()) } @@ -172,25 +177,19 @@ async fn h3_payload(sender: quinn_h3::server::Sender, len: usize) -> Result<()> return Ok(()); } + let mut buf = TEXT.repeat(len / TEXT.len() + 1); + buf.truncate(len); + let response = Response::builder() .status(StatusCode::OK) - .body(()) + .body(Bytes::from(buf)) .expect("failed to build response"); - let mut body_writer = sender + sender .send_response(response) .await .map_err(|e| anyhow!("failed to send response: {:?}", e))?; - let mut remaining = len; - while remaining > 0 { - let size = cmp::min(remaining, TEXT.len()); - body_writer.write_all(&TEXT[..size]).await?; - remaining -= size; - } - body_writer.flush().await?; - body_writer.close().await?; - Ok(()) } @@ -310,6 +309,98 @@ fn parse_size(literal: &str) -> Result { Ok(num * scale) } +fn h2_home() -> hyper::Response { + Response::builder() + .status(StatusCode::OK) + .body(HOME.into()) + .expect("failed to build response") +} + +fn h2_payload(len: usize) -> hyper::Response { + if len > 1_000_000_000 { + let response = Response::builder() + .status(StatusCode::BAD_REQUEST) + .body(Bytes::from(format!("requested {}: too large", len)).into()) + .expect("failed to build response"); + return response; + } + + let mut buf = TEXT.repeat(len / TEXT.len() + 1); + buf.truncate(len); + Response::builder() + .status(StatusCode::OK) + .body(buf.into()) + .expect("failed to build response") +} + +async fn h2_handle(request: hyper::Request) -> Result> { + Ok(match request.uri().path() { + "/" => h2_home(), + path => match parse_size(path) { + Ok(n) => h2_payload(n), + Err(_) => h2_home(), + }, + }) +} + +async fn h2_server(server_config: quinn::ServerConfigBuilder) -> Result<()> { + let mut tls_cfg = (*server_config.build().crypto).clone(); + tls_cfg.set_protocols(&[b"h2".to_vec()]); + let tls_acceptor = TlsAcceptor::from(sync::Arc::new(tls_cfg)); + + let tcp = TcpListener::bind(&SocketAddr::new([0, 0, 0, 0].into(), 443)).await?; + + let service = make_service_fn(|_conn| async { Ok::<_, anyhow::Error>(service_fn(h2_handle)) }); + let server = hyper::Server::builder(HyperAcceptor::new(tcp, tls_acceptor)) + .http2_only(true) + .serve(service); + + if let Err(e) = server.await { + error!("server error: {}", e); + } + Ok(()) +} + +struct HyperAcceptor { + tcp: TcpListener, + tls: TlsAcceptor, + handshake: Option>, +} + +impl HyperAcceptor { + pub fn new(tcp: TcpListener, tls: TlsAcceptor) -> Self { + Self { + tls, + tcp, + handshake: None, + } + } +} + +impl hyper::server::accept::Accept for HyperAcceptor { + type Conn = TlsStream; + type Error = anyhow::Error; + + fn poll_accept( + mut self: Pin<&mut Self>, + cx: &mut Context, + ) -> Poll>> { + loop { + match self.handshake { + Some(ref mut h) => { + let conn = ready!(Pin::new(h).poll(cx))?; + std::mem::replace(&mut self.handshake, None); + return Poll::Ready(Some(Ok(conn))); + } + None => { + let (stream, _) = ready!(self.tcp.poll_accept(cx))?; + self.handshake = Some(self.tls.accept(stream)); + } + } + } + } +} + const TEXT: &[u8] = b"It would be different if we could not step back and reflect on the process,\n\ but were merely led from impulse to impulse without self- consciousness. But human\n\