feat: iteration 35 — net seams + the serving slice (fiber-per-connection)
- 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 <noreply@anthropic.com>
This commit is contained in:
parent
40127bf53e
commit
bafd532afa
19 changed files with 737 additions and 31 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 |
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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 |
|
||||
|
|
|
|||
|
|
@ -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/) —
|
||||
|
|
|
|||
|
|
@ -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
|
||||
102
docs/superpowers/specs/2026-08-23-net-seams-park-design.md
Normal file
102
docs/superpowers/specs/2026-08-23-net-seams-park-design.md
Normal file
|
|
@ -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.
|
||||
|
|
@ -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).
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -28,6 +28,7 @@
|
|||
#include <string.h>
|
||||
#include <sys/socket.h>
|
||||
#include <sys/stat.h>
|
||||
#include <sys/un.h>
|
||||
#include <sys/wait.h>
|
||||
#include <time.h>
|
||||
#include <unistd.h>
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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 */
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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"'
|
||||
|
|
|
|||
Loading…
Reference in a new issue