11 Commits
Author SHA1 Message Date
Claude (sandbox) 8bdec97842 feat(config): optional TOML config loading behind config-file feature
A file-based way to set tuning knobs without recompiling. urus is a
library, so it never presumes a config path or reads the environment — the
binary hands the text in:

- Config::with_toml_str(&str) -> Result<Config, ConfigError>: sparse overlay
  onto an existing Config (built with the addr the binary chose). Only keys
  present are applied; durations are integer seconds; unknown keys are a
  hard error (deny_unknown_fields) so a typo is loud, not a silent no-op.
- Scope: the four slowloris knobs (head_timeout_secs, body_timeout_secs,
  body_burst_bytes, body_stall_timeout_secs). Migrating the rest of Config
  into the file is a separate, additive job — TomlOverrides just grows.
- Feature `config-file = ["dep:serde", "dep:toml"]`; the optional serde dep
  gains the derive feature. The default build is unchanged (deps + code are
  all gated).
- examples/serve_toml.rs (required-features = ["config-file"]): a `--config
  PATH` demo with no presumed default location. plain_serve and its env
  vars are left untouched.

Tests (feature-gated): empty keeps defaults, partial overrides only named,
full overrides all, unknown key errors, malformed errors. 84 lib with the
feature / 79 without; clippy --lib clean both ways; e2e smoke serves 200
from a file and rejects an unknown key loudly.
2026-08-12 13:41:57 +00:00
Claude (sandbox) f3ccb6e468 feat(serve): burst-gated body stall eviction
The body budget from the prior commit is a generous absolute cap; alone it
just hands a body-phase slowloris a bigger window. Add a stall gate under
that cap that distinguishes a slowloris trickle from a slow-but-legit
client by requiring BURSTS, not a mere average rate:

- BodyStallGate: each body read is bounded by min(body cap, mark + stall).
  The stall mark advances only when body_burst_bytes accumulate since the
  last advance, so a steady sub-burst trickle never moves it and is evicted
  at ~body_stall_timeout, while a bursty slow client keeps resetting it.
- Two words of state; one add + one compare per read. Raw socket bytes are
  counted, so chunked framing counts and an MSS-fragmented burst still
  accumulates. Reuses the existing read_some deadline plumbing.
- Wired into read_body (fixed CL) and read_chunked_body (via fill_to, the
  single choke point all chunked reads pass through).
- New knobs body_burst_bytes (4 KiB) + body_stall_timeout (20s); effective
  floor ~205 B/s enforced in bursts.

Tests: body_smooth_trickle_evicted_at_stall_timeout (chunked; active
sub-burst trickle evicted at ~stall while the cap is far away) and
bursty_slow_body_survives_stall_gate (fixed CL; real bursts with sub-stall
gaps complete intact). 79 lib + 45 integration green; clippy --lib clean.
2026-08-12 13:38:16 +00:00
Claude (sandbox) 4f06265338 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.
2026-08-12 13:00:22 +00:00
Claude (sandbox) 1b1ea124c8 feat(serve): plumb Config.max_actors through to smarm init
Config gained a max_actors: Option<usize> (None = smarm's DEFAULT_MAX_ACTORS
of 16_384). serve_with now applies it to the smarm runtime config, so the
per-connection actor slab can be sized to the deployment's peak concurrent
connections. Since each connection is one actor, the slab was the hard cap
on concurrent connections (previously an un-raisable 16_384) regardless of
RAM/fds; slots are ~256 B so raising it is cheap next to per-conn stacks.

