diff --git a/linkerd/app/src/env/http2.rs b/linkerd/app/src/env/http2.rs index 42c5806ca6..18393606e0 100644 --- a/linkerd/app/src/env/http2.rs +++ b/linkerd/app/src/env/http2.rs @@ -9,6 +9,11 @@ pub(super) fn parse_server( Ok(ServerParams { flow_control: Some(parse_flow_control(strings, base)?), keep_alive: parse_keep_alive(strings, &format!("{base}_KEEP_ALIVE"))?, + max_connection_age: parse( + strings, + &format!("{base}_MAX_CONNECTION_AGE"), + parse_duration, + )?, max_concurrent_streams: parse( strings, &format!("{base}_MAX_CONCURRENT_STREAMS"), @@ -107,6 +112,7 @@ mod tests { env.insert("TEST_KEEP_ALIVE_INTERVAL", "2s"); env.insert("TEST_INITIAL_STREAM_WINDOW_SIZE", "1"); env.insert("TEST_INITIAL_CONNECTION_WINDOW_SIZE", "2"); + env.insert("TEST_MAX_CONNECTION_AGE", "900s"); let expected = h2::ServerParams { flow_control: Some(h2::FlowControl::Fixed { initial_stream_window_size: 1, @@ -116,6 +122,7 @@ mod tests { interval: Duration::from_secs(2), timeout: Duration::from_secs(1), }), + max_connection_age: Some(Duration::from_secs(900)), max_concurrent_streams: Some(3), max_frame_size: Some(4), max_header_list_size: Some(5), @@ -147,4 +154,31 @@ mod tests { } ); } + + /// Like every other SERVER_HTTP2 setting, max connection age is parsed + /// under both the inbound and outbound server prefixes. + #[test] + fn max_connection_age_parses_for_both_servers() { + let mut env = HashMap::default(); + env.insert( + "LINKERD2_PROXY_INBOUND_SERVER_HTTP2_MAX_CONNECTION_AGE", + "15m", + ); + env.insert( + "LINKERD2_PROXY_OUTBOUND_SERVER_HTTP2_MAX_CONNECTION_AGE", + "1h", + ); + + let inbound = parse_server(&env, "LINKERD2_PROXY_INBOUND_SERVER_HTTP2").unwrap(); + assert_eq!( + inbound.max_connection_age, + Some(Duration::from_secs(15 * 60)) + ); + + let outbound = parse_server(&env, "LINKERD2_PROXY_OUTBOUND_SERVER_HTTP2").unwrap(); + assert_eq!( + outbound.max_connection_age, + Some(Duration::from_secs(60 * 60)) + ); + } } diff --git a/linkerd/http/h2/src/lib.rs b/linkerd/http/h2/src/lib.rs index ad1b109595..5bd844c2e4 100644 --- a/linkerd/http/h2/src/lib.rs +++ b/linkerd/http/h2/src/lib.rs @@ -5,6 +5,7 @@ pub struct ServerParams { pub flow_control: Option, pub keep_alive: Option, pub max_concurrent_streams: Option, + pub max_connection_age: Option, // Internals pub max_frame_size: Option, diff --git a/linkerd/proxy/http/src/server.rs b/linkerd/proxy/http/src/server.rs index 5e4fcdc2d0..6823985c8d 100644 --- a/linkerd/proxy/http/src/server.rs +++ b/linkerd/proxy/http/src/server.rs @@ -8,6 +8,7 @@ use std::{ future::Future, pin::Pin, task::{Context, Poll}, + time::Duration, }; use tower::Service; use tracing::{debug, Instrument}; @@ -30,12 +31,26 @@ pub struct NewServeHttp { params: X, } +/// How a served connection is being torn down, as decided by the serve loop's +/// completion branches. +enum Teardown { + /// The client closed the connection; nothing left to do. + ClientClosed, + /// The process is draining: shut down gracefully, holding the drain + /// release guard until the connection closes. + Drain(drain::ReleaseShutdown), + /// Stack teardown or max connection age: shut down gracefully, then wait + /// for the connection to close. + Graceful, +} + /// Serves HTTP connections with an inner service. #[derive(Clone, Debug)] pub struct ServeHttp { version: Variant, http1: hyper::server::conn::http1::Builder, http2: hyper::server::conn::http2::Builder, + max_connection_age: Option, inner: N, drain: drain::Watch, } @@ -70,6 +85,7 @@ where keep_alive, flow_control, max_concurrent_streams, + max_connection_age, max_frame_size, max_header_list_size, max_send_buf_size, @@ -124,6 +140,7 @@ where drain, http1, http2, + max_connection_age, } } } @@ -154,17 +171,18 @@ where let drain = self.drain.clone(); let http1 = self.http1.clone(); let http2 = self.http2.clone(); + let max_connection_age = self.max_connection_age; - let res = io.peer_addr().map(|pa| { - let (handle, closed) = ClientHandle::new(pa); + let res = io.peer_addr().map(|peer_addr| { + let (handle, closed) = ClientHandle::new(peer_addr); let svc = self.inner.new_service(handle.clone()); let svc = SetClientHandle::new(handle, svc); - (svc, closed) + (svc, closed, peer_addr) }); Box::pin( async move { - let (svc, closed) = res?; + let (svc, closed, peer_addr) = res?; debug!(?version, "Handling as HTTP"); match version { Variant::Http1 => { @@ -201,18 +219,51 @@ where let io = hyper_util::rt::TokioIo::new(io); let mut conn = http2.serve_connection(io, svc); - tokio::select! { + // Bound the connection's lifetime if a max age is + // configured, so long-lived pooled peer connections + // (and the buffer high-water memory they retain) are + // periodically recycled. A deterministic per-connection + // jitter of up to +10%, derived from the client's + // ephemeral port, avoids synchronized shutdowns. + let max_age = async move { + match max_connection_age { + Some(age) => { + let jitter = + age.mul_f64(f64::from(peer_addr.port() % 128) / 1280.0); + tokio::time::sleep(age + jitter).await; + } + None => std::future::pending().await, + } + }; + tokio::pin!(max_age); + + let teardown = tokio::select! { res = &mut conn => { debug!(?res, "The client is shutting down the connection"); - res? + res?; + Teardown::ClientClosed } shutdown = drain.signaled() => { debug!("The process is shutting down the connection"); - Pin::new(&mut conn).graceful_shutdown(); - shutdown.release_after(conn).await?; + Teardown::Drain(shutdown) } () = closed => { debug!("The stack is tearing down the connection"); + Teardown::Graceful + } + () = &mut max_age => { + debug!(client.addr = %peer_addr, "Max connection age reached; gracefully shutting down the connection"); + Teardown::Graceful + } + }; + + match teardown { + Teardown::ClientClosed => {} + Teardown::Drain(shutdown) => { + Pin::new(&mut conn).graceful_shutdown(); + shutdown.release_after(conn).await?; + } + Teardown::Graceful => { Pin::new(&mut conn).graceful_shutdown(); conn.await?; } diff --git a/linkerd/proxy/http/src/server/tests.rs b/linkerd/proxy/http/src/server/tests.rs index a7b806144b..28ee40d429 100644 --- a/linkerd/proxy/http/src/server/tests.rs +++ b/linkerd/proxy/http/src/server/tests.rs @@ -159,6 +159,76 @@ async fn h2_stream_window_exhaustion() { */ } +/// Tests that a server connection is gracefully shut down after its max +/// connection age elapses: in-flight streams complete, but the connection is +/// then closed so the client must reconnect. +#[tokio::test(flavor = "current_thread", start_paused = true)] +async fn h2_max_connection_age() { + let _trace = linkerd_tracing::test::with_default_filter(LOG_LEVEL); + + const MAX_AGE: time::Duration = time::Duration::from_secs(10); + // The effective age includes a deterministic per-connection jitter of up + // to +10%; sleep comfortably past it. + const PAST_AGE: time::Duration = time::Duration::from_secs(12); + + let mut server = TestServer::connect_h2( + h2::ServerParams { + max_connection_age: Some(MAX_AGE), + ..Default::default() + }, + hyper::client::conn::http2::Builder::new(TokioExecutor::new()) + .timer(hyper_util::rt::TokioTimer::new()), + ) + .await; + + let bytes = Bytes::from_static(b"hello"); + + tracing::info!("Before the age elapses, requests are served normally"); + let rx = timeout(server.respond(bytes.clone())) + .await + .expect("timed out"); + let body = timeout(rx.collect()) + .await + .expect("response timed out") + .expect("response"); + assert_eq!(body.to_bytes(), bytes); + + tracing::info!("A stream that is in flight when the age elapses still completes"); + let (mut tx, mut body) = timeout(server.get()).await.expect("timed out"); + time::sleep(PAST_AGE).await; + tx.send_data(bytes.clone()) + .await + .expect("in-flight stream remains writable after the age elapses"); + drop(tx); + let data = body + .frame() + .await + .expect("yields a result") + .expect("yields a frame") + .into_data() + .expect("yields data"); + assert_eq!(data, bytes); + // Drain the body to its end; only an empty end-of-stream frame may remain. + while let Some(frame) = timeout(body.frame()).await.expect("body ends") { + let frame = frame.expect("body yields frames until end"); + if let Ok(data) = frame.into_data() { + assert!(data.is_empty(), "no more data expected"); + } + } + + tracing::info!("Once in-flight streams complete, the connection is closed"); + timeout(async { + loop { + if server.client.ready().await.is_err() { + return; + } + time::sleep(time::Duration::from_millis(100)).await; + } + }) + .await + .expect("client should observe the connection closing"); +} + // === Utilities === const LOG_LEVEL: &str = "h2::proto=trace,hyper=trace,linkerd=trace,info";