Compare commits
7
Commits
535f7bcc68
..
v0.2.3
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8bdec97842 | ||
|
|
f3ccb6e468 | ||
|
|
4f06265338 | ||
|
|
1b1ea124c8 | ||
|
|
b86c64d490 | ||
|
|
394e9b962a | ||
|
|
6f02cec261 |
+8
-1
@@ -14,14 +14,17 @@ sha1_smol = "1"
|
||||
|
||||
# dep #4, ratified 2026-06-12: serde/serde_json behind the opt-in
|
||||
# "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 }
|
||||
# config-file feature: TOML loader for tuning knobs (dep #5, 2026-08-12)
|
||||
toml = { version = "0.8", optional = true }
|
||||
|
||||
[features]
|
||||
smarm-trace = ["smarm/smarm-trace"]
|
||||
smarm-causal = ["smarm/smarm-causal"]
|
||||
channels = []
|
||||
phoenix = ["channels", "dep:serde", "dep:serde_json"]
|
||||
config-file = ["dep:serde", "dep:toml"]
|
||||
|
||||
[dev-dependencies]
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
@@ -47,3 +50,7 @@ path = "examples/crud.rs"
|
||||
name = "channels_chat"
|
||||
path = "examples/channels_chat.rs"
|
||||
required-features = ["phoenix"]
|
||||
|
||||
[[example]]
|
||||
name = "serve_toml"
|
||||
required-features = ["config-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
@@ -46,14 +46,34 @@ 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,
|
||||
/// 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
|
||||
@@ -76,7 +96,10 @@ 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),
|
||||
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,
|
||||
@@ -111,7 +134,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 +145,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 +165,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 +197,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 +219,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, &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
|
||||
@@ -330,9 +358,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 +375,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<u8>,
|
||||
limits: &ConnLimits,
|
||||
) -> Result<(parser::ParsedHead, Instant), ReadHeadErr> {
|
||||
) -> Result<parser::ParsedHead, ReadHeadErr> {
|
||||
let entry = Instant::now();
|
||||
let idle_deadline = entry + limits.keep_alive_timeout;
|
||||
// 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
|
||||
} else {
|
||||
Some(entry + limits.request_timeout)
|
||||
Some(entry + limits.head_timeout)
|
||||
};
|
||||
|
||||
loop {
|
||||
@@ -375,11 +403,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 +414,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
|
||||
});
|
||||
@@ -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
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -421,7 +497,8 @@ fn read_body(
|
||||
buf: &mut Vec<u8>,
|
||||
head_len: usize,
|
||||
body_len: usize,
|
||||
deadline: Instant,
|
||||
limits: &ConnLimits,
|
||||
cap: Instant,
|
||||
) -> io::Result<Vec<u8>> {
|
||||
// Bytes already in `buf` past the head belong to the body.
|
||||
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());
|
||||
}
|
||||
|
||||
// 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 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),
|
||||
}
|
||||
}
|
||||
@@ -451,8 +532,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.
|
||||
//
|
||||
@@ -479,23 +561,25 @@ fn read_chunked_body(
|
||||
limits: &ConnLimits,
|
||||
deadline: Instant,
|
||||
) -> Result<(Vec<u8>, 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<u8>,
|
||||
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)),
|
||||
}
|
||||
}
|
||||
@@ -510,7 +594,7 @@ fn read_chunked_body(
|
||||
buf: &mut Vec<u8>,
|
||||
from: usize,
|
||||
max_line: usize,
|
||||
deadline: Instant,
|
||||
gate: &mut BodyStallGate,
|
||||
) -> Result<usize, ChunkedBodyErr> {
|
||||
let mut scan = from;
|
||||
loop {
|
||||
@@ -523,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<u8> = 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
|
||||
@@ -549,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 {
|
||||
@@ -566,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);
|
||||
@@ -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",
|
||||
ParseError::Unsupported =>
|
||||
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
|
||||
// appear (read_head loops on it).
|
||||
_ =>
|
||||
|
||||
@@ -44,3 +44,5 @@ pub use ws::{Message, WsClosed, WsHandler, WsSender};
|
||||
pub use serve::{
|
||||
serve, serve_with, serve_with_shutdown, shutdown_handle, Config, Handle, ShutdownSignal,
|
||||
};
|
||||
#[cfg(feature = "config-file")]
|
||||
pub use serve::ConfigError;
|
||||
|
||||
+231
-12
@@ -9,8 +9,10 @@
|
||||
//! - No body header — empty body.
|
||||
//! - `Transfer-Encoding: chunked` (HTTP/1.1) — flagged in `ParsedHead`;
|
||||
//! the connection actor decodes incrementally (`read_chunked_body`).
|
||||
//! Chunked + Content-Length together, or chunked on HTTP/1.0, is
|
||||
//! Malformed (request-smuggling ambiguity; RFC 7230 §3.3.3).
|
||||
//! TE is 1.1-only and overrides Content-Length: TE on HTTP/1.0, or TE
|
||||
//! 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};
|
||||
|
||||
@@ -33,6 +35,10 @@ pub enum ParseError {
|
||||
/// (chunked decoding landed in v0.3); kept for future unsupported
|
||||
/// framings. Connection actor responds 411 + close.
|
||||
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 chunked = 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() {
|
||||
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() {
|
||||
"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(
|
||||
value.trim()
|
||||
.parse::<usize>()
|
||||
@@ -114,10 +129,17 @@ pub fn parse_head(buf: &[u8], max_headers: usize) -> Result<ParsedHead, ParseErr
|
||||
);
|
||||
}
|
||||
"transfer-encoding" => {
|
||||
// We only care whether it includes "chunked". Multiple codings
|
||||
// can appear; chunked is the only one we'd need to decode.
|
||||
if value.to_ascii_lowercase().split(',').any(|t| t.trim() == "chunked") {
|
||||
chunked = true;
|
||||
// Collect the ordered coding list across any number of TE
|
||||
// headers; finality/known-ness is decided post-loop. Empty
|
||||
// list elements (legacy `#rule`, e.g. a trailing comma) are
|
||||
// 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" => {
|
||||
@@ -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_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());
|
||||
}
|
||||
|
||||
if chunked {
|
||||
// Transfer-Encoding is an HTTP/1.1 mechanism; a 1.0 request
|
||||
// carrying it is malformed. And a request carrying BOTH a
|
||||
// Content-Length and TE: chunked is the classic request-smuggling
|
||||
// ambiguity — RFC 7230 §3.3.3 lets a server reject it, and we do.
|
||||
if version == HttpVersion::Http10 || content_length.is_some() {
|
||||
// Host (RFC 9112 §3.2): an HTTP/1.1 request MUST carry exactly one valid
|
||||
// Host; a missing, duplicate, or malformed Host is a 400. HTTP/1.0 may
|
||||
// omit Host, but a duplicate or invalid one is still rejected on any
|
||||
// version (ambiguous / malformed authority).
|
||||
if host_count > 1 || !host_ok {
|
||||
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);
|
||||
}
|
||||
// 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:
|
||||
@@ -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
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -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]
|
||||
fn serialise_basic_200() {
|
||||
let conn = Conn::new().put_status(200).put_body("hi");
|
||||
|
||||
+165
-4
@@ -37,7 +37,21 @@ 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,
|
||||
/// 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,
|
||||
@@ -65,6 +79,12 @@ pub struct Config {
|
||||
/// 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`]).
|
||||
@@ -82,7 +102,10 @@ 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),
|
||||
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),
|
||||
@@ -90,6 +113,7 @@ impl Config {
|
||||
max_message_bytes: 4 * 1024 * 1024,
|
||||
scheduler_threads: None,
|
||||
conn_stack_reserve: DEFAULT_CONN_STACK_RESERVE,
|
||||
max_actors: None,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -100,7 +124,10 @@ 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,
|
||||
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,
|
||||
@@ -108,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
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -326,10 +426,13 @@ pub fn serve_with_shutdown(
|
||||
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),
|
||||
None => smarm::Config::default(),
|
||||
};
|
||||
if let Some(m) = config.max_actors {
|
||||
smarm_cfg = smarm_cfg.max_actors(m);
|
||||
}
|
||||
let rt = smarm::init(smarm_cfg);
|
||||
// Listener self-termination flag — see the shutdown sequence below.
|
||||
let shutdown_flag = Arc::new(AtomicBool::new(false));
|
||||
@@ -440,3 +543,61 @@ pub fn serve(addr: impl ToSocketAddrs, pipeline: Pipeline) -> io::Result<()> {
|
||||
.ok_or_else(|| io::Error::new(ErrorKind::InvalidInput, "no addresses resolved"))?;
|
||||
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
@@ -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};
|
||||
|
||||
|
||||
+184
-12
@@ -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,41 @@ 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();
|
||||
});
|
||||
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)
|
||||
};
|
||||
serve_with(cfg, pipeline).unwrap();
|
||||
@@ -486,7 +521,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 +547,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 +950,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 +975,139 @@ 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");
|
||||
}
|
||||
|
||||
/// 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)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
Reference in New Issue
Block a user