From bafd532afa687eaa7d6ca9d19181ecaacbed6a37 Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Sun, 23 Aug 2026 08:04:32 +0200 Subject: [PATCH] =?UTF-8?q?feat:=20iteration=2035=20=E2=80=94=20net=20seam?= =?UTF-8?q?s=20+=20the=20serving=20slice=20(fiber-per-connection)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - runtime ids 91-95: net.read_dl/accept_dl/write_dl (per-call deadline, nil/false = the EXPECTED timeout; ms<=0 = old behavior bit for bit), net.listen_unix (unlink-before-bind, O_NONBLOCK on the listener — probe-found: accept4's flag covers accepted sockets only), net.peer - plane: one-op-per-park stays law — deadlines ride one per-shard TIMEOUT tick (sentinel user_data) + post-CQE expiry sweep + POLL_REMOVE tombstone; epoll's deadline scan grew the fd-park case; fibers POOL instead of freeing mid-run (stale-CQE UAF); plain parks zero park_deadline (no stale sleep deadlines) - probe: all five seams verified on BOTH WO_IO backends (timeout timing exact, peer round-trip, unix rebind) - framework: parse_request grows first_ms/read_ms; serve_conn — the keep-alive loop with deadlines where parked idle conns are LEGAL (close-when-idle RETIRED); App.handle_conn exposes it; plain serve() unchanged for simple apps - web-app: app-owned accept_dl loop + ConnWorker actor per connection (each builds its own App; cross-shard placement rides the DB actor); WA_IDLE_MS knob; gate grows to 41 checks — two slow requests served in PARALLEL, stalled client evicted at the idle deadline, slow-loris torn at the read deadline (400) - docs: story 35 -> done with banner; SQE/CQE design spec LANDED (was the review doc); ledger rows (timeouts/unix/keep-alive/peer), graph (NETSEAM cleared, KEEPAL done), builtin-surface rows, runtime CODE-LOGIC section, board entry - battery 13/13 fresh-built Co-Authored-By: Claude Fable 5 --- compiler/src/types.ml | 7 + docs/00-dependency-graph.md | 10 +- docs/examples/web-app/main.wo | 67 +++++- docs/examples/writeonce-framework/README.md | 18 +- docs/examples/writeonce-framework/app.wo | 10 + .../writeonce-framework/internal/parse.wo | 23 +- .../writeonce-framework/internal/serve.wo | 50 +++- docs/plan/oop-vm/08-builtin-surface.md | 5 + docs/stories/00-status.md | 1 + .../{refine => done}/35-net-runtime-seams.md | 13 +- .../specs/2026-08-23-net-seams-park-design.md | 102 ++++++++ runtime/src/CODE-LOGIC.md | 21 ++ runtime/src/builtin.c | 3 +- runtime/src/park.c | 73 +++++- runtime/src/sysio.c | 222 ++++++++++++++++++ runtime/src/vm.c | 32 ++- runtime/src/vm.h | 21 ++ runtime/src/wob.h | 18 +- scripts/web-app-accept.sh | 72 +++++- 19 files changed, 737 insertions(+), 31 deletions(-) rename docs/stories/language-runtime-database/{refine => done}/35-net-runtime-seams.md (87%) create mode 100644 docs/superpowers/specs/2026-08-23-net-seams-park-design.md diff --git a/compiler/src/types.ml b/compiler/src/types.ml index dec3da7..4afc6cd 100644 --- a/compiler/src/types.ml +++ b/compiler/src/types.ml @@ -302,6 +302,13 @@ let stdlib_members : stdlib_member list = m "net" "read" 2 53 (Some (TScalar "Text")) None; m "net" "write" 2 54 None None; m "net" "close" 1 55 None None; + (* iteration 35: per-call deadlines (nil/false = the EXPECTED timeout), + unix listeners, the peer's address *) + m "net" "read_dl" 3 91 (Some (TNullable (TScalar "Text"))) None; + m "net" "accept_dl" 2 92 (Some (TNullable (TScalar "Int"))) None; + m "net" "write_dl" 3 93 (Some (TScalar "Bool")) None; + m "net" "listen_unix" 1 94 (Some (TScalar "Int")) None; + m "net" "peer" 1 95 (Some (TScalar "Text")) None; (* proc *) m "proc" "run" 2 56 (Some (TNullable (TScalar proc_record_name))) (Some proc_record_name); (* json — both members are lowered specially (emit.ml): encode needs its diff --git a/docs/00-dependency-graph.md b/docs/00-dependency-graph.md index 291d072..a610c6d 100644 --- a/docs/00-dependency-graph.md +++ b/docs/00-dependency-graph.md @@ -99,7 +99,7 @@ flowchart TD I11x["11 fibers: reduction-budget preemption, blocking builtins park"]:::rt I9ex["22 baseline (numbers 8/23 sign against)"]:::rt - KEEPAL["keep-alive parking retired (close-when-idle policy dies; parked fds)"]:::gated + KEEPAL["keep-alive parking retired ✅ iteration 35 (app-owned fiber-per-connection + idle deadline)"]:::rt H2C2["h2c HTTP/2 cleartext (spec §C: also needs 23)"]:::gated STREAM2["request body streaming + backpressure"]:::gated SRESP2["streaming responses + explicit commit point"]:::gated @@ -161,10 +161,10 @@ flowchart TD XFF["client_ip: X-Forwarded-For parsing ✅ slice 2 (peer VERIFY stays gated)"]:::done ACCEPT["accepts(): response-side negotiation ✅ slice 2"]:::done - NETSEAM["GATE: net runtime seams (timeouts, unix socket, peer address) — story 35 owns"]:::gate - TMOUT["read/write/idle timeouts"]:::blocked - UNIX["unix socket binding"]:::blocked - PEERV["trusted-proxy PEER verification"]:::blocked + NETSEAM["GATE CLEARED: iteration 35 landed the net seams — _dl deadlines, listen_unix, peer (ids 91-95)"]:::done + TMOUT["read/write/idle timeouts ✅ iteration 35 (serve_conn read_ms/idle_ms)"]:::done + UNIX["unix socket binding ✅ iteration 35"]:::done + PEERV["trusted-proxy PEER verification — net.peer landed; the verify middleware is a ready framework slice"]:::ready CRYPTO["GATE CLEARED: iteration 34 landed C builtins — sha1/sha256/hmac_sha256 (ids 85-87)"]:::done SHA["sha1/sha256/hmac_sha256 ✅ iteration 34; SHA-512/CRC32 wait for a consumer"]:::done diff --git a/docs/examples/web-app/main.wo b/docs/examples/web-app/main.wo index 0c4e260..76a9d81 100644 --- a/docs/examples/web-app/main.wo +++ b/docs/examples/web-app/main.wo @@ -4,6 +4,8 @@ -- separate process, one binary. use env use json +use net +use time use framework use framework/http use framework/router @@ -118,6 +120,39 @@ class DeleteProduct { } } +-- ---- iteration 35, the serving slice: fiber-per-connection ------------ +-- The app owns the accept loop and spawns ONE ConnWorker actor per +-- accepted connection (spawn takes a class literal, so this ten-line +-- pattern lives app-side by doctrine). Each worker builds its OWN App — +-- route tables are small, and per-shard placement means no shared state +-- crosses heaps — then runs the framework's keep-alive loop with +-- deadlines. Concurrency: a slow request no longer blocks the next one; +-- a stalled client is evicted at the read deadline; an idle keep-alive +-- connection parks (blocking nobody) until the idle deadline. + +class Conn { + fd: net.Conn +} + +class ConnWorker { + token: Text + read_ms: Int + idle_ms: Int + fn receive(msg: Conn) { + let app = build_app(self.token); + app.handle_conn(msg.fd, self.read_ms, self.idle_ms); + } +} + +-- a deliberately slow route: the concurrency proof's workload +class Slow { + pad: Int + fn handle(req: Req) -> Resp { + time.sleep(400); + return ok_text("slow done"); + } +} + -- ---- framework v1 slice 2: the storefront exercises the new surface ---- -- wildcard capture: GET /files/*path echoes the rest @@ -185,13 +220,40 @@ fn main(args: multi Text) -> Int { return 2; } + -- iteration 35: the app-owned accept loop. accept_dl's tick keeps the + -- stop flag honored within 250ms; each accepted fd moves into its own + -- ConnWorker actor (placed round-robin across shards — DB access from + -- any shard rides the transparent DB actor). + let read_ms = 5000; + let idle_ms = 5000; + let ie = env.get("WA_IDLE_MS"); + if ie != nil { + let iv = parse_int(ie); + if iv != nil { idle_ms = iv; read_ms = iv; } + } + let srv = net.listen("127.0.0.1", port); + print("listening on 127.0.0.1:${port}"); + while true { + if env.stopping() { net.close(srv); return 0; } + let c = net.accept_dl(srv, 250); + if c != nil { + let w: actor Conn = spawn ConnWorker { token: "${token}", read_ms: read_ms, idle_ms: idle_ms }; + send(w, Conn { fd: c }); + } + } +} + +-- Every worker's own App: registration is idempotent code, and per-worker +-- construction is what lets connections live on ANY shard without shared +-- state (route tables are small; measured before optimized, per doctrine). +fn build_app(token: Text) -> App { let app = App { middleware: [], routes: [] }; -- v1 slice 2: host gate first (421 before anything runs), CORS preflight -- next, then the framework's Bearer mechanism (constant-time compare, -- principal attached to req.principal for handlers that want "who") app.use_mw(Mw { m: HostAllow { host: "a" } }); app.use_mw(Mw { m: Cors { allow_origin: "*" } }); - app.use_mw(Mw { m: BearerAuth { token: token, principal: "api" } }); + app.use_mw(Mw { m: BearerAuth { token: "${token}", principal: "api" } }); -- the response half: security headers + the CORS origin stamp on every -- response that leaves dispatch (404/405/401 included) app.use_after(Aw { a: SecurityHeaders { pad: 0 } }); @@ -204,9 +266,10 @@ fn main(args: multi Text) -> Int { app.get("/files/*path", EchoPath { pad: 0 }); app.get("/etag-probe", EtagProbe { pad: 0 }); app.get("/nego", NegoProbe { pad: 0 }); + app.get("/slow", Slow { pad: 0 }); let g = Group { prefix: "/api" }; g.use_mw(Mw { m: StampCtx { pad: 0 } }); g.get("/ping", ApiPing { pad: 0 }); app.mount(g); - return app.serve("127.0.0.1", port); + return app; } diff --git a/docs/examples/writeonce-framework/README.md b/docs/examples/writeonce-framework/README.md index 79677a2..c799f9b 100644 --- a/docs/examples/writeonce-framework/README.md +++ b/docs/examples/writeonce-framework/README.md @@ -52,9 +52,13 @@ writeonce-framework = { git = "https://github.com/shoneyj/writeonce-framework", ## Honest limits (v1, all deliberate) -- **Single-threaded, blocking** — one request at a time. Concurrency arrives - underneath this same surface now that the arc (8/11) has landed - (2026-08-21); the switch itself rides iteration 24's serving slice. +- **Concurrency is the APP's ten lines** (iteration 35's serving slice): + the framework ships `serve_conn` — the keep-alive loop with read/idle + deadlines — and the app owns accept + one spawned ConnWorker actor per + connection (web-app's pattern; `spawn` takes a class literal, so this + cannot live in the library). Parallel requests, stalled-client + eviction and parked idle keep-alive are gate-proven. The plain + `serve()` stays single-threaded for simple apps. - **TLS: none, anywhere.** Deploy behind nginx/caddy; the proxy terminates TLS+ALPN and gives browsers HTTP/2 while this backend speaks HTTP/1.1 keep-alive. See the web-app sample's README for the nginx sketch. @@ -81,10 +85,10 @@ first (pure `.wo` cannot express it yet). | Item | State | | --- | --- | | HTTP/1.1 parsing | ✅ parses + 400-and-survive; duplicate `Content-Length` rejected outright (RFC 9112 §6.3, slice 2); BODY_MAX bounds headers and body | -| Keep-alive | ✅ pipelined-serve / close-when-idle (arc landed 2026-08-21; retirement of close-when-idle rides iteration 24's fiber-per-connection slice) | -| Read/write/idle timeouts | 🔧 `net` has no timeout surface — story 35 owns the seam (park_deadline infra already exists for sleeps), then a framework knob | +| Keep-alive | ✅ RETIRED close-when-idle (iteration 35's serving slice): under the app-owned fiber-per-connection pattern, idle connections PARK until the idle deadline; the sequential `serve()` keeps the old policy for simple apps | +| Read/write/idle timeouts | ✅ iteration 35: per-call deadlines (`net.read_dl`/`accept_dl`/`write_dl`, nil/false = the expected timeout); `serve_conn(read_ms, idle_ms)` bounds slow-loris AND idle keep-alive | | Request size limits | ✅ BODY_MAX bounds headers AND body | -| Unix socket binding | 🔧 `net.listen` is TCP-only — story 35 owns the seam | +| Unix socket binding | ✅ `net.listen_unix(path)` (iteration 35) — stale sockets unlinked before bind, same accept/read/write after | | Graceful SIGTERM | ✅ in-flight request completes (blocking model), listener + fds closed, storage is per-commit durable (WAL fdatasync — nothing to checkpoint) | ### Routing @@ -104,7 +108,7 @@ first (pure `.wo` cannot express it yet). | Case-insensitive headers · query parsing | ✅ (names lowercased on read) | | JSON · form-urlencoded · multipart | ✅ all three hooks (`json.decode`, `form_values`, `multipart_parts`) | | Content negotiation | ✅ `media_type(req)` request-side; `accepts(req, mtype)` response-side (exact, type/*, */*; q-values stripped not ranked — ranking waits for an app serving alternates) — slice 2 | -| Trusted-proxy client IP | 🔶 `client_ip(req)` parses X-Forwarded-For (slice 2); VERIFYING the peer is the trusted proxy still needs the peer-address seam 🔧 — story 35 owns it | +| Trusted-proxy client IP | 🔶 `client_ip(req)` parses X-Forwarded-For; `net.peer(fd)` (iteration 35) exposes the peer — the verify middleware is now a pure-`.wo` candidate slice | | Status/header setting · redirects | ✅ builders + `set_header` | | Lazy body streaming + backpressure · streaming responses · explicit commit point | ⏸ UNBLOCKED by the arc (8/11 landed 2026-08-21) — stays parked until its own slice | | ETag + conditional requests | ✅ `etag_for` (quoted base64 SHA-256) + `with_etag` (If-None-Match → 304) over iteration 34's digest builtins — slice 2 | diff --git a/docs/examples/writeonce-framework/app.wo b/docs/examples/writeonce-framework/app.wo index a8ab72e..78ec0b4 100644 --- a/docs/examples/writeonce-framework/app.wo +++ b/docs/examples/writeonce-framework/app.wo @@ -115,4 +115,14 @@ pub class App { fn serve(host: Text, port: Int) -> Int { return internal.serve(host, port, self); } + + -- iteration 35, the serving slice: serve ONE accepted connection to + -- completion (the keep-alive loop with deadlines) — the body of an + -- app-spawned per-connection actor. The app owns the accept loop and + -- the spawn (a class literal, so the framework cannot spawn it); + -- each worker builds its own App and calls this. See web-app's + -- ConnWorker for the ten-line pattern. + fn handle_conn(c: net.Conn, read_ms: Int, idle_ms: Int) { + internal.serve_conn(c, self, read_ms, idle_ms); + } } \ No newline at end of file diff --git a/docs/examples/writeonce-framework/internal/parse.wo b/docs/examples/writeonce-framework/internal/parse.wo index 2d83aef..68be0c3 100644 --- a/docs/examples/writeonce-framework/internal/parse.wo +++ b/docs/examples/writeonce-framework/internal/parse.wo @@ -83,12 +83,25 @@ fn malformed(rest: Text) -> Parsed { } -- One request off the connection. `carry` = leftover bytes from the same --- connection's previous request (keep-alive). -pub fn parse_request(c: net.Conn, carry: Text) -> Parsed { +-- connection's previous request (keep-alive). Deadlines (iteration 35): +-- `first_ms` bounds the wait for a request's FIRST bytes (the keep-alive +-- idle window — expiry is a CLEAN close, not an error), `read_ms` bounds +-- every later read (a slow-loris mid-request is torn = 400-and-close). +-- ms <= 0 = wait forever, the pre-35 behavior bit for bit. +pub fn parse_request(c: net.Conn, carry: Text, first_ms: Int, read_ms: Int) -> Parsed { let buf = carry; let header_end = index_of(buf, "\r\n\r\n"); while header_end == -1 { - let got = net.read(c, 8192); + let dl = read_ms; + if buf == "" { dl = first_ms; } + let r = net.read_dl(c, 8192, dl); + if r == nil { + -- deadline expired: idle (nothing arrived) closes clean; a stalled + -- peer MID-request is torn + if trim(buf) == "" { return Parsed { closed: true, ok: true, req: nil, rest: "" }; } + return malformed(""); + } + let got = "${r}"; if len(got) == 0 { -- peer closed: clean between requests (empty buffer), torn otherwise if trim(buf) == "" { return Parsed { closed: true, ok: true, req: nil, rest: "" }; } @@ -148,7 +161,9 @@ pub fn parse_request(c: net.Conn, carry: Text) -> Parsed { let body = substr(buf, header_end + 4, len(buf) - header_end - 4); while len(body) < want { - let got = net.read(c, 8192); + let r2 = net.read_dl(c, 8192, read_ms); + if r2 == nil { return malformed(""); } -- stalled mid-body + let got = "${r2}"; if len(got) == 0 { return malformed(""); } -- peer died mid-body body = body .. got; } diff --git a/docs/examples/writeonce-framework/internal/serve.wo b/docs/examples/writeonce-framework/internal/serve.wo index 38f0079..56e5983 100644 --- a/docs/examples/writeonce-framework/internal/serve.wo +++ b/docs/examples/writeonce-framework/internal/serve.wo @@ -48,6 +48,54 @@ pub fn serialize(resp: Resp, keep: Bool, head_only: Bool) -> Text { return head .. resp.body; } +-- iteration 35, the serving slice: ONE connection served to completion — +-- the keep-alive loop with per-read deadlines. Meant to run INSIDE an +-- app-spawned per-connection actor (fiber): a parked idle connection is +-- legal there (it blocks nobody), so keep-alive stays OPEN until the +-- idle deadline evicts it — close-when-idle retires. read_ms bounds a +-- slow peer mid-request (torn = 400-and-close); idle_ms bounds the wait +-- for a request's first bytes (expiry = clean close). ms <= 0 = forever. +-- Closes the fd on every path except a WS hijack (status 101). +pub fn serve_conn(c: net.Conn, d: Dispatcher, read_ms: Int, idle_ms: Int) { + let carry = ""; + let alive = true; + let hijacked = false; + while alive { + if env.stopping() { alive = false; continue; } + let p = try parse_request(c, carry, idle_ms, read_ms) catch (e) nil; + if p == nil { alive = false; continue; } + if p.closed { alive = false; continue; } + if p.ok == false { + try net.write(c, serialize(bad_request("malformed request"), false, false)) catch (e) {} + alive = false; + continue; + } + let r = p.req; + if r == nil { alive = false; continue; } + carry = p.rest; + let is_head = r.method == "HEAD"; + if is_head { r.method = "GET"; } + -- fiber-per-connection: keep-alive stays OPEN (the idle deadline is + -- the eviction policy), unless the client asks to close + let keep = true; + let conn = r.headers["connection"]; + if conn != nil { + if to_lower(conn) == "close" { keep = false; } + } + let resp = try d.dispatch(r) catch (e) server_error(); + if resp.status == 101 { + hijacked = true; + alive = false; + continue; + } + try net.write(c, serialize(resp, keep, is_head)) catch (e) { alive = false; } + if keep == false { alive = false; } + } + if hijacked == false { + net.close(c); + } +} + pub fn serve(host: Text, port: Int, d: Dispatcher) -> Int { let srv = net.listen(host, port); print("listening on ${host}:${port}"); @@ -59,7 +107,7 @@ pub fn serve(host: Text, port: Int, d: Dispatcher) -> Int { let hijacked = false; while alive { if env.stopping() { alive = false; continue; } - let p = try parse_request(c, carry) catch (e) nil; -- an IO trap = gone + let p = try parse_request(c, carry, 0, 0) catch (e) nil; -- an IO trap = gone if p == nil { alive = false; continue; } if p.closed { alive = false; continue; } if p.ok == false { diff --git a/docs/plan/oop-vm/08-builtin-surface.md b/docs/plan/oop-vm/08-builtin-surface.md index 574c1e0..4dd458b 100644 --- a/docs/plan/oop-vm/08-builtin-surface.md +++ b/docs/plan/oop-vm/08-builtin-surface.md @@ -281,6 +281,11 @@ unset `env.get` are nil. | `env.get(name)` | `-> ?Text` | unset is nil | | `env.stopping()` | `-> Bool` | SIGTERM/SIGINT latch, handlers installed on first use | | `net.listen(host, port)` | `-> Int` | IPv4, SO_REUSEADDR, backlog 64; returns an fd | +| `net.read_dl(fd, max, ms)` | `-> ?Text` | iteration 35 (id 91): read with a per-call deadline — nil = expired (an EXPECTED outcome, never a trap), "" = EOF; ms <= 0 = wait forever | +| `net.accept_dl(fd, ms)` | `-> ?Int` | iteration 35 (id 92): accept with a deadline — nil = nothing arrived | +| `net.write_dl(fd, t, ms)` | `-> Bool` | iteration 35 (id 93): false = deadline mid-write — the stream is torn, close it | +| `net.listen_unix(path)` | `-> Int` | iteration 35 (id 94): AF_UNIX listener, stale socket unlinked before bind | +| `net.peer(fd)` | `-> Text` | iteration 35 (id 95): "ip:port" (TCP), "unix", "" on error | | `net.accept(fd)` | `-> Int` | | | `net.read(fd, max)` | `-> Text` | one read; the empty Text is EOF | | `net.write(fd, text)` | — | writes all of it | diff --git a/docs/stories/00-status.md b/docs/stories/00-status.md index 0d3d76a..c38305c 100644 --- a/docs/stories/00-status.md +++ b/docs/stories/00-status.md @@ -265,6 +265,7 @@ that sequences its tasks. Read one, approve, then the next starts. | -------- | --------------------------------------------------------------------------- | ---------------------------------------------------------- | | Language | 🔄 [iteration 36 — operator parity](language-runtime-database/in-progress/36-operator-parity.md): `not`, bitwise `& \| ^ << >>`, hex/binary/`_` literals, compound assigns — CODE LANDED 2026-08-22 (branch operator-parity, `.wob` v6, all gates green; reference project `.dev/reference/go` drove the design). Awaiting the developer's MANUAL pass on `docs/examples/operators/` (no test fixtures by directive); unblocks story 34's pure-`.wo` HMAC question | [plan](../superpowers/plans/2026-08-22-operator-parity.md) | | Language | the framework v1-polish slice landed 2026-08-20 (branch framework-v1, awaiting merge); next per the order: brainstorm 20/21's forks | [order](#implementation-order-re-sequenced-2026-08-21--concurrency-chain) | +| Runtime | ✅ **iteration 35 landed 2026-08-23** (branch `framework-v1b`, with framework v1 slice 2 + the serving slice): net deadlines/unix/peer (ids 91–95), fiber pooling, serve_conn + web-app fiber-per-connection — web-app gate 41/0, both WO_IO backends | [design](../superpowers/specs/2026-08-23-net-seams-park-design.md) | | Runtime | 🔄 **iteration 24 (absorbing 31 + 34): chat + actor lifecycle** — spec + plan approved 2026-08-23 (24 absorbs 31 by directive; 34 resolved C-builtins); executing on branch `chat-ws-lifecycle` | [marker](../in-progress/2026-08-23-chat-ws-lifecycle.md) · [plan](../superpowers/plans/2026-08-23-chat-ws-lifecycle.md) | The active slice's marker doc lives in [`in-progress/`](../in-progress/) — diff --git a/docs/stories/language-runtime-database/refine/35-net-runtime-seams.md b/docs/stories/language-runtime-database/done/35-net-runtime-seams.md similarity index 87% rename from docs/stories/language-runtime-database/refine/35-net-runtime-seams.md rename to docs/stories/language-runtime-database/done/35-net-runtime-seams.md index 98ad7c2..8f6e795 100644 --- a/docs/stories/language-runtime-database/refine/35-net-runtime-seams.md +++ b/docs/stories/language-runtime-database/done/35-net-runtime-seams.md @@ -1,6 +1,6 @@ --- iteration: "35" -status: refine +status: done --- # Iteration 35 — `net` runtime seams: timeouts, Unix sockets, peer address @@ -8,6 +8,17 @@ status: refine > Format: fiberloom `product/story-iteration-template`. Part of > [Story — one language, one runtime, one database, one binary](../00-story.md). > +> **LANDED 2026-08-23** (branch `framework-v1b`, with the serving slice +> riding it): per-call `_dl` deadlines (nil/false = the expected +> timeout; ids 91–93), `net.listen_unix` (unlink-before-bind, id 94), +> `net.peer` (id 95). Plane: shard-tick TIMEOUT + expiry sweep + +> POLL_REMOVE tombstone on uring, extended deadline scan on epoll, +> fibers POOLED against stale-CQE UAF — the full design in +> [the review spec](../../../superpowers/specs/2026-08-23-net-seams-park-design.md). +> Proof: all five seams probe-verified on BOTH `WO_IO` backends; +> `serve_conn` + web-app's fiber-per-connection pattern gate parallel +> requests, idle eviction, and slow-loris tearing (web-app 41 checks). +> > **Inserted 2026-08-22** — the framework ledger's three 🔧 rows get one > owner: "Read/write/idle timeouts — `net` has no timeout surface", > "Unix socket binding — `net.listen` is TCP-only", and "Trusted-proxy diff --git a/docs/superpowers/specs/2026-08-23-net-seams-park-design.md b/docs/superpowers/specs/2026-08-23-net-seams-park-design.md new file mode 100644 index 0000000..c76a412 --- /dev/null +++ b/docs/superpowers/specs/2026-08-23-net-seams-park-design.md @@ -0,0 +1,102 @@ +# Iteration 35 — net seams: the SQE/CQE plane design (for review) + +> **Status: LANDED 2026-08-23** (branch `framework-v1b`) — implemented +> as designed; one addition found by the probe: `listen_unix` must set +> O_NONBLOCK on the LISTENER (accept4's flag covers only accepted +> sockets). Covers story 35's fork 2 (deadline plumbing on the plane) +> plus the surface decisions taken with it. Normative park-protocol +> home once approved: +> [`../../plan/oop-vm/03-concurrency-coroutines.md`](../../plan/oop-vm/03-concurrency-coroutines.md). + +## Decisions taken (the story's four forks) + +1. **Timeout result**: nil/false, never a trap — a timeout is an + EXPECTED outcome (`parse_int` doctrine). `net.read_dl -> ?Text` + (nil = deadline, "" = EOF), `net.accept_dl -> ?Int`, + `net.write_dl -> Bool` (false = torn mid-write, close the fd). +2. **Plane plumbing**: shard-tick TIMEOUT + expiry sweep + POLL_REMOVE + tombstone (the table below) — NOT per-fiber second ops, NOT linked + ops. +3. **Surface**: per-CALL deadline argument (`_dl` builtin variants, + ids 91–95; 89/90 stay reserved for monitor/time.after). No hidden + fd state; `ms <= 0` = the old blocking behavior bit for bit. +4. **Unix sockets**: `net.listen_unix(path)` unlinks a stale socket + file before bind — a restart never needs manual cleanup. + `net.peer(fd) -> Text`: `"ip:port"` (TCP), `"unix"`, `""` on error. + +## The ring today (arc T4, landed) and the additions + +| SQE | user_data | CQE consumer action | +| --- | --- | --- | +| POLL_ADD fd, oneshot — net accept/read/write park | fiber pointer | `state == PARKED` → wake; else ignore (stale, benign) | +| TIMEOUT from `fb->park_ts` — time.sleep (`park_done=1`: resume continues PAST the builtin) | fiber pointer | same wake path; the op IS the waker — 1 op, 1 CQE, consumed exactly at wake | +| POLL_ADD wake_efd, oneshot — inbox envelopes arrived | `EFD_SENTINEL` (1) | drain eventfd, re-arm, return "adopt-needed" | +| **NEW** TIMEOUT, one per SHARD ("tick"), armed for the NEAREST fd-park deadline (`vm->tick_ts`) | `TICK_SENTINEL` (2) | `tick_armed = 0`; the sweep after the CQE batch wakes every expired fd-park | +| **NEW** POLL_REMOVE, `addr` = the expired park's fiber pointer (matches its POLL's user_data) | `CANCEL_SENTINEL` (3) | nothing — the tombstone's own completion | + +## Why this shape + +- **Standing invariant (arc T4):** one park = one SQE, and its CQE is + consumed precisely when the fiber wakes — no op ever outlives its + fiber. +- **The problem a deadline'd read creates:** two racing wait sources + (fd readiness, timer). Two per-fiber ops would let the LOSER's CQE + land after the fiber is freed — a use-after-free on the `user_data` + dereference one wait later. +- **The fix, three parts:** + 1. fd-parks keep exactly ONE op (their POLL_ADD); deadlines ride the + shard tick, whose user_data is a sentinel and can never dangle; + 2. the post-CQE sweep wakes expired fd-parks and submits POLL_REMOVE + to tombstone the orphaned poll (its -ECANCELED CQE arrives later + with the fiber's user_data and is dropped by the + `state == PARKED` guard); + 3. dead fibers are POOLED, never freed mid-run (`vm->fib_pool`, + freed at vm teardown) — even a post-mortem `state` read hits live + memory; the worst outcome anywhere is a SPURIOUS wake, which the + re-execute protocol absorbs (the builtin re-checks EAGAIN and the + deadline). Steady-state pool size = peak live fiber count. +- **epoll fallback:** zero ops — the existing parked-list deadline scan + (previously sleeps only) now also covers fd-parks with + `park_deadline > 0`; expiry = `EPOLL_CTL_DEL` + wake. +- **Rejected:** `IOSQE_IO_LINK` POLL→TIMEOUT chains (kernel cancels the + loser) — tighter, but linked-op error semantics are subtle and the + ring stays on 5.1-safe ops by doctrine. Also rejected: per-fd + deadline setting (`net.set_deadline`) — hidden fd state, against the + no-coloring lean the story records. + +## The `_dl` builtin state machine + +- FIRST entry stamps the absolute deadline into the fiber + (`fb->dl_active`, `fb->dl_at` — wall ms). The park protocol + RE-EXECUTES a parked builtin, and these fields are how the retry + remembers the original deadline. +- Every entry retries the syscall. Success/EOF/error → clear + `dl_active`, answer as the plain builtin would. +- EAGAIN with the deadline passed → clear `dl_active`, answer the + timeout result (nil / false). +- EAGAIN before the deadline → park with `park_fd` AND + `park_deadline = dl_at` both set; whichever fires first resumes the + builtin, which loops back to "every entry retries". +- Existing plain parks (`net.read`/`accept`/`write`) now explicitly + zero `park_deadline` — the sweep must never read a stale sleep + deadline off a reused fiber. + +## Consumers this unblocks (the serving slice, same branch) + +- framework `serve_conn(fd, dispatcher, read_ms, idle_ms)`: the + keep-alive loop where an idle parked connection is finally LEGAL + (idle deadline evicts it — close-when-idle retires); slow-client + reads bounded by `read_ms`. +- App-owned acceptor pattern: `accept_dl` loop (stop-flag checked per + tick) + the app spawns ITS conn-worker actor per connection, each + building its own App (interface-typed actor state — spike-proven). +- `client_ip` verification (`net.peer`) and unix-socket deployment + become pure framework/app slices. + +## Proof plan (pending) + +- Full battery (byte-identical: `ms <= 0` paths and untouched + builtins). +- New gate: stalled client evicted at the deadline (fds + RSS flat + over a soak), parallel requests complete out of order, unix listener + + peer round-trip — web-app gate additions + both `WO_IO` backends. diff --git a/runtime/src/CODE-LOGIC.md b/runtime/src/CODE-LOGIC.md index 6379457..c00223f 100644 --- a/runtime/src/CODE-LOGIC.md +++ b/runtime/src/CODE-LOGIC.md @@ -232,3 +232,24 @@ layout. - Proof: `just db-actor` (docs/examples/db-actor — multi-shard set ×3, both forced backends, single-shard byte-exact, WO_DATA replay pair); ASan/TSan clean on the RPC path. + +## Net deadlines + the deadline tick (iteration 35) + +- **`_dl` builtins (91–95) are per-call**: the fiber carries the absolute + deadline (`dl_active`/`dl_at`) across the park protocol's re-execution; + a timeout is the EXPECTED nil/false result, never a trap. `ms <= 0` is + the pre-35 behavior bit for bit. +- **One op per fd-park stays the law.** Deadlines ride ONE per-shard + TIMEOUT ("tick", sentinel user_data) armed for the nearest fd-park + deadline; the post-CQE sweep wakes expired parks and POLL_REMOVE + tombstones their poll. epoll needs no ops — its deadline scan grew the + fd-park case. Full design + rejected alternatives: + docs/superpowers/specs/2026-08-23-net-seams-park-design.md. +- **Fibers pool, never free mid-run** (`vm->fib_pool`): the loser of a + readiness-vs-deadline race can complete one wait late, and its + user_data must never point at freed memory. Worst case anywhere is a + spurious wake, absorbed by re-execution. Pool dies with the vm; + steady-state size = peak live fibers. +- **`listen_unix` sets O_NONBLOCK on the listener itself** — accept4's + SOCK_NONBLOCK flags the ACCEPTED socket only; a blocking listener + would block the whole shard (found by the seam probe, both backends). diff --git a/runtime/src/builtin.c b/runtime/src/builtin.c index 74c199d..095367e 100644 --- a/runtime/src/builtin.c +++ b/runtime/src/builtin.c @@ -171,7 +171,8 @@ int wo_builtin(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { both ranges (WO_B_MAP_GET_OPT and anything added after it) stay here */ if (C == WO_B_JSON_ENCODE || C == WO_B_JSON_DECODE) return wo_builtin_json(vm, R, ins, msg); - if ((C >= WO_B_SYS_FIRST && C <= WO_B_PROC_RUN) || C == WO_B_TIME_TICKS) + if ((C >= WO_B_SYS_FIRST && C <= WO_B_PROC_RUN) || C == WO_B_TIME_TICKS + || (C >= WO_B_NET_READ_DL && C <= WO_B_NET_PEER)) return wo_builtin_sys(vm, R, ins, msg); if (C >= WO_B_SHA1 && C <= WO_B_HMAC_SHA256) return wo_builtin_crypto(vm, R, ins, msg); diff --git a/runtime/src/park.c b/runtime/src/park.c index e13770d..6fbe11a 100644 --- a/runtime/src/park.c +++ b/runtime/src/park.c @@ -257,6 +257,63 @@ int wo_io_arm(wo_vm *vm, wo_fiber *fb) { /* user_data sentinel for the wake-eventfd's own readiness (fibers are * heap pointers, never 1) */ #define EFD_SENTINEL 1ull +/* iteration 35: the shard deadline tick (one TIMEOUT op armed for the + * nearest fd-park deadline) and the tombstone POLL_REMOVE's own CQE. + * Sentinels, never pointers — a late completion can never dangle. */ +#define TICK_SENTINEL 2ull +#define CANCEL_SENTINEL 3ull +#define IORING_OP_POLL_REMOVE 7 + +/* iteration 35: wake every fd-park whose deadline passed and tombstone + * its POLL op (the resumed builtin answers nil — the timeout result). + * The removed poll's CQE (-ECANCELED, user_data = the fiber) arrives + * later and is ignored: the fiber is RUNNABLE by then, and even a + * recycled fiber just takes a benign spurious wake (the park protocol + * re-executes the builtin, which re-checks). Returns woke-count. */ +static int deadline_sweep_uring(wo_vm *vm, int64_t now) { + int woke = 0; + for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext) { + if (fb->state == WO_FIB_PARKED && fb->park_fd >= 0 + && fb->park_deadline > 0 && fb->park_deadline <= now) { + struct io_uring_sqe sqe; + memset(&sqe, 0, sizeof sqe); + sqe.opcode = IORING_OP_POLL_REMOVE; + sqe.fd = -1; + sqe.addr = (uint64_t)(uintptr_t)fb; /* match the poll's user_data */ + sqe.user_data = CANCEL_SENTINEL; + (void)uring_submit(vm, &sqe); + wake(vm, fb); + woke++; + } + } + return woke; +} + +/* Arm (or re-arm) the tick for the nearest fd-park deadline. Cheap + * over-arming is fine: a tick firing with nothing expired just re-arms. */ +static void tick_arm_uring(wo_vm *vm, int64_t now) { + int64_t next = 0; + for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext) + if (fb->state == WO_FIB_PARKED && fb->park_fd >= 0 && fb->park_deadline > 0) + if (next == 0 || fb->park_deadline < next) next = fb->park_deadline; + if (next == 0) return; + if (vm->tick_armed && vm->tick_at <= next) return; + int64_t rel = next - now; + if (rel < 0) rel = 0; + vm->tick_ts.sec = rel / 1000; + vm->tick_ts.nsec = (rel % 1000) * 1000000LL; + struct io_uring_sqe sqe; + memset(&sqe, 0, sizeof sqe); + sqe.opcode = IORING_OP_TIMEOUT; + sqe.fd = -1; + sqe.addr = (uint64_t)(uintptr_t)&vm->tick_ts; + sqe.len = 1; + sqe.user_data = TICK_SENTINEL; + if (uring_submit(vm, &sqe) == 0) { + vm->tick_armed = 1; + vm->tick_at = next; + } +} static void efd_drain(wo_vm *vm) { uint64_t v = 0; @@ -279,6 +336,7 @@ int wo_io_wait(wo_vm *vm) { sqe.user_data = EFD_SENTINEL; if (uring_submit(vm, &sqe) == 0) vm->efd_armed = 1; } + tick_arm_uring(vm, now_ms()); rings r = ring_ptrs(vm); uint32_t head = *r.cq_head; uint32_t tail = __atomic_load_n(r.cq_tail, __ATOMIC_ACQUIRE); @@ -299,6 +357,10 @@ int wo_io_wait(wo_vm *vm) { vm->efd_armed = 0; efd_drain(vm); woke = 2; /* inbox wake: the caller adopts */ + } else if (cqe->user_data == TICK_SENTINEL) { + vm->tick_armed = 0; /* the sweep below decides who expired */ + } else if (cqe->user_data == CANCEL_SENTINEL) { + /* the tombstone's own completion: nothing to do */ } else { wo_fiber *fb = (wo_fiber *)(uintptr_t)cqe->user_data; if (fb && fb->state == WO_FIB_PARKED) { @@ -309,6 +371,7 @@ int wo_io_wait(wo_vm *vm) { head++; } __atomic_store_n(r.cq_head, head, __ATOMIC_RELEASE); + if (deadline_sweep_uring(vm, now_ms()) && woke != 2) woke = 1; if (woke == 2) return 1; /* adopt-needed */ if (woke) return 0; continue; @@ -327,7 +390,9 @@ int wo_io_wait(wo_vm *vm) { int timeout = -1; int64_t now = now_ms(); for (wo_fiber *fb = vm->parked; fb; fb = fb->pnext) - if (fb->park_fd == -1) { /* deadline waits only, never INBOX */ + if (fb->park_fd == -1 + || (fb->park_fd >= 0 && fb->park_deadline > 0)) { + /* sleeps AND deadline'd fd-parks (iteration 35); never INBOX */ int64_t rel = fb->park_deadline - now; if (rel < 0) rel = 0; if (timeout < 0 || rel < timeout) timeout = (int)rel; @@ -355,7 +420,11 @@ int wo_io_wait(wo_vm *vm) { wo_fiber *fb = vm->parked; while (fb) { wo_fiber *nx = fb->pnext; - if (fb->park_fd == -1 && fb->park_deadline <= now) { + if ((fb->park_fd == -1 + || (fb->park_fd >= 0 && fb->park_deadline > 0)) + && fb->park_deadline <= now) { + if (fb->park_fd >= 0) + epoll_ctl(vm->io_fd, EPOLL_CTL_DEL, fb->park_fd, NULL); wake(vm, fb); woke = 1; } diff --git a/runtime/src/sysio.c b/runtime/src/sysio.c index add11db..1a60501 100644 --- a/runtime/src/sysio.c +++ b/runtime/src/sysio.c @@ -28,6 +28,7 @@ #include #include #include +#include #include #include #include @@ -392,6 +393,7 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { if (fd < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { /* arc T4: park until the listener is readable, then retry */ vm->cur->park_fd = (int)R[B]; + vm->cur->park_deadline = 0; vm->cur->park_events = POLLIN; vm->cur->park_done = 0; return WO_SYS_PARKED; @@ -426,6 +428,7 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { * re-allocates) and park until the fd is readable */ wo_str_free(rt, s); vm->cur->park_fd = (int)R[B]; + vm->cur->park_deadline = 0; vm->cur->park_events = POLLIN; vm->cur->park_done = 0; return WO_SYS_PARKED; @@ -471,6 +474,7 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { if (errno == EAGAIN || errno == EWOULDBLOCK) { vm->cur->park_wr_at = at; vm->cur->park_fd = (int)R[B]; + vm->cur->park_deadline = 0; vm->cur->park_events = POLLOUT; vm->cur->park_done = 0; return WO_SYS_PARKED; @@ -488,6 +492,224 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { R[A] = 0; return 0; } + /* ---- iteration 35: per-call deadlines + unix sockets + peer ------- + * The _dl protocol: the FIRST entry computes the absolute deadline + * into the fiber (dl_active/dl_at — the park/retry re-executes the + * builtin, and this is how the retry remembers it); every entry + * re-tries the syscall; EAGAIN past the deadline answers the timeout + * result (nil/false — an EXPECTED outcome, never a trap); EAGAIN + * before it parks with BOTH the fd and the deadline armed (park.c's + * sweep wakes whichever fires first). ms <= 0 = no deadline. */ + case WO_B_NET_READ_DL: { + wo_fiber *fb = vm->cur; + struct timespec dts; + clock_gettime(CLOCK_REALTIME, &dts); + int64_t dnow = (int64_t)dts.tv_sec * 1000 + dts.tv_nsec / 1000000; + if (!fb->dl_active) { + int64_t ms = (int64_t)R[B + 2]; + fb->dl_active = 1; + fb->dl_at = ms > 0 ? dnow + ms : 0; + } + int64_t max = (int64_t)R[B + 1]; + if (max < 0) max = 0; + wo_str *s = wo_str_alloc(rt, (uint32_t)max); + if (!s) { + fb->dl_active = 0; + *msg = "out of memory"; + return WO_T_OOM; + } + ssize_t n; + for (;;) { + n = read((int)R[B], s->data, (size_t)max); + if (n >= 0 || errno != EINTR) break; + if (stop_pending()) { + wo_str_free(rt, s); + fb->dl_active = 0; + return WO_SYS_STOPPED; + } + } + if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { + wo_str_free(rt, s); + if (fb->dl_at > 0 && dnow >= fb->dl_at) { + fb->dl_active = 0; + R[A] = 0; /* ?Text nil: the deadline expired */ + return 0; + } + fb->park_fd = (int)R[B]; + fb->park_deadline = fb->dl_at; /* 0 = wait forever, like read */ + fb->park_events = POLLIN; + fb->park_done = 0; + return WO_SYS_PARKED; + } + fb->dl_active = 0; + if (n < 0) { + wo_str_free(rt, s); + *msg = strerror(errno); + return WO_T_IO; + } + if ((size_t)n == (size_t)max) { + R[A] = (uint64_t)(uintptr_t)s; + return 0; + } + wo_str *exact = wo_str_new(rt, s->data, (uint32_t)n); + wo_str_free(rt, s); + if (!exact) { + *msg = "out of memory"; + return WO_T_OOM; + } + R[A] = (uint64_t)(uintptr_t)exact; + return 0; + } + case WO_B_NET_ACCEPT_DL: { + wo_fiber *fb = vm->cur; + struct timespec dts; + clock_gettime(CLOCK_REALTIME, &dts); + int64_t dnow = (int64_t)dts.tv_sec * 1000 + dts.tv_nsec / 1000000; + if (!fb->dl_active) { + int64_t ms = (int64_t)R[B + 1]; + fb->dl_active = 1; + fb->dl_at = ms > 0 ? dnow + ms : 0; + } + int fd; + for (;;) { + fd = accept4((int)R[B], NULL, NULL, SOCK_NONBLOCK); + if (fd >= 0 || errno != EINTR) break; + if (stop_pending()) { + fb->dl_active = 0; + return WO_SYS_STOPPED; + } + } + if (fd < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { + if (fb->dl_at > 0 && dnow >= fb->dl_at) { + fb->dl_active = 0; + R[A] = WO_NIL_SCALAR; /* ?Int nil: nothing arrived */ + return 0; + } + fb->park_fd = (int)R[B]; + fb->park_deadline = fb->dl_at; + fb->park_events = POLLIN; + fb->park_done = 0; + return WO_SYS_PARKED; + } + fb->dl_active = 0; + if (fd < 0) { + *msg = strerror(errno); + return WO_T_IO; + } + R[A] = (uint64_t)fd; + return 0; + } + case WO_B_NET_WRITE_DL: { + wo_fiber *fb = vm->cur; + const wo_str *body = (const wo_str *)(uintptr_t)R[B + 1]; + if (!body || body->h.class_id != WO_CLS_STR) { + *msg = "not a text value"; + return WO_T_BOUNDS; + } + struct timespec dts; + clock_gettime(CLOCK_REALTIME, &dts); + int64_t dnow = (int64_t)dts.tv_sec * 1000 + dts.tv_nsec / 1000000; + if (!fb->dl_active) { + int64_t ms = (int64_t)R[B + 2]; + fb->dl_active = 1; + fb->dl_at = ms > 0 ? dnow + ms : 0; + } + uint32_t at = fb->park_wr_at; + fb->park_wr_at = 0; + while (at < body->len) { + ssize_t n = write((int)R[B], body->data + at, body->len - at); + if (n < 0) { + if (errno == EINTR) { + if (stop_pending()) { + fb->dl_active = 0; + return WO_SYS_STOPPED; + } + continue; + } + if (errno == EAGAIN || errno == EWOULDBLOCK) { + if (fb->dl_at > 0 && dnow >= fb->dl_at) { + fb->dl_active = 0; + R[A] = 0; /* false: torn mid-write — close the fd */ + return 0; + } + fb->park_wr_at = at; + fb->park_fd = (int)R[B]; + fb->park_deadline = fb->dl_at; + fb->park_events = POLLOUT; + fb->park_done = 0; + return WO_SYS_PARKED; + } + fb->dl_active = 0; + *msg = strerror(errno); + return WO_T_IO; + } + at += (uint32_t)n; + } + fb->dl_active = 0; + R[A] = 1; + return 0; + } + case WO_B_NET_LISTEN_UNIX: { /* unlink-before-bind: a restart never + * needs manual socket-file cleanup */ + if (cstr_of(R[B], path, sizeof path, msg)) return WO_T_BOUNDS; + struct sockaddr_un ua; + if (strlen(path) >= sizeof(ua.sun_path)) { + *msg = "unix socket path too long"; + return WO_T_BOUNDS; + } + int fd = socket(AF_UNIX, SOCK_STREAM, 0); + if (fd < 0) { + *msg = strerror(errno); + return WO_T_IO; + } + unlink(path); + memset(&ua, 0, sizeof ua); + ua.sun_family = AF_UNIX; + strncpy(ua.sun_path, path, sizeof(ua.sun_path) - 1); + if (bind(fd, (struct sockaddr *)&ua, sizeof ua) != 0 || listen(fd, 64) != 0) { + *msg = strerror(errno); + close(fd); + return WO_T_IO; + } + /* the listener must be NONBLOCKING like net.listen's (arc T4): + * accept4's SOCK_NONBLOCK flags the ACCEPTED socket, not this one — + * a blocking listener would block the whole shard in the syscall */ + fcntl(fd, F_SETFL, fcntl(fd, F_GETFL, 0) | O_NONBLOCK); + R[A] = (uint64_t)fd; + return 0; + } + case WO_B_NET_PEER: { /* "ip:port" (TCP), "unix" (unix peers), "" error */ + struct sockaddr_storage ss; + socklen_t sl = sizeof ss; + if (getpeername((int)R[B], (struct sockaddr *)&ss, &sl) != 0) { + wo_str *e = wo_str_new(rt, "", 0); + if (!e) { + *msg = "out of memory"; + return WO_T_OOM; + } + R[A] = (uint64_t)(uintptr_t)e; + return 0; + } + char pbuf[64]; + if (ss.ss_family == AF_INET) { + struct sockaddr_in *in = (struct sockaddr_in *)&ss; + uint32_t ip = ntohl(in->sin_addr.s_addr); + snprintf(pbuf, sizeof pbuf, "%u.%u.%u.%u:%u", (ip >> 24) & 255, + (ip >> 16) & 255, (ip >> 8) & 255, ip & 255, + (unsigned)ntohs(in->sin_port)); + } else if (ss.ss_family == AF_UNIX) { + snprintf(pbuf, sizeof pbuf, "unix"); + } else { + pbuf[0] = 0; + } + wo_str *out = wo_str_new(rt, pbuf, (uint32_t)strlen(pbuf)); + if (!out) { + *msg = "out of memory"; + return WO_T_OOM; + } + R[A] = (uint64_t)(uintptr_t)out; + return 0; + } /* ---- proc -------------------------------------------------------- */ case WO_B_PROC_RUN: { /* Proc: 0 code, 1 out, 2 err. argv[0] is the * command itself; the `multi Text` argument diff --git a/runtime/src/vm.c b/runtime/src/vm.c index ca83c74..1ca9271 100644 --- a/runtime/src/vm.c +++ b/runtime/src/vm.c @@ -550,6 +550,12 @@ int wo_vm_init(wo_vm *vm, const wo_module *mod, size_t heap_cap) { } void wo_vm_destroy(wo_vm *vm) { + /* iteration 35: the fiber pool dies with the vm */ + while (vm->fib_pool) { + wo_fiber *fb = vm->fib_pool; + vm->fib_pool = fb->next; + free(fb); + } /* actors first — dropping their state and queued messages needs the * runtime alive */ wo_actor *a = vm->actors; @@ -591,12 +597,28 @@ static wo_fiber *fib_dequeue(wo_vm *vm) { return fb; } +/* iteration 35: dead fibers pool instead of freeing (vm.h's UAF note). + * next links the pool; a pooled fiber's state is DONE, so a stale plane + * completion reading it is harmless. */ +static void fib_retire(wo_vm *vm, wo_fiber *fb) { + fb->state = WO_FIB_DONE; + fb->next = vm->fib_pool; + vm->fib_pool = fb; +} + wo_fiber *wo_vm_spawn_fiber(wo_vm *vm, uint32_t method_idx, const uint64_t *args, uint32_t argc) { if (method_idx >= vm->mod->method_cnt) return NULL; const wo_methodrec *sme = &vm->mod->methods[method_idx]; if (argc != sme->arg_cnt) return NULL; - wo_fiber *fb = calloc(1, sizeof(*fb)); + wo_fiber *fb; + if (vm->fib_pool) { + fb = vm->fib_pool; + vm->fib_pool = fb->next; + memset(fb, 0, sizeof(*fb)); + } else { + fb = calloc(1, sizeof(*fb)); + } if (!fb) return NULL; fb->depth = 1; fb->frames[0].method = method_idx; @@ -625,7 +647,7 @@ static void fib_reap(wo_vm *vm, wo_fiber *fb) { } if (fb != &vm->f0) { vm->nfibers--; - free(fb); + fib_retire(vm, fb); } } @@ -1175,7 +1197,7 @@ static int vm_run(wo_vm *vm, uint64_t *ret, wo_err *err) { * dangled — the actor's mailbox rotted forever). */ \ if (dead->actor) actor_die(vm, dead->actor, dead); \ vm->nfibers--; \ - free(dead); \ + fib_retire(vm, dead); \ NEXT_RUNNABLE(); \ RELOAD(); \ NEXT(); \ @@ -1543,7 +1565,7 @@ dispatch: a->active = NULL; \ } \ vm->nfibers--; \ - free(dead); \ + fib_retire(vm, dead); \ NEXT_RUNNABLE(); \ RELOAD(); \ NEXT(); \ @@ -1713,7 +1735,7 @@ dispatch: wo_fiber *dead = vm->cur; vm->cur = &vm->f0; vm->nfibers--; - free(dead); + fib_retire(vm, dead); if (vm->f0.depth) { /* main was queued mid-run: release its frames too */ wo_fiber *q = vm->qhead, *prev = NULL; diff --git a/runtime/src/vm.h b/runtime/src/vm.h index b86616a..78675cb 100644 --- a/runtime/src/vm.h +++ b/runtime/src/vm.h @@ -93,6 +93,12 @@ typedef struct wo_fiber { * a plain send) — where FIBER_DONE ships the receive's return value. */ struct wo_fiber *msg_caller; uint32_t msg_caller_shard; + /* iteration 35: the in-flight per-CALL deadline (_dl builtins). Set on + * the builtin's first entry, cleared when it answers — the park/retry + * protocol re-executes the builtin, and this is how the retry knows + * the original deadline. */ + int dl_active; + int64_t dl_at; /* wall ms */ } wo_fiber; /* arc stage 3: park_fd sentinel — PARKED with NO plane wait; the wake is @@ -164,6 +170,21 @@ typedef struct wo_vm { /* the I/O plane (arc T4, park.c): io_uring primary, epoll fallback */ wo_fiber *parked; /* fibers waiting on the plane */ uint32_t nparked; + /* iteration 35: dead fibers are POOLED, never freed mid-run — a stale + * plane completion (the loser of a poll-vs-deadline race, consumed one + * wait later) may still read the fiber's `state` word, and reading + * freed memory is the UAF this prevents. Steady-state pool size = the + * peak live fiber count; the pool dies with the vm. */ + wo_fiber *fib_pool; + /* iteration 35, uring backend: the shard's ONE deadline tick — a + * TIMEOUT op with a sentinel user_data armed for the nearest fd-park + * deadline (fd parks keep exactly one POLL op each; expiry wakes them + * from the scan and POLL_REMOVE tombstones the poll). */ + int tick_armed; + int64_t tick_at; + struct { + long long sec, nsec; + } tick_ts; int io_kind; /* 0 = uring, 1 = epoll */ int efd_armed; /* wake_efd registered on the plane (uring oneshot) */ int io_fd; /* ring fd or epoll fd */ diff --git a/runtime/src/wob.h b/runtime/src/wob.h index 36d44a0..506d59a 100644 --- a/runtime/src/wob.h +++ b/runtime/src/wob.h @@ -468,9 +468,25 @@ enum { * return value arrives. R is a SCALAR (v1, * compiler-enforced WO-E226). Dead callee = * WO_T_ACTOR, immediately or mid-call. */ + /* ids 89 (monitor) and 90 (time.after) are RESERVED for the rest of + * the lifecycle slice — do not reuse. */ + /* ---- iteration 35: net seams (sysio.c). Deadlines are per-CALL (no + * hidden fd state); a timeout is an EXPECTED outcome, so it answers + * nil/false, never a trap. ms <= 0 = no deadline (the old behavior, + * bit for bit). ---- */ + WO_B_NET_READ_DL = 91, /* (fd, max, ms) -> ?Text: nil = deadline + * expired with nothing read; "" = EOF */ + WO_B_NET_ACCEPT_DL = 92, /* (fd, ms) -> ?Int: nil = nothing arrived */ + WO_B_NET_WRITE_DL = 93, /* (fd, text, ms) -> Bool: false = deadline + * mid-write — the stream is torn, close it */ + WO_B_NET_LISTEN_UNIX = 94,/* (path) -> Int: AF_UNIX listener; a stale + * socket file is unlinked first (a restart + * never needs manual cleanup) */ + WO_B_NET_PEER = 95, /* (fd) -> Text: "ip:port" for TCP peers, + * "unix" for unix-socket peers, "" on error */ }; -#define WO_B_MAX 88u +#define WO_B_MAX 95u /* ids at or above this one live in sysio.c, not builtin.c */ #define WO_B_SYS_FIRST WO_B_FS_EXISTS diff --git a/scripts/web-app-accept.sh b/scripts/web-app-accept.sh index 40835a3..28ca4e1 100755 --- a/scripts/web-app-accept.sh +++ b/scripts/web-app-accept.sh @@ -82,7 +82,7 @@ else fi DATA="$W/data"; mkdir -p "$DATA" -WA_TOKEN=s3cr3t WO_DATA="$DATA" "$W/app/target/web-app" "$PORT" >"$W/srv.out" 2>&1 & +WA_TOKEN=s3cr3t WA_IDLE_MS=600 WO_DATA="$DATA" "$W/app/target/web-app" "$PORT" >"$W/srv.out" 2>&1 & SRV=$! for _ in $(seq 1 40); do grep -q listening "$W/srv.out" 2>/dev/null && break; sleep 0.1; done @@ -280,6 +280,74 @@ r="$(hraw "$AUTH origin: http://x" GET /products)" [[ "$r" == 200\|*"access-control-allow-origin: *"* ]] \ && ok "CORS origin stamped on real responses" || bad "cors-after" "$r" +# ---- 12c. the serving slice: fiber-per-connection + deadlines ---- +r="$(timeout 10 python3 - "$PORT" <<'PYEOF' +import socket, sys, time, threading +port = int(sys.argv[1]) +REQ = b"GET /slow HTTP/1.1\r\nhost: a\r\nauthorization: Bearer s3cr3t\r\nconnection: close\r\ncontent-length: 0\r\n\r\n" +def one(res, i): + s = socket.create_connection(("127.0.0.1", port), timeout=8) + s.sendall(REQ) + d = b"" + while True: + c = s.recv(4000) + if not c: break + d += c + res[i] = b"slow done" in d +t0 = time.time() +res = [False, False] +ts = [threading.Thread(target=one, args=(res, i)) for i in (0, 1)] +[t.start() for t in ts]; [t.join() for t in ts] +el = int((time.time() - t0) * 1000) +print(f"{res[0] and res[1]}|{el}") +PYEOF +)" +pw="${r%%|*}"; pe="${r#*|}" +[[ "$pw" == "True" && "$pe" -lt 700 ]] \ + && ok "two slow requests served in PARALLEL (${pe}ms, serial would be 800+)" \ + || bad "parallel" "$r" +r="$(timeout 10 python3 - "$PORT" <<'PYEOF' +import socket, sys, time +port = int(sys.argv[1]) +# a client that connects and sends NOTHING: the idle deadline must evict it +s = socket.create_connection(("127.0.0.1", port), timeout=8) +t0 = time.time() +s.settimeout(5) +try: + d = s.recv(100) + print(f"closed|{int((time.time()-t0)*1000)}" if d == b"" else f"data|{d[:20]}") +except socket.timeout: + print("still-open|5000") +PYEOF +)" +sw="${r%%|*}"; se="${r#*|}" +[[ "$sw" == "closed" && "$se" -lt 2500 ]] \ + && ok "stalled client evicted at the idle deadline (${se}ms)" \ + || bad "stalled" "$r" +r="$(timeout 10 python3 - "$PORT" <<'PYEOF' +import socket, sys, time +port = int(sys.argv[1]) +# half a request then silence: the READ deadline tears it (400-and-close) +s = socket.create_connection(("127.0.0.1", port), timeout=8) +s.sendall(b"GET /products HTTP/1.1\r\nhost: a\r\nauthor") +t0 = time.time() +d = b"" +s.settimeout(5) +try: + while True: + c = s.recv(400) + if not c: break + d += c +except socket.timeout: pass +status = d.decode(errors="replace").split(" ")[1] if d else "closed" +print(f"{status}|{int((time.time()-t0)*1000)}") +PYEOF +)" +tw="${r%%|*}"; te="${r#*|}" +[[ "$tw" == "400" && "$te" -lt 2500 ]] \ + && ok "slow-loris torn at the read deadline (400, ${te}ms)" \ + || bad "slowloris" "$r" + r="$(timeout 5 python3 - "$PORT" <<'PYEOF' import socket, sys s = socket.create_connection(("127.0.0.1", int(sys.argv[1])), timeout=5) @@ -325,7 +393,7 @@ for _ in $(seq 1 30); do kill -0 "$SRV" 2>/dev/null || { stopped=0; break; }; sl SRV="" # ---- 15. restart persistence (WAL replay) ---- -WA_TOKEN=s3cr3t WO_DATA="$DATA" "$W/app/target/web-app" "$PORT" >>"$W/srv.out" 2>&1 & +WA_TOKEN=s3cr3t WA_IDLE_MS=600 WO_DATA="$DATA" "$W/app/target/web-app" "$PORT" >>"$W/srv.out" 2>&1 & SRV=$! sleep 0.5 expect "product survives a restart (WAL)" "$(hit GET /products)" 200 '"name":"mug"'