diff --git a/src/conn_actor.rs b/src/conn_actor.rs index 8a7b584..6ca1633 100644 --- a/src/conn_actor.rs +++ b/src/conn_actor.rs @@ -61,6 +61,19 @@ pub struct ConnLimits { /// covered (that's the handler's business); the write phase has its own /// per-write budget (`write_timeout`). pub body_timeout: Duration, + /// Burst-gated body stall eviction: the bytes that must accumulate + /// since the last advance to count as a "burst" and reset the stall + /// clock. A body that dribbles fewer than this per `body_stall_timeout` + /// window is evicted — the discriminator between a slowloris trickle + /// (near-zero, smooth) and a slow-but-legit client (delivers real + /// bursts). The pair implies an effective floor of + /// body_burst_bytes / body_stall_timeout, enforced in bursts. + pub body_burst_bytes: usize, + /// Max time since the last qualifying burst (`body_burst_bytes`) + /// before a stalled body read is evicted. Must comfortably exceed a + /// legit client's worst quiet gap (e.g. cellular RRC/handover/DRX + /// stalls). The absolute `body_timeout` always backstops it. + pub body_stall_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 @@ -85,6 +98,8 @@ impl Default for ConnLimits { keep_alive_timeout: Duration::from_secs(60), head_timeout: Duration::from_secs(30), body_timeout: Duration::from_secs(300), + body_burst_bytes: 4 * 1024, + body_stall_timeout: Duration::from_secs(20), write_timeout: Duration::from_secs(30), max_frame_payload: 1024 * 1024, max_message_bytes: 4 * 1024 * 1024, @@ -204,7 +219,7 @@ pub fn run_connection( Err(ChunkedBodyErr::Io(_)) => return, } } else { - match read_body(raw, &mut buf, parsed.head_len, body_len, body_deadline) { + match read_body(raw, &mut buf, parsed.head_len, body_len, &limits, 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 @@ -421,6 +436,58 @@ fn read_head( } } +// --------------------------------------------------------------------------- +// BodyStallGate — burst-gated stall eviction for body reads +// --------------------------------------------------------------------------- +// +// Each body read is bounded by the SOONER of two deadlines: the absolute +// body cap (`body_timeout`, passed in as `cap`) and a sliding stall window +// (`mark + body_stall_timeout`). The stall mark only advances when the +// client delivers a full burst (`body_burst_bytes` accumulated since the +// last advance) — so a steady sub-burst trickle never moves the mark and +// is evicted at ~body_stall_timeout, while a bursty slow-but-legit client +// keeps resetting it and survives up to the absolute cap. +// +// State is two words (`mark`, `since_mark`); the per-read cost is one add +// and one compare. Bytes counted are RAW socket bytes (progress = the +// client is sending *something*), so chunked framing counts too, and a +// burst that the kernel fragments into several reads still accumulates. +struct BodyStallGate { + cap: Instant, + stall_timeout: Duration, + burst_bytes: usize, + mark: Instant, + since_mark: usize, +} + +impl BodyStallGate { + fn new(cap: Instant, limits: &ConnLimits, now: Instant) -> Self { + Self { + cap, + stall_timeout: limits.body_stall_timeout, + burst_bytes: limits.body_burst_bytes, + mark: now, + since_mark: 0, + } + } + + /// Deadline for the next read: the sooner of the absolute cap and the + /// current stall window. + fn deadline(&self) -> Instant { + (self.mark + self.stall_timeout).min(self.cap) + } + + /// Record `n` freshly-read raw body bytes; advance the stall mark if a + /// full burst has accumulated since the last advance. + fn record(&mut self, n: usize, now: Instant) { + self.since_mark += n; + if self.since_mark >= self.burst_bytes { + self.mark = now; + self.since_mark = 0; + } + } +} + // --------------------------------------------------------------------------- // read_body // --------------------------------------------------------------------------- @@ -430,7 +497,8 @@ fn read_body( buf: &mut Vec, head_len: usize, body_len: usize, - deadline: Instant, + limits: &ConnLimits, + cap: Instant, ) -> io::Result> { // Bytes already in `buf` past the head belong to the body. let already = buf.len().saturating_sub(head_len); @@ -442,13 +510,17 @@ fn read_body( return Ok(buf[head_len..head_len + body_len].to_vec()); } - // Read until we have the rest, bounded by the body budget (anchored - // by the caller when the head finished parsing). + // Read until we have the rest, bounded by the body cap AND the + // burst-gated stall window (whichever is sooner). + let mut gate = BodyStallGate::new(cap, limits, Instant::now()); let mut total_read = already; while total_read < body_len { - match read_some(fd, buf, 8 * 1024, deadline) { + match read_some(fd, buf, 8 * 1024, gate.deadline()) { Ok(0) => return Err(io::Error::new(ErrorKind::UnexpectedEof, "client closed during body")), - Ok(n) => total_read += n, + Ok(n) => { + total_read += n; + gate.record(n, Instant::now()); + } Err(e) => return Err(e), } } @@ -489,23 +561,25 @@ fn read_chunked_body( limits: &ConnLimits, deadline: Instant, ) -> Result<(Vec, usize), ChunkedBodyErr> { - // Ensure `buf` holds at least `until` bytes, reading on the request - // deadline. Io(TimedOut) on expiry, UnexpectedEof on early close. + // Ensure `buf` holds at least `until` bytes, reading under the body + // stall gate (absolute cap AND burst-gated stall window). Io(TimedOut) + // on expiry, UnexpectedEof on early close. All chunked socket reads + // funnel through here, so recording bytes here covers the whole path. fn fill_to( fd: RawFd, buf: &mut Vec, until: usize, - deadline: Instant, + gate: &mut BodyStallGate, ) -> Result<(), ChunkedBodyErr> { while buf.len() < until { - match read_some(fd, buf, 8 * 1024, deadline) { + match read_some(fd, buf, 8 * 1024, gate.deadline()) { Ok(0) => { return Err(ChunkedBodyErr::Io(io::Error::new( ErrorKind::UnexpectedEof, "client closed during chunked body", ))) } - Ok(_) => {} + Ok(n) => gate.record(n, Instant::now()), Err(e) => return Err(ChunkedBodyErr::Io(e)), } } @@ -520,7 +594,7 @@ fn read_chunked_body( buf: &mut Vec, from: usize, max_line: usize, - deadline: Instant, + gate: &mut BodyStallGate, ) -> Result { let mut scan = from; loop { @@ -533,16 +607,19 @@ fn read_chunked_body( return Err(ChunkedBodyErr::Malformed); } } - fill_to(fd, buf, buf.len() + 1, deadline)?; + fill_to(fd, buf, buf.len() + 1, gate)?; } } + // `deadline` is the absolute body cap; the gate layers the burst-gated + // stall window under it. All reads below go through find_crlf/fill_to. + let mut gate = BodyStallGate::new(deadline, limits, Instant::now()); let mut pos = head_len; let mut decoded: Vec = Vec::new(); loop { // ----- size line: HEX[;extensions]\r\n ----- - let line_end = find_crlf(fd, buf, pos, MAX_SIZE_LINE, deadline)?; + let line_end = find_crlf(fd, buf, pos, MAX_SIZE_LINE, &mut gate)?; let line = &buf[pos..line_end]; let size_str = match line.iter().position(|&b| b == b';') { Some(i) => &line[..i], // chunk extensions: ignored @@ -559,7 +636,7 @@ fn read_chunked_body( // ----- trailer section: zero or more header lines, then CRLF ----- let trailer_start = pos; loop { - let t_end = find_crlf(fd, buf, pos, MAX_SIZE_LINE.max(1024), deadline)?; + let t_end = find_crlf(fd, buf, pos, MAX_SIZE_LINE.max(1024), &mut gate)?; let empty = t_end == pos; pos = t_end + 2; if empty { @@ -576,7 +653,7 @@ fn read_chunked_body( } // ----- chunk payload + trailing CRLF ----- - fill_to(fd, buf, pos + size + 2, deadline)?; + fill_to(fd, buf, pos + size + 2, &mut gate)?; decoded.extend_from_slice(&buf[pos..pos + size]); if &buf[pos + size..pos + size + 2] != b"\r\n" { return Err(ChunkedBodyErr::Malformed); diff --git a/src/serve.rs b/src/serve.rs index 615ee56..2a0a75f 100644 --- a/src/serve.rs +++ b/src/serve.rs @@ -45,6 +45,13 @@ pub struct Config { /// to full body). Sized for slow links, so much larger than /// `head_timeout`. See `ConnLimits::body_timeout`. pub body_timeout: Duration, + /// Burst size that resets the body stall clock. A body dribbling fewer + /// than this per `body_stall_timeout` window is evicted — the slowloris + /// / slow-legit discriminator. See `ConnLimits::body_burst_bytes`. + pub body_burst_bytes: usize, + /// Max time since the last qualifying body burst before eviction; + /// backstopped by `body_timeout`. See `ConnLimits::body_stall_timeout`. + pub body_stall_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, @@ -97,6 +104,8 @@ impl Config { read_buf_size: 8 * 1024, head_timeout: Duration::from_secs(30), body_timeout: Duration::from_secs(300), + body_burst_bytes: 4 * 1024, + body_stall_timeout: Duration::from_secs(20), write_timeout: Duration::from_secs(30), max_body_bytes: 16 * 1024 * 1024, drain_timeout: Duration::from_secs(30), @@ -117,6 +126,8 @@ impl Config { keep_alive_timeout: self.keep_alive_timeout, head_timeout: self.head_timeout, body_timeout: self.body_timeout, + body_burst_bytes: self.body_burst_bytes, + body_stall_timeout: self.body_stall_timeout, write_timeout: self.write_timeout, max_frame_payload: self.max_frame_payload, max_message_bytes: self.max_message_bytes, diff --git a/tests/integration.rs b/tests/integration.rs index 93fd171..65b3c15 100644 --- a/tests/integration.rs +++ b/tests/integration.rs @@ -462,6 +462,39 @@ fn spawn_server_with_timeouts( panic!("server didn't come up on {addr}"); } +/// Spawn a server with the body stall-gate knobs under test; keep-alive +/// and head budgets are set out of the way so only the body path matters. +fn spawn_server_with_body_gate( + pipeline: Pipeline, + head: Duration, + body: Duration, + burst_bytes: usize, + stall: Duration, +) -> u16 { + let port = free_port(); + let addr: SocketAddr = format!("127.0.0.1:{port}").parse().unwrap(); + std::thread::spawn(move || { + let cfg = Config { + listener_pool: 2, + scheduler_threads: Some(2), + keep_alive_timeout: Duration::from_secs(30), + head_timeout: head, + body_timeout: body, + body_burst_bytes: burst_bytes, + body_stall_timeout: stall, + ..Config::new(addr) + }; + serve_with(cfg, pipeline).unwrap(); + }); + for _ in 0..50 { + if TcpStream::connect(addr).is_ok() { + return port; + } + std::thread::sleep(Duration::from_millis(50)); + } + panic!("server didn't come up on {addr}"); +} + /// Read from `s` until the response head is complete (double CRLF). Only /// suitable for responses with an empty body. fn read_response_head(s: &mut TcpStream) -> Vec { @@ -1003,6 +1036,78 @@ fn fixed_body_stall_killed_at_body_timeout() { assert!(start.elapsed() < Duration::from_secs(3), "close took too long"); } +/// Burst gate, NEGATIVE (chunked path): a client that ACTIVELY but SMOOTHLY +/// trickles sub-burst bytes is evicted at ~body_stall_timeout — even though +/// the absolute body_timeout is far away and the client never goes fully +/// silent. This is the slowloris-body case the gate exists to catch, and +/// exercises the fill_to gate in read_chunked_body. +#[test] +fn body_smooth_trickle_evicted_at_stall_timeout() { + let port = spawn_server_with_body_gate( + echo_pipeline(), + Duration::from_secs(30), // head_timeout out of the way + Duration::from_secs(30), // body_timeout out of the way (prove it's the STALL gate) + 4096, // body_burst_bytes + Duration::from_millis(800), // body_stall_timeout under test + ); + let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap(); + s.set_read_timeout(Some(Duration::from_secs(5))).unwrap(); + // Head + a chunk-size line announcing a 4096-byte chunk, then trickle + // its payload one byte at a time: never a full burst, so the stall mark + // never advances. + s.write_all(b"POST /echo HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: chunked\r\n\r\n1000\r\n") + .unwrap(); + let start = std::time::Instant::now(); + let mut evicted = false; + for _ in 0..200 { // up to ~20s; eviction expected at ~800ms + if s.write_all(&[b'x']).is_err() { + evicted = true; // server closed on us -> write failed + break; + } + std::thread::sleep(Duration::from_millis(100)); + } + assert!(evicted, "server never evicted the smooth sub-burst trickle"); + assert!( + start.elapsed() < Duration::from_secs(3), + "eviction took {:?}, expected ~800ms (stall gate, not the 30s cap)", start.elapsed() + ); +} + +/// Burst gate, POSITIVE (fixed-CL path): a slow-but-legit client that +/// delivers real bursts with gaps SHORTER than body_stall_timeout keeps +/// resetting the stall mark and completes intact. This is the slow IoT +/// upload the gate must NOT punish; exercises the read_body gate. +#[test] +fn bursty_slow_body_survives_stall_gate() { + let port = spawn_server_with_body_gate( + echo_pipeline(), + Duration::from_secs(30), // head_timeout out of the way + Duration::from_secs(30), // body_timeout out of the way + 4096, // body_burst_bytes + Duration::from_secs(2), // body_stall_timeout: gaps stay under this + ); + let mut s = TcpStream::connect(("127.0.0.1", port)).unwrap(); + s.set_read_timeout(Some(Duration::from_secs(10))).unwrap(); + // Promise 3 * 4096 bytes, Connection: close so read_to_end lands the + // full echo. + let burst = vec![b'x'; 4096]; + s.write_all(b"POST /echo HTTP/1.1\r\nHost: x\r\nContent-Length: 12288\r\nConnection: close\r\n\r\n") + .unwrap(); + for i in 0..3 { + s.write_all(&burst).unwrap(); + if i < 2 { + std::thread::sleep(Duration::from_millis(500)); // < 2s stall window + } + } + let mut resp = Vec::new(); + s.read_to_end(&mut resp).expect("expected full response"); + assert_eq!(http_status(&resp), 200, "resp head: {:?}", String::from_utf8_lossy(&resp[..resp.len().min(120)])); + let body_at = resp.windows(4).position(|w| w == b"\r\n\r\n").expect("no head terminator") + 4; + let body = &resp[body_at..]; + assert_eq!(body.len(), 12288, "echoed body truncated: {} bytes", body.len()); + assert!(body.iter().all(|&b| b == b'x'), "echoed body corrupted"); +} + // --------------------------------------------------------------------------- // SSE (v0.3 chunk 3) // ---------------------------------------------------------------------------