Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 34 additions & 0 deletions linkerd/app/src/env/http2.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,11 @@ pub(super) fn parse_server<S: Strings>(
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"),
Comment thread
ergdevops marked this conversation as resolved.
parse_duration,
)?,
max_concurrent_streams: parse(
strings,
&format!("{base}_MAX_CONCURRENT_STREAMS"),
Expand Down Expand Up @@ -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,
Expand All @@ -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),
Expand Down Expand Up @@ -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))
);
}
}
1 change: 1 addition & 0 deletions linkerd/http/h2/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ pub struct ServerParams {
pub flow_control: Option<FlowControl>,
pub keep_alive: Option<KeepAlive>,
pub max_concurrent_streams: Option<u32>,
pub max_connection_age: Option<Duration>,

// Internals
pub max_frame_size: Option<u32>,
Expand Down
67 changes: 59 additions & 8 deletions linkerd/proxy/http/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ use std::{
future::Future,
pin::Pin,
task::{Context, Poll},
time::Duration,
};
use tower::Service;
use tracing::{debug, Instrument};
Expand All @@ -30,12 +31,26 @@ pub struct NewServeHttp<X, N> {
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<N> {
version: Variant,
http1: hyper::server::conn::http1::Builder,
http2: hyper::server::conn::http2::Builder<TokioExecutor>,
max_connection_age: Option<Duration>,
inner: N,
drain: drain::Watch,
}
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -124,6 +140,7 @@ where
drain,
http1,
http2,
max_connection_age,
}
}
}
Expand Down Expand Up @@ -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 => {
Expand Down Expand Up @@ -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?;
}
Comment thread
ergdevops marked this conversation as resolved.
Expand Down
70 changes: 70 additions & 0 deletions linkerd/proxy/http/src/server/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down