From 4f06265338448bb34a61805bc78b07037568ea72 Mon Sep 17 00:00:00 2001 From: "Claude (sandbox)" Date: Wed, 12 Aug 2026 13:00:22 +0000 Subject: [PATCH] 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. --- src/conn_actor.rs | 88 +++++++++++++++++++++++------------------- src/serve.rs | 15 ++++++-- src/sse.rs | 3 +- tests/integration.rs | 91 ++++++++++++++++++++++++++++++++++++++------ 4 files changed, 142 insertions(+), 55 deletions(-) diff --git a/src/conn_actor.rs b/src/conn_actor.rs index 87b3c7e..8a7b584 100644 --- a/src/conn_actor.rs +++ b/src/conn_actor.rs @@ -46,14 +46,21 @@ pub struct ConnLimits { /// connection). Expiry closes the connection silently — nothing is /// owed to a client that isn't talking. pub keep_alive_timeout: Duration, - /// Per-request wall-clock budget, measured from the first byte of a - /// request until the request (head + body) is fully read. Expiry - /// mid-head gets a best-effort 408; expiry mid-body just closes. - /// Pipeline run time is NOT covered — that's the handler's business. - /// Covers the READ phase only; the write phase has its own - /// per-write budget (`write_timeout`) so a streaming response can - /// legitimately outlive any whole-request clock. - pub request_timeout: Duration, + /// Wall-clock budget for reading the request HEAD, measured from the + /// first byte of a request until the head is fully parsed. Expiry + /// mid-head gets a best-effort 408. Kept short: an incomplete head is + /// the classic slowloris, and a legitimate client sends its head in a + /// single burst. The BODY has its own, larger budget (`body_timeout`) + /// so a slow-but-legit upload is not judged by the head clock. + pub head_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 /// head+body, and each streamed chunk) must complete within this. /// 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_body_bytes: 16 * 1024 * 1024, 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), max_frame_payload: 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. ----- // We are idle until a head parses: stoppable by a draining // 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, Err(ReadHeadErr::ClientClosed) => { // Clean EOF between requests (or before any request). Normal. @@ -122,8 +130,8 @@ pub fn run_connection( // a request. Nothing is owed; close silently. return; } - Err(ReadHeadErr::RequestTimeout) => { - // request_timeout expired mid-head (slowloris and friends). + Err(ReadHeadErr::HeadTimeout) => { + // head_timeout expired mid-head (slowloris and friends). // Best-effort 408 WITHOUT parking on writability — a client // that stalls reads must not defeat the timeout by making // the 408 write park forever. @@ -142,6 +150,11 @@ pub fn run_connection( let _ = registry.cast(Cast::ConnBusy(me)); // ----- 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 // body is bounded incrementally by the decoder. 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 // next pipelined request. 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, Err(ChunkedBodyErr::TooLarge) => { let _ = write_all( @@ -191,7 +204,7 @@ pub fn run_connection( Err(ChunkedBodyErr::Io(_)) => return, } } 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), // Timeout mid-body (and any other body io error) -> just // 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 /// request. Close silently. IdleTimeout, - /// request_timeout expired after the request had started arriving. - /// Best-effort 408. - RequestTimeout, + /// head_timeout expired after the request had started arriving but + /// before the head finished parsing. Best-effort 408. + HeadTimeout, Io(io::Error), Parse(ParseError), } @@ -347,22 +360,22 @@ enum ReadHeadErr { /// - while `buf` is empty and nothing has arrived, we are *idle* and the /// wait is bounded by `keep_alive_timeout`; /// - the instant the request has started (first byte read, or pipelined -/// bytes already in `buf` at entry), the *request* clock starts: an -/// `Instant` deadline of `request_timeout` from that moment, which also -/// covers body reads — it is returned alongside the parsed head so the -/// caller can thread it into `read_body`. +/// bytes already in `buf` at entry), the *head* clock starts: an +/// `Instant` deadline of `head_timeout` from that moment. This budget +/// covers the HEAD only; the body has its own budget (`body_timeout`), +/// which the caller anchors once the head has parsed. fn read_head( fd: RawFd, buf: &mut Vec, limits: &ConnLimits, -) -> Result<(parser::ParsedHead, Instant), ReadHeadErr> { +) -> Result { let entry = Instant::now(); let idle_deadline = entry + limits.keep_alive_timeout; // Pipelined leftovers count as a started request. - let mut request_deadline: Option = if buf.is_empty() { + let mut head_deadline: Option = if buf.is_empty() { None } else { - Some(entry + limits.request_timeout) + Some(entry + limits.head_timeout) }; loop { @@ -375,11 +388,7 @@ fn read_head( parser::parse_head(buf, limits.max_headers) }; match head { - Ok(h) => { - let deadline = request_deadline - .unwrap_or_else(|| Instant::now() + limits.request_timeout); - return Ok((h, deadline)); - } + Ok(h) => return Ok(h), Err(ParseError::Incomplete) => {} // need more bytes Err(e) => return Err(ReadHeadErr::Parse(e)), } @@ -390,19 +399,19 @@ fn read_head( } // 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) { Ok(0) => return Err(ReadHeadErr::ClientClosed), Ok(_) => { - if request_deadline.is_none() { - // First byte(s) of this request: the request clock + if head_deadline.is_none() { + // First byte(s) of this request: the head clock // 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 => { - return Err(if request_deadline.is_some() { - ReadHeadErr::RequestTimeout + return Err(if head_deadline.is_some() { + ReadHeadErr::HeadTimeout } else { ReadHeadErr::IdleTimeout }); @@ -433,8 +442,8 @@ fn read_body( return Ok(buf[head_len..head_len + body_len].to_vec()); } - // Read until we have the rest, on the same request budget that the - // head was read under. + // Read until we have the rest, bounded by the body budget (anchored + // by the caller when the head finished parsing). let mut total_read = already; while total_read < body_len { 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 -// from the socket as needed on the SAME request deadline the head was read -// under. Returns (decoded_body, raw_bytes_consumed_past_head) — the raw +// from the socket as needed on the body deadline (anchored by the caller +// 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 // keep-alive drain lands exactly on the next pipelined request. // diff --git a/src/serve.rs b/src/serve.rs index 2124ad2..615ee56 100644 --- a/src/serve.rs +++ b/src/serve.rs @@ -37,7 +37,14 @@ pub struct Config { pub keep_alive_timeout: Duration, pub max_header_count: 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 /// each streamed chunk). See `ConnLimits::write_timeout`. pub write_timeout: Duration, @@ -88,7 +95,8 @@ impl Config { keep_alive_timeout: Duration::from_secs(60), max_header_count: 64, 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), max_body_bytes: 16 * 1024 * 1024, drain_timeout: Duration::from_secs(30), @@ -107,7 +115,8 @@ impl Config { max_head_bytes: 64 * 1024, max_body_bytes: self.max_body_bytes, 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, max_frame_payload: self.max_frame_payload, max_message_bytes: self.max_message_bytes, diff --git a/src/sse.rs b/src/sse.rs index ea599c2..5a60527 100644 --- a/src/sse.rs +++ b/src/sse.rs @@ -33,7 +33,8 @@ //! write — event or heartbeat — stalls past `write_timeout`; the conn //! actor then drops the stream and the producer's next [`EventSender`] //! 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}; diff --git a/tests/integration.rs b/tests/integration.rs index babc2a2..93fd171 100644 --- a/tests/integration.rs +++ b/tests/integration.rs @@ -437,7 +437,8 @@ fn shutdown_force_stops_at_drain_deadline() { fn spawn_server_with_timeouts( pipeline: Pipeline, keep_alive: Duration, - request: Duration, + head: Duration, + body: Duration, ) -> u16 { let port = free_port(); let addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap(); @@ -446,7 +447,8 @@ fn spawn_server_with_timeouts( listener_pool: 2, scheduler_threads: Some(2), keep_alive_timeout: keep_alive, - request_timeout: request, + head_timeout: head, + body_timeout: body, ..Config::new(addr) }; serve_with(cfg, pipeline).unwrap(); @@ -486,7 +488,8 @@ fn idle_keepalive_reaped_at_keep_alive_timeout() { let port = spawn_server_with_timeouts( pipe, 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(); @@ -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 -/// at request_timeout with a best-effort 408, even though the (large) -/// keep-alive budget hasn't expired. +/// at head_timeout with a best-effort 408, even though the (large) +/// keep-alive and body budgets haven't expired. #[test] -fn slowloris_partial_head_killed_at_request_timeout() { +fn slowloris_partial_head_killed_at_head_timeout() { let pipe = Pipeline::new().plug( Router::new().get("/", |c: Conn, _n: Next| c.put_status(200)) ); let port = spawn_server_with_timeouts( pipe, 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(); @@ -913,14 +917,16 @@ fn chunked_plus_content_length_400() { assert_eq!(http_status(&resp), 400); } -/// A chunked body that stalls mid-stream is killed by the request -/// deadline: the connection just closes (no response owed mid-body). +/// A chunked body that stalls mid-stream is killed by the BODY deadline +/// (head budget generous): the connection just closes (no response owed +/// mid-body). #[test] -fn chunked_request_stall_killed_at_request_timeout() { +fn chunked_body_stall_killed_at_body_timeout() { let port = spawn_server_with_timeouts( echo_pipeline(), - Duration::from_secs(30), - Duration::from_millis(400), // request_timeout + 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(); @@ -936,6 +942,67 @@ fn chunked_request_stall_killed_at_request_timeout() { 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) // ---------------------------------------------------------------------------