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

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

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

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

Tests: body_smooth_trickle_evicted_at_stall_timeout (chunked; active
sub-burst trickle evicted at ~stall while the cap is far away) and
bursty_slow_body_survives_stall_gate (fixed CL; real bursts with sub-stall
gaps complete intact). 79 lib + 45 integration green; clippy --lib clean.
2026-08-12 13:38:16 +00:00
Claude (sandbox) 4f06265338 feat(serve): split request read budget into head_timeout + body_timeout
The single request_timeout covered head + body under one wall clock, so a
slow-but-legit body upload (e.g. a trickling cellular IoT client) was
judged by the short head deadline and killed mid-body. Split into:

- head_timeout (default 30s): first byte -> full head parse; the classic
  slowloris surface, kept short.
- body_timeout (default 300s): head parse -> full body; an absolute cap
  sized for slow links, anchored independently once the head has parsed.

read_head no longer returns a shared deadline; run_connection anchors the
body deadline itself. ReadHeadErr::RequestTimeout -> HeadTimeout. Config
and ConnLimits gain body_timeout; request_timeout renamed to head_timeout
(breaking, but this axis is unreleased).

Tests: slow_body_outlives_head_timeout (positive: body survives past the
head clock), fixed_/chunked_body_stall_killed_at_body_timeout (body cap
still bites), slowloris_partial_head_killed_at_head_timeout (head clock
unchanged). 79 lib + 43 integration green; clippy --lib clean.
2026-08-12 13:00:22 +00:00
Claude (sandbox) 1b1ea124c8 feat(serve): plumb Config.max_actors through to smarm init
Config gained a max_actors: Option<usize> (None = smarm's DEFAULT_MAX_ACTORS
of 16_384). serve_with now applies it to the smarm runtime config, so the
per-connection actor slab can be sized to the deployment's peak concurrent
connections. Since each connection is one actor, the slab was the hard cap
on concurrent connections (previously an un-raisable 16_384) regardless of
RAM/fds; slots are ~256 B so raising it is cheap next to per-conn stacks.

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

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

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

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

Verified: 90 lib + 45 integration + 2 doc tests green, examples build,
under default and --all-features.
2026-08-08 11:50:51 +02:00
Claude 34730930e7 feat(bench): E1 scheduler-count sweep orchestrator
Discriminates herd-contention vs stranding for the ka low-c latency
floor: sweeps URUS_SCHED_THREADS {1,2,4,8} x CONNS {4,8} over the plain
matrix (PLAIN_ONLY; close rides along as control). Per-cell audit
asserts the knob reached the server via plain_serve's stderr line.
Summary parser (req/s + p50/75/90/99, unit-normalized) validated
against the 96b40ad5 sweep output.
2026-07-23 15:00:29 +00:00
Claude 72064c7a79 feat(bench): URUS_SCHED_THREADS knob in plain_serve
E1 (herd-vs-stranding discrimination) sweeps smarm scheduler count at
fixed low concurrency. Strict parse — a malformed value panics instead
of silently reverting to default, and the effective value is logged to
stderr so every bench cell's server.log records it (same per-cell
verification discipline as mode-verify).
2026-07-23 14:59:45 +00:00
claude 14be21db15 feat(bench): PLAIN_ONLY knob — skip causal cells for concurrency sweeps
The c=64 matrix run (job 7dbc1ede) showed plain-mode parity but a 2.2x ka
latency penalty; the decisive test is a CONNS sweep, which only needs the
plain cells. PLAIN_ONLY=1 skips the two causal cells; the summary already
tolerates their absence.
2026-07-20 20:43:32 +00:00
claude 182f4fe602 feat(bench): ka/close A/B matrix — plain_serve baseline + box orchestrator
The jar's ka-vs-close throughput inversion was observed on the box but its
run configuration died with the job workspaces, so it can't be re-derived —
this commit makes the comparison a committed, controlled experiment instead
of a lost one-off.

- examples/plain_serve: the load_profile request path with zero causal
  machinery (same route/handler/Config), serving until killed. The plain
  half of the A/B; throughput measured externally by wrk.
