feat(serve): split request read budget into head_timeout + body_timeout

The single request_timeout covered head + body under one wall clock, so a
slow-but-legit body upload (e.g. a trickling cellular IoT client) was
judged by the short head deadline and killed mid-body. Split into:

- head_timeout (default 30s): first byte -> full head parse; the classic
  slowloris surface, kept short.
- body_timeout (default 300s): head parse -> full body; an absolute cap
  sized for slow links, anchored independently once the head has parsed.

read_head no longer returns a shared deadline; run_connection anchors the
body deadline itself. ReadHeadErr::RequestTimeout -> HeadTimeout. Config
and ConnLimits gain body_timeout; request_timeout renamed to head_timeout
(breaking, but this axis is unreleased).

Tests: slow_body_outlives_head_timeout (positive: body survives past the
head clock), fixed_/chunked_body_stall_killed_at_body_timeout (body cap
still bites), slowloris_partial_head_killed_at_head_timeout (head clock
unchanged). 79 lib + 43 integration green; clippy --lib clean.
This commit is contained in:
Claude (sandbox)
2026-08-12 13:00:22 +00:00
parent 1b1ea124c8
commit 4f06265338
4 changed files with 142 additions and 55 deletions
+49 -39
View File
@@ -46,14 +46,21 @@ pub struct ConnLimits {
/// connection). Expiry closes the connection silently — nothing is /// connection). Expiry closes the connection silently — nothing is
/// owed to a client that isn't talking. /// owed to a client that isn't talking.
pub keep_alive_timeout: Duration, pub keep_alive_timeout: Duration,
/// Per-request wall-clock budget, measured from the first byte of a /// Wall-clock budget for reading the request HEAD, measured from the
/// request until the request (head + body) is fully read. Expiry /// first byte of a request until the head is fully parsed. Expiry
/// mid-head gets a best-effort 408; expiry mid-body just closes. /// mid-head gets a best-effort 408. Kept short: an incomplete head is
/// Pipeline run time is NOT covered — that's the handler's business. /// the classic slowloris, and a legitimate client sends its head in a
/// Covers the READ phase only; the write phase has its own /// single burst. The BODY has its own, larger budget (`body_timeout`)
/// per-write budget (`write_timeout`) so a streaming response can /// so a slow-but-legit upload is not judged by the head clock.
/// legitimately outlive any whole-request clock. pub head_timeout: Duration,
pub request_timeout: Duration, /// Absolute wall-clock cap on reading the request BODY, measured from
/// the moment the head finished parsing until the body is fully read.
/// Sized for slow links (e.g. a trickling cellular IoT client), so it
/// is much larger than `head_timeout`. Expiry mid-body just closes —
/// nothing is owed to a client this far gone. Pipeline run time is NOT
/// covered (that's the handler's business); the write phase has its own
/// per-write budget (`write_timeout`).
pub body_timeout: Duration,
/// Per-write budget for response bytes: every `write_all` (the fixed /// Per-write budget for response bytes: every `write_all` (the fixed
/// head+body, and each streamed chunk) must complete within this. /// head+body, and each streamed chunk) must complete within this.
/// A client that stops reading mid-response is dropped when its /// A client that stops reading mid-response is dropped when its
@@ -76,7 +83,8 @@ impl Default for ConnLimits {
max_head_bytes: 64 * 1024, max_head_bytes: 64 * 1024,
max_body_bytes: 16 * 1024 * 1024, max_body_bytes: 16 * 1024 * 1024,
keep_alive_timeout: Duration::from_secs(60), keep_alive_timeout: Duration::from_secs(60),
request_timeout: Duration::from_secs(30), head_timeout: Duration::from_secs(30),
body_timeout: Duration::from_secs(300),
write_timeout: Duration::from_secs(30), write_timeout: Duration::from_secs(30),
max_frame_payload: 1024 * 1024, max_frame_payload: 1024 * 1024,
max_message_bytes: 4 * 1024 * 1024, max_message_bytes: 4 * 1024 * 1024,
@@ -111,7 +119,7 @@ pub fn run_connection(
// ----- 1. Read until we have a full request head. ----- // ----- 1. Read until we have a full request head. -----
// We are idle until a head parses: stoppable by a draining // We are idle until a head parses: stoppable by a draining
// registry while parked here. // registry while parked here.
let (parsed, request_deadline) = match read_head(raw, &mut buf, &limits) { let parsed = match read_head(raw, &mut buf, &limits) {
Ok(p) => p, Ok(p) => p,
Err(ReadHeadErr::ClientClosed) => { Err(ReadHeadErr::ClientClosed) => {
// Clean EOF between requests (or before any request). Normal. // Clean EOF between requests (or before any request). Normal.
@@ -122,8 +130,8 @@ pub fn run_connection(
// a request. Nothing is owed; close silently. // a request. Nothing is owed; close silently.
return; return;
} }
Err(ReadHeadErr::RequestTimeout) => { Err(ReadHeadErr::HeadTimeout) => {
// request_timeout expired mid-head (slowloris and friends). // head_timeout expired mid-head (slowloris and friends).
// Best-effort 408 WITHOUT parking on writability — a client // Best-effort 408 WITHOUT parking on writability — a client
// that stalls reads must not defeat the timeout by making // that stalls reads must not defeat the timeout by making
// the 408 write park forever. // the 408 write park forever.
@@ -142,6 +150,11 @@ pub fn run_connection(
let _ = registry.cast(Cast::ConnBusy(me)); let _ = registry.cast(Cast::ConnBusy(me));
// ----- 2. Read body. ----- // ----- 2. Read body. -----
// The body has its OWN absolute budget, anchored here (head just
// parsed) and independent of the head clock — a slow-but-legit
// upload must not be judged by the short head deadline. Expiry
// closes the connection (nothing owed mid-body).
let body_deadline = Instant::now() + limits.body_timeout;
// Content-Length pre-check only applies to fixed bodies; a chunked // Content-Length pre-check only applies to fixed bodies; a chunked
// body is bounded incrementally by the decoder. // body is bounded incrementally by the decoder.
let body_len = parsed.content_length.unwrap_or(0); let body_len = parsed.content_length.unwrap_or(0);
@@ -169,7 +182,7 @@ pub fn run_connection(
// bottom of the loop must drop exactly this much to land on the // bottom of the loop must drop exactly this much to land on the
// next pipelined request. // next pipelined request.
let (body, consumed_past_head) = if parsed.chunked { let (body, consumed_past_head) = if parsed.chunked {
match read_chunked_body(raw, &mut buf, parsed.head_len, &limits, request_deadline) { match read_chunked_body(raw, &mut buf, parsed.head_len, &limits, body_deadline) {
Ok(ok) => ok, Ok(ok) => ok,
Err(ChunkedBodyErr::TooLarge) => { Err(ChunkedBodyErr::TooLarge) => {
let _ = write_all( let _ = write_all(
@@ -191,7 +204,7 @@ pub fn run_connection(
Err(ChunkedBodyErr::Io(_)) => return, Err(ChunkedBodyErr::Io(_)) => return,
} }
} else { } else {
match read_body(raw, &mut buf, parsed.head_len, body_len, request_deadline) { match read_body(raw, &mut buf, parsed.head_len, body_len, body_deadline) {
Ok(b) => (b, body_len), Ok(b) => (b, body_len),
// Timeout mid-body (and any other body io error) -> just // Timeout mid-body (and any other body io error) -> just
// close; there's no point talking HTTP to a client this far // close; there's no point talking HTTP to a client this far
@@ -330,9 +343,9 @@ enum ReadHeadErr {
/// keep_alive_timeout expired while waiting for the first byte of a /// keep_alive_timeout expired while waiting for the first byte of a
/// request. Close silently. /// request. Close silently.
IdleTimeout, IdleTimeout,
/// request_timeout expired after the request had started arriving. /// head_timeout expired after the request had started arriving but
/// Best-effort 408. /// before the head finished parsing. Best-effort 408.
RequestTimeout, HeadTimeout,
Io(io::Error), Io(io::Error),
Parse(ParseError), Parse(ParseError),
} }
@@ -347,22 +360,22 @@ enum ReadHeadErr {
/// - while `buf` is empty and nothing has arrived, we are *idle* and the /// - while `buf` is empty and nothing has arrived, we are *idle* and the
/// wait is bounded by `keep_alive_timeout`; /// wait is bounded by `keep_alive_timeout`;
/// - the instant the request has started (first byte read, or pipelined /// - the instant the request has started (first byte read, or pipelined
/// bytes already in `buf` at entry), the *request* clock starts: an /// bytes already in `buf` at entry), the *head* clock starts: an
/// `Instant` deadline of `request_timeout` from that moment, which also /// `Instant` deadline of `head_timeout` from that moment. This budget
/// covers body reads — it is returned alongside the parsed head so the /// covers the HEAD only; the body has its own budget (`body_timeout`),
/// caller can thread it into `read_body`. /// which the caller anchors once the head has parsed.
fn read_head( fn read_head(
fd: RawFd, fd: RawFd,
buf: &mut Vec<u8>, buf: &mut Vec<u8>,
limits: &ConnLimits, limits: &ConnLimits,
) -> Result<(parser::ParsedHead, Instant), ReadHeadErr> { ) -> Result<parser::ParsedHead, ReadHeadErr> {
let entry = Instant::now(); let entry = Instant::now();
let idle_deadline = entry + limits.keep_alive_timeout; let idle_deadline = entry + limits.keep_alive_timeout;
// Pipelined leftovers count as a started request. // Pipelined leftovers count as a started request.
let mut request_deadline: Option<Instant> = if buf.is_empty() { let mut head_deadline: Option<Instant> = if buf.is_empty() {
None None
} else { } else {
Some(entry + limits.request_timeout) Some(entry + limits.head_timeout)
}; };
loop { loop {
@@ -375,11 +388,7 @@ fn read_head(
parser::parse_head(buf, limits.max_headers) parser::parse_head(buf, limits.max_headers)
}; };
match head { match head {
Ok(h) => { Ok(h) => return Ok(h),
let deadline = request_deadline
.unwrap_or_else(|| Instant::now() + limits.request_timeout);
return Ok((h, deadline));
}
Err(ParseError::Incomplete) => {} // need more bytes Err(ParseError::Incomplete) => {} // need more bytes
Err(e) => return Err(ReadHeadErr::Parse(e)), Err(e) => return Err(ReadHeadErr::Parse(e)),
} }
@@ -390,19 +399,19 @@ fn read_head(
} }
// Read more, bounded by whichever budget is active. // Read more, bounded by whichever budget is active.
let deadline = request_deadline.unwrap_or(idle_deadline); let deadline = head_deadline.unwrap_or(idle_deadline);
match read_some(fd, buf, limits.initial_read_buf, deadline) { match read_some(fd, buf, limits.initial_read_buf, deadline) {
Ok(0) => return Err(ReadHeadErr::ClientClosed), Ok(0) => return Err(ReadHeadErr::ClientClosed),
Ok(_) => { Ok(_) => {
if request_deadline.is_none() { if head_deadline.is_none() {
// First byte(s) of this request: the request clock // First byte(s) of this request: the head clock
// starts now. // starts now.
request_deadline = Some(Instant::now() + limits.request_timeout); head_deadline = Some(Instant::now() + limits.head_timeout);
} }
} }
Err(e) if e.kind() == ErrorKind::TimedOut => { Err(e) if e.kind() == ErrorKind::TimedOut => {
return Err(if request_deadline.is_some() { return Err(if head_deadline.is_some() {
ReadHeadErr::RequestTimeout ReadHeadErr::HeadTimeout
} else { } else {
ReadHeadErr::IdleTimeout ReadHeadErr::IdleTimeout
}); });
@@ -433,8 +442,8 @@ fn read_body(
return Ok(buf[head_len..head_len + body_len].to_vec()); return Ok(buf[head_len..head_len + body_len].to_vec());
} }
// Read until we have the rest, on the same request budget that the // Read until we have the rest, bounded by the body budget (anchored
// head was read under. // by the caller when the head finished parsing).
let mut total_read = already; let mut total_read = already;
while total_read < body_len { while total_read < body_len {
match read_some(fd, buf, 8 * 1024, deadline) { match read_some(fd, buf, 8 * 1024, deadline) {
@@ -451,8 +460,9 @@ fn read_body(
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// //
// Decodes `Transfer-Encoding: chunked` from `buf[head_len..]`, reading more // Decodes `Transfer-Encoding: chunked` from `buf[head_len..]`, reading more
// from the socket as needed on the SAME request deadline the head was read // from the socket as needed on the body deadline (anchored by the caller
// under. Returns (decoded_body, raw_bytes_consumed_past_head) — the raw // when the head finished parsing, independent of the head clock).
// Returns (decoded_body, raw_bytes_consumed_past_head) — the raw
// count includes all framing and the trailer section, so the caller's // count includes all framing and the trailer section, so the caller's
// keep-alive drain lands exactly on the next pipelined request. // keep-alive drain lands exactly on the next pipelined request.
// //
+12 -3
View File
@@ -37,7 +37,14 @@ pub struct Config {
pub keep_alive_timeout: Duration, pub keep_alive_timeout: Duration,
pub max_header_count: usize, pub max_header_count: usize,
pub read_buf_size: usize, pub read_buf_size: usize,
pub request_timeout: Duration, /// Wall-clock budget for reading the request HEAD (from first byte to
/// full head parse). Kept short — an incomplete head is the classic
/// slowloris. See `ConnLimits::head_timeout`.
pub head_timeout: Duration,
/// Absolute wall-clock cap on reading the request BODY (from head-parse
/// to full body). Sized for slow links, so much larger than
/// `head_timeout`. See `ConnLimits::body_timeout`.
pub body_timeout: Duration,
/// Per-write budget for response bytes (the fixed head+body write, and /// Per-write budget for response bytes (the fixed head+body write, and
/// each streamed chunk). See `ConnLimits::write_timeout`. /// each streamed chunk). See `ConnLimits::write_timeout`.
pub write_timeout: Duration, pub write_timeout: Duration,
@@ -88,7 +95,8 @@ impl Config {
keep_alive_timeout: Duration::from_secs(60), keep_alive_timeout: Duration::from_secs(60),
max_header_count: 64, max_header_count: 64,
read_buf_size: 8 * 1024, read_buf_size: 8 * 1024,
request_timeout: Duration::from_secs(30), head_timeout: Duration::from_secs(30),
body_timeout: Duration::from_secs(300),
write_timeout: Duration::from_secs(30), write_timeout: Duration::from_secs(30),
max_body_bytes: 16 * 1024 * 1024, max_body_bytes: 16 * 1024 * 1024,
drain_timeout: Duration::from_secs(30), drain_timeout: Duration::from_secs(30),
@@ -107,7 +115,8 @@ impl Config {
max_head_bytes: 64 * 1024, max_head_bytes: 64 * 1024,
max_body_bytes: self.max_body_bytes, max_body_bytes: self.max_body_bytes,
keep_alive_timeout: self.keep_alive_timeout, keep_alive_timeout: self.keep_alive_timeout,
request_timeout: self.request_timeout, head_timeout: self.head_timeout,
body_timeout: self.body_timeout,
write_timeout: self.write_timeout, write_timeout: self.write_timeout,
max_frame_payload: self.max_frame_payload, max_frame_payload: self.max_frame_payload,
max_message_bytes: self.max_message_bytes, max_message_bytes: self.max_message_bytes,
+2 -1
View File
@@ -33,7 +33,8 @@
//! write — event or heartbeat — stalls past `write_timeout`; the conn //! write — event or heartbeat — stalls past `write_timeout`; the conn
//! actor then drops the stream and the producer's next [`EventSender`] //! actor then drops the stream and the producer's next [`EventSender`]
//! call returns `Err(SseClosed)`. There is no request clock on an SSE //! call returns `Err(SseClosed)`. There is no request clock on an SSE
//! response: `request_timeout` covers only the read phase, by design. //! response: the head/body read budgets cover only the read phase, by
//! design.
use crate::conn::{Conn, RespBody, StreamBody}; use crate::conn::{Conn, RespBody, StreamBody};
+79 -12
View File
@@ -437,7 +437,8 @@ fn shutdown_force_stops_at_drain_deadline() {
fn spawn_server_with_timeouts( fn spawn_server_with_timeouts(
pipeline: Pipeline, pipeline: Pipeline,
keep_alive: Duration, keep_alive: Duration,
request: Duration, head: Duration,
body: Duration,
) -> u16 { ) -> u16 {
let port = free_port(); let port = free_port();
let addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap(); let addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap();
@@ -446,7 +447,8 @@ fn spawn_server_with_timeouts(
listener_pool: 2, listener_pool: 2,
scheduler_threads: Some(2), scheduler_threads: Some(2),
keep_alive_timeout: keep_alive, keep_alive_timeout: keep_alive,
request_timeout: request, head_timeout: head,
body_timeout: body,
..Config::new(addr) ..Config::new(addr)
}; };
serve_with(cfg, pipeline).unwrap(); serve_with(cfg, pipeline).unwrap();
@@ -486,7 +488,8 @@ fn idle_keepalive_reaped_at_keep_alive_timeout() {
let port = spawn_server_with_timeouts( let port = spawn_server_with_timeouts(
pipe, pipe,
Duration::from_millis(300), // keep_alive_timeout under test Duration::from_millis(300), // keep_alive_timeout under test
Duration::from_secs(10), // request_timeout out of the way Duration::from_secs(10), // head_timeout out of the way
Duration::from_secs(10), // body_timeout out of the way
); );
let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap(); let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap();
@@ -511,17 +514,18 @@ fn idle_keepalive_reaped_at_keep_alive_timeout() {
} }
/// A slowloris client that sends a partial head and then stalls is killed /// A slowloris client that sends a partial head and then stalls is killed
/// at request_timeout with a best-effort 408, even though the (large) /// at head_timeout with a best-effort 408, even though the (large)
/// keep-alive budget hasn't expired. /// keep-alive and body budgets haven't expired.
#[test] #[test]
fn slowloris_partial_head_killed_at_request_timeout() { fn slowloris_partial_head_killed_at_head_timeout() {
let pipe = Pipeline::new().plug( let pipe = Pipeline::new().plug(
Router::new().get("/", |c: Conn, _n: Next| c.put_status(200)) Router::new().get("/", |c: Conn, _n: Next| c.put_status(200))
); );
let port = spawn_server_with_timeouts( let port = spawn_server_with_timeouts(
pipe, pipe,
Duration::from_secs(10), // keep_alive_timeout out of the way Duration::from_secs(10), // keep_alive_timeout out of the way
Duration::from_millis(300), // request_timeout under test Duration::from_millis(300), // head_timeout under test
Duration::from_secs(10), // body_timeout out of the way
); );
let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap(); let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap();
@@ -913,14 +917,16 @@ fn chunked_plus_content_length_400() {
assert_eq!(http_status(&resp), 400); assert_eq!(http_status(&resp), 400);
} }
/// A chunked body that stalls mid-stream is killed by the request /// A chunked body that stalls mid-stream is killed by the BODY deadline
/// deadline: the connection just closes (no response owed mid-body). /// (head budget generous): the connection just closes (no response owed
/// mid-body).
#[test] #[test]
fn chunked_request_stall_killed_at_request_timeout() { fn chunked_body_stall_killed_at_body_timeout() {
let port = spawn_server_with_timeouts( let port = spawn_server_with_timeouts(
echo_pipeline(), echo_pipeline(),
Duration::from_secs(30), Duration::from_secs(30), // keep_alive_timeout out of the way
Duration::from_millis(400), // request_timeout Duration::from_secs(30), // head_timeout out of the way
Duration::from_millis(400), // body_timeout under test
); );
let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap(); let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap();
s.set_read_timeout(Some(Duration::from_secs(5))).unwrap(); s.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
@@ -936,6 +942,67 @@ fn chunked_request_stall_killed_at_request_timeout() {
assert!(start.elapsed() < Duration::from_secs(3), "close took too long"); assert!(start.elapsed() < Duration::from_secs(3), "close took too long");
} }
/// The core of the head/body split: a client that sends a COMPLETE head
/// promptly and then trickles its (small) body over a span LONGER than
/// head_timeout still succeeds, because the body runs on its own, larger
/// budget. Under the old shared request clock this would have been killed
/// mid-body at head_timeout. This is the slow-but-legit IoT upload we must
/// not punish.
#[test]
fn slow_body_outlives_head_timeout() {
let port = spawn_server_with_timeouts(
echo_pipeline(),
Duration::from_secs(30), // keep_alive_timeout out of the way
Duration::from_millis(500), // head_timeout: SHORT
Duration::from_secs(8), // body_timeout: generous
);
let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap();
s.set_read_timeout(Some(Duration::from_secs(10))).unwrap();
// Full head at once (parses well within head_timeout), Connection:
// close so the server closes after responding and read_to_end lands
// the whole response.
s.write_all(
b"POST /echo HTTP/1.1\r\nHost: x\r\nContent-Length: 4\r\nConnection: close\r\n\r\n",
)
.unwrap();
// Trickle the 4-byte body at 250ms/byte => ~1s total, well past the
// 500ms head_timeout but inside the 8s body_timeout.
for b in b"test" {
std::thread::sleep(Duration::from_millis(250));
s.write_all(&[*b]).unwrap();
}
let mut resp = Vec::new();
s.read_to_end(&mut resp).expect("expected full response");
assert_eq!(http_status(&resp), 200, "resp: {:?}", String::from_utf8_lossy(&resp));
assert!(
resp.ends_with(b"test"),
"expected echoed body 'test', got: {:?}", String::from_utf8_lossy(&resp)
);
}
/// A fixed-Content-Length body that stalls before completing is killed by
/// the BODY deadline (head budget generous): silent close, nothing owed
/// mid-body. The fixed-path twin of chunked_body_stall_killed_at_body_timeout.
#[test]
fn fixed_body_stall_killed_at_body_timeout() {
let port = spawn_server_with_timeouts(
echo_pipeline(),
Duration::from_secs(30), // keep_alive_timeout out of the way
Duration::from_secs(30), // head_timeout out of the way
Duration::from_millis(400), // body_timeout under test
);
let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap();
s.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
// Promises 100 bytes, sends a few, then stalls forever.
s.write_all(b"POST /echo HTTP/1.1\r\nHost: x\r\nContent-Length: 100\r\n\r\npartial")
.unwrap();
let start = std::time::Instant::now();
let mut resp = Vec::new();
s.read_to_end(&mut resp).unwrap(); // server closes; EOF
assert!(resp.is_empty(), "expected silent close, got: {:?}", String::from_utf8_lossy(&resp));
assert!(start.elapsed() < Duration::from_secs(3), "close took too long");
}
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// SSE (v0.3 chunk 3) // SSE (v0.3 chunk 3)
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------