Compare commits
21
Commits
86c8d31b93
..
v0.2.3
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
8bdec97842 | ||
|
|
f3ccb6e468 | ||
|
|
4f06265338 | ||
|
|
1b1ea124c8 | ||
|
|
b86c64d490 | ||
|
|
394e9b962a | ||
|
|
6f02cec261 | ||
|
|
535f7bcc68 | ||
|
|
b137c646b0 | ||
|
|
792897d3e4 | ||
|
|
b77448191e | ||
|
|
34730930e7 | ||
|
|
72064c7a79 | ||
|
|
14be21db15 | ||
|
|
182f4fe602 | ||
|
|
17cd4a5ceb | ||
|
|
288b52d89c | ||
|
|
ecaddc579a | ||
|
|
072ee126f9 | ||
|
|
0f824635d1 | ||
|
|
078072b527 |
+19
-6
@@ -1,30 +1,39 @@
|
||||
[package]
|
||||
name = "urus"
|
||||
version = "0.1.0"
|
||||
version = "0.2.2"
|
||||
edition = "2021"
|
||||
rust-version = "1.95"
|
||||
description = "Cowboy/bandit-style HTTP library for the smarm actor runtime"
|
||||
license = "MIT"
|
||||
|
||||
[dependencies]
|
||||
smarm = { path = "../smarm" }
|
||||
smarm = { git = "https://git.kalsbeek.dev/Markk116/smarm", tag = "v0.6.0" }
|
||||
httparse = "1.9"
|
||||
libc = "0.2"
|
||||
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"]
|
||||
channels = []
|
||||
phoenix = ["channels", "dep:serde", "dep:serde_json"]
|
||||
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"] }
|
||||
serde_json = "1"
|
||||
|
||||
[[example]]
|
||||
name = "causal_bench"
|
||||
required-features = ["smarm-causal"]
|
||||
|
||||
[profile.dev]
|
||||
panic = "unwind"
|
||||
|
||||
@@ -41,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,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.
|
||||
@@ -378,3 +378,9 @@ reconnect cycles.
|
||||
- **Bench suite** — `urus-bench-spec.md` exists in the artefact store;
|
||||
wire it up once v0.2 lands (supervision changes the hot path not at all,
|
||||
but prove it).
|
||||
- **Operator introspection via typed names** — the v0.2-era
|
||||
`urus.server` / `urus.listener.{i}` pid tags were dropped in the
|
||||
RFC 014 port (names are messageable endpoints now, self-registered
|
||||
with a real `Sender<M>`). If wanted back, do it properly: register the
|
||||
serve loop's shutdown/control channel under a typed `urus.server`
|
||||
name instead of faking pid tags with unit channels.
|
||||
|
||||
@@ -0,0 +1,322 @@
|
||||
//! Causal-profiling bench (RFC 007): a real urus webserver as the workload,
|
||||
//! with a planted, known-answer bottleneck.
|
||||
//!
|
||||
//! Where smarm's `causal_pipeline` demo is synthetic, this exercises the
|
||||
//! exact machinery the causal commits touched, for real:
|
||||
//!
|
||||
//! - epoll parks (`wait_readable_timeout` / `wait_writable_timeout`)
|
||||
//! → park-gated resume credit
|
||||
//! - keep-alive / request / write deadlines → timer-heap virtual time
|
||||
//! - the experiment controller's fixed windows → wall-anchored timers
|
||||
//! - mixed CPU (parse/serialize) and IO (socket) → non-trivial ranking
|
||||
//!
|
||||
//! Shape: `LOAD_CONNS` plain OS threads hammer `GET /order/:id` over
|
||||
//! blocking loopback TCP (keep-alive). OS threads on the client side are
|
||||
//! deliberate: they can't absorb injected virtual delay, so the only
|
||||
//! delayable code is the server's — the same trick the smarm-side timer
|
||||
//! tests use. The handler calls a single store actor (crud's pattern:
|
||||
//! one actor owns the data, serialization is structural) which burns
|
||||
//! `STORE_US` of calibrated *work* — work-shaped, not timed, because a
|
||||
//! timed wait absorbs injected delay and flattens every experiment.
|
||||
//!
|
||||
//! Known answer: the store is serialized and saturated, so it is the
|
||||
//! throughput ceiling (1e6/STORE_US rps). A virtual speedup of `store`
|
||||
//! by p% predicts throughput ×1/(1-p): +33% @25%, +100% @50%. The other
|
||||
//! five sites (parse, router, pipeline, serialize, socket-write — the
|
||||
//! urus lib instrumentation) are tens of µs and parallel across
|
||||
//! connections: predicted impact ≈0. Per the smarm-side validation runs,
|
||||
//! absolute magnitudes are trustworthy to roughly ±15% while the
|
||||
//! *ranking* is solid — the verdict thresholds reflect that.
|
||||
//!
|
||||
//! Run (needs cores; the verdict is SKIPPED below 4):
|
||||
//! cargo run --release --example causal_bench --features smarm-causal
|
||||
//!
|
||||
//! Prints the summary, writes `profile.coz` next to the CWD, and exits
|
||||
//! nonzero if the expected separation doesn't hold.
|
||||
|
||||
use std::io::{Read, Write};
|
||||
use std::net::{SocketAddr, TcpListener, TcpStream};
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
use std::sync::{Arc, OnceLock};
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use smarm::{channel, Sender};
|
||||
use urus::{serve_with_shutdown, shutdown_handle, Config, Conn, Next, Pipeline, Router};
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Tunables
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Planted bottleneck cost: calibrated work per store request, µs.
|
||||
const STORE_US: u64 = 400;
|
||||
/// Handler-side render work per request, µs. Runs under the `pipeline`
|
||||
/// site (no inner guard), parallel across connection actors.
|
||||
const RENDER_US: u64 = 50;
|
||||
/// Client connections = OS load threads. Plenty to keep a 400µs store
|
||||
/// saturated even at a 50% virtual speedup (see saturation note in main).
|
||||
const LOAD_CONNS: usize = 16;
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Calibrated work (lifted from smarm's causal_pipeline demo)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// LCG-mix `iters` times in dependent sequence (unvectorizable, un-elidable),
|
||||
/// staying preemptible — and causal-sampleable/delayable — via `check!()`.
|
||||
fn work_iters(iters: u64) {
|
||||
let mut acc = 0x2545_f491_4f6c_dd1du64;
|
||||
let mut i = 0u64;
|
||||
while i < iters {
|
||||
let chunk_end = (i + 256).min(iters);
|
||||
while i < chunk_end {
|
||||
acc = acc.wrapping_mul(6364136223846793005).wrapping_add(i);
|
||||
i += 1;
|
||||
}
|
||||
std::hint::black_box(acc);
|
||||
smarm::check!();
|
||||
}
|
||||
}
|
||||
|
||||
/// Iterations per microsecond on this machine, so stage costs are
|
||||
/// meaningful in time while staying work-shaped.
|
||||
fn calibrate_iters_per_us() -> u64 {
|
||||
let n = 8_000_000u64;
|
||||
let t = Instant::now();
|
||||
work_iters(n);
|
||||
(n / (t.elapsed().as_micros().max(1) as u64)).max(1)
|
||||
}
|
||||
|
||||
static PER_US: OnceLock<u64> = OnceLock::new();
|
||||
|
||||
fn work_us(us: u64) {
|
||||
work_iters(us * PER_US.get().copied().unwrap_or(1));
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Store actor — the planted bottleneck (crud's once-cell pattern)
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
enum Req {
|
||||
Get { id: u64, reply: Sender<u64> },
|
||||
}
|
||||
|
||||
/// Set at teardown; the store polls it (a cross-thread `Sender::send`
|
||||
/// can't wake a parked actor, so a bare `recv` would hang AllDone —
|
||||
/// same limitation crud's store works around).
|
||||
static STORE_STOP: AtomicBool = AtomicBool::new(false);
|
||||
|
||||
static STORE_TX: OnceLock<Sender<Req>> = OnceLock::new();
|
||||
|
||||
fn store_loop(rx: smarm::Receiver<Req>) {
|
||||
loop {
|
||||
let req = match rx.recv_timeout(Duration::from_millis(250)) {
|
||||
Ok(r) => r,
|
||||
Err(smarm::channel::RecvTimeoutError::Timeout) => {
|
||||
if STORE_STOP.load(Ordering::Relaxed) {
|
||||
return;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
Err(smarm::channel::RecvTimeoutError::Disconnected) => return,
|
||||
};
|
||||
let Req::Get { id, reply } = req;
|
||||
{
|
||||
// The known-answer site: serialized by construction (one
|
||||
// actor, one request at a time), saturated by the load.
|
||||
let _g = smarm::causal_site!("store");
|
||||
work_us(STORE_US);
|
||||
}
|
||||
let _ = reply.send(id);
|
||||
}
|
||||
}
|
||||
|
||||
/// Spawned from the first handler that runs — connection actors are
|
||||
/// inside the runtime, which `spawn` requires; the static then shares
|
||||
/// the Sender with every later handler.
|
||||
fn store() -> &'static Sender<Req> {
|
||||
STORE_TX.get_or_init(|| {
|
||||
let (tx, rx) = channel::<Req>();
|
||||
smarm::spawn(move || store_loop(rx));
|
||||
tx
|
||||
})
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Handler
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn order(conn: Conn, _next: Next) -> Conn {
|
||||
let id: u64 = conn.params.get("id").and_then(|s| s.parse().ok()).unwrap_or(0);
|
||||
let (tx, rx) = channel::<u64>();
|
||||
store().send(Req::Get { id, reply: tx }).ok();
|
||||
let got = rx.recv().expect("store dropped");
|
||||
// Render work: attributed to the enclosing `pipeline` site.
|
||||
work_us(RENDER_US);
|
||||
conn.put_status(200).put_body(got.to_string())
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Load generation — plain OS threads, blocking loopback TCP, keep-alive
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Read one HTTP/1.1 response (headers + content-length body) into `buf`,
|
||||
/// draining exactly what was consumed. Errors mean: reconnect.
|
||||
fn read_response(s: &mut TcpStream, buf: &mut Vec<u8>) -> std::io::Result<()> {
|
||||
let mut tmp = [0u8; 4096];
|
||||
let header_end = loop {
|
||||
if let Some(i) = buf.windows(4).position(|w| w == b"\r\n\r\n") {
|
||||
break i + 4;
|
||||
}
|
||||
let n = s.read(&mut tmp)?;
|
||||
if n == 0 {
|
||||
return Err(std::io::ErrorKind::UnexpectedEof.into());
|
||||
}
|
||||
buf.extend_from_slice(&tmp[..n]);
|
||||
};
|
||||
let head = String::from_utf8_lossy(&buf[..header_end]).to_ascii_lowercase();
|
||||
let len: usize = head
|
||||
.lines()
|
||||
.find_map(|l| l.strip_prefix("content-length:"))
|
||||
.and_then(|v| v.trim().parse().ok())
|
||||
.unwrap_or(0);
|
||||
while buf.len() < header_end + len {
|
||||
let n = s.read(&mut tmp)?;
|
||||
if n == 0 {
|
||||
return Err(std::io::ErrorKind::UnexpectedEof.into());
|
||||
}
|
||||
buf.extend_from_slice(&tmp[..n]);
|
||||
}
|
||||
buf.drain(..header_end + len);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn load_loop(port: u16, stop: Arc<AtomicBool>) {
|
||||
let mut n: u64 = 1;
|
||||
'outer: while !stop.load(Ordering::Relaxed) {
|
||||
let mut s = match TcpStream::connect(("127.0.0.1", port)) {
|
||||
Ok(s) => s,
|
||||
Err(_) => {
|
||||
std::thread::sleep(Duration::from_millis(10));
|
||||
continue;
|
||||
}
|
||||
};
|
||||
s.set_read_timeout(Some(Duration::from_secs(5))).ok();
|
||||
s.set_nodelay(true).ok();
|
||||
let mut buf = Vec::with_capacity(4096);
|
||||
while !stop.load(Ordering::Relaxed) {
|
||||
let req = format!("GET /order/{n} HTTP/1.1\r\nhost: bench\r\n\r\n");
|
||||
if s.write_all(req.as_bytes()).is_err() {
|
||||
continue 'outer;
|
||||
}
|
||||
if read_response(&mut s, &mut buf).is_err() {
|
||||
continue 'outer;
|
||||
}
|
||||
n += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// main — boot, load, sweep, verdict
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn main() {
|
||||
let per_us = calibrate_iters_per_us();
|
||||
PER_US.set(per_us).expect("PER_US set twice");
|
||||
println!("calibration: {per_us} work iters/µs");
|
||||
println!(
|
||||
"planted: store {STORE_US}µs serialized (ceiling ≈{} rps), render {RENDER_US}µs, {LOAD_CONNS} conns",
|
||||
1_000_000 / STORE_US
|
||||
);
|
||||
println!("theory: store +33% @25%, +100% @50%; every other site ≈0");
|
||||
|
||||
// Ephemeral port via bind-and-release (integration-test pattern; the
|
||||
// race window before the server rebinds is far below flake threshold).
|
||||
let port = {
|
||||
let l = TcpListener::bind("127.0.0.1:0").expect("bind probe");
|
||||
l.local_addr().expect("local_addr").port()
|
||||
};
|
||||
let addr: SocketAddr = format!("127.0.0.1:{port}").parse().expect("addr");
|
||||
|
||||
let (handle, signal) = shutdown_handle();
|
||||
let server = std::thread::spawn(move || {
|
||||
let pipe = Pipeline::new().plug(Router::new().get("/order/:id", order));
|
||||
serve_with_shutdown(Config::new(addr), pipe, signal).expect("serve");
|
||||
});
|
||||
|
||||
// Wait until it's accepting.
|
||||
let up = (0..100).any(|_| {
|
||||
std::thread::sleep(Duration::from_millis(50));
|
||||
TcpStream::connect(addr).is_ok()
|
||||
});
|
||||
assert!(up, "server didn't come up on {addr}");
|
||||
|
||||
let stop = Arc::new(AtomicBool::new(false));
|
||||
let loaders: Vec<_> = (0..LOAD_CONNS)
|
||||
.map(|_| {
|
||||
let stop = stop.clone();
|
||||
std::thread::spawn(move || load_loop(port, stop))
|
||||
})
|
||||
.collect();
|
||||
|
||||
// Warm up: queues to steady state, every site + the progress point
|
||||
// registered by real traffic before the sweep enumerates them.
|
||||
std::thread::sleep(Duration::from_millis(500));
|
||||
|
||||
// The controller runs on this plain OS thread: its windows are wall
|
||||
// time by construction here; inside a runtime they'd be wall-anchored
|
||||
// via the efbc254 timer path (the thing the jobrunner sweep confirms).
|
||||
let results = smarm::causal::run_experiments(&smarm::causal::ExperimentPlan {
|
||||
speedups_pct: vec![0, 25, 50],
|
||||
experiment: Duration::from_millis(700),
|
||||
cooldown: Duration::from_millis(150),
|
||||
});
|
||||
|
||||
stop.store(true, Ordering::Relaxed);
|
||||
for l in loaders {
|
||||
l.join().expect("loader panicked");
|
||||
}
|
||||
STORE_STOP.store(true, Ordering::Relaxed);
|
||||
handle.shutdown();
|
||||
server.join().expect("server panicked");
|
||||
|
||||
print!("{}", smarm::causal::render_summary(&results));
|
||||
match std::fs::write("profile.coz", smarm::causal::render_coz(&results)) {
|
||||
Ok(()) => println!("\nwrote profile.coz"),
|
||||
Err(e) => eprintln!("\nfailed to write profile.coz: {e}"),
|
||||
}
|
||||
|
||||
// Verdict — same discipline as the smarm demo: skip where the
|
||||
// separation physically can't exist.
|
||||
let cores = std::thread::available_parallelism().map(|n| n.get()).unwrap_or(1);
|
||||
if cores < 4 {
|
||||
println!("verdict: SKIPPED ({cores} cores; separation needs real parallelism)");
|
||||
return;
|
||||
}
|
||||
let mut failures: Vec<String> = Vec::new();
|
||||
let impact =
|
||||
|site: &str| smarm::causal::impact_pct(&results, site, 25, "responses");
|
||||
let mut expect = |site: &str, ok: &dyn Fn(f64) -> bool, want: &str| match impact(site) {
|
||||
Some(p) => {
|
||||
let verdict = if ok(p) { "ok" } else { "FAIL" };
|
||||
println!("verdict: {site} @25% -> {p:+.1}% (want {want}) {verdict}");
|
||||
if !ok(p) {
|
||||
failures.push(format!("{site}: {p:+.1}% (want {want})"));
|
||||
}
|
||||
}
|
||||
None => {
|
||||
println!("verdict: {site} @25% -> missing cell FAIL");
|
||||
failures.push(format!("{site}: missing cell"));
|
||||
}
|
||||
};
|
||||
expect("store", &|p| p > 15.0, "> +15%");
|
||||
for site in ["parse", "router", "pipeline", "serialize", "socket-write"] {
|
||||
expect(site, &|p| p < 10.0, "< +10%");
|
||||
}
|
||||
|
||||
if failures.is_empty() {
|
||||
println!("verdict: PASS — planted bottleneck found, cold sites quiet");
|
||||
} else {
|
||||
println!("verdict: FAIL — {}", failures.join("; "));
|
||||
std::process::exit(1);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,106 @@
|
||||
//! Causal load-profiling target (RFC 007): the real urus hot path under
|
||||
//! EXTERNAL load — no planted bottleneck, no in-process load generation.
|
||||
//!
|
||||
//! Where `causal_bench` validates the machinery against a planted,
|
||||
//! known-answer store, this is the discovery tool: boot a bare server,
|
||||
//! let an external generator (wrk, pinned to different cores) drive it,
|
||||
//! sweep virtual speedups over the five lib sites (parse, router,
|
||||
//! pipeline, serialize, socket-write), and report which one causally
|
||||
//! limits throughput. External load is deliberate: the generator lives
|
||||
//! outside the smarm runtime, so it cannot absorb injected virtual
|
||||
//! delay — the same property `causal_bench` gets from plain OS threads.
|
||||
//!
|
||||
//! No pass/fail verdict — there is no known answer here. Trust gates
|
||||
//! only: traffic actually flowed through the sweep, and the ledger
|
||||
//! audit (printed) balances.
|
||||
//!
|
||||
//! Env:
|
||||
//! URUS_PORT listen port (default 8080; binds 0.0.0.0)
|
||||
//! COZ_OUT coz profile output path (default profile.coz)
|
||||
//! WARMUP_MS settle time after first traffic, ms (default 2000)
|
||||
//!
|
||||
//! Orchestration contract (the jobrunner run.sh): start this, wait for
|
||||
//! it to accept, point wrk at `GET /json/:id` with a duration that
|
||||
//! outlives the sweep, and stop wrk when this process exits.
|
||||
//!
|
||||
//! cargo run --release --example load_profile --features smarm-causal
|
||||
|
||||
use std::net::SocketAddr;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use urus::{serve_with_shutdown, shutdown_handle, Config, Conn, Next, Pipeline, Router};
|
||||
|
||||
fn env_u64(name: &str, default: u64) -> u64 {
|
||||
std::env::var(name).ok().and_then(|v| v.parse().ok()).unwrap_or(default)
|
||||
}
|
||||
|
||||
/// The whole handler: param parse + JSON render. Deliberately thin — the
|
||||
/// subject is the lib path around it, not application work.
|
||||
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 port = env_u64("URUS_PORT", 8080) as u16;
|
||||
let warmup = Duration::from_millis(env_u64("WARMUP_MS", 2000));
|
||||
let coz_out = std::env::var("COZ_OUT").unwrap_or_else(|_| "profile.coz".into());
|
||||
|
||||
let addr: SocketAddr = format!("0.0.0.0:{port}").parse().expect("addr");
|
||||
let (handle, signal) = shutdown_handle();
|
||||
let server = std::thread::spawn(move || {
|
||||
let pipe = Pipeline::new().plug(Router::new().get("/json/:id", json_id));
|
||||
serve_with_shutdown(Config::new(addr), pipe, signal).expect("serve");
|
||||
});
|
||||
|
||||
// Readiness is the orchestrator's job (TCP probe); ours is to not
|
||||
// sweep before real traffic has registered every site and the
|
||||
// progress point — `run_experiments` enumerates *registered* sites.
|
||||
// 100 completed responses guarantees full end-to-end coverage.
|
||||
let responses = |snap: &[(String, u64)]| {
|
||||
snap.iter().find(|(n, _)| n == "responses").map(|(_, c)| *c).unwrap_or(0)
|
||||
};
|
||||
let t0 = Instant::now();
|
||||
let seen = loop {
|
||||
let n = responses(&smarm::causal::progress_snapshot());
|
||||
if n >= 100 {
|
||||
break n;
|
||||
}
|
||||
assert!(
|
||||
t0.elapsed() < Duration::from_secs(120),
|
||||
"no load after 120s ({n} responses) — is the generator running?"
|
||||
);
|
||||
std::thread::sleep(Duration::from_millis(100));
|
||||
};
|
||||
println!("traffic up: {seen} responses; settling {warmup:?}");
|
||||
std::thread::sleep(warmup);
|
||||
|
||||
// Same plan as causal_bench, for comparability across workloads.
|
||||
let results = smarm::causal::run_experiments(&smarm::causal::ExperimentPlan {
|
||||
speedups_pct: vec![0, 25, 50],
|
||||
experiment: Duration::from_millis(700),
|
||||
cooldown: Duration::from_millis(150),
|
||||
});
|
||||
|
||||
handle.shutdown();
|
||||
server.join().expect("server panicked");
|
||||
|
||||
print!("{}", smarm::causal::render_summary(&results));
|
||||
print!("{}", smarm::causal::render_ledger_audit(&results));
|
||||
match std::fs::write(&coz_out, smarm::causal::render_coz(&results)) {
|
||||
Ok(()) => println!("\nwrote {coz_out}"),
|
||||
Err(e) => eprintln!("\nfailed to write {coz_out}: {e}"),
|
||||
}
|
||||
|
||||
// Trust gate: the sweep is meaningless if load didn't flow through it.
|
||||
let total: u64 = results
|
||||
.iter()
|
||||
.flat_map(|r| r.deltas.iter())
|
||||
.filter(|(n, _)| n == "responses")
|
||||
.map(|(_, c)| c)
|
||||
.sum();
|
||||
println!("total responses across windows: {total}");
|
||||
assert!(total > 0, "sweep saw zero progress");
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
//! Plain benchmark server: the `load_profile` request path with zero causal
|
||||
//! machinery — the baseline half of the ka/close A/B matrix (the
|
||||
//! throughput-inversion chase).
|
||||
//!
|
||||
//! Identical route, handler, and `Config` construction to `load_profile`;
|
||||
//! differs only in having no `smarm-causal` feature, no sweep, and no
|
||||
//! shutdown path — it serves until killed. Throughput is measured
|
||||
//! externally (wrk).
|
||||
//!
|
||||
//! Env:
|
||||
//! URUS_PORT listen port (default 8080; binds 0.0.0.0)
|
||||
//! URUS_SCHED_THREADS smarm scheduler OS threads (default: smarm's own
|
||||
//! default). Malformed values panic rather than
|
||||
//! silently falling back — a benchmark knob that
|
||||
//! quietly reverts to default poisons the cell.
|
||||
//!
|
||||
//! cargo run --release --example plain_serve
|
||||
|
||||
use std::net::SocketAddr;
|
||||
|
||||
use urus::{serve_with, Config, Conn, Next, Pipeline, Router};
|
||||
|
||||
/// Mirrors `load_profile`'s handler byte for byte: param parse + JSON render.
|
||||
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 port: u16 = std::env::var("URUS_PORT")
|
||||
.ok()
|
||||
.and_then(|v| v.parse().ok())
|
||||
.unwrap_or(8080);
|
||||
let addr: SocketAddr = format!("0.0.0.0:{port}").parse().expect("addr");
|
||||
let sched_threads: Option<usize> = std::env::var("URUS_SCHED_THREADS").ok().map(|v| {
|
||||
v.parse()
|
||||
.unwrap_or_else(|_| panic!("URUS_SCHED_THREADS not a usize: {v:?}"))
|
||||
});
|
||||
// Audit line: lands in each bench cell's server.log so the effective
|
||||
// scheduler count is recorded per cell, same discipline as mode-verify.
|
||||
eprintln!("plain_serve: scheduler_threads={sched_threads:?}");
|
||||
let mut cfg = Config::new(addr);
|
||||
cfg.scheduler_threads = sched_threads;
|
||||
let pipe = Pipeline::new().plug(Router::new().get("/json/:id", json_id));
|
||||
serve_with(cfg, pipe).expect("serve");
|
||||
}
|
||||
@@ -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");
|
||||
}
|
||||
+1
-1
@@ -92,7 +92,7 @@ impl WsHandler for ChatHandler {
|
||||
// subscribed as.
|
||||
let _ = self
|
||||
.bus()
|
||||
.broadcast_from(smarm::self_pid(), &self.topic(), text);
|
||||
.broadcast_from(smarm::self_pid(), self.topic(), text);
|
||||
}
|
||||
|
||||
fn on_close(&mut self, _code: Option<u16>, _reason: &str) {
|
||||
|
||||
Executable
+89
@@ -0,0 +1,89 @@
|
||||
#!/usr/bin/env bash
|
||||
# E1: does the ka low-concurrency latency floor track idle-scheduler count?
|
||||
#
|
||||
# Mechanism under test (smarm): one shared level-triggered wake pipe; every
|
||||
# completion byte wakes ALL idle schedulers; one drain-lock winner; losers
|
||||
# stampede the timers/io/queue mutexes and re-sleep; `enqueue` is silent.
|
||||
# Prediction if right: at c <= 8, ka latency tail shrinks and throughput
|
||||
# rises as scheduler count drops (fewer idle pollers -> smaller herd);
|
||||
# close mode (control) stays flat or worsens as threads drop.
|
||||
#
|
||||
# Sweeps URUS_SCHED_THREADS x CONNS over the plain ka/close matrix,
|
||||
# PLAIN_ONLY — causal cells are irrelevant to E1. wrk side is byte-identical
|
||||
# to the 96b40ad5 conns sweep (THREADS=4, --latency) for comparability.
|
||||
#
|
||||
# Knobs: THREADS_SET CONNS_SET DUR REPS OUT (+ everything the matrix takes)
|
||||
set -euo pipefail
|
||||
cd "$(dirname "$0")/.."
|
||||
|
||||
THREADS_SET="${THREADS_SET:-1 2 4 8}"
|
||||
CONNS_SET="${CONNS_SET:-4 8}"
|
||||
DUR="${DUR:-15}"
|
||||
REPS="${REPS:-2}"
|
||||
OUT="${OUT:-/workspace/results}"
|
||||
|
||||
mkdir -p "$OUT"
|
||||
say() { echo "[$(date +%H:%M:%S)] $*" | tee -a "$OUT/e1.log"; }
|
||||
|
||||
for t in $THREADS_SET; do
|
||||
for c in $CONNS_SET; do
|
||||
cell="$OUT/t${t}-c${c}"
|
||||
say "=== E1 cell: sched_threads=$t conns=$c -> $cell ==="
|
||||
URUS_SCHED_THREADS="$t" PLAIN_ONLY=1 \
|
||||
DUR="$DUR" REPS="$REPS" CONNS="$c" OUT="$cell" \
|
||||
bash scripts/ka-close-matrix.sh
|
||||
# Cell audit: the knob must have actually reached the server. A cell
|
||||
# whose server silently ran at default threads poisons the sweep — the
|
||||
# exact failure class per-cell verification exists to catch.
|
||||
for m in ka close; do
|
||||
if ! grep -q "scheduler_threads=Some($t)" "$cell/plain-$m/server.log"; then
|
||||
say "FATAL: t=$t c=$c mode=$m server.log lacks scheduler_threads=Some($t)"
|
||||
exit 1
|
||||
fi
|
||||
done
|
||||
say "cell t=$t c=$c audit OK"
|
||||
done
|
||||
done
|
||||
|
||||
say "=== E1 SUMMARY ==="
|
||||
python3 - "$OUT" <<'PY' | tee -a "$OUT/e1.log"
|
||||
import re, sys, pathlib, statistics
|
||||
|
||||
out = pathlib.Path(sys.argv[1])
|
||||
|
||||
def ms(tok):
|
||||
m = re.match(r"([\d.]+)(us|ms|s)$", tok)
|
||||
if not m:
|
||||
return None
|
||||
v = float(m.group(1))
|
||||
return {"us": v / 1000, "ms": v, "s": v * 1000}[m.group(2)]
|
||||
|
||||
def cell_stats(d):
|
||||
reps, pcts = [], {}
|
||||
for rep in sorted(d.glob("rep*.txt")):
|
||||
t = rep.read_text()
|
||||
m = re.search(r"Requests/sec:\s+([\d.]+)", t)
|
||||
if m:
|
||||
reps.append(float(m.group(1)))
|
||||
for p, tok in re.findall(r"^\s+(50|75|90|99)%\s+(\S+)$", t, re.M):
|
||||
pcts.setdefault(p, []).append(ms(tok))
|
||||
if not reps:
|
||||
return None
|
||||
return (statistics.mean(reps),
|
||||
{p: statistics.mean([v for v in vs if v is not None])
|
||||
for p, vs in pcts.items()})
|
||||
|
||||
cells = sorted(out.glob("t*-c*"))
|
||||
hdr = f"{'cell':>10} {'mode':>6} {'req/s':>9} {'p50ms':>7} {'p75ms':>7} {'p90ms':>7} {'p99ms':>7}"
|
||||
print(hdr); print("-" * len(hdr))
|
||||
for cell in cells:
|
||||
for mode in ("ka", "close"):
|
||||
s = cell_stats(cell / f"plain-{mode}")
|
||||
if s is None:
|
||||
continue
|
||||
rps, p = s
|
||||
print(f"{cell.name:>10} {mode:>6} {rps:>9.0f} "
|
||||
f"{p.get('50', float('nan')):>7.3f} {p.get('75', float('nan')):>7.3f} "
|
||||
f"{p.get('90', float('nan')):>7.3f} {p.get('99', float('nan')):>7.3f}")
|
||||
PY
|
||||
say "=== E1 DONE ==="
|
||||
Executable
+181
@@ -0,0 +1,181 @@
|
||||
#!/usr/bin/env bash
|
||||
# ka-close-matrix.sh — controlled ka/close A/B for the throughput-inversion
|
||||
# chase (handoff jar item), with the forgiveness-fix box validation folded
|
||||
# into the causal cells (handoff v13 PENDING, pull-forward agreed 2026-07-20).
|
||||
#
|
||||
# Cells, run in order on fresh ports:
|
||||
# plain-ka, plain-close examples/plain_serve (no causal feature).
|
||||
# Metric: wrk Requests/sec, REPS reps after a
|
||||
# discarded warmup.
|
||||
# causal-ka, causal-close examples/load_profile (smarm-causal). Metric:
|
||||
# the sweep summary + ledger audit (forgiveness
|
||||
# column, books balance); wrk is backdrop load
|
||||
# whose own numbers are injection-contaminated.
|
||||
#
|
||||
# Controls: byte-identical wrk invocation per mode pair except the
|
||||
# `Connection: close` header; server and wrk pinned to disjoint core sets
|
||||
# (SMT siblings left idle by default on the 5900X); the negotiated
|
||||
# connection behavior is verified via curl per cell and logged, so a
|
||||
# loadgen-config asymmetry can never silently explain a result again.
|
||||
#
|
||||
# Knobs (env, defaults for the 5900X box):
|
||||
# DUR=30 REPS=2 CONNS=64 THREADS=4 PIN=1
|
||||
# SERVER_CPUS=0-7 WRK_CPUS=8-11
|
||||
# CAUSAL_WRK_DUR=300 BASE_PORT=8080
|
||||
# OUT=/workspace/results
|
||||
# PLAIN_BIN, CAUSAL_BIN binary paths (default: target/release/examples/*)
|
||||
set -euo pipefail
|
||||
|
||||
cd "$(dirname "$0")/.."
|
||||
|
||||
DUR="${DUR:-30}"
|
||||
REPS="${REPS:-2}"
|
||||
CONNS="${CONNS:-64}"
|
||||
THREADS="${THREADS:-4}"
|
||||
PIN="${PIN:-1}"
|
||||
SERVER_CPUS="${SERVER_CPUS:-0-7}"
|
||||
WRK_CPUS="${WRK_CPUS:-8-11}"
|
||||
CAUSAL_WRK_DUR="${CAUSAL_WRK_DUR:-300}"
|
||||
BASE_PORT="${BASE_PORT:-8080}"
|
||||
OUT="${OUT:-/workspace/results}"
|
||||
PLAIN_BIN="${PLAIN_BIN:-target/release/examples/plain_serve}"
|
||||
CAUSAL_BIN="${CAUSAL_BIN:-target/release/examples/load_profile}"
|
||||
|
||||
mkdir -p "$OUT"
|
||||
say() { echo "[$(date +%H:%M:%S)] $*" | tee -a "$OUT/matrix.log"; }
|
||||
|
||||
pin_server=(); pin_wrk=()
|
||||
if [ "$PIN" = 1 ]; then
|
||||
pin_server=(taskset -c "$SERVER_CPUS")
|
||||
pin_wrk=(taskset -c "$WRK_CPUS")
|
||||
fi
|
||||
|
||||
# ---- environment record --------------------------------------------------
|
||||
{
|
||||
date -u
|
||||
uname -a
|
||||
echo "nproc: $(nproc)"
|
||||
lscpu -e 2>/dev/null || true
|
||||
echo "port range: $(cat /proc/sys/net/ipv4/ip_local_port_range 2>/dev/null)"
|
||||
echo "somaxconn: $(cat /proc/sys/net/core/somaxconn 2>/dev/null)"
|
||||
echo "wrk: $(wrk --version 2>&1 | head -1 || true)"
|
||||
echo "PIN=$PIN SERVER_CPUS=$SERVER_CPUS WRK_CPUS=$WRK_CPUS"
|
||||
echo "DUR=$DUR REPS=$REPS CONNS=$CONNS THREADS=$THREADS CAUSAL_WRK_DUR=$CAUSAL_WRK_DUR"
|
||||
} > "$OUT/env.txt"
|
||||
say "env recorded -> env.txt"
|
||||
|
||||
# ---- helpers -------------------------------------------------------------
|
||||
run_wrk() { # $1=mode $2=duration_s $3=port
|
||||
if [ "$1" = close ]; then
|
||||
"${pin_wrk[@]}" wrk -t "$THREADS" -c "$CONNS" -d "${2}s" --latency \
|
||||
-H "Connection: close" "http://127.0.0.1:$3/json/7"
|
||||
else
|
||||
"${pin_wrk[@]}" wrk -t "$THREADS" -c "$CONNS" -d "${2}s" --latency \
|
||||
"http://127.0.0.1:$3/json/7"
|
||||
fi
|
||||
}
|
||||
|
||||
wait_port() { # $1=port
|
||||
for _ in $(seq 1 150); do
|
||||
curl -s -o /dev/null "http://127.0.0.1:$1/json/1" && return 0
|
||||
sleep 0.2
|
||||
done
|
||||
return 1
|
||||
}
|
||||
|
||||
verify_mode() { # $1=mode $2=port — record what actually goes over the wire
|
||||
echo "--- single request, mode=$1 ---"
|
||||
if [ "$1" = close ]; then
|
||||
curl -sv --http1.1 -H "Connection: close" -o /dev/null \
|
||||
"http://127.0.0.1:$2/json/7" 2>&1 | grep -iE "^(> |< )(GET|HTTP|connection)" || true
|
||||
else
|
||||
curl -sv --http1.1 -o /dev/null "http://127.0.0.1:$2/json/7" 2>&1 \
|
||||
| grep -iE "^(> |< )(GET|HTTP|connection)" || true
|
||||
fi
|
||||
echo "--- reuse probe (two requests, one curl) ---"
|
||||
curl -sv --http1.1 -o /dev/null -o /dev/null \
|
||||
"http://127.0.0.1:$2/json/1" "http://127.0.0.1:$2/json/2" 2>&1 \
|
||||
| grep -icE "re-us(ed|ing)" || true
|
||||
}
|
||||
|
||||
# ---- plain cells ---------------------------------------------------------
|
||||
run_plain() { # $1=mode $2=port
|
||||
local mode="$1" port="$2" d="$OUT/plain-$1"
|
||||
mkdir -p "$d"
|
||||
say "=== plain / $mode (port $port) ==="
|
||||
URUS_PORT="$port" "${pin_server[@]}" "$PLAIN_BIN" > "$d/server.log" 2>&1 &
|
||||
local spid=$!
|
||||
wait_port "$port" || { say "FATAL: plain server never came up"; cat "$d/server.log"; exit 1; }
|
||||
verify_mode "$mode" "$port" > "$d/mode-verify.txt" 2>&1
|
||||
ss -s > "$d/ss-before.txt" 2>/dev/null || true
|
||||
run_wrk "$mode" 5 "$port" > "$d/warmup.txt" 2>&1
|
||||
for r in $(seq 1 "$REPS"); do
|
||||
run_wrk "$mode" "$DUR" "$port" > "$d/rep$r.txt" 2>&1
|
||||
grep -E "Requests/sec|Latency |requests in|Socket errors|Non-2xx" "$d/rep$r.txt" \
|
||||
| sed "s/^/ [plain-$mode r$r] /" | tee -a "$OUT/matrix.log" || true
|
||||
done
|
||||
ss -s > "$d/ss-after.txt" 2>/dev/null || true
|
||||
awk '{printf "server cpu jiffies (utime+stime): %d\n", $14+$15}' \
|
||||
"/proc/$spid/stat" > "$d/server-cpu.txt" 2>/dev/null || true
|
||||
kill "$spid" 2>/dev/null || true
|
||||
wait "$spid" 2>/dev/null || true
|
||||
}
|
||||
|
||||
# ---- causal cells --------------------------------------------------------
|
||||
run_causal() { # $1=mode $2=port
|
||||
local mode="$1" port="$2" d="$OUT/causal-$1"
|
||||
mkdir -p "$d"
|
||||
say "=== causal / $mode (port $port) ==="
|
||||
URUS_PORT="$port" COZ_OUT="$d/profile.coz" WARMUP_MS=2000 \
|
||||
"${pin_server[@]}" "$CAUSAL_BIN" > "$d/sweep.log" 2>&1 &
|
||||
local spid=$!
|
||||
wait_port "$port" || { say "FATAL: causal server never came up"; cat "$d/sweep.log"; exit 1; }
|
||||
verify_mode "$mode" "$port" > "$d/mode-verify.txt" 2>&1
|
||||
run_wrk "$mode" "$CAUSAL_WRK_DUR" "$port" > "$d/wrk-backdrop.txt" 2>&1 &
|
||||
local wpid=$!
|
||||
local rc=0
|
||||
wait "$spid" || rc=$?
|
||||
kill "$wpid" 2>/dev/null || true
|
||||
wait "$wpid" 2>/dev/null || true
|
||||
say "causal/$mode server exit=$rc (sweep + audit in sweep.log)"
|
||||
if [ "$rc" -ne 0 ]; then say "WARNING: causal/$mode exited nonzero"; fi
|
||||
}
|
||||
|
||||
# ---- matrix --------------------------------------------------------------
|
||||
run_plain ka "$BASE_PORT"
|
||||
run_plain close "$((BASE_PORT + 1))"
|
||||
if [ "${PLAIN_ONLY:-0}" != 1 ]; then
|
||||
run_causal ka "$((BASE_PORT + 2))"
|
||||
run_causal close "$((BASE_PORT + 3))"
|
||||
fi
|
||||
|
||||
# ---- summary -------------------------------------------------------------
|
||||
say "=== SUMMARY ==="
|
||||
python3 - "$OUT" <<'PY' | tee -a "$OUT/matrix.log"
|
||||
import re, sys, pathlib, statistics
|
||||
out = pathlib.Path(sys.argv[1])
|
||||
means = {}
|
||||
for mode in ("ka", "close"):
|
||||
vals = []
|
||||
for rep in sorted((out / f"plain-{mode}").glob("rep*.txt")):
|
||||
m = re.search(r"Requests/sec:\s+([\d.]+)", rep.read_text())
|
||||
if m:
|
||||
vals.append(float(m.group(1)))
|
||||
if vals:
|
||||
means[mode] = statistics.mean(vals)
|
||||
print(f"plain-{mode}: reps={[f'{v:.0f}' for v in vals]} mean={means[mode]:.0f} req/s")
|
||||
if len(means) == 2:
|
||||
r = means["close"] / means["ka"]
|
||||
print(f"close/ka ratio: {r:.3f}", "-> INVERSION PRESENT (close beats ka)" if r > 1.05
|
||||
else "-> no inversion (ka >= close)" if r < 0.95 else "-> parity")
|
||||
for mode in ("ka", "close"):
|
||||
log = out / f"causal-{mode}" / "sweep.log"
|
||||
if log.exists():
|
||||
t = log.read_text()
|
||||
keep = [l for l in t.splitlines()
|
||||
if re.search(r"forgiv|books|audit|balance|total responses", l, re.I)]
|
||||
print(f"-- causal-{mode} audit lines --")
|
||||
for l in keep[:14]:
|
||||
print(" " + l)
|
||||
PY
|
||||
say "=== MATRIX DONE ==="
|
||||
@@ -48,8 +48,8 @@
|
||||
|
||||
use super::*;
|
||||
|
||||
use smarm::gen_server::{self, GenServer, ServerCtx};
|
||||
use smarm::{Down, ServerRef, Watcher};
|
||||
use smarm::gen_server::{self, GenServer, GenServerCtx};
|
||||
use smarm::{Down, GenServerRef, Watcher};
|
||||
|
||||
use std::collections::VecDeque;
|
||||
use std::hash::Hash;
|
||||
@@ -87,7 +87,7 @@ pub trait ChannelSession<P>: Send + 'static {
|
||||
pub(super) struct SessionFactory<P: Encode + Decode + Send + Sync + 'static, K: SessionKey> {
|
||||
inner: Arc<dyn ChannelFactory<P>>,
|
||||
keyfn: fn(&str, &P) -> K,
|
||||
registry: ServerRef<Registry<P, K>>,
|
||||
registry: GenServerRef<Registry<P, K>>,
|
||||
}
|
||||
|
||||
/// The registry's working bounds for a session key.
|
||||
@@ -140,14 +140,14 @@ pub(super) struct Join<P: Send + Sync + 'static, K> {
|
||||
inbound: Receiver<In<P>>,
|
||||
}
|
||||
|
||||
struct Registry<P: Send + Sync + 'static, K> {
|
||||
struct Registry<P: Encode + Decode + Send + Sync + 'static, K: SessionKey> {
|
||||
factory: Arc<dyn ChannelFactory<P>>,
|
||||
cap: usize,
|
||||
ttl: Duration,
|
||||
sessions: HashMap<K, (Pid, Sender<Ctl<P>>)>,
|
||||
/// Reverse index for `handle_down` — the reason `Key: Clone`.
|
||||
pid_key: HashMap<Pid, K>,
|
||||
watcher: Option<Watcher>,
|
||||
watcher: Option<Watcher<Registry<P, K>>>,
|
||||
}
|
||||
|
||||
impl<P: Encode + Decode + Send + Sync + 'static, K: SessionKey> GenServer for Registry<P, K> {
|
||||
@@ -155,8 +155,9 @@ impl<P: Encode + Decode + Send + Sync + 'static, K: SessionKey> GenServer for Re
|
||||
type Reply = ();
|
||||
type Cast = Join<P, K>;
|
||||
type Info = ();
|
||||
type Timer = ();
|
||||
|
||||
fn init(&mut self, ctx: &ServerCtx) {
|
||||
fn init(&mut self, ctx: &GenServerCtx<Self>) {
|
||||
self.watcher = Some(ctx.watcher());
|
||||
}
|
||||
|
||||
|
||||
+2
-4
@@ -158,7 +158,9 @@ impl Body {
|
||||
// Empty `Vec`s are skipped by the pump (a zero-length chunk would
|
||||
// terminate chunked framing early), so they are safe to send but useless.
|
||||
|
||||
#[derive(Default)]
|
||||
pub enum RespBody {
|
||||
#[default]
|
||||
Empty,
|
||||
Bytes(Vec<u8>),
|
||||
Stream(StreamBody),
|
||||
@@ -217,10 +219,6 @@ impl RespBody {
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for RespBody {
|
||||
fn default() -> Self { RespBody::Empty }
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Params — path parameters extracted by the router.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
+170
-60
@@ -19,7 +19,7 @@ use crate::net::OwnedFd;
|
||||
use crate::parser::{self, ParseError};
|
||||
use crate::plug::Pipeline;
|
||||
|
||||
use smarm::ServerRef;
|
||||
use smarm::GenServerRef;
|
||||
|
||||
use std::io::{self, ErrorKind};
|
||||
use std::os::fd::RawFd;
|
||||
@@ -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,
|
||||
@@ -92,7 +115,7 @@ pub fn run_connection(
|
||||
fd: OwnedFd,
|
||||
pipeline: Pipeline,
|
||||
limits: ConnLimits,
|
||||
registry: ServerRef<ConnRegistry>,
|
||||
registry: GenServerRef<ConnRegistry>,
|
||||
) {
|
||||
// The OwnedFd cleans up via Drop on any exit path (panic, error, or
|
||||
// normal close). No explicit close calls below.
|
||||
@@ -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);
|
||||
@@ -157,10 +185,10 @@ pub fn run_connection(
|
||||
// If client sent `Expect: 100-continue`, emit it before reading the
|
||||
// body. RFC 7231 §5.1.1. We don't gate on app logic here; v1 always
|
||||
// accepts.
|
||||
if parsed.expect_100 {
|
||||
if write_all(raw, b"HTTP/1.1 100 Continue\r\n\r\n", Instant::now() + limits.write_timeout).is_err() {
|
||||
return;
|
||||
}
|
||||
if parsed.expect_100
|
||||
&& write_all(raw, b"HTTP/1.1 100 Continue\r\n\r\n", Instant::now() + limits.write_timeout).is_err()
|
||||
{
|
||||
return;
|
||||
}
|
||||
|
||||
// `consumed_past_head`: how many RAW bytes of `buf` past the head
|
||||
@@ -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
|
||||
@@ -209,6 +237,8 @@ pub fn run_connection(
|
||||
// Catch panics at the actor boundary — a panicking handler should
|
||||
// not take down the whole connection silently with no response.
|
||||
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
|
||||
#[cfg(feature = "smarm-causal")]
|
||||
let _g = smarm::causal_site!("pipeline");
|
||||
pipeline.run(conn)
|
||||
}));
|
||||
|
||||
@@ -273,7 +303,11 @@ pub fn run_connection(
|
||||
let keep_alive = keep_alive
|
||||
&& !(version == HttpVersion::Http10 && (is_stream || is_error));
|
||||
|
||||
let head_bytes = parser::serialise_response(&response_conn, keep_alive);
|
||||
let head_bytes = {
|
||||
#[cfg(feature = "smarm-causal")]
|
||||
let _g = smarm::causal_site!("serialize");
|
||||
parser::serialise_response(&response_conn, keep_alive)
|
||||
};
|
||||
if write_all(raw, &head_bytes, Instant::now() + limits.write_timeout).is_err() {
|
||||
return;
|
||||
}
|
||||
@@ -285,10 +319,15 @@ pub fn run_connection(
|
||||
}
|
||||
if !chunked {
|
||||
// EOF delimits the HTTP/1.0 stream body.
|
||||
smarm::progress!("responses");
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
// A full response (head + body, streamed or not) is on the wire:
|
||||
// the unit of useful work for causal profiling (RFC 007).
|
||||
smarm::progress!("responses");
|
||||
|
||||
// ----- 5. Loop or close. -----
|
||||
if !keep_alive {
|
||||
return;
|
||||
@@ -319,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),
|
||||
}
|
||||
@@ -336,34 +375,35 @@ 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 {
|
||||
// Try to parse what we already have. On the first iteration of a
|
||||
// fresh keep-alive cycle, `buf` may already hold the next request.
|
||||
if !buf.is_empty() {
|
||||
match parser::parse_head(buf, limits.max_headers) {
|
||||
Ok(h) => {
|
||||
let deadline = request_deadline
|
||||
.unwrap_or_else(|| Instant::now() + limits.request_timeout);
|
||||
return Ok((h, deadline));
|
||||
}
|
||||
let head = {
|
||||
#[cfg(feature = "smarm-causal")]
|
||||
let _g = smarm::causal_site!("parse");
|
||||
parser::parse_head(buf, limits.max_headers)
|
||||
};
|
||||
match head {
|
||||
Ok(h) => return Ok(h),
|
||||
Err(ParseError::Incomplete) => {} // need more bytes
|
||||
Err(e) => return Err(ReadHeadErr::Parse(e)),
|
||||
}
|
||||
@@ -374,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
|
||||
});
|
||||
@@ -396,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
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -405,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);
|
||||
@@ -417,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),
|
||||
}
|
||||
}
|
||||
@@ -435,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.
|
||||
//
|
||||
@@ -463,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)),
|
||||
}
|
||||
}
|
||||
@@ -494,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 {
|
||||
@@ -507,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
|
||||
@@ -533,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 {
|
||||
@@ -550,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);
|
||||
@@ -634,6 +737,11 @@ fn try_write_once(fd: RawFd, buf: &[u8]) {
|
||||
// wait_writable forever (the write-side twin of slowloris).
|
||||
|
||||
pub(crate) fn write_all(fd: RawFd, mut buf: &[u8], deadline: Instant) -> io::Result<()> {
|
||||
// The whole loop — writability parks included — runs under the
|
||||
// `socket-write` causal site: a park inside a site is exactly what
|
||||
// RFC 007's park-gated resume credit exists to attribute.
|
||||
#[cfg(feature = "smarm-causal")]
|
||||
let _g = smarm::causal_site!("socket-write");
|
||||
while !buf.is_empty() {
|
||||
// Park on writability before each syscall, bounded by the budget.
|
||||
let remaining = deadline.saturating_duration_since(Instant::now());
|
||||
@@ -752,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).
|
||||
_ =>
|
||||
|
||||
@@ -33,7 +33,7 @@
|
||||
//! (roadmap v0.4+), which is why it exists as its own module rather than
|
||||
//! being inlined into `serve`.
|
||||
|
||||
use smarm::{GenServer, Pid, ServerBuilder, ServerRef};
|
||||
use smarm::{GenServer, Pid, GenServerBuilder, GenServerRef};
|
||||
|
||||
use std::collections::HashMap;
|
||||
|
||||
@@ -86,6 +86,7 @@ impl GenServer for ConnRegistry {
|
||||
type Reply = Reply;
|
||||
type Cast = Cast;
|
||||
type Info = ();
|
||||
type Timer = ();
|
||||
|
||||
fn handle_call(&mut self, request: Call) -> Reply {
|
||||
match request {
|
||||
@@ -137,8 +138,8 @@ impl GenServer for ConnRegistry {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn start() -> ServerRef<ConnRegistry> {
|
||||
ServerBuilder::new(ConnRegistry::default()).start()
|
||||
pub fn start() -> GenServerRef<ConnRegistry> {
|
||||
GenServerBuilder::new(ConnRegistry::default()).start()
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -149,13 +150,13 @@ pub fn start() -> ServerRef<ConnRegistry> {
|
||||
/// unwind, and panic unwind alike; the cast is infallible from the
|
||||
/// guard's perspective (a dead registry just returns an ignored Err).
|
||||
pub struct DeregisterGuard {
|
||||
registry: ServerRef<ConnRegistry>,
|
||||
registry: GenServerRef<ConnRegistry>,
|
||||
pid: Pid,
|
||||
make: fn(Pid) -> Cast,
|
||||
}
|
||||
|
||||
impl DeregisterGuard {
|
||||
pub fn new(registry: ServerRef<ConnRegistry>, pid: Pid, make: fn(Pid) -> Cast) -> Self {
|
||||
pub fn new(registry: GenServerRef<ConnRegistry>, pid: Pid, make: fn(Pid) -> Cast) -> Self {
|
||||
Self { registry, pid, make }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
+232
-15
@@ -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,18 +129,30 @@ 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" => {
|
||||
connection_hdr = Some(value.to_ascii_lowercase());
|
||||
}
|
||||
"expect" => {
|
||||
if value.eq_ignore_ascii_case("100-continue") {
|
||||
expect_100 = true;
|
||||
"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;
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
@@ -133,14 +160,56 @@ pub fn parse_head(buf: &[u8], max_headers: usize) -> Result<ParsedHead, ParseErr
|
||||
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:
|
||||
@@ -165,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
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -405,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");
|
||||
|
||||
+7
-6
@@ -58,7 +58,7 @@
|
||||
//! # Why there is no `register(name)` helper (yet)
|
||||
//!
|
||||
//! smarm's registry maps `name → Pid`, but a `Pid` cannot be turned back
|
||||
//! into a `ServerRef` (the ref *is* the inbox sender). A useful named
|
||||
//! into a `GenServerRef` (the ref *is* the inbox sender). A useful named
|
||||
//! lookup therefore needs either smarm support (registry-held senders)
|
||||
//! or a process-global type-erased map here — both against the grain of
|
||||
//! the ratified design. Deferred; pass the handle.
|
||||
@@ -67,8 +67,8 @@ use std::collections::{HashMap, HashSet};
|
||||
use std::fmt;
|
||||
use std::sync::Arc;
|
||||
|
||||
use smarm::gen_server::{self, GenServer, ServerCtx};
|
||||
use smarm::{channel, Down, Pid, Receiver, Sender, ServerRef, Watcher};
|
||||
use smarm::gen_server::{self, GenServer, GenServerCtx};
|
||||
use smarm::{channel, Down, Pid, Receiver, Sender, GenServerRef, Watcher};
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Public handle
|
||||
@@ -92,7 +92,7 @@ impl std::error::Error for PubSubDown {}
|
||||
/// Payloads are broadcast as `Arc<M>`: one allocation per broadcast, not
|
||||
/// per subscriber.
|
||||
pub struct PubSub<M: Send + Sync + 'static> {
|
||||
server: ServerRef<Table<M>>,
|
||||
server: GenServerRef<Table<M>>,
|
||||
}
|
||||
|
||||
impl<M: Send + Sync + 'static> Clone for PubSub<M> {
|
||||
@@ -226,7 +226,7 @@ struct Table<M: Send + Sync + 'static> {
|
||||
/// retires the entry so a reused-slot pid (fresh generation) gets a
|
||||
/// fresh monitor.
|
||||
monitored: HashSet<Pid>,
|
||||
watcher: Option<Watcher>,
|
||||
watcher: Option<Watcher<Table<M>>>,
|
||||
}
|
||||
|
||||
impl<M: Send + Sync + 'static> Table<M> {
|
||||
@@ -240,8 +240,9 @@ impl<M: Send + Sync + 'static> GenServer for Table<M> {
|
||||
type Reply = Reply;
|
||||
type Cast = Cast<M>;
|
||||
type Info = ();
|
||||
type Timer = ();
|
||||
|
||||
fn init(&mut self, ctx: &ServerCtx) {
|
||||
fn init(&mut self, ctx: &GenServerCtx<Self>) {
|
||||
self.watcher = Some(ctx.watcher());
|
||||
}
|
||||
|
||||
|
||||
+19
-6
@@ -135,16 +135,29 @@ impl Plug for Router {
|
||||
// path matches but a same-path-different-method does, return 405.
|
||||
// Otherwise pass through to `next` so outer pipelines can layer a
|
||||
// 404 handler (or skip and let the connection actor emit a default).
|
||||
//
|
||||
// Matching runs under the `router` causal site (RFC 007); the
|
||||
// winning handler and the `next` fall-through run outside it, so
|
||||
// the site measures dispatch, not what it dispatches to.
|
||||
let mut path_seen = false;
|
||||
for route in &self.routes {
|
||||
if let Some(params) = route.pattern.match_path(&conn.path) {
|
||||
if route.method == conn.method {
|
||||
let conn = conn.put_params(params);
|
||||
return route.handler.call(conn, next);
|
||||
let mut hit = None;
|
||||
{
|
||||
#[cfg(feature = "smarm-causal")]
|
||||
let _g = smarm::causal_site!("router");
|
||||
for (i, route) in self.routes.iter().enumerate() {
|
||||
if let Some(params) = route.pattern.match_path(&conn.path) {
|
||||
if route.method == conn.method {
|
||||
hit = Some((i, params));
|
||||
break;
|
||||
}
|
||||
path_seen = true;
|
||||
}
|
||||
path_seen = true;
|
||||
}
|
||||
}
|
||||
if let Some((i, params)) = hit {
|
||||
let conn = conn.put_params(params);
|
||||
return self.routes[i].handler.call(conn, next);
|
||||
}
|
||||
if path_seen {
|
||||
// Path is known, method isn't — RFC 7231 §6.5.5.
|
||||
conn.put_status(405)
|
||||
|
||||
+197
-25
@@ -17,7 +17,7 @@ use crate::conn_registry::{self, Call, Cast, ConnRegistry, Reply};
|
||||
use crate::net::{accept_nonblocking, bind_and_listen, OwnedFd};
|
||||
use crate::plug::Pipeline;
|
||||
|
||||
use smarm::{ChildSpec, OneForOne, Restart, ServerRef, Strategy};
|
||||
use smarm::{ChildSpec, OneForOne, Restart, GenServerRef, Strategy};
|
||||
|
||||
use std::io::{self, ErrorKind};
|
||||
use std::net::{SocketAddr, ToSocketAddrs};
|
||||
@@ -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,
|
||||
@@ -57,8 +71,25 @@ pub struct Config {
|
||||
/// (one per CPU). Set this to a small fixed number in tests so multiple
|
||||
/// concurrent test servers don't oversubscribe the host.
|
||||
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 {
|
||||
pub fn new(addr: SocketAddr) -> Self {
|
||||
let pool = std::thread::available_parallelism()
|
||||
@@ -71,13 +102,18 @@ 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),
|
||||
max_frame_payload: 1024 * 1024,
|
||||
max_message_bytes: 4 * 1024 * 1024,
|
||||
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_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,
|
||||
@@ -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
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -125,7 +237,8 @@ fn listener_loop(
|
||||
listener: Arc<OwnedFd>,
|
||||
pipeline: Pipeline,
|
||||
limits: ConnLimits,
|
||||
registry: ServerRef<ConnRegistry>,
|
||||
conn_stack_reserve: usize,
|
||||
registry: GenServerRef<ConnRegistry>,
|
||||
shutdown: Arc<AtomicBool>,
|
||||
) {
|
||||
let fd = listener.as_raw();
|
||||
@@ -157,7 +270,11 @@ fn listener_loop(
|
||||
let p = pipeline.clone();
|
||||
let l = limits;
|
||||
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 => {
|
||||
// No pending connection. Park until the listener is
|
||||
@@ -305,13 +422,17 @@ pub fn serve_with_shutdown(
|
||||
listener_fds.push(Arc::new(dup));
|
||||
}
|
||||
|
||||
let limits = config.to_conn_limits();
|
||||
let drain_timeout = config.drain_timeout;
|
||||
let limits = config.to_conn_limits();
|
||||
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));
|
||||
@@ -327,27 +448,20 @@ pub fn serve_with_shutdown(
|
||||
let sf = shutdown_flag.clone();
|
||||
sup = sup.child(ChildSpec::new(Restart::Transient, move || {
|
||||
println!("urus: listener {} starting", i);
|
||||
// Named for whereis-style introspection. On a restart the
|
||||
// old binding points at a dead pid; smarm's registry
|
||||
// evicts stale bindings lazily, so re-registering the
|
||||
// same name is fine. Ignore the result — a registry
|
||||
// hiccup must not take the listener down.
|
||||
let _ = smarm::register(
|
||||
format!("urus.listener.{i}"),
|
||||
smarm::self_pid(),
|
||||
);
|
||||
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
|
||||
// faster than that trips the cap and tears the pool down — loud
|
||||
// failure over a zombie server.
|
||||
let sup_h = smarm::spawn(move || sup.run());
|
||||
// Register via the JoinHandle's pid rather than inside the
|
||||
// closure: the binding exists before the supervisor body runs a
|
||||
// single instruction, so an early `whereis("urus.server")` can't
|
||||
// race a None. Result ignored for the same reason as listeners.
|
||||
let _ = smarm::register("urus.server", sup_h.pid());
|
||||
// The old `urus.server` / `urus.listener.{i}` name registrations
|
||||
// are gone with smarm's RFC 014 registry rework: `register` is now
|
||||
// `(Name<M>, Sender<M>)`, self-only — a name is a typed messaging
|
||||
// endpoint, not a pid tag. urus's bindings were introspection-only
|
||||
// with no channel behind them, so they were dropped rather than
|
||||
// faked with a unit channel. A real messageable `urus.server`
|
||||
// name is in the icebox (ROADMAP.md).
|
||||
|
||||
// Block until told to shut down. We poll `try_recv` + `sleep`
|
||||
// rather than parking in `recv`: a smarm `Sender::send` from a
|
||||
@@ -399,7 +513,7 @@ pub fn serve_with_shutdown(
|
||||
}
|
||||
}
|
||||
|
||||
// 5. Our ServerRef drops here. The registry's inbox closes once
|
||||
// 5. Our GenServerRef drops here. The registry's inbox closes once
|
||||
// the last conn's clone drops with it, and the runtime winds
|
||||
// down when the last actor exits.
|
||||
});
|
||||
@@ -429,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};
|
||||
|
||||
|
||||
+187
-30
@@ -89,24 +89,9 @@ fn hello_world() {
|
||||
assert_eq!(http_body(&resp), b"hello urus");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn server_and_listeners_are_registered() {
|
||||
// `whereis` must run inside the runtime (the test itself is a foreign
|
||||
// OS thread with no runtime in its TLS), so probe from a handler —
|
||||
// connection actors live in the runtime by construction.
|
||||
let pipe = Pipeline::new().plug(
|
||||
Router::new().get("/whereis", |c: Conn, _n: Next| {
|
||||
let ok = smarm::whereis("urus.server").is_some()
|
||||
&& smarm::whereis("urus.listener.0").is_some()
|
||||
&& smarm::whereis("urus.listener.1").is_some(); // pool of 2
|
||||
c.put_status(200).put_body(if ok { "registered" } else { "missing" })
|
||||
})
|
||||
);
|
||||
let port = spawn_server(pipe);
|
||||
let resp = send_request(port, b"GET /whereis HTTP/1.1\r\nHost: x\r\nConnection: close\r\n\r\n");
|
||||
assert_eq!(http_status(&resp), 200);
|
||||
assert_eq!(http_body(&resp), b"registered");
|
||||
}
|
||||
// `server_and_listeners_are_registered` was deleted with the RFC 014 port:
|
||||
// urus no longer binds `urus.server` / `urus.listener.{i}` names (see the
|
||||
// note in serve.rs and the icebox entry in ROADMAP.md).
|
||||
|
||||
#[test]
|
||||
fn echo_body() {
|
||||
@@ -452,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();
|
||||
@@ -461,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();
|
||||
@@ -501,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();
|
||||
@@ -526,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();
|
||||
@@ -928,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();
|
||||
@@ -951,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