- scripts/ka-close-matrix.sh: {ka,close} x {plain,causal} on fresh ports,
  byte-identical wrk args per mode pair except the Connection: close
  header, disjoint-core pinning for server vs wrk, per-cell curl
  verification of the negotiated connection behavior, ss/env capture, and
  a parsed summary with a close/ka ratio verdict. The causal cells double
  as the forgiveness-fix box validation (v13 pending item, pull-forward
  agreed): key on forgiven staying ~ injected x parked-actors in
  close-mode churn, not process-lifetime phantoms.
2026-07-20 07:01:14 +00:00
Claude 17cd4a5ceb feat(causal): load_profile example — external-loadgen causal target
Server-only discovery tool: real hot path, no planted bottleneck, no
in-process loadgen. wrk (or any external generator, pinned to other
cores) drives GET /json/:id; the sweep runs the causal_bench plan over
the five lib sites. Trust gates only (traffic flowed, ledger printed) —
no known-answer verdict, this one asks the question instead.
2026-07-17 19:36:04 +00:00
Claude 288b52d89c feat(causal): webserver bench with a planted, known-answer bottleneck
examples/causal_bench.rs (required-features smarm-causal): the Phase 2
RFC 007 test case — a real urus server replacing the synthetic demo as
validation workload.

Shape: 16 plain-OS-thread clients hammer GET /order/:id over blocking
loopback keep-alive TCP (OS threads can't absorb injected virtual delay
— the established trick, so only server code is delayable). The handler
calls a single store actor (crud's once-cell pattern, serialization
structural) burning 400µs of calibrated *work* under
causal_site!("store") — work-shaped, not timed, or injection reads as a
no-op. 50µs handler render lands under the enclosing pipeline site.

Known answer: store serialized + saturated => ceiling 2500 rps; speedup
p => x1/(1-p): +33% @25, +100% @50. The five lib sites are tens of µs
and parallel across conns => ~0. Verdict mirrors the demo: store @25 >
+15%, others < +10%, SKIPPED under 4 cores; magnitudes trusted to ~±15%
per the smarm-side validation notes, ranking is the hard check.

1-core sandbox smoke (verdict SKIPPED but numbers indicative): store
+26.2% @25 / +74.0% @50, all other sites within ±2%. Unlike the demo
this workload keeps its bottleneck structure on one core — everything
but the store is IO-parked — and the theory shortfall there is plain
core contention (render/parse share the store's CPU), which the 24-core
sweep should close.
2026-07-13 08:59:49 +00:00
Claude ecaddc579a feat(causal): instrument the request hot path — five sites and a progress point
RFC 007 causal-site coverage for the Phase 2 bench (and any downstream
profiling): parse (head parsing attempts in read_head), router (dispatch
matching only — the winning handler and the next fall-through run outside
the guard, so the site measures dispatch, not what it dispatches to;
call() restructured match-then-dispatch for that, semantics preserved
incl. first-path+method-match-wins and the 405 two-pass), pipeline (the
plug chain inside catch_unwind; nested sites like router take over
attribution during their span, so this reads as chain overhead + handler
code outside inner sites), serialize (response head+bytes-body
serialisation; chunked stream framing is not covered — it interleaves
with writes in pump_stream and the bench doesn't stream), and
socket-write (all of write_all, writability parks included: a park
inside a site is exactly what the park-gated resume credit attributes).

progress!("responses") fires once per response fully on the wire, at
both completion points (the common path and the HTTP/1.0 EOF-stream
early return).

Site guards are #[cfg(feature = "smarm-causal")]-gated: the macros are
no-ops featureless, but binding the unit expansion would trip clippy's
let_unit_value. progress! is bare — its featureless expansion is an
empty block, free and lint-clean. Featureless build byte-behavior is
unchanged; both configs clippy-clean, full suite green both ways.
2026-07-13 08:30:53 +00:00
Claude 072ee126f9 feat: smarm-causal passthrough feature
Mirrors the smarm-trace precedent: urus itself gains no causal code
paths yet — this just lets downstream binaries (the Phase 2 causal
bench lives in examples/) flip smarm's instrumentation on through the
urus dep. Whole suite is green with the feature enabled: the causal
runtime is behaviorally inert while no experiment is active.
2026-07-13 08:25:49 +00:00
Claude 0f824635d1 chore(clippy): appease 1.97 lints — derive Default, collapse ifs, drop needless borrow
Pre-existing, surfaced by the toolchain bump (rust-version is 1.95; the
sandbox gates with stable 1.97). All four are cargo clippy --fix output
with the mechanical-collapse indentation hand-tidied to house style;
the parser change is semantics-preserving (a non-100-continue Expect
value now falls to the _ arm instead of an empty if — headers.append
still runs after the match either way). No fmt pass: rustfmt would
clobber the aligned-assignment style, so only the touched lines moved.
2026-07-13 08:24:31 +00:00
Claude 078072b527 port(smarm): track HEAD efbc254 — RFC 014/015 API sync
- gen_server rename (RFC 015, 3e31606): ServerRef/ServerCtx/ServerBuilder
  -> GenServerRef/GenServerCtx/GenServerBuilder across conn_actor,
  conn_registry, serve, pubsub, channels::session.
- type Timer = () on the three GenServer impls (RFC 015, 57eadb5);
  handle_timer/handle_idle/tick_every stay defaulted — opt-in later.
- Watcher and GenServerCtx are now generic over the server type: watcher
  fields typed Watcher<Table<M>> / Watcher<Registry<P, K>>; Registry's
  struct bounds strengthened to its GenServer impl bounds
  (P: Encode + Decode + Send + Sync + 'static, K: SessionKey) so the
  Watcher field's G: GenServer bound is satisfiable at the declaration.
- RFC 014 (a866e34) registry: register is (Name<M>, Sender<M>), self-only
  — a name is a typed messaging endpoint, not a pid tag. The
  introspection-only urus.server / urus.listener.{i} bindings are dropped
  rather than faked with unit channels; the whereis integration test is
  deleted; a proper messageable-name design is icebox'd in ROADMAP.md.
- Audit vs smarm 6c2b7e9 (queued messages dropped when Receiver drops):
  pubsub's prune-on-send-failure retain still holds — send Errs once
  receiver_alive is false, so a dropped rx prunes on next broadcast,
  exactly what subscriber_count's doc already promised. Freeing stranded
  Arc<M> broadcasts is strictly good. No change needed.
- The 7x E0283 in channels/mod.rs were cascade fallout of the generics
  changes; dissolved with the port, as discovery predicted.

Suite: 90 lib + 45 integration + 2 doc, green under default, smarm-trace,
phoenix, and all-features. 3 pre-existing clippy lints (conn.rs,
parser.rs, conn_actor.rs; clippy 1.97 strictness) deferred to a follow-up
chore commit to keep this diff pure.
2026-07-13 08:22:40 +00:00
21 changed files with 1672 additions and 165 deletions
+16 -3
View File
@@ -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"]
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"]
+21
View File
@@ -0,0 +1,21 @@
MIT License
Copyright (c) 2026 Mark Kalsbeek
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
+6
View File
@@ -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.
+322
View File
@@ -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);
}
}
+106
View File
@@ -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");
}
+48
View File
@@ -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");
}
+48
View File
@@ -0,0 +1,48 @@
//! Serve with an optional TOML config overlay.
//!
//! urus is a library and never presumes a config path or reads the
//! environment for one — the binary decides where the file lives and hands
//! the text to `Config::with_toml_str`. Here that's a `--config PATH` flag;
//! with no flag, the compiled defaults are used unchanged.
//!
//! Requires the `config-file` feature:
//! cargo run --example serve_toml --features config-file -- --config urus.toml
//!
//! Example urus.toml (all keys optional, sparse override; seconds):
//! head_timeout_secs = 15
//! body_timeout_secs = 300
//! body_burst_bytes = 4096
//! body_stall_timeout_secs = 20
use std::net::SocketAddr;
use urus::{serve_with, Config, Conn, Next, Pipeline, Router};
fn json_id(conn: Conn, _next: Next) -> Conn {
let id: u64 = conn.params.get("id").and_then(|s| s.parse().ok()).unwrap_or(0);
conn.put_status(200)
.put_header("content-type", "application/json")
.put_body(format!("{{\"id\":{id}}}"))
}
fn main() {
let addr: SocketAddr = "0.0.0.0:8080".parse().expect("addr");
let mut cfg = Config::new(addr);
// Minimal flag scan: `--config PATH`. No presumed default location.
let mut args = std::env::args().skip(1);
while let Some(arg) = args.next() {
if arg == "--config" {
let path = args.next().expect("--config needs a PATH");
let toml = std::fs::read_to_string(&path)
.unwrap_or_else(|e| panic!("reading config {path}: {e}"));
cfg = cfg
.with_toml_str(&toml)
.unwrap_or_else(|e| panic!("invalid config {path}: {e}"));
eprintln!("serve_toml: loaded config from {path}");
}
}
let pipe = Pipeline::new().plug(Router::new().get("/json/:id", json_id));
serve_with(cfg, pipe).expect("serve");
}
+1 -1
View File
@@ -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) {
+89
View File
@@ -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 ==="
+181
View File
@@ -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 ==="
+7 -6
View File
@@ -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
View File
@@ -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.
// ---------------------------------------------------------------------------
+169 -59
View File
@@ -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,11 +185,11 @@ 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() {
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
// this request's body occupied — for chunked bodies that is framing
@@ -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).
_ =>
+6 -5
View File
@@ -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 }
}
}
+2
View File
@@ -44,3 +44,5 @@ pub use ws::{Message, WsClosed, WsHandler, WsSender};
pub use serve::{
serve, serve_with, serve_with_shutdown, shutdown_handle, Config, Handle, ShutdownSignal,
};
#[cfg(feature = "config-file")]
pub use serve::ConfigError;
+231 -14
View File
@@ -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,33 +129,87 @@ 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" if value.eq_ignore_ascii_case("100-continue") => {
expect_100 = true;
}
"host" => {
// Presence/uniqueness enforced post-loop; validity here.
host_count += 1;
if !valid_host(value) {
host_ok = false;
}
}
_ => {}
}
headers.append(&name_lower, value.to_string());
}
if chunked {
// Transfer-Encoding is an HTTP/1.1 mechanism; a 1.0 request
// carrying it is malformed. And a request carrying BOTH a
// Content-Length and TE: chunked is the classic request-smuggling
// ambiguity — RFC 7230 §3.3.3 lets a server reject it, and we do.
if version == HttpVersion::Http10 || content_length.is_some() {
// Host (RFC 9112 §3.2): an HTTP/1.1 request MUST carry exactly one valid
// Host; a missing, duplicate, or malformed Host is a 400. HTTP/1.0 may
// omit Host, but a duplicate or invalid one is still rejected on any
// version (ambiguous / malformed authority).
if host_count > 1 || !host_ok {
return Err(ParseError::Malformed);
}
if version == HttpVersion::Http11 && host_count == 0 {
return Err(ParseError::Malformed);
}
// Content-Length (RFC 9112 §6.3): more than one Content-Length is an
// unrecoverable framing ambiguity (CL.CL request smuggling). We are
// strict — reject any duplicate, not only differing values.
if cl_count > 1 {
return Err(ParseError::BadContentLength);
}
// Transfer-Encoding (RFC 9112 §6.1/§6.3). TE is a 1.1 mechanism and
// overrides Content-Length; only `chunked` is implemented here.
if te_present {
// TE on HTTP/1.0 is malformed (no 1.0 chunked).
if version == HttpVersion::Http10 {
return Err(ParseError::Malformed);
}
// TE together with Content-Length is the classic smuggling
// ambiguity; TE overrides CL and we reject rather than forward.
if content_length.is_some() {
return Err(ParseError::Malformed);
}
// A Transfer-Encoding header that carries no coding frames nothing.
if te_codings.is_empty() {
return Err(ParseError::Malformed);
}
let last_is_chunked = te_codings.last().map(String::as_str) == Some("chunked");
let has_chunked = te_codings.iter().any(|c| c == "chunked");
if has_chunked && !last_is_chunked {
// chunked present but not final: body length isn't reliably
// determinable -> 400.
return Err(ParseError::Malformed);
}
if te_codings.iter().any(|c| c != "chunked") {
// Some coding we don't implement (chunked is the only decodable
// one). Whether or not chunked is final, we can't apply it -> 501.
return Err(ParseError::UnknownTransferCoding);
}
// Sole, final `chunked`: the connection actor decodes the body.
chunked = true;
}
// Keep-alive logic, RFC 7230 §6.3:
@@ -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
View File
@@ -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());
}
+16 -3
View File
@@ -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 {
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 {
let conn = conn.put_params(params);
return route.handler.call(conn, next);
hit = Some((i, params));
break;
}
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)
+195 -23
View File
@@ -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
@@ -306,12 +423,16 @@ pub fn serve_with_shutdown(
}
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
View File
@@ -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
View File
@@ -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)
// ---------------------------------------------------------------------------