Verified on the GPU box: default caps at 16_383 held connections; with
max_actors raised, a paced ramp holds 100_000 concurrent slow-header
connections at 12.4 KB RSS / 2 VMAs each (1.29 GB total) on one pinned core.
2026-08-10 05:47:22 +00:00
Claude (sandbox) b86c64d490 feat(parser): strict Transfer-Encoding framing; unknown coding -> 501
The TE arm set chunked whenever the token appeared anywhere in the value,
so 'chunked, gzip' (chunked not final) was accepted and an unknown coding
like 'bogus' was treated as no-body (h1spec #18/#19 -> 404). Collect the
ordered coding list across all TE headers and decide post-loop: TE on
HTTP/1.0 or TE+Content-Length -> 400 (the CL check now covers ANY TE, not
just chunked, closing the old TE:unknown + CL smuggling gap); chunked
present but not final -> 400; any coding other than chunked -> 501 via a
new UnknownTransferCoding variant (emit_error_response gains the 501 arm);
only a sole final chunked sets the flag. Tests cover each branch.

Note: the dead Unsupported/411 variant is left as-is (separate cleanup).
2026-08-09 08:02:04 +00:00
Claude (sandbox) 394e9b962a feat(parser): reject duplicate Content-Length (RFC 9112 §6.3)
The content-length arm ran content_length = Some(parse) per header, so a
second Content-Length silently overwrote the first with no conflict check
(CL.CL request smuggling; h1spec #21 -> 404 instead of 400). Count
occurrences and reject any duplicate post-loop, strictly (even equal
values), reusing BadContentLength (400). A single value is still required
to be one decimal integer, so a comma-list or non-numeric keeps failing at
parse as before. Tests: differing dup, equal dup, single-CL regression.
2026-08-09 07:59:25 +00:00
Claude (sandbox) 6f02cec261 feat(parser): reject missing/duplicate/invalid Host (RFC 9112 §3.2)
parse_head never inspected Host, so a missing (HTTP/1.1), duplicate, or
syntactically invalid Host all passed through to the router (h1spec #8/#9/
#10 -> 404 instead of 400). Add per-header validity (RFC 3986 host[:port]
charset via valid_host) plus a post-loop presence/uniqueness check: 1.1
MUST carry exactly one valid Host; 1.0 may omit it but a duplicate/invalid
one is still 400. Unit matrix mirrors the three h1spec cases with reg-name/
port/IPv6-literal positive controls.
2026-08-09 07:58:16 +00:00
Markk116 535f7bcc68 feat(serve): give connection actors a 256 KiB stack via smarm SpawnOpts
Connection actors were still spawned with a bare smarm::spawn(), which
gets the runtime's fixed 64 KiB default stack regardless of smarm
v0.6.0's RFC 019 SpawnOpts/stack_reserve work landing one crate down.
Any handler that leans on app code with real stack needs (DB drivers,
(de)compression, ...) blows the guard page and the connection just
dies with no response - reproduced with a CCC handler that decompresses
gzip on the identity-encoding path.

Add Config::conn_stack_reserve (default DEFAULT_CONN_STACK_RESERVE =
256 KiB) and thread it through listener_loop into a
smarm::spawn_with(SpawnOpts { stack_reserve: Some(_), .. }, ...) call
for every accepted connection. Existing Config { ..Config::new(addr) }
call sites (tests/integration.rs) pick up the new field automatically
via struct-update syntax; no call-site churn beyond that.

Bump to 0.2.2.
2026-08-08 22:39:37 +02:00
Claude (sandbox)andClaude b137c646b0 release: v0.2.1 — switch smarm to the pinned v0.6.0 git tag
RFC 019 lands upstream: per-actor stacks, park-path shrink, recycle zap,
SIGSEGV diagnostics, introspect surface. No urus code changes required —
E1 interleaved A/B on the box shows every ka cell within +0.3..+2.9% of
the v0.5.0 pin (t8-c4 close-mode control is bistable either side; see
smarm v0.6.0 release notes). Tag must exist upstream before this builds:
push smarm master + v0.6.0 first.
2026-08-08 22:17:32 +02:00
Markk116 792897d3e4 License under MIT
- LICENSE: MIT text
- Cargo.toml: license = "MIT"
2026-08-08 16:37:57 +02:00
Markk116 b77448191e release: v0.2.0 — switch smarm to the pinned v0.5.0 git tag
The gen_server API port itself already landed upstream (078072b, tracking
smarm HEAD efbc254 pre-git-dep). This just moves the dependency off the
local path checkout onto the git remote, pinned to the smarm v0.5.0
release tag (a03a7ca) rather than a floating branch HEAD.

Verified: 90 lib + 45 integration + 2 doc tests green, examples build,
under default and --all-features.
2026-08-08 11:50:51 +02:00
9 changed files with 827 additions and 88 deletions
+11 -3
View File
@@ -1,26 +1,30 @@
[package] [package]
name = "urus" name = "urus"
version = "0.1.0" version = "0.2.2"
edition = "2021" edition = "2021"
rust-version = "1.95" rust-version = "1.95"
description = "Cowboy/bandit-style HTTP library for the smarm actor runtime" description = "Cowboy/bandit-style HTTP library for the smarm actor runtime"
license = "MIT"
[dependencies] [dependencies]
smarm = { path = "../smarm" } smarm = { git = "https://git.kalsbeek.dev/Markk116/smarm", tag = "v0.6.0" }
httparse = "1.9" httparse = "1.9"
libc = "0.2" libc = "0.2"
sha1_smol = "1" sha1_smol = "1"
# dep #4, ratified 2026-06-12: serde/serde_json behind the opt-in # dep #4, ratified 2026-06-12: serde/serde_json behind the opt-in
# "phoenix" feature only — the "channels" core stays dependency-free. # "phoenix" feature only — the "channels" core stays dependency-free.
serde = { version = "1", optional = true } serde = { version = "1", optional = true, features = ["derive"] }
serde_json = { version = "1", optional = true } serde_json = { version = "1", optional = true }
# config-file feature: TOML loader for tuning knobs (dep #5, 2026-08-12)
toml = { version = "0.8", optional = true }
[features] [features]
smarm-trace = ["smarm/smarm-trace"] smarm-trace = ["smarm/smarm-trace"]
smarm-causal = ["smarm/smarm-causal"] smarm-causal = ["smarm/smarm-causal"]
channels = [] channels = []
phoenix = ["channels", "dep:serde", "dep:serde_json"] phoenix = ["channels", "dep:serde", "dep:serde_json"]
config-file = ["dep:serde", "dep:toml"]
[dev-dependencies] [dev-dependencies]
serde = { version = "1", features = ["derive"] } serde = { version = "1", features = ["derive"] }
@@ -46,3 +50,7 @@ path = "examples/crud.rs"
name = "channels_chat" name = "channels_chat"
path = "examples/channels_chat.rs" path = "examples/channels_chat.rs"
required-features = ["phoenix"] required-features = ["phoenix"]
[[example]]
name = "serve_toml"
required-features = ["config-file"]
+21
View File
@@ -0,0 +1,21 @@
MIT License
Copyright (c) 2026 Mark Kalsbeek
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
+48
View File
@@ -0,0 +1,48 @@
//! Serve with an optional TOML config overlay.
//!
//! urus is a library and never presumes a config path or reads the
//! environment for one — the binary decides where the file lives and hands
//! the text to `Config::with_toml_str`. Here that's a `--config PATH` flag;
//! with no flag, the compiled defaults are used unchanged.
//!
//! Requires the `config-file` feature:
//! cargo run --example serve_toml --features config-file -- --config urus.toml
//!
//! Example urus.toml (all keys optional, sparse override; seconds):
//! head_timeout_secs = 15
//! body_timeout_secs = 300
//! body_burst_bytes = 4096
//! body_stall_timeout_secs = 20
use std::net::SocketAddr;
use urus::{serve_with, Config, Conn, Next, Pipeline, Router};
fn json_id(conn: Conn, _next: Next) -> Conn {
let id: u64 = conn.params.get("id").and_then(|s| s.parse().ok()).unwrap_or(0);
conn.put_status(200)
.put_header("content-type", "application/json")
.put_body(format!("{{\"id\":{id}}}"))
}
fn main() {
let addr: SocketAddr = "0.0.0.0:8080".parse().expect("addr");
let mut cfg = Config::new(addr);
// Minimal flag scan: `--config PATH`. No presumed default location.
let mut args = std::env::args().skip(1);
while let Some(arg) = args.next() {
if arg == "--config" {
let path = args.next().expect("--config needs a PATH");
let toml = std::fs::read_to_string(&path)
.unwrap_or_else(|e| panic!("reading config {path}: {e}"));
cfg = cfg
.with_toml_str(&toml)
.unwrap_or_else(|e| panic!("invalid config {path}: {e}"));
eprintln!("serve_toml: loaded config from {path}");
}
}
let pipe = Pipeline::new().plug(Router::new().get("/json/:id", json_id));
serve_with(cfg, pipe).expect("serve");
}
+141 -52
View File
@@ -46,14 +46,34 @@ 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,
/// 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 /// 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 +96,10 @@ 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),
body_burst_bytes: 4 * 1024,
body_stall_timeout: Duration::from_secs(20),
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 +134,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 +145,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 +165,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 +197,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 +219,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, &limits, 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 +358,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 +375,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 +403,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 +414,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
}); });
@@ -412,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 // read_body
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
@@ -421,7 +497,8 @@ fn read_body(
buf: &mut Vec<u8>, buf: &mut Vec<u8>,
head_len: usize, head_len: usize,
body_len: usize, body_len: usize,
deadline: Instant, limits: &ConnLimits,
cap: Instant,
) -> io::Result<Vec<u8>> { ) -> io::Result<Vec<u8>> {
// Bytes already in `buf` past the head belong to the body. // Bytes already in `buf` past the head belong to the body.
let already = buf.len().saturating_sub(head_len); let already = buf.len().saturating_sub(head_len);
@@ -433,13 +510,17 @@ 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 cap AND the
// head was read under. // burst-gated stall window (whichever is sooner).
let mut gate = BodyStallGate::new(cap, limits, Instant::now());
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, gate.deadline()) {
Ok(0) => return Err(io::Error::new(ErrorKind::UnexpectedEof, "client closed during body")), 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), Err(e) => return Err(e),
} }
} }
@@ -451,8 +532,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.
// //
@@ -479,23 +561,25 @@ fn read_chunked_body(
limits: &ConnLimits, limits: &ConnLimits,
deadline: Instant, deadline: Instant,
) -> Result<(Vec<u8>, usize), ChunkedBodyErr> { ) -> Result<(Vec<u8>, usize), ChunkedBodyErr> {
// Ensure `buf` holds at least `until` bytes, reading on the request // Ensure `buf` holds at least `until` bytes, reading under the body
// deadline. Io(TimedOut) on expiry, UnexpectedEof on early close. // 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( fn fill_to(
fd: RawFd, fd: RawFd,
buf: &mut Vec<u8>, buf: &mut Vec<u8>,
until: usize, until: usize,
deadline: Instant, gate: &mut BodyStallGate,
) -> Result<(), ChunkedBodyErr> { ) -> Result<(), ChunkedBodyErr> {
while buf.len() < until { while buf.len() < until {
match read_some(fd, buf, 8 * 1024, deadline) { match read_some(fd, buf, 8 * 1024, gate.deadline()) {
Ok(0) => { Ok(0) => {
return Err(ChunkedBodyErr::Io(io::Error::new( return Err(ChunkedBodyErr::Io(io::Error::new(
ErrorKind::UnexpectedEof, ErrorKind::UnexpectedEof,
"client closed during chunked body", "client closed during chunked body",
))) )))
} }
Ok(_) => {} Ok(n) => gate.record(n, Instant::now()),
Err(e) => return Err(ChunkedBodyErr::Io(e)), Err(e) => return Err(ChunkedBodyErr::Io(e)),
} }
} }
@@ -510,7 +594,7 @@ fn read_chunked_body(
buf: &mut Vec<u8>, buf: &mut Vec<u8>,
from: usize, from: usize,
max_line: usize, max_line: usize,
deadline: Instant, gate: &mut BodyStallGate,
) -> Result<usize, ChunkedBodyErr> { ) -> Result<usize, ChunkedBodyErr> {
let mut scan = from; let mut scan = from;
loop { loop {
@@ -523,16 +607,19 @@ fn read_chunked_body(
return Err(ChunkedBodyErr::Malformed); 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 pos = head_len;
let mut decoded: Vec<u8> = Vec::new(); let mut decoded: Vec<u8> = Vec::new();
loop { loop {
// ----- size line: HEX[;extensions]\r\n ----- // ----- 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 line = &buf[pos..line_end];
let size_str = match line.iter().position(|&b| b == b';') { let size_str = match line.iter().position(|&b| b == b';') {
Some(i) => &line[..i], // chunk extensions: ignored Some(i) => &line[..i], // chunk extensions: ignored
@@ -549,7 +636,7 @@ fn read_chunked_body(
// ----- trailer section: zero or more header lines, then CRLF ----- // ----- trailer section: zero or more header lines, then CRLF -----
let trailer_start = pos; let trailer_start = pos;
loop { 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; let empty = t_end == pos;
pos = t_end + 2; pos = t_end + 2;
if empty { if empty {
@@ -566,7 +653,7 @@ fn read_chunked_body(
} }
// ----- chunk payload + trailing CRLF ----- // ----- 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]); decoded.extend_from_slice(&buf[pos..pos + size]);
if &buf[pos + size..pos + size + 2] != b"\r\n" { if &buf[pos + size..pos + size + 2] != b"\r\n" {
return Err(ChunkedBodyErr::Malformed); return Err(ChunkedBodyErr::Malformed);
@@ -773,6 +860,8 @@ fn emit_error_response(fd: RawFd, err: &ParseError, deadline: Instant) {
b"HTTP/1.1 400 Bad Request\r\ncontent-length: 0\r\nconnection: close\r\n\r\n", b"HTTP/1.1 400 Bad Request\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
ParseError::Unsupported => ParseError::Unsupported =>
b"HTTP/1.1 411 Length Required\r\ncontent-length: 0\r\nconnection: close\r\n\r\n", b"HTTP/1.1 411 Length Required\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
ParseError::UnknownTransferCoding =>
b"HTTP/1.1 501 Not Implemented\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
// Incomplete and Malformed both lead here; Incomplete shouldn't // Incomplete and Malformed both lead here; Incomplete shouldn't
// appear (read_head loops on it). // appear (read_head loops on it).
_ => _ =>
+2
View File
@@ -44,3 +44,5 @@ pub use ws::{Message, WsClosed, WsHandler, WsSender};
pub use serve::{ pub use serve::{
serve, serve_with, serve_with_shutdown, shutdown_handle, Config, Handle, ShutdownSignal, serve, serve_with, serve_with_shutdown, shutdown_handle, Config, Handle, ShutdownSignal,
}; };
#[cfg(feature = "config-file")]
pub use serve::ConfigError;
+231 -12
View File
@@ -9,8 +9,10 @@
//! - No body header — empty body. //! - No body header — empty body.
//! - `Transfer-Encoding: chunked` (HTTP/1.1) — flagged in `ParsedHead`; //! - `Transfer-Encoding: chunked` (HTTP/1.1) — flagged in `ParsedHead`;
//! the connection actor decodes incrementally (`read_chunked_body`). //! the connection actor decodes incrementally (`read_chunked_body`).
//! Chunked + Content-Length together, or chunked on HTTP/1.0, is //! TE is 1.1-only and overrides Content-Length: TE on HTTP/1.0, or TE
//! Malformed (request-smuggling ambiguity; RFC 7230 §3.3.3). //! together with a Content-Length, is Malformed (400). `chunked` must be
//! the final coding (non-final -> 400); any other coding is unimplemented
//! (-> 501). Only a sole final `chunked` sets the flag (RFC 9112 §6.1/§6.3).
use crate::conn::{Body, Conn, HeaderMap, HttpVersion, Method, RespBody}; use crate::conn::{Body, Conn, HeaderMap, HttpVersion, Method, RespBody};
@@ -33,6 +35,10 @@ pub enum ParseError {
/// (chunked decoding landed in v0.3); kept for future unsupported /// (chunked decoding landed in v0.3); kept for future unsupported
/// framings. Connection actor responds 411 + close. /// framings. Connection actor responds 411 + close.
Unsupported, Unsupported,
/// `Transfer-Encoding` names a transfer coding we don't implement
/// (`chunked` is the only one urus decodes). Connection actor responds
/// 501 Not Implemented + close (RFC 9112 §6.1, §7).
UnknownTransferCoding,
} }
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
@@ -100,6 +106,11 @@ pub fn parse_head(buf: &[u8], max_headers: usize) -> Result<ParsedHead, ParseErr
let mut connection_hdr = None; let mut connection_hdr = None;
let mut chunked = false; let mut chunked = false;
let mut expect_100 = false; let mut expect_100 = false;
let mut host_count = 0usize;
let mut host_ok = true;
let mut cl_count = 0usize;
let mut te_present = false;
let mut te_codings: Vec<String> = Vec::new();
for h in req.headers.iter() { for h in req.headers.iter() {
let name_lower = h.name.to_ascii_lowercase(); let name_lower = h.name.to_ascii_lowercase();
@@ -107,6 +118,10 @@ pub fn parse_head(buf: &[u8], max_headers: usize) -> Result<ParsedHead, ParseErr
match name_lower.as_str() { match name_lower.as_str() {
"content-length" => { "content-length" => {
// Count occurrences; duplicates (even equal) are rejected
// post-loop. A single value must be one decimal integer —
// a comma-list ("5, 5") or non-numeric fails parse here.
cl_count += 1;
content_length = Some( content_length = Some(
value.trim() value.trim()
.parse::<usize>() .parse::<usize>()
@@ -114,10 +129,17 @@ pub fn parse_head(buf: &[u8], max_headers: usize) -> Result<ParsedHead, ParseErr
); );
} }
"transfer-encoding" => { "transfer-encoding" => {
// We only care whether it includes "chunked". Multiple codings // Collect the ordered coding list across any number of TE
// can appear; chunked is the only one we'd need to decode. // headers; finality/known-ness is decided post-loop. Empty
if value.to_ascii_lowercase().split(',').any(|t| t.trim() == "chunked") { // list elements (legacy `#rule`, e.g. a trailing comma) are
chunked = true; // skipped; a wholly empty value leaves te_codings empty and
// is caught below.
te_present = true;
for coding in value.split(',') {
let c = coding.trim().to_ascii_lowercase();
if !c.is_empty() {
te_codings.push(c);
}
} }
} }
"connection" => { "connection" => {
@@ -126,19 +148,68 @@ pub fn parse_head(buf: &[u8], max_headers: usize) -> Result<ParsedHead, ParseErr
"expect" if value.eq_ignore_ascii_case("100-continue") => { "expect" if value.eq_ignore_ascii_case("100-continue") => {
expect_100 = true; expect_100 = true;
} }
"host" => {
// Presence/uniqueness enforced post-loop; validity here.
host_count += 1;
if !valid_host(value) {
host_ok = false;
}
}
_ => {} _ => {}
} }
headers.append(&name_lower, value.to_string()); headers.append(&name_lower, value.to_string());
} }
if chunked { // Host (RFC 9112 §3.2): an HTTP/1.1 request MUST carry exactly one valid
// Transfer-Encoding is an HTTP/1.1 mechanism; a 1.0 request // Host; a missing, duplicate, or malformed Host is a 400. HTTP/1.0 may
// carrying it is malformed. And a request carrying BOTH a // omit Host, but a duplicate or invalid one is still rejected on any
// Content-Length and TE: chunked is the classic request-smuggling // version (ambiguous / malformed authority).
// ambiguity — RFC 7230 §3.3.3 lets a server reject it, and we do. if host_count > 1 || !host_ok {
if version == HttpVersion::Http10 || content_length.is_some() { return Err(ParseError::Malformed);
}
if version == HttpVersion::Http11 && host_count == 0 {
return Err(ParseError::Malformed);
}
// Content-Length (RFC 9112 §6.3): more than one Content-Length is an
// unrecoverable framing ambiguity (CL.CL request smuggling). We are
// strict — reject any duplicate, not only differing values.
if cl_count > 1 {
return Err(ParseError::BadContentLength);
}
// Transfer-Encoding (RFC 9112 §6.1/§6.3). TE is a 1.1 mechanism and
// overrides Content-Length; only `chunked` is implemented here.
if te_present {
// TE on HTTP/1.0 is malformed (no 1.0 chunked).
if version == HttpVersion::Http10 {
return Err(ParseError::Malformed); return Err(ParseError::Malformed);
} }
// TE together with Content-Length is the classic smuggling
// ambiguity; TE overrides CL and we reject rather than forward.
if content_length.is_some() {
return Err(ParseError::Malformed);
}
// A Transfer-Encoding header that carries no coding frames nothing.
if te_codings.is_empty() {
return Err(ParseError::Malformed);
}
let last_is_chunked = te_codings.last().map(String::as_str) == Some("chunked");
let has_chunked = te_codings.iter().any(|c| c == "chunked");
if has_chunked && !last_is_chunked {
// chunked present but not final: body length isn't reliably
// determinable -> 400.
return Err(ParseError::Malformed);
}
if te_codings.iter().any(|c| c != "chunked") {
// Some coding we don't implement (chunked is the only decodable
// one). Whether or not chunked is final, we can't apply it -> 501.
return Err(ParseError::UnknownTransferCoding);
}
// Sole, final `chunked`: the connection actor decodes the body.
chunked = true;
} }
// Keep-alive logic, RFC 7230 §6.3: // Keep-alive logic, RFC 7230 §6.3:
@@ -163,6 +234,32 @@ pub fn parse_head(buf: &[u8], max_headers: usize) -> Result<ParsedHead, ParseErr
}) })
} }
/// Conservative RFC 3986 check for a `Host` field-value: non-empty and every
/// byte drawn from the `host[:port]` productions (reg-name / IP-literal
/// brackets / port colon). This is charset-level, not full structural
/// validation (no bracket matching, no pct-encoding well-formedness) — enough
/// to reject the smuggling-relevant garbage (whitespace, controls, `@`, `/`,
/// `?`, `#`) while accepting every legitimate host. Tighter structural checks
/// (bracketed IPv6, single port colon) are a possible follow-up.
fn valid_host(value: &str) -> bool {
!value.is_empty()
&& value.bytes().all(|b| {
b.is_ascii_alphanumeric()
|| matches!(
b,
// unreserved punctuation
b'-' | b'.' | b'_' | b'~'
// sub-delims
| b'!' | b'$' | b'&' | b'\'' | b'(' | b')'
| b'*' | b'+' | b',' | b';' | b'='
// pct-encoded lead
| b'%'
// IP-literal brackets + port separator
| b'[' | b']' | b':'
)
})
}
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// Conn assembly // Conn assembly
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
@@ -403,6 +500,128 @@ mod tests {
} }
} }
// --- Host (RFC 9112 §3.2) -------------------------------------------
#[test]
fn parse_missing_host_http11_is_malformed() {
let req = b"GET / HTTP/1.1\r\n\r\n";
match parse_head(req, 64) {
Err(ParseError::Malformed) => {}
_ => panic!("expected Malformed for missing Host on 1.1"),
}
}
#[test]
fn parse_missing_host_http10_is_allowed() {
// Host is optional in HTTP/1.0.
let req = b"GET / HTTP/1.0\r\n\r\n";
assert!(parse_head(req, 64).is_ok(), "1.0 may omit Host");
}
#[test]
fn parse_duplicate_host_is_malformed() {
let req = b"GET / HTTP/1.1\r\nHost: a\r\nHost: b\r\n\r\n";
match parse_head(req, 64) {
Err(ParseError::Malformed) => {}
_ => panic!("expected Malformed for duplicate Host"),
}
}
#[test]
fn parse_invalid_host_value_is_malformed() {
// Embedded whitespace — invalid in an RFC 3986 authority.
let req = b"GET / HTTP/1.1\r\nHost: bad host\r\n\r\n";
match parse_head(req, 64) {
Err(ParseError::Malformed) => {}
_ => panic!("expected Malformed for invalid Host"),
}
}
#[test]
fn parse_valid_hosts_accepted() {
// Positive controls: reg-name, reg-name:port, and IPv6-literal:port.
for req in [
b"GET / HTTP/1.1\r\nHost: example.com\r\n\r\n".as_slice(),
b"GET / HTTP/1.1\r\nHost: example.com:8080\r\n\r\n".as_slice(),
b"GET / HTTP/1.1\r\nHost: [::1]:443\r\n\r\n".as_slice(),
] {
assert!(parse_head(req, 64).is_ok(), "should accept a valid Host");
}
}
// --- Content-Length (RFC 9112 §6.3) ---------------------------------
#[test]
fn parse_conflicting_content_length_is_rejected() {
// Two differing Content-Length values — classic CL.CL smuggling.
let req = b"POST / HTTP/1.1\r\nHost: x\r\nContent-Length: 5\r\nContent-Length: 7\r\n\r\nhello!!";
match parse_head(req, 64) {
Err(ParseError::BadContentLength) => {}
_ => panic!("expected BadContentLength for conflicting CL"),
}
}
#[test]
fn parse_duplicate_equal_content_length_is_rejected() {
// Strict: even identical duplicates are rejected.
let req = b"POST / HTTP/1.1\r\nHost: x\r\nContent-Length: 5\r\nContent-Length: 5\r\n\r\nhello";
match parse_head(req, 64) {
Err(ParseError::BadContentLength) => {}
_ => panic!("expected BadContentLength for duplicate CL"),
}
}
#[test]
fn parse_single_content_length_still_ok() {
// Regression: the ordinary single-CL path is unchanged.
let req = b"POST / HTTP/1.1\r\nHost: x\r\nContent-Length: 5\r\n\r\nhello";
let head = parse_head(req, 64).unwrap();
assert_eq!(head.content_length, Some(5));
}
// --- Transfer-Encoding (RFC 9112 §6.1/§6.3) -------------------------
#[test]
fn parse_non_final_chunked_is_malformed() {
// chunked must be the FINAL coding.
let req = b"POST / HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: chunked, gzip\r\n\r\n";
match parse_head(req, 64) {
Err(ParseError::Malformed) => {}
_ => panic!("expected Malformed for non-final chunked"),
}
}
#[test]
fn parse_unknown_transfer_coding_is_unimplemented() {
// A coding urus doesn't implement, no chunked at all -> 501.
let req = b"POST / HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: nonsense\r\n\r\n";
match parse_head(req, 64) {
Err(ParseError::UnknownTransferCoding) => {}
_ => panic!("expected UnknownTransferCoding for unknown coding"),
}
}
#[test]
fn parse_gzip_then_chunked_is_unimplemented() {
// chunked IS final, but gzip is still a coding we can't apply -> 501.
let req = b"POST / HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: gzip, chunked\r\n\r\n";
match parse_head(req, 64) {
Err(ParseError::UnknownTransferCoding) => {}
_ => panic!("expected UnknownTransferCoding for gzip,chunked"),
}
}
#[test]
fn parse_te_with_content_length_is_malformed() {
// ANY Transfer-Encoding + Content-Length -> reject (smuggling),
// not only chunked+CL. This closes the old TE:unknown + CL gap.
let req = b"POST / HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: bogus\r\nContent-Length: 5\r\n\r\nhello";
match parse_head(req, 64) {
Err(ParseError::Malformed) => {}
_ => panic!("expected Malformed for TE + CL"),
}
}
#[test] #[test]
fn serialise_basic_200() { fn serialise_basic_200() {
let conn = Conn::new().put_status(200).put_body("hi"); let conn = Conn::new().put_status(200).put_body("hi");
+187 -8
View File
@@ -37,7 +37,21 @@ 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,
/// 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 /// 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,
@@ -57,8 +71,25 @@ pub struct Config {
/// (one per CPU). Set this to a small fixed number in tests so multiple /// (one per CPU). Set this to a small fixed number in tests so multiple
/// concurrent test servers don't oversubscribe the host. /// concurrent test servers don't oversubscribe the host.
pub scheduler_threads: Option<usize>, pub scheduler_threads: Option<usize>,
/// Stack reserve (RFC 019 `smarm::SpawnOpts::stack_reserve`) given to
/// each per-connection actor. Request handlers routinely pull in
/// application code — DB drivers, (de)compression, templating — whose
/// stack needs comfortably exceed smarm's bare-actor default of 64 KiB
/// (the exact shape of bug this exists to head off; see smarm RFC 019).
/// Default: 256 KiB. The reserve is virtual/demand-paged, so raising it
/// costs address space, not RSS, until a handler actually uses it.
pub conn_stack_reserve: usize,
/// Maximum concurrently-live actors — smarm's fixed slot slab, allocated
/// once at init. Each connection is one actor, so this is also the hard
/// cap on concurrent connections. `None` uses smarm's default (16_384).
/// Slots are ~256 B, so raising this is cheap relative to per-connection
/// stacks; size it to peak concurrent connections.
pub max_actors: Option<usize>,
} }
/// Default per-connection actor stack reserve (see [`Config::conn_stack_reserve`]).
pub const DEFAULT_CONN_STACK_RESERVE: usize = 256 * 1024;
impl Config { impl Config {
pub fn new(addr: SocketAddr) -> Self { pub fn new(addr: SocketAddr) -> Self {
let pool = std::thread::available_parallelism() let pool = std::thread::available_parallelism()
@@ -71,13 +102,18 @@ 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),
body_burst_bytes: 4 * 1024,
body_stall_timeout: Duration::from_secs(20),
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),
max_frame_payload: 1024 * 1024, max_frame_payload: 1024 * 1024,
max_message_bytes: 4 * 1024 * 1024, max_message_bytes: 4 * 1024 * 1024,
scheduler_threads: None, scheduler_threads: None,
conn_stack_reserve: DEFAULT_CONN_STACK_RESERVE,
max_actors: None,
} }
} }
@@ -88,7 +124,10 @@ 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,
body_burst_bytes: self.body_burst_bytes,
body_stall_timeout: self.body_stall_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,
@@ -96,6 +135,79 @@ impl Config {
} }
} }
// ---------------------------------------------------------------------------
// config-file: TOML overlay for tuning knobs
// ---------------------------------------------------------------------------
//
// urus is a library, so it never presumes a config-file path or reads the
// environment — the embedding binary decides where a file lives and hands
// the text here. This overlays a sparse TOML document onto an existing
// `Config` (built with an addr the binary chose): only the keys present are
// applied, everything else keeps the compiled default. Durations are
// integer seconds. Unknown keys are a hard error so a typo is loud, not a
// silent no-op.
//
// Scope for now: the slowloris-tuning knobs only. Migrating the rest of the
// Config surface into the file is a separate, additive job (the loader
// mechanism is general — it just extends `TomlOverrides`).
/// Error from [`Config::with_toml_str`]: the TOML failed to parse or carried
/// an unknown/mistyped key.
#[cfg(feature = "config-file")]
#[derive(Debug)]
pub enum ConfigError {
Toml(String),
}
#[cfg(feature = "config-file")]
impl std::fmt::Display for ConfigError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ConfigError::Toml(m) => write!(f, "config TOML error: {m}"),
}
}
}
#[cfg(feature = "config-file")]
impl std::error::Error for ConfigError {}
#[cfg(feature = "config-file")]
#[derive(serde::Deserialize)]
#[serde(deny_unknown_fields)]
struct TomlOverrides {
head_timeout_secs: Option<u64>,
body_timeout_secs: Option<u64>,
body_burst_bytes: Option<usize>,
body_stall_timeout_secs: Option<u64>,
}
#[cfg(feature = "config-file")]
impl Config {
/// Overlay a TOML document of tuning knobs onto this config (sparse:
/// only the keys present are applied). Durations are integer seconds.
///
/// Recognized keys: `head_timeout_secs`, `body_timeout_secs`,
/// `body_burst_bytes`, `body_stall_timeout_secs`. Unknown keys error.
/// Other `Config` knobs are not yet file-configurable.
pub fn with_toml_str(mut self, s: &str) -> Result<Self, ConfigError> {
let o: TomlOverrides =
toml::from_str(s).map_err(|e| ConfigError::Toml(e.to_string()))?;
if let Some(v) = o.head_timeout_secs {
self.head_timeout = Duration::from_secs(v);
}
if let Some(v) = o.body_timeout_secs {
self.body_timeout = Duration::from_secs(v);
}
if let Some(v) = o.body_burst_bytes {
self.body_burst_bytes = v;
}
if let Some(v) = o.body_stall_timeout_secs {
self.body_stall_timeout = Duration::from_secs(v);
}
Ok(self)
}
}
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// dup helper // dup helper
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
@@ -125,6 +237,7 @@ fn listener_loop(
listener: Arc<OwnedFd>, listener: Arc<OwnedFd>,
pipeline: Pipeline, pipeline: Pipeline,
limits: ConnLimits, limits: ConnLimits,
conn_stack_reserve: usize,
registry: GenServerRef<ConnRegistry>, registry: GenServerRef<ConnRegistry>,
shutdown: Arc<AtomicBool>, shutdown: Arc<AtomicBool>,
) { ) {
@@ -157,7 +270,11 @@ fn listener_loop(
let p = pipeline.clone(); let p = pipeline.clone();
let l = limits; let l = limits;
let r = registry.clone(); let r = registry.clone();
smarm::spawn(move || run_connection(client, p, l, r)); let opts = smarm::SpawnOpts {
stack_reserve: Some(conn_stack_reserve),
..smarm::SpawnOpts::default()
};
smarm::spawn_with(opts, move || run_connection(client, p, l, r));
} }
Err(e) if e.kind() == ErrorKind::WouldBlock => { Err(e) if e.kind() == ErrorKind::WouldBlock => {
// No pending connection. Park until the listener is // No pending connection. Park until the listener is
@@ -305,13 +422,17 @@ pub fn serve_with_shutdown(
listener_fds.push(Arc::new(dup)); listener_fds.push(Arc::new(dup));
} }
let limits = config.to_conn_limits(); let limits = config.to_conn_limits();
let drain_timeout = config.drain_timeout; let conn_stack_reserve = config.conn_stack_reserve;
let drain_timeout = config.drain_timeout;
let smarm_cfg = match config.scheduler_threads { let mut smarm_cfg = match config.scheduler_threads {
Some(n) => smarm::Config::exact(n), Some(n) => smarm::Config::exact(n),
None => smarm::Config::default(), None => smarm::Config::default(),
}; };
if let Some(m) = config.max_actors {
smarm_cfg = smarm_cfg.max_actors(m);
}
let rt = smarm::init(smarm_cfg); let rt = smarm::init(smarm_cfg);
// Listener self-termination flag — see the shutdown sequence below. // Listener self-termination flag — see the shutdown sequence below.
let shutdown_flag = Arc::new(AtomicBool::new(false)); let shutdown_flag = Arc::new(AtomicBool::new(false));
@@ -327,7 +448,7 @@ pub fn serve_with_shutdown(
let sf = shutdown_flag.clone(); let sf = shutdown_flag.clone();
sup = sup.child(ChildSpec::new(Restart::Transient, move || { sup = sup.child(ChildSpec::new(Restart::Transient, move || {
println!("urus: listener {} starting", i); println!("urus: listener {} starting", i);
listener_loop(lfd.clone(), p.clone(), limits, r.clone(), sf.clone()); listener_loop(lfd.clone(), p.clone(), limits, conn_stack_reserve, r.clone(), sf.clone());
})); }));
} }
// Default intensity (3 per 5s) applies; a listener crash-looping // Default intensity (3 per 5s) applies; a listener crash-looping
@@ -422,3 +543,61 @@ pub fn serve(addr: impl ToSocketAddrs, pipeline: Pipeline) -> io::Result<()> {
.ok_or_else(|| io::Error::new(ErrorKind::InvalidInput, "no addresses resolved"))?; .ok_or_else(|| io::Error::new(ErrorKind::InvalidInput, "no addresses resolved"))?;
serve_with(Config::new(addr), pipeline) serve_with(Config::new(addr), pipeline)
} }
#[cfg(all(test, feature = "config-file"))]
mod config_file_tests {
use super::*;
fn base() -> Config {
Config::new("127.0.0.1:0".parse().unwrap())
}
#[test]
fn toml_empty_keeps_defaults() {
let d = base();
let c = base().with_toml_str("").unwrap();
assert_eq!(c.head_timeout, d.head_timeout);
assert_eq!(c.body_timeout, d.body_timeout);
assert_eq!(c.body_burst_bytes, d.body_burst_bytes);
assert_eq!(c.body_stall_timeout, d.body_stall_timeout);
}
#[test]
fn toml_partial_overrides_only_named() {
let d = base();
let c = base().with_toml_str("head_timeout_secs = 5").unwrap();
assert_eq!(c.head_timeout, Duration::from_secs(5)); // overridden
assert_eq!(c.body_timeout, d.body_timeout); // default kept
assert_eq!(c.body_burst_bytes, d.body_burst_bytes); // default kept
assert_eq!(c.body_stall_timeout, d.body_stall_timeout);
}
#[test]
fn toml_full_overrides_all() {
let c = base()
.with_toml_str(
"head_timeout_secs = 10\n\
body_timeout_secs = 120\n\
body_burst_bytes = 8192\n\
body_stall_timeout_secs = 15\n",
)
.unwrap();
assert_eq!(c.head_timeout, Duration::from_secs(10));
assert_eq!(c.body_timeout, Duration::from_secs(120));
assert_eq!(c.body_burst_bytes, 8192);
assert_eq!(c.body_stall_timeout, Duration::from_secs(15));
}
#[test]
fn toml_unknown_key_errors() {
// A mistyped/unknown key is a hard error, not a silent no-op.
let e = base().with_toml_str("body_timeout_sec = 120"); // typo: missing 's'
assert!(e.is_err(), "unknown key should error");
}
#[test]
fn toml_malformed_errors() {
let e = base().with_toml_str("this is not = valid = toml");
assert!(e.is_err(), "malformed TOML should error");
}
}
+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};
+184 -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,41 @@ 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)
};
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}");
}
/// 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) ..Config::new(addr)
}; };
serve_with(cfg, pipeline).unwrap(); serve_with(cfg, pipeline).unwrap();
@@ -486,7 +521,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 +547,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 +950,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 +975,139 @@ 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");
}
/// 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) // SSE (v0.3 chunk 3)
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------