The v0.3 shape from the spec: an app owns its runtime and root supervisor
and places urus in it as one ordered child among its own.
your root sup
└── ChildSpec(Permanent, urus::endpoint(cfg, pipeline)?) <- Endpoint
└── listener_sup OneForOne over N listeners
└── plain connection actors
- src/conn_registry.rs -> src/endpoint.rs. The registry gains the listener
pool it registers for and becomes the Endpoint gen_server; ConnRegistry
-> Endpoint. It runs inline as the ChildSpec's actor
(NamedGenServerBuilder::run), so supervisor shutdown arrives as
handle_shutdown and a restart re-runs init on the same still-open fds.
- Endpoint spawns its OWN listener sup in init rather than being its
sibling: smarm's supervisor start order is not start *readiness* (spawn
is fire-and-forget), so a sibling listener could whereis the name before
the registry actor ran. Registrar-spawns-consumers makes it program
order inside one init. Gap filed in smarm ROADMAP (readiness ack);
making spawn blocking would only shrink the window, not close it —
'has begun executing' is not 'has bound its name'.
- Listener sup is monitored: death outside shutdown = panic (loud, the
app's supervisor decides) instead of a zombie on a dead port. Death
during shutdown is the 'no new connections' barrier.
- DELETED: the shutdown AtomicBool, LISTENER_TICK (250ms wake per listener
per tick, now an untimed wait_readable park), SHUTDOWN_POLL (100ms root
poll — the root parks on the signal channel now), Restart::Transient
(listeners are Permanent: they only exit by supervisor action, so a
self-exit always means breakage). Verified against current smarm:
request_stop unwinds an untimed wait_readable park, is no longer lossy
against a QUEUED actor, and supervisor shutdown joins in ~200us.
- Config: scheduler_threads/max_actors removed (runtime knobs an
endpoint-as-child cannot honour) -> serve_with(cfg, smarm::Config, pipe)
and serve_with_shutdown(cfg, smarm::Config, pipe, signal). Added
Config.name (default 'urus'): the endpoint's registry name, unique per
endpoint, and the introspection handle via endpoint::whereis(name).
- serve* keep their meaning as the batteries-included path: they build a
one-child tree around endpoint() with Shutdown::Infinity. Handle stays
(a serve* caller has no RuntimeHandle to reach for) and now backs a real
park instead of a poll.
Tests: 4 drain tests ported onto a real endpoint (bound socket, supervised
child, request_shutdown driven); new integration test boots an app tree
with an ordered sibling and asserts serve-then-drain, reverse-order
teardown and a closed port. 106 lib + 50 integration green.
872 lines
36 KiB
Rust
872 lines
36 KiB
Rust
//! The connection actor — one per accepted TCP connection.
|
|
//!
|
|
//! Runs the HTTP/1.1 request loop:
|
|
//!
|
|
//! loop {
|
|
//! read bytes → parse → build Conn
|
|
//! pipeline.run(conn) // inline; no spawn
|
|
//! write response
|
|
//! if !keep_alive { break }
|
|
//! }
|
|
//!
|
|
//! Everything in here happens in one smarm green thread. The actor parks on
|
|
//! `wait_readable` between bytes and `wait_writable` during slow writes;
|
|
//! during those parks, other connection actors progress freely.
|
|
|
|
use crate::conn::{Body, Conn, HttpVersion, RespBody, StreamBody};
|
|
use crate::endpoint::{Cast, DeregisterGuard, Endpoint};
|
|
use crate::net::OwnedFd;
|
|
use crate::parser::{self, ParseError};
|
|
use crate::plug::Pipeline;
|
|
|
|
use smarm::GenServerRef;
|
|
|
|
use std::io::{self, ErrorKind};
|
|
use std::os::fd::RawFd;
|
|
use std::time::{Duration, Instant};
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Limits
|
|
// ---------------------------------------------------------------------------
|
|
|
|
/// Per-connection settings the connection actor needs to honour.
|
|
#[derive(Clone, Copy, Debug)]
|
|
pub struct ConnLimits {
|
|
pub max_headers: usize,
|
|
pub initial_read_buf: usize,
|
|
/// Hard cap on the request head to bound buffer growth. 64 KiB is
|
|
/// well over Apache's 8 KiB default; protects against pathological
|
|
/// clients streaming headers forever.
|
|
pub max_head_bytes: usize,
|
|
/// Hard cap on Content-Length we'll accept. 16 MiB is enough for a CRUD
|
|
/// example; configurable in `Config`.
|
|
pub max_body_bytes: usize,
|
|
/// Idle budget between requests: how long we'll park waiting for the
|
|
/// FIRST byte of a request (including the first request on a fresh
|
|
/// connection). Expiry closes the connection silently — nothing is
|
|
/// owed to a client that isn't talking.
|
|
pub keep_alive_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
|
|
/// socket buffer fills and a write stalls past the budget.
|
|
pub write_timeout: Duration,
|
|
/// WebSocket: hard cap on a single frame's payload, enforced from
|
|
/// the frame header BEFORE the payload is buffered. Violation closes
|
|
/// with 1009.
|
|
pub max_frame_payload: usize,
|
|
/// WebSocket: hard cap on a complete (reassembled) message; spans
|
|
/// fragments. Violation closes with 1009.
|
|
pub max_message_bytes: usize,
|
|
}
|
|
|
|
impl Default for ConnLimits {
|
|
fn default() -> Self {
|
|
Self {
|
|
max_headers: 64,
|
|
initial_read_buf: 8 * 1024,
|
|
max_head_bytes: 64 * 1024,
|
|
max_body_bytes: 16 * 1024 * 1024,
|
|
keep_alive_timeout: Duration::from_secs(60),
|
|
head_timeout: Duration::from_secs(30),
|
|
body_timeout: Duration::from_secs(300),
|
|
body_burst_bytes: 4 * 1024,
|
|
body_stall_timeout: Duration::from_secs(20),
|
|
write_timeout: Duration::from_secs(30),
|
|
max_frame_payload: 1024 * 1024,
|
|
max_message_bytes: 4 * 1024 * 1024,
|
|
}
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// run_connection — entry point spawned by the listener actor.
|
|
// ---------------------------------------------------------------------------
|
|
|
|
pub fn run_connection(
|
|
fd: OwnedFd,
|
|
pipeline: Pipeline,
|
|
limits: ConnLimits,
|
|
registry: GenServerRef<Endpoint>,
|
|
) {
|
|
// The OwnedFd cleans up via Drop on any exit path (panic, error, or
|
|
// normal close). No explicit close calls below.
|
|
let raw = fd.as_raw();
|
|
let mut buf: Vec<u8> = Vec::with_capacity(limits.initial_read_buf);
|
|
|
|
// Self-register (initially idle: no request head parsed yet) and arm
|
|
// the deregistration guard. Both casts come from this actor, so
|
|
// Started always precedes Ended in the registry's inbox — see
|
|
// endpoint module docs for why the listener must not do this.
|
|
let me = smarm::self_pid();
|
|
let _ = registry.cast(Cast::ConnStarted(me));
|
|
let _guard = DeregisterGuard::new(registry.clone(), me, Cast::ConnEnded);
|
|
|
|
loop {
|
|
// ----- 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 = match read_head(raw, &mut buf, &limits) {
|
|
Ok(p) => p,
|
|
Err(ReadHeadErr::ClientClosed) => {
|
|
// Clean EOF between requests (or before any request). Normal.
|
|
return;
|
|
}
|
|
Err(ReadHeadErr::IdleTimeout) => {
|
|
// keep_alive_timeout expired waiting for the first byte of
|
|
// a request. Nothing is owed; close silently.
|
|
return;
|
|
}
|
|
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.
|
|
try_write_once(raw, b"HTTP/1.1 408 Request Timeout\r\ncontent-length: 0\r\nconnection: close\r\n\r\n");
|
|
return;
|
|
}
|
|
Err(ReadHeadErr::Io(_)) => {
|
|
// Network error. Best-effort close; we're done.
|
|
return;
|
|
}
|
|
Err(ReadHeadErr::Parse(e)) => {
|
|
emit_error_response(raw, &e, Instant::now() + limits.write_timeout);
|
|
return;
|
|
}
|
|
};
|
|
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);
|
|
if body_len > limits.max_body_bytes {
|
|
let _ = write_all(
|
|
raw,
|
|
b"HTTP/1.1 413 Payload Too Large\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
|
|
Instant::now() + limits.write_timeout,
|
|
);
|
|
return;
|
|
}
|
|
|
|
// 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
|
|
&& 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
|
|
// included, NOT the decoded length. The keep-alive drain at the
|
|
// 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, body_deadline) {
|
|
Ok(ok) => ok,
|
|
Err(ChunkedBodyErr::TooLarge) => {
|
|
let _ = write_all(
|
|
raw,
|
|
b"HTTP/1.1 413 Payload Too Large\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
|
|
Instant::now() + limits.write_timeout,
|
|
);
|
|
return;
|
|
}
|
|
Err(ChunkedBodyErr::Malformed) => {
|
|
let _ = write_all(
|
|
raw,
|
|
b"HTTP/1.1 400 Bad Request\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
|
|
Instant::now() + limits.write_timeout,
|
|
);
|
|
return;
|
|
}
|
|
// Timeout mid-body (and any other io error) -> just close.
|
|
Err(ChunkedBodyErr::Io(_)) => return,
|
|
}
|
|
} else {
|
|
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
|
|
// gone.
|
|
Err(_) => return,
|
|
}
|
|
};
|
|
|
|
let keep_alive = parsed.keep_alive;
|
|
let version = parsed.version;
|
|
let head_len = parsed.head_len;
|
|
let conn = parser::build_conn(parsed, Body::from_bytes(body));
|
|
|
|
// ----- 3. Run the pipeline. -----
|
|
// 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)
|
|
}));
|
|
|
|
let mut response_conn = match result {
|
|
Ok(c) => c,
|
|
Err(_) => {
|
|
// Distinguish a genuine handler panic from smarm's stop
|
|
// sentinel, which is also a panic payload and which this
|
|
// catch_unwind would otherwise swallow — turning a
|
|
// graceful stop into a 500-and-keep-running. The stop
|
|
// flag is persistent (not consumed by raising the
|
|
// sentinel), so if we were stopped this re-raises it
|
|
// here, outside the catch, and we unwind properly (the
|
|
// fd and registry guards clean up).
|
|
smarm::preempt::check_cancelled();
|
|
|
|
// Compose a 500 manually; the original Conn was moved into
|
|
// the closure.
|
|
let mut c = Conn::new();
|
|
c.version = version;
|
|
c.put_status(500).put_header("content-length", "0")
|
|
}
|
|
};
|
|
|
|
// If no plug touched status, that's a configuration error (no router
|
|
// matched, no default handler). Emit 404.
|
|
if response_conn.status.is_none() {
|
|
response_conn = response_conn.put_status(404)
|
|
.put_body(RespBody::Empty);
|
|
}
|
|
|
|
// ----- 3.5 WebSocket upgrade. -----
|
|
// An accepted handshake (payload + 101) ends HTTP on this socket:
|
|
// write the 101 head, hand the fd to the duplex loop with the
|
|
// boxed handler and any bytes already read past this request (a
|
|
// client may pipeline its first frame behind the handshake — those
|
|
// bytes are ws bytes now). The registry entry stays Busy for the
|
|
// whole ws lifetime: an open WebSocket is in-flight work, not
|
|
// reapable idle HTTP; graceful shutdown force-stops it out of the
|
|
// select park at the drain deadline. The status check is
|
|
// defensive: a post-handler plug that clobbered the 101 forfeits
|
|
// the upgrade and falls through to plain HTTP.
|
|
if response_conn.upgrade.is_some() && response_conn.status == Some(101) {
|
|
let head = parser::serialise_response(&response_conn, true);
|
|
if write_all(raw, &head, Instant::now() + limits.write_timeout).is_err() {
|
|
return;
|
|
}
|
|
let upgrade = response_conn.upgrade.take().expect("checked above");
|
|
buf.drain(..head_len + consumed_past_head);
|
|
crate::ws::duplex::run_duplex(raw, buf, upgrade.handler, &limits);
|
|
return;
|
|
}
|
|
|
|
// ----- 4. Write the response. -----
|
|
// A Stream body on HTTP/1.0 has no chunked framing: the body is
|
|
// delimited by EOF, so keep-alive is forced off for this response
|
|
// (and `connection: close` goes on the wire). Error statuses on
|
|
// 1.0 also force close: we don't advertise keep-alive on a
|
|
// 4xx/5xx to a protocol generation where reuse is opt-in.
|
|
let is_stream = matches!(response_conn.resp_body, RespBody::Stream(_));
|
|
let is_error = response_conn.status.unwrap_or(200) >= 400;
|
|
let keep_alive = keep_alive
|
|
&& !(version == HttpVersion::Http10 && (is_stream || is_error));
|
|
|
|
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;
|
|
}
|
|
|
|
if let RespBody::Stream(stream) = response_conn.resp_body {
|
|
let chunked = version == HttpVersion::Http11;
|
|
if pump_stream(raw, stream, chunked, limits.write_timeout).is_err() {
|
|
return;
|
|
}
|
|
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;
|
|
}
|
|
|
|
// Response is on the wire; nothing is owed. Going idle here makes
|
|
// us stoppable by a draining registry while we park for the next
|
|
// keep-alive request. (If the next request is already pipelined in
|
|
// `buf`, the very next read_head parses it without parking and we
|
|
// go Busy again — a draining registry's request_stop may still
|
|
// catch us, which is acceptable: drain means no new work.)
|
|
let _ = registry.cast(Cast::ConnIdle(me));
|
|
|
|
// Drop the request bytes (head + raw body framing) from `buf`;
|
|
// anything past them is the start of the next pipelined request.
|
|
let consumed = head_len + consumed_past_head;
|
|
buf.drain(..consumed);
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// read_head
|
|
// ---------------------------------------------------------------------------
|
|
|
|
#[allow(dead_code)] // io::Error is captured for future logging
|
|
enum ReadHeadErr {
|
|
ClientClosed,
|
|
/// keep_alive_timeout expired while waiting for the first byte of a
|
|
/// request. Close silently.
|
|
IdleTimeout,
|
|
/// head_timeout expired after the request had started arriving but
|
|
/// before the head finished parsing. Best-effort 408.
|
|
HeadTimeout,
|
|
Io(io::Error),
|
|
Parse(ParseError),
|
|
}
|
|
|
|
/// Read until `parse_head` succeeds or fails definitively. `buf` may already
|
|
/// contain leftover bytes from a previous keep-alive cycle; we try to parse
|
|
/// those before reading more from the socket.
|
|
///
|
|
/// Two wall-clock budgets govern the waits (each wait uses whichever budget
|
|
/// is currently active):
|
|
///
|
|
/// - 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 *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, ReadHeadErr> {
|
|
let entry = Instant::now();
|
|
let idle_deadline = entry + limits.keep_alive_timeout;
|
|
// Pipelined leftovers count as a started request.
|
|
let mut head_deadline: Option<Instant> = if buf.is_empty() {
|
|
None
|
|
} else {
|
|
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() {
|
|
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)),
|
|
}
|
|
}
|
|
|
|
if buf.len() >= limits.max_head_bytes {
|
|
return Err(ReadHeadErr::Parse(ParseError::TooManyHeaders));
|
|
}
|
|
|
|
// Read more, bounded by whichever budget is active.
|
|
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 head_deadline.is_none() {
|
|
// First byte(s) of this request: the head clock
|
|
// starts now.
|
|
head_deadline = Some(Instant::now() + limits.head_timeout);
|
|
}
|
|
}
|
|
Err(e) if e.kind() == ErrorKind::TimedOut => {
|
|
return Err(if head_deadline.is_some() {
|
|
ReadHeadErr::HeadTimeout
|
|
} else {
|
|
ReadHeadErr::IdleTimeout
|
|
});
|
|
}
|
|
Err(e) => return Err(ReadHeadErr::Io(e)),
|
|
}
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// 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
|
|
// ---------------------------------------------------------------------------
|
|
|
|
fn read_body(
|
|
fd: RawFd,
|
|
buf: &mut Vec<u8>,
|
|
head_len: usize,
|
|
body_len: usize,
|
|
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);
|
|
let need = body_len.saturating_sub(already);
|
|
|
|
if need == 0 {
|
|
// We have the full body in `buf` already. Extract a copy; `buf` is
|
|
// drained later in the connection loop.
|
|
return Ok(buf[head_len..head_len + body_len].to_vec());
|
|
}
|
|
|
|
// 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, gate.deadline()) {
|
|
Ok(0) => return Err(io::Error::new(ErrorKind::UnexpectedEof, "client closed during body")),
|
|
Ok(n) => {
|
|
total_read += n;
|
|
gate.record(n, Instant::now());
|
|
}
|
|
Err(e) => return Err(e),
|
|
}
|
|
}
|
|
Ok(buf[head_len..head_len + body_len].to_vec())
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// read_chunked_body — incremental chunked transfer-decoding (request side).
|
|
// ---------------------------------------------------------------------------
|
|
//
|
|
// Decodes `Transfer-Encoding: chunked` from `buf[head_len..]`, reading more
|
|
// 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.
|
|
//
|
|
// Bounds: the DECODED size is capped at max_body_bytes (-> TooLarge/413);
|
|
// a single size line (incl. chunk extensions, which are ignored) is capped
|
|
// at MAX_SIZE_LINE and the trailer section at MAX_TRAILER_BYTES (->
|
|
// Malformed/400) so framing spam can't grow `buf` unboundedly. Trailers
|
|
// are consumed and discarded — nothing in the pipeline wants them yet.
|
|
|
|
const MAX_SIZE_LINE: usize = 128;
|
|
const MAX_TRAILER_BYTES: usize = 8 * 1024;
|
|
|
|
#[allow(dead_code)] // io::Error is captured for future logging
|
|
enum ChunkedBodyErr {
|
|
Io(io::Error),
|
|
Malformed,
|
|
TooLarge,
|
|
}
|
|
|
|
fn read_chunked_body(
|
|
fd: RawFd,
|
|
buf: &mut Vec<u8>,
|
|
head_len: usize,
|
|
limits: &ConnLimits,
|
|
deadline: Instant,
|
|
) -> Result<(Vec<u8>, usize), ChunkedBodyErr> {
|
|
// 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,
|
|
gate: &mut BodyStallGate,
|
|
) -> Result<(), ChunkedBodyErr> {
|
|
while buf.len() < until {
|
|
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(n) => gate.record(n, Instant::now()),
|
|
Err(e) => return Err(ChunkedBodyErr::Io(e)),
|
|
}
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
// Find "\r\n" in buf[from..], reading more as needed; the line may be
|
|
// at most `max_line` bytes (terminator excluded). Returns the index of
|
|
// the '\r'.
|
|
fn find_crlf(
|
|
fd: RawFd,
|
|
buf: &mut Vec<u8>,
|
|
from: usize,
|
|
max_line: usize,
|
|
gate: &mut BodyStallGate,
|
|
) -> Result<usize, ChunkedBodyErr> {
|
|
let mut scan = from;
|
|
loop {
|
|
while scan + 1 < buf.len() {
|
|
if buf[scan] == b'\r' && buf[scan + 1] == b'\n' {
|
|
return Ok(scan);
|
|
}
|
|
scan += 1;
|
|
if scan - from > max_line {
|
|
return Err(ChunkedBodyErr::Malformed);
|
|
}
|
|
}
|
|
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, &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
|
|
None => line,
|
|
};
|
|
let size_str = std::str::from_utf8(size_str)
|
|
.map_err(|_| ChunkedBodyErr::Malformed)?
|
|
.trim();
|
|
let size = usize::from_str_radix(size_str, 16)
|
|
.map_err(|_| ChunkedBodyErr::Malformed)?;
|
|
pos = line_end + 2;
|
|
|
|
if size == 0 {
|
|
// ----- 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), &mut gate)?;
|
|
let empty = t_end == pos;
|
|
pos = t_end + 2;
|
|
if empty {
|
|
return Ok((decoded, pos - head_len));
|
|
}
|
|
if pos - trailer_start > MAX_TRAILER_BYTES {
|
|
return Err(ChunkedBodyErr::Malformed);
|
|
}
|
|
}
|
|
}
|
|
|
|
if decoded.len() + size > limits.max_body_bytes {
|
|
return Err(ChunkedBodyErr::TooLarge);
|
|
}
|
|
|
|
// ----- chunk payload + trailing CRLF -----
|
|
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);
|
|
}
|
|
pos += size + 2;
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// read_some — single epoll-park + read loop, bounded by a deadline.
|
|
// ---------------------------------------------------------------------------
|
|
//
|
|
// Appends what it reads onto `buf`. Returns bytes read, 0 for EOF,
|
|
// `ErrorKind::TimedOut` when `deadline` passes before the fd turns
|
|
// readable, or the last io error.
|
|
|
|
pub(crate) fn read_some(
|
|
fd: RawFd,
|
|
buf: &mut Vec<u8>,
|
|
chunk: usize,
|
|
deadline: Instant,
|
|
) -> io::Result<usize> {
|
|
// Loop to absorb EAGAIN: a readable wakeup followed by EAGAIN is
|
|
// possible (signal race, etc). Re-park and retry rather than returning
|
|
// 0 (which would be confused with EOF by callers). The deadline is an
|
|
// Instant, so spurious wakes don't reset the budget.
|
|
loop {
|
|
let remaining = deadline.saturating_duration_since(Instant::now());
|
|
if remaining.is_zero() {
|
|
return Err(io::Error::new(ErrorKind::TimedOut, "read deadline elapsed"));
|
|
}
|
|
if !smarm::wait_readable_timeout(fd, remaining)? {
|
|
return Err(io::Error::new(ErrorKind::TimedOut, "read deadline elapsed"));
|
|
}
|
|
|
|
let start = buf.len();
|
|
buf.resize(start + chunk, 0);
|
|
|
|
let n = unsafe {
|
|
libc::read(fd, buf.as_mut_ptr().add(start) as *mut _, chunk)
|
|
};
|
|
|
|
if n < 0 {
|
|
let err = io::Error::last_os_error();
|
|
buf.truncate(start);
|
|
if err.kind() == ErrorKind::WouldBlock || err.kind() == ErrorKind::Interrupted {
|
|
continue;
|
|
}
|
|
return Err(err);
|
|
}
|
|
|
|
let n = n as usize;
|
|
buf.truncate(start + n);
|
|
return Ok(n); // n == 0 here is real EOF
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// try_write_once — single non-parking write attempt, result ignored.
|
|
// ---------------------------------------------------------------------------
|
|
//
|
|
// For best-effort farewells (the 408) to clients we've decided to drop:
|
|
// one non-blocking write syscall, no wait_writable park. A client that
|
|
// stalls its read side must not be able to keep this actor alive past its
|
|
// own timeout. The socket send buffer almost always has room for a
|
|
// one-liner, so in practice the 408 lands.
|
|
|
|
fn try_write_once(fd: RawFd, buf: &[u8]) {
|
|
unsafe {
|
|
let _ = libc::write(fd, buf.as_ptr() as *const _, buf.len());
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// write_all — robust write loop, bounded by a deadline.
|
|
// ---------------------------------------------------------------------------
|
|
//
|
|
// Mirrors read_some: each writability park is bounded by the remaining
|
|
// budget. `ErrorKind::TimedOut` when the deadline passes before the bytes
|
|
// are down — a client that stops reading must not pin this actor in
|
|
// 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());
|
|
if remaining.is_zero() {
|
|
return Err(io::Error::new(ErrorKind::TimedOut, "write deadline elapsed"));
|
|
}
|
|
if !smarm::wait_writable_timeout(fd, remaining)? {
|
|
return Err(io::Error::new(ErrorKind::TimedOut, "write deadline elapsed"));
|
|
}
|
|
|
|
let n = unsafe {
|
|
libc::write(fd, buf.as_ptr() as *const _, buf.len())
|
|
};
|
|
if n < 0 {
|
|
let err = io::Error::last_os_error();
|
|
if err.kind() == ErrorKind::WouldBlock {
|
|
continue; // spurious wake; retry
|
|
}
|
|
return Err(err);
|
|
}
|
|
if n == 0 {
|
|
return Err(io::Error::new(ErrorKind::WriteZero, "write returned 0"));
|
|
}
|
|
buf = &buf[n as usize..];
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// pump_stream — drive a RespBody::Stream onto the wire.
|
|
// ---------------------------------------------------------------------------
|
|
//
|
|
// The pull side of the v0.3 streaming design: the handler's producer actor
|
|
// owns the Sender; this conn actor owns the socket and every write
|
|
// deadline. We park in `recv()` between chunks — that park is stoppable
|
|
// (a draining registry's `request_stop` unwinds us out of `park_current`
|
|
// via the stop sentinel; fd + registry guards clean up), so an infinite
|
|
// stream is force-stoppable at the drain deadline like any other in-flight
|
|
// request. End of stream is the channel closing: every Sender dropped.
|
|
//
|
|
// Each chunk gets a FRESH write_timeout budget — a stream is expected to
|
|
// outlive any whole-response clock; what is not tolerated is a single
|
|
// write stalling. On write failure we return Err: the conn loop drops the
|
|
// Receiver, and the producer's next `send` observes the closed channel and
|
|
// should exit (that is the documented producer contract).
|
|
//
|
|
// `chunked` selects HTTP/1.1 chunked framing (hex-length CRLF payload
|
|
// CRLF, terminated by a 0-chunk) vs HTTP/1.0 raw writes (EOF-delimited;
|
|
// caller closes). Empty chunks are skipped — a zero-length chunk would
|
|
// terminate the framing early.
|
|
|
|
fn pump_stream(
|
|
fd: RawFd,
|
|
stream: StreamBody,
|
|
chunked: bool,
|
|
write_timeout: Duration,
|
|
) -> io::Result<()> {
|
|
let write_chunk = |payload: &[u8]| -> io::Result<()> {
|
|
let deadline = Instant::now() + write_timeout;
|
|
if chunked {
|
|
let mut framed = Vec::with_capacity(payload.len() + 20);
|
|
framed.extend_from_slice(format!("{:x}\r\n", payload.len()).as_bytes());
|
|
framed.extend_from_slice(payload);
|
|
framed.extend_from_slice(b"\r\n");
|
|
write_all(fd, &framed, deadline)
|
|
} else {
|
|
write_all(fd, payload, deadline)
|
|
}
|
|
};
|
|
|
|
loop {
|
|
// With a heartbeat configured (SSE), the wait between chunks is the
|
|
// heartbeat interval; expiry emits the ping and keeps waiting. A
|
|
// ping write that stalls past write_timeout errors out below —
|
|
// that is the dead-client detector. Without one, plain recv():
|
|
// stoppable by the draining registry either way.
|
|
let msg = match &stream.heartbeat {
|
|
Some((interval, payload)) => match stream.rx.recv_timeout(*interval) {
|
|
Ok(chunk) => Some(chunk),
|
|
Err(smarm::RecvTimeoutError::Timeout) => {
|
|
write_chunk(payload)?;
|
|
continue;
|
|
}
|
|
Err(smarm::RecvTimeoutError::Disconnected) => None,
|
|
},
|
|
None => stream.rx.recv().ok(),
|
|
};
|
|
|
|
match msg {
|
|
Some(chunk) => {
|
|
if chunk.is_empty() {
|
|
continue;
|
|
}
|
|
write_chunk(&chunk)?;
|
|
}
|
|
None => {
|
|
// All senders dropped: end of stream.
|
|
if chunked {
|
|
write_all(fd, b"0\r\n\r\n", Instant::now() + write_timeout)?;
|
|
}
|
|
return Ok(());
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Error responses for unparseable / malformed requests.
|
|
// ---------------------------------------------------------------------------
|
|
|
|
fn emit_error_response(fd: RawFd, err: &ParseError, deadline: Instant) {
|
|
let resp: &[u8] = match err {
|
|
ParseError::TooManyHeaders =>
|
|
b"HTTP/1.1 431 Request Header Fields Too Large\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
|
|
ParseError::BadContentLength =>
|
|
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).
|
|
_ =>
|
|
b"HTTP/1.1 400 Bad Request\r\ncontent-length: 0\r\nconnection: close\r\n\r\n",
|
|
};
|
|
let _ = write_all(fd, resp, deadline);
|
|
}
|