8bdec9784223fb3e95bb482b64cc7934bdd85cc3
10
Commits
| Author | SHA1 | Message | Date | |
|---|---|---|---|---|
|
|
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. |
||
|
|
0de2baa72d |
feat(channels): opt-in session persistence — v0.6 chunk 3
ChannelSession<P> (src/channels/session.rs): a channel actor that outlives its transport, buffering outbound broadcasts as the same Arc<Broadcast<P>>s the relay path carries (zero re-allocation) until rejoin, buffer cap, or TTL. Registered per pattern via PrefixRouter::channel_session::<S>(pattern, factory) / channel_session_default::<C>(pattern); one registry gen_server per registration, monomorphic over S::Key (no type erasure). Machinery: deploy() computes the key on the conn actor and casts the join handshake to the registry; the registry maps key -> (pid, control sender), forwards reattaches, spawns-and-monitors fresh actors, and prunes on Down (pid-guarded against replaced entries). The session actor selects [inbound, bus, ctl] attached and select_timeout([ctl, bus], ttl-remaining) detached; detach triggers are the closed inbound arm and ws send failure (the failed broadcast is the buffer's first entry — nothing lost). Drain re-encodes under the new join generation. As-landed decisions (each in module docs): - ChannelFactory gained a #[doc(hidden)] deploy() seam (default = ephemeral linked actor, the chunk-1 behavior verbatim); JoinHandshake and ChannelInbox are opaque pub structs. Channel construction moved from run_channel into deploy — same linked blast radius. - Every attach calls ch.join() again (rejoin acks need a payload and channels get to re-auth); session channels must treat join as re-entrant. - A rejected rejoin ENDS the session (no zombie sessions held open for unauthorized clients); next join is a cold start. - A second transport EVICTS the first (best-effort phx_close); a session is single-transport. - Explicit leave ends the session, not just the attachment. - TTL per detach episode; subscribe once, on first successful join; cap.max(1) semantics (cap 0 = first buffered message tears down, but an empty buffer still waits out the TTL). - Key: Clone beyond the ratified Eq+Hash+Send — the registry keeps a pid -> key reverse index for Down pruning. - The registry spawns eagerly at registration (channel_session is in-runtime-only, same law as ChannelHub::new): a lazy OnceLock would block std-sync under green threads on first-join races. - Shutdown composes WITHOUT links: conns die -> router Arcs drop -> registry inbox closes -> registry exits -> ctl senders drop -> detached sessions wake on the closed ctl arm, terminate, exit. shutdown_with_detached_session_terminates is the proof (ttl 30s vs 5s budget — only the chain can reap it). Tests: 6 session unit tests (buffer-and-drain ordering across reconnect, TTL expiry, cap teardown, leave, eviction incl. stale conn-entry self-prune, rejected rejoin) + the shutdown integration test. Shared toy codec + conn harness extracted to channels/testkit.rs (mod.rs tests now import it; no behavior change). Hammer: 35x lifecycle subset (+channels +session) + 3x full (phoenix) + 1x full (phoenix+smarm-trace) + featureless + channels-only, all green; dbg-grep clean. Suite: 90 unit + 46 integration + 2 doc under --features phoenix. |
||
|
|
0531390613 |
feat(channels): core protocol machinery — v0.6 chunk 1
Feature 'channels' (zero new deps): wire-neutral ChannelFrame/FrameRef + Encode/Decode codec traits (one method each, whole-envelope — the codec owns the wire format end to end); Channel trait; ChannelFactory + TopicRouter + shipped PrefixRouter (exact map + head:* prefix map); ChannelSocket (reply/push/broadcast/broadcast_from/broadcast_to); ChannelHub (router + PubSub<Broadcast<P>>, in-runtime + non-static per the v0.5 pattern, hub.upgrade(conn) entry, server-side hub.broadcast). Topology as ratified: one channel actor per (socket, topic), linked to the conn actor, select(&[&inbound, &bus_rx]) with inbound at index 0; subscription pinned to the channel actor pid. Subscribe-before-ok-reply makes the join ack first-class (the v0.5 on_open race, answered structurally). Cleanup is the sender-drop chain (conn handler drops -> inbound closes -> closed arm wakes -> terminate); bus arm dropped from the select set once observed closed (closed arms stay ready forever). Heartbeat answered conn-side on topic 'phoenix'; rejoin replaces the old generation with a phx_close to the old join_ref; lazy prune of dead channel entries on first failed forward (the pubsub discipline). As-landed deviations from the ratified sketch (module docs, veto by diff): join takes &mut self (trait-object callable; construction moved to ChannelFactory which still gets the full raw topic); reply(status, payload) with the ref threaded invisibly (the sketch listed both); broadcast gained the event parameter (unencodable without it); route() returns Arc not Box. Outbound payloads encode from borrows (FrameRef) so bus fan-out never clones P; payload-less control frames (phx_close, error/heartbeat replies) encode as the codec's empty payload rather than conjuring a P. 11 runtime-backed unit tests with a toy pipe codec (core provably JSON-free). Suites: 78 w/ feature, 67 featureless, both green. |
||
|
|
341d20ab45 |
feat(pubsub): topic table gen_server — v0.5 chunk 1
urus::pubsub, phoenix_pubsub-shaped, local node only. Independent of HTTP (imports smarm only); crate-extraction candidate. Design as ratified (handoff Q1-Q5, user said "push" on the leans): - Q1 generic: PubSub<M> per instance, zero-cost, no downcasts. Heterogeneous events = app-side enum M. v0.6 channels build on it. - Q2 Receiver: subscribe(topic) -> Receiver<Arc<M>>, fresh channel per subscription, subscriber owns its loop. Fan-in via a Sender-passing subscribe_with is a compatible later addition, deliberately not v1. subscribe_as(pid, topic) pins the subscription lifetime to an explicit pid for relay patterns (session actor subscribes, spawned relay receives) — without it the monitor would watch the wrong actor. - Q3 unbounded: broadcast never blocks the table (bounded-block = head-of-line across ALL topics; bounded-drop = silent loss). Memory risk documented; TWO cleanup paths: monitor Down on subscriber death (eager) + prune-on-send-failure for dropped receivers (lazy, on next broadcast). Bench before sharding. - Q4 user-owned ServerRef wrapper. DISCOVERY: the ratified optional register-by-name helper is NOT implementable against smarm's registry — it maps name -> Pid, and a Pid cannot be turned back into a ServerRef (the ref is the inbox sender). Needs smarm support or a global type-erased map; deferred, noted in module docs. pid() exposed for introspection. - Q5 unique per (pid, topic): HashMap<Topic, HashMap<Pid, Sender>>, subscribe idempotent (replaces sender; stale receiver's channel closes). One monitor per live pid (monitored: HashSet<Pid>, retired in handle_down so slot-reuse pids re-monitor). handle_down does the full-scan cleanup as ratified; pid->topics reverse index is the documented optimization if deaths ever measure hot. Subscribe is a call (table updated before return: a broadcast issued right after by the same caller is seen); unsubscribe/broadcast are casts. subscriber_count(topic) added as a call — needed by the tests to observe async cleanup, legitimate API anyway. 8 runtime-backed unit tests (smarm::run + collect-outside-assert-after pattern from smarm's own gen_server tests), including: Arc payload ptr-equality across subscribers, monitor-path pruning with NO broadcast issued (isolates it from the lazy path), broadcast_from self-skip, resubscribe replacement. Suite: 67 unit + 40 integration + 2 doc, green. smarm PRISTINE. |
||
|
|
ffc579c500 |
feat(ws): duplex wiring + handler API + echo example (v0.4 chunk 3)
Topology (b) as agreed: ONE actor per connection via smarm RFC 008 fd
arms. After the 101, ws::duplex::run_duplex replaces the HTTP loop:
try_select(&[&outbound_rx, &FdArm::readable(fd)]) per iteration,
outbound at index 0 (priority order = owed writes drain before reads;
honest backpressure). Accepted gap documented: mid-write of a large
outbound frame the actor isn't reading.
Handler API (deferred questions, as answered this session):
- WsHandler { on_message, on_close } runs IN the conn actor's select
loop; concurrency = spawn an actor with a WsSender clone (the SSE
producer pattern). Conn::upgrade(self, handler) flat signature;
WsUpgrade grows from marker to Box<dyn WsHandler> payload.
- WsSender: Clone; send/text/binary/ping/close -> Err(WsClosed).
Unbounded channel like SSE; slow client bounded by write_timeout.
- Control frames invisible v1: auto-pong inline (5.5.2), peer close
auto-echoed (5.5.1; server closes TCP first per 7.1.1) surfacing
only as on_close(Some(code), reason); wordless endings (EOF, write
failure, drain timeout) = on_close(None, ""). on_ping/on_pong can
land later as default methods, non-breaking.
- Server-initiated close (WsSender::close, FrameError->1002/1009/1007,
handler panic->1011): close out, bounded drain (one write_timeout
budget) for the peer echo, data discarded (1.4). Handler panics
re-raise smarm's stop sentinel first (the pipeline catch_unwind
dance, replicated).
- Caps: Config.max_frame_payload (1 MiB) / max_message_bytes (4 MiB).
Conn actor: 101 branch now hands fd + leftover buf (pipelined first
frame carries over; tested) + boxed handler to run_duplex; registry
entry stays Busy for the ws lifetime, so graceful shutdown force-stops
the conn out of the select park at the drain deadline (tested — the
in-runtime request_stop wake covers select parks too).
Tests: chunk-1 EOF test rewritten into a 9-test duplex suite (echo,
pipelined first frame, ping/pong, both close directions, 1002 unmasked,
1009 header-cap, fragmentation with interleaved ping, shutdown
force-stop). Suite 59u+40i+2d. Hammer: 35x lifecycle+ws subset + 3 full
+ 1 full under smarm-trace, all green; no AlreadyExists out of
try_select (the fresh eager-cleanup path held). Validated against
python websocket-client (echo, pong payload, close 1000).
examples/ws_echo.rs: /echo in-actor + /clock producer-spawning.
|
||
|
|
49538ccc83 |
feat(ws): RFC 6455 opening handshake — Conn::upgrade, 101 short-circuit (v0.4 chunk 1)
Design (as agreed in the v0.4 pass): - Dep #3 spent: sha1_smol 1.0.1 (zero transitive deps) instead of a vendored SHA-1 — user call: no hand-rolled crypto. Consequence noted in ROADMAP: the v0.6 wire-format JSON question now costs dep #4. - Base64 stays in-tree but ENCODE-ONLY (RFC 4648 vectors tested): it is an encoding, not crypto, and the decode direction (where parsing bugs live) is deliberately not implemented. - src/ws/{mod,handshake}.rs: validate() implements §4.2.1 (GET, 1.1, Host, upgrade/connection token lists case-insensitively, key shape = base64 of exactly 16 bytes, version 13). Version mismatch -> 426 + sec-websocket-version: 13 (§4.2.2); everything else -> 400. Rejected upgrades stay plain HTTP with keep-alive intact (tested: 400 then 101 on the same connection). - Conn grows pub upgrade: Option<WsUpgrade> (opaque marker; chunk 3 turns it into the duplex/handler handoff payload — shape deliberately uncommitted while the handler API is open). Conn::upgrade(self) is the commit point: 101 + accept key + marker + halt, or the rejection. - conn_actor: upgrade marker + status 101 -> write head, leave the HTTP loop. Until chunk 3 that drops the fd (clients see 101 then EOF); status check is defensive against a post-handler plug clobbering 101. - serialise_response: 1xx no longer get content-length injected (RFC 7230 §3.3.2). 204 left alone on purpose — separate conversation. Accept-key path pinned to the §1.3 worked example end-to-end (unit + wire). Suite: 34 unit + 33 integration + 2 doc. |
||
|
|
d14fb7c331 |
feat(sse): Conn::sse() sugar — EventSender + pump heartbeats (v0.3 chunk 3)
New src/sse.rs. Conn::sse() (and sse_with_heartbeat) turns the response
into a text/event-stream: status 200 if unset, content-type +
cache-control: no-cache, and a RespBody::Stream whose heartbeat is wired
to the new StreamBody.heartbeat field. Returns (Conn, EventSender);
EventSender is Clone (stream ends when the LAST clone drops) with
send(event, data) / data(data) / comment(text), all returning
Err(SseClosed) once the conn actor has dropped the stream — the producer
exit signal. Multi-line data becomes one data: line per line (pure
format_event, unit-tested).
Heartbeat lives in the pump, as the roadmap leaned: with
StreamBody.heartbeat set, pump_stream waits recv_timeout(interval) and on
expiry writes the payload (": keep-alive" comment chunk) and keeps
waiting. Liveness inversion is deliberate: no request clock governs an
SSE response (request_timeout is read-phase only since chunk 1); a dead
client is detected when an event or heartbeat write stalls past
write_timeout, and a draining shutdown force-stops the recv park like any
stream.
Tests: 4 unit (3 framing, 1 sse() wiring) + 2 integration (event round
trip over decoded chunked stream incl. named + unnamed events and clean
0-chunk EOS on sender drop; >=2 heartbeats observed across 450ms of
producer silence at a 100ms interval).
|
||
|
|
42a2464743 |
feat(conn): streaming response bodies — RespBody::Stream + chunked TE (v0.3 chunk 1)
Design decisions (per handoff Q1 user answer + Q3 proposal): - Producer shape is PULL: the handler returns a smarm Receiver<Vec<u8>> (RespBody::Stream / StreamBody, From<Receiver<Vec<u8>>> for ergonomics); the producer is an actor the handler spawned. The conn actor pumps in pump_stream(): it keeps sole ownership of the socket and of write deadlines, and the recv() park between chunks is stoppable by the draining registry's request_stop (stop sentinel unwinds out of park_current; fd + registry guards clean up) — so an infinite stream is force-stoppable at the drain deadline like any in-flight request, with no polling. End of stream = every Sender dropped = terminating 0-chunk. Producer contract: a send() error means the conn died; exit. - Framing: HTTP/1.1 gets transfer-encoding: chunked and the connection stays reusable after the terminator (keep-alive after chunked). HTTP/1.0 has no chunked TE: bytes go raw, keep-alive is forced off, EOF delimits. User-set content-length/transfer-encoding headers are DROPPED for stream bodies — we own the framing, and CL+chunked is a smuggling vector. Empty producer chunks are skipped (a 0-length chunk would terminate the framing early). - request_timeout stays READ-phase only (unchanged). NEW Config/ConnLimits field write_timeout (default 30s) gives every response write a per-write budget: the fixed head+body write, each streamed chunk, and the error paths (413/100-continue/4xx) all go through the now deadline-bounded write_all (wait_writable_timeout, mirroring read_some). Each chunk gets a FRESH budget — streams may outlive any whole-response clock; a single stalled write may not (write-side slowloris). Naming/default were flagged as a user call in the handoff: veto here if write_timeout(30s) isn't it. - RespBody loses derive(Clone) (Receiver isn't Clone; Clone was unused) and gets a manual Debug. Tests: 3 serialiser unit tests (chunked head shape, user-framing-header stripping, 1.0 fallback) + 5 integration (chunked round-trip with decoder, keep-alive after chunked, 1.0 EOF-delimited, shutdown force-stops an infinite stream at drain deadline, stalled reader killed at write_timeout with producer observing the closed channel). |
||
|
|
c658ac06b2 |
feat(serve): graceful shutdown + connection registry (v0.2 chunk 2)
Shutdown is drain-then-force-stop: flag.store(true) -> sup_h.join() (the no-more-accepts barrier) -> cast BeginDrain -> poll ConnCount every 50ms until 0 or drain_timeout; past the deadline each tick casts ForceStopConns and re-sweeps, so conns that registered after the deadline are still caught. Listener shutdown is redesigned vs chunk 1: Restart::Permanent -> Transient, wait-failure return -> panic (Transient restarts it), and a shared Arc<AtomicBool> flag checked at loop top with a 250ms wait_readable_timeout tick instead of request_stop — no request_stop on the supervisor or listeners anywhere. Historical note: the redesign was originally forced by smarm's then-lossy stop against QUEUED actors (fixed upstream in 7bab4d2); it is kept because flag-based listener shutdown is simpler and stop-semantics-free. Conns self-register with the registry (ConnStarted from the conn actor, not a listener cast): per-sender FIFO ordering is the only thing that prevents ConnEnded overtaking ConnStarted; see conn_registry module docs. conn_actor: in the catch_unwind(pipeline.run) Err branch, smarm::preempt::check_cancelled() runs BEFORE composing the 500 — smarm's stop is an undowncastable panic_any(StopSentinel), so the catch_unwind would otherwise swallow a stop and 500-and-keep-running. The stop flag is persistent, so check_cancelled re-raises cleanly. Handle.shutdown() wakes the root via a 100ms try_recv poll (SHUTDOWN_POLL): a cross-thread Sender::send's unpark is a try_with_runtime no-op without runtime TLS on the sending thread (still true on smarm 8e5b754; cross-thread unpark is a recorded smarm roadmap candidate — this poll dies with it). Two smarm bugs were found during this chunk and fixed upstream: lossy stop against QUEUED actors (7bab4d2) and the terminal-wake shutdown stall/hang (eddf3fe); post-mortem in artefact smarm-bug-terminal-wake.md. |
||
|
|
3b6c466210 | Initial commit |