From 735fd270db9efb02e982f4dc2283a4a8e4d0916a Mon Sep 17 00:00:00 2001 From: "shoney.arickathil" Date: Sun, 23 Aug 2026 09:53:21 +0200 Subject: [PATCH] feat: chat sample + gate (T8/T9, IN PROGRESS) + stop-drain semantics MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - docs/examples/chat: registry (call consumer) / room / reader+writer actor pair per connection over ws_accept + wsframe; presence, broadcast, cross-room isolation, mailbox-full = drop-from-room; reader tail sends hardened (a full writer no longer orphans the fd) - RUNTIME SEMANTICS CHANGE (the drain): SIGTERM no longer kills parked fibers from outside — the plane WAKES them and each wait RESOLVES (deadline'd waits answer their timeout result, sleeps return early, plain waits answer WO_SYS_STOPPED and unwind THAT fiber alone; main's STOPPED still ends the program). Workers keep adopting their inboxes after stop until eng_shutdown. This is what lets a program drain: chat's close frames now reach clients (byte-verified 0x88), then main returns and the reap runs - also: SIGPIPE ignored process-wide (EPIPE trap instead of death); two-phase engine teardown (real drops while arenas+routing live, settle passes for routed frees) — fixes the registry-map leak and the drain UAF ASan found - gate scripts/chat-accept.sh + just chat: handshake independently verified, functional matrix on BOTH backends, 1k-hot-room soak (1000/1000 in ~35ms), drain close-frames, SIGTERM exit 0, ASan leg clean. OPEN: soak-fds check (18 fds settle slower than the window) + full battery after the semantics change — NOT yet run - committed for manual testing at the user's request Co-Authored-By: Claude Fable 5 --- docs/examples/chat/main.wo | 335 +++++++++++++++++++++++++++++++++++++ docs/examples/chat/wo.toml | 9 + justfile | 7 + runtime/src/main.c | 5 + runtime/src/park.c | 22 ++- runtime/src/sysio.c | 13 +- runtime/src/vm.c | 140 +++++++++++++++- scripts/chat-accept.sh | 290 ++++++++++++++++++++++++++++++++ 8 files changed, 809 insertions(+), 12 deletions(-) create mode 100644 docs/examples/chat/main.wo create mode 100644 docs/examples/chat/wo.toml create mode 100755 scripts/chat-accept.sh diff --git a/docs/examples/chat/main.wo b/docs/examples/chat/main.wo new file mode 100644 index 0000000..2721878 --- /dev/null +++ b/docs/examples/chat/main.wo @@ -0,0 +1,335 @@ +-- chat — iteration 24's acceptance workload. Rooms, presence and +-- broadcast over WebSocket: every connection is a reader actor (sole fd +-- reader) plus a writer actor (sole fd writer); rooms and the registry +-- are actors; delivery between them is ownership-moving sends, across +-- shards when placement lands them there. One binary, no broker. +-- +-- CHAT_TOKEN is not needed — chat is open; the framework serves it +-- through [deps] exactly like web-app: +-- woc . && ./target/chat 8080 +-- ws://127.0.0.1:8080/ws?room=lobby&name=alice +-- +-- The actor split exists because an actor takes ONE message at a time: +-- a single per-connection actor blocked in net read could never hear a +-- broadcast. The reader owns the socket's inbound half and the carry +-- buffer; the writer owns the outbound half so frames never interleave. +use env +use net +use time +use framework +use framework/http +use framework/router + +-- ---- message types (one per actor) -------------------------------------- + +-- To a writer: 1 = text frame, 2 = close (frame + fd close), 3 = pong. +class WriterMsg { + kind: Int + text: Text +} + +-- To a room: 1 = join, 2 = leave, 3 = text, 4 = shutdown (drain). +class RoomMsg { + kind: Int + name: Text + text: Text + writer: actor WriterMsg +} + +-- To the registry: 1 = lookup (a `call` — the reply is the room's +-- address), 2 = shutdown every room (a `send` on SIGTERM). +class Lookup { + kind: Int + room: Text +} + +-- To a reader: everything the connection's inbound loop needs. +class ReaderMsg { + fd: net.Conn + room: actor RoomMsg + writer: actor WriterMsg + name: Text +} + +-- One connection accepted, one worker: builds its own App and runs the +-- framework's keep-alive loop (the serving-slice pattern). +class Conn { + fd: net.Conn +} + +-- ---- the writer: sole owner of the outbound half ------------------------- + +class Writer { + fd: net.Conn + dead: Int + fn receive(msg: WriterMsg) { + if self.dead == 1 { return; } + if msg.kind == 1 { + let ok = try net.write_dl(self.fd, ws_text(msg.text), 2000) catch (e) false; + if ok == false { + -- a stalled or gone client: tear the fd; the reader will see EOF + -- and route the leave through the room + self.dead = 1; + net.close(self.fd); + } + return; + } + if msg.kind == 3 { + let ok2 = try net.write_dl(self.fd, ws_pong(msg.text), 2000) catch (e) false; + if ok2 == false { + self.dead = 1; + net.close(self.fd); + } + return; + } + -- close: the drain path (room shutdown or reader-detected close) + self.dead = 1; + let ig = try net.write_dl(self.fd, ws_close(), 1000) catch (e) false; + net.close(self.fd); + } +} + +-- ---- the room: members, presence, fan-out -------------------------------- + +class Mem { + w: actor WriterMsg + name: Text +} + +class Room { + members: multi Mem + fn receive(msg: RoomMsg) { + if msg.kind == 1 { + push(self.members, Mem { w: msg.writer, name: "${msg.name}" }); + self.say("* ${msg.name} joined"); + return; + } + if msg.kind == 2 { + let keep: multi Mem = []; + while len(self.members) > 0 { + let m = shift(self.members); + if m.name != msg.name { push(keep, m); } + } + self.members = keep; + self.say("* ${msg.name} left"); + return; + } + if msg.kind == 3 { + self.say("${msg.name}: ${msg.text}"); + return; + } + -- shutdown: every member gets a close frame; the list empties + while len(self.members) > 0 { + let m = shift(self.members); + let r = try send_close(m.w) catch (e) 0; + } + } + + -- fan-out one line; a member whose mailbox is FULL is a slow client — + -- the fail-fast cap turns it into a drop-from-the-room (the backpressure + -- policy earning its keep) + fn say(line: Text) { + let keep: multi Mem = []; + while len(self.members) > 0 { + let m = shift(self.members); + let ok = try send_text(m.w, "${line}") catch (e) 0; + if ok == 1 { + push(keep, m); + } else { + let r = try send_close(m.w) catch (e) 0; + } + } + self.members = keep; + } +} + +-- send wrappers: `try` is an expression, so give it Int results +fn send_text(w: actor WriterMsg, line: Text) -> Int { + send(w, WriterMsg { kind: 1, text: line }); + return 1; +} + +fn send_close(w: actor WriterMsg) -> Int { + send(w, WriterMsg { kind: 2, text: "" }); + return 1; +} + +-- ---- the registry: name -> room, spawn on demand -------------------------- + +class RoomRef { + r: actor RoomMsg +} + +class Registry { + rooms: map + fallback: actor RoomMsg + fn receive(msg: Lookup) -> actor RoomMsg { + if msg.kind == 2 { + for k, v in self.rooms { + send(v.r, RoomMsg { kind: 4, name: "", text: "", writer: dummy_writer() }); + } + return self.fallback; + } + if has(self.rooms, msg.room) == 1 { + let have = self.rooms[msg.room]; + if have != nil { + return have.r; + } + } + let room: actor RoomMsg = spawn Room { members: [] }; + self.rooms[msg.room] = RoomRef { r: room }; + return room; + } +} + +-- RoomMsg requires a writer field on every construction; the shutdown +-- message has no meaningful one, so a throwaway satisfies the shape (it +-- never receives anything — kind 4 reads no fields). +fn dummy_writer() -> actor WriterMsg { + let w: actor WriterMsg = spawn Writer { fd: 0 - 1, dead: 1 }; + return w; +} + +-- ---- the reader: sole owner of the inbound half --------------------------- + +class Reader { + pad: Int + fn receive(msg: ReaderMsg) { + let carry = ""; + let alive = true; + while alive { + if env.stopping() { alive = false; continue; } + let got = try net.read_dl(msg.fd, 4096, 30000) catch (e) nil; + if got == nil { + -- idle deadline or I/O trap: this client is done + alive = false; + continue; + } + let bytes = "${got}"; + if len(bytes) == 0 { + alive = false; + continue; + } + carry = carry .. bytes; + let more = true; + while more { + let f = ws_parse(carry); + if f.kind == 0 { + more = false; + continue; + } + carry = f.rest; + if f.kind == 1 { + send(msg.room, RoomMsg { kind: 3, name: "${msg.name}", text: f.payload, writer: msg.writer }); + continue; + } + if f.kind == 9 { + send(msg.writer, WriterMsg { kind: 3, text: f.payload }); + continue; + } + if f.kind == 10 or f.kind == 2 { + continue; -- pongs ignored; binary tolerated (echo is not chat) + } + -- close frame or protocol error: stop reading + alive = false; + more = false; + } + } + -- the tail sends must survive full mailboxes (a leave storm after a + -- mass close): a trap here would kill the reader and orphan the fd + let r1 = try send_leave(msg.room, "${msg.name}", msg.writer) catch (e) 0; + let r2 = try send_close(msg.writer) catch (e) 0; + if r2 == 0 { + -- the writer is unreachable (full/dead): close the fd ourselves + net.close(msg.fd); + } + } +} + +fn send_leave(room: actor RoomMsg, name: Text, w: actor WriterMsg) -> Int { + send(room, RoomMsg { kind: 2, name: name, text: "", writer: w }); + return 1; +} + +-- ---- HTTP: the upgrade route + usage -------------------------------------- + +class WsRoute { + reg: actor Lookup + fn handle(req: Req) -> Resp { + if ws_upgrade_valid(req) == false { + return bad_request("expected a websocket upgrade"); + } + let rname = req.query["room"]; + if rname == nil { return bad_request("expected ?room=&name="); } + let who = req.query["name"]; + if who == nil { return bad_request("expected ?room=&name="); } + -- the cross-shard call: this handler runs on the connection worker's + -- shard, the registry lives wherever placement put it + let room = call(self.reg, Lookup { kind: 1, room: "${rname}" }); + let fd = ws_accept(req); + let w: actor WriterMsg = spawn Writer { fd: fd, dead: 0 }; + let rd: actor ReaderMsg = spawn Reader { pad: 0 }; + send(room, RoomMsg { kind: 1, name: "${who}", text: "", writer: w }); + send(rd, ReaderMsg { fd: fd, room: room, writer: w, name: "${who}" }); + return hijacked(); + } +} + +class Usage { + pad: Int + fn handle(req: Req) -> Resp { + return ok_json("{\"ws\":\"/ws?room=&name=\"}"); + } +} + +fn build_app(reg: actor Lookup) -> App { + let app = App { middleware: [], routes: [] }; + app.get("/", Usage { pad: 0 }); + app.get("/ws", WsRoute { reg: reg }); + return app; +} + +class ConnWorker { + reg: actor Lookup + fn receive(msg: Conn) { + let app = build_app(self.reg); + app.handle_conn(msg.fd, 10000, 10000); + } +} + +fn main(args: multi Text) -> Int { + if len(args) < 1 { + print_err("usage: chat "); + return 2; + } + let port = parse_int(args[0]); + if port == nil { + print_err("chat: must be a number"); + return 2; + } + let fb: actor RoomMsg = spawn Room { members: [] }; + let reg: actor Lookup = spawn Registry { rooms: {}, fallback: fb }; + let srv = net.listen("127.0.0.1", port); + print("listening on 127.0.0.1:${port}"); + while true { + if env.stopping() { + -- the drain: every room broadcasts a close frame and writers flush. + -- main must NOT park here (a park after the stop flag unwinds), so + -- it SPINS — each loop back-edge pays a reduction, and the budget + -- hands the shard to the draining actors between slices; worker + -- shards keep adopting their inboxes until the engine stops. + send(reg, Lookup { kind: 2, room: "" }); + let spin = 0; + while spin < 20000000 { + spin = spin + 1; + } + net.close(srv); + return 0; + } + let c = net.accept_dl(srv, 250); + if c != nil { + let w: actor Conn = spawn ConnWorker { reg: reg }; + send(w, Conn { fd: c }); + } + } +} diff --git a/docs/examples/chat/wo.toml b/docs/examples/chat/wo.toml new file mode 100644 index 0000000..c58a51a --- /dev/null +++ b/docs/examples/chat/wo.toml @@ -0,0 +1,9 @@ +name = "chat" +version = "0.1.0" +description = "Iteration 24's acceptance workload: rooms + presence + broadcast over WebSocket — actors on fibers across shards, one binary, no broker" + +[runtime] +wo = ">= 0.1" + +[deps] +framework = { git = "https://github.com/shoneyj/writeonce-framework", rev = "v0.1.0" } diff --git a/justfile b/justfile index 9efc7c4..67a80f4 100644 --- a/justfile +++ b/justfile @@ -60,6 +60,13 @@ site: fibers: ./scripts/fibers-accept.sh +# chat: iteration 24's gate (docs/examples/chat) — rooms/presence/broadcast +# over WebSocket via actors: functional on both WO_IO backends, the +# 1k-clients-one-hot-room soak (fds/RSS accounted), SIGTERM drain with +# close frames, and an ASan leg. `just chat` runs it (CHAT_SOAK=N trims). +chat: + ./scripts/chat-accept.sh + # db-actor: arc stage 3's gate (docs/examples/db-actor) — worker-shard # actors read/write the database through the transparent DB actor; WAL # replay pair included. `just db-actor` runs it. diff --git a/runtime/src/main.c b/runtime/src/main.c index 28e8b6f..01dd4e5 100644 --- a/runtime/src/main.c +++ b/runtime/src/main.c @@ -4,6 +4,7 @@ * 2 = usage or load failure (loader's message on stderr) * Heap cap defaults to 64 MiB, overridable via WO_HEAP_MB. */ #include +#include #include #include #include @@ -178,6 +179,10 @@ int main(int argc, char **argv) { return 2; } wo_tls_set(&VM); + /* iteration 24: a write to a peer-closed socket must be EPIPE (a + * catchable WO_T_IO), never a process-killing SIGPIPE — every + * serving program writes to sockets whose peers vanish. */ + signal(SIGPIPE, SIG_IGN); /* The database engine boots with the VM: every class IS a table. * Durability is opt-in — WO_DATA= opens /shard-0.wal, * replays it before the entry runs (boot-before-listeners doctrine), diff --git a/runtime/src/park.c b/runtime/src/park.c index 599a4b1..af4345a 100644 --- a/runtime/src/park.c +++ b/runtime/src/park.c @@ -325,7 +325,27 @@ static void efd_drain(wo_vm *vm) { int wo_io_wait(wo_vm *vm) { for (;;) { - if (wo_sys_stop_pending()) return WO_IO_STOP; + if (wo_sys_stop_pending()) { + /* iteration 24 (the drain): a STOP does not kill parked fibers + * from the outside — it WAKES them all, and each blocking + * builtin resolves per its own stop contract (deadline'd waits + * answer their timeout result, sleeps return early, plain + * waits answer WO_SYS_STOPPED and that fiber unwinds). The + * program's own code then drains and returns. Nothing parked + * = nothing to resolve: the old immediate-stop answer. */ + int woke = 0; + wo_fiber *fb = vm->parked; + while (fb) { + wo_fiber *nx = fb->pnext; + if (fb->state == WO_FIB_PARKED) { + wake(vm, fb); + woke = 1; + } + fb = nx; + } + if (woke) return 0; + return WO_IO_STOP; + } if (vm->io_kind == 0) { /* keep the wake eventfd armed (oneshot POLL_ADD, re-armed * after each firing) so inbox pushes interrupt the wait */ diff --git a/runtime/src/sysio.c b/runtime/src/sysio.c index 5379963..7d0ca08 100644 --- a/runtime/src/sysio.c +++ b/runtime/src/sysio.c @@ -395,6 +395,7 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { if (stop_pending()) return WO_SYS_STOPPED; } if (fd < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { + if (stop_pending()) return WO_SYS_STOPPED; /* arc T4: park until the listener is readable, then retry */ vm->cur->park_fd = (int)R[B]; vm->cur->park_deadline = 0; @@ -431,6 +432,7 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { /* arc T4: nothing readable yet — free the buffer (the retry * re-allocates) and park until the fd is readable */ wo_str_free(rt, s); + if (stop_pending()) return WO_SYS_STOPPED; vm->cur->park_fd = (int)R[B]; vm->cur->park_deadline = 0; vm->cur->park_events = POLLIN; @@ -476,6 +478,7 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { continue; } if (errno == EAGAIN || errno == EWOULDBLOCK) { + if (stop_pending()) return WO_SYS_STOPPED; vm->cur->park_wr_at = at; vm->cur->park_fd = (int)R[B]; vm->cur->park_deadline = 0; @@ -534,7 +537,9 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { } if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { wo_str_free(rt, s); - if (fb->dl_at > 0 && dnow >= fb->dl_at) { + if (stop_pending() || (fb->dl_at > 0 && dnow >= fb->dl_at)) { + /* iteration 24: a STOP resolves the wait as its timeout + * result — the program's own drain code decides what next */ fb->dl_active = 0; R[A] = 0; /* ?Text nil: the deadline expired */ return 0; @@ -584,9 +589,9 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { } } if (fd < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) { - if (fb->dl_at > 0 && dnow >= fb->dl_at) { + if (stop_pending() || (fb->dl_at > 0 && dnow >= fb->dl_at)) { fb->dl_active = 0; - R[A] = WO_NIL_SCALAR; /* ?Int nil: nothing arrived */ + R[A] = WO_NIL_SCALAR; /* ?Int nil: nothing arrived (or stop) */ return 0; } fb->park_fd = (int)R[B]; @@ -631,7 +636,7 @@ int wo_builtin_sys(wo_vm *vm, uint64_t *R, uint32_t ins, const char **msg) { continue; } if (errno == EAGAIN || errno == EWOULDBLOCK) { - if (fb->dl_at > 0 && dnow >= fb->dl_at) { + if (stop_pending() || (fb->dl_at > 0 && dnow >= fb->dl_at)) { fb->dl_active = 0; R[A] = 0; /* false: torn mid-write — close the fd */ return 0; diff --git a/runtime/src/vm.c b/runtime/src/vm.c index 344ede9..86244d0 100644 --- a/runtime/src/vm.c +++ b/runtime/src/vm.c @@ -467,6 +467,85 @@ int wo_engine_primary_inbox(int wake_efd) { return 0; } +/* iteration 24 teardown phase 1 (single-threaded, BEFORE eng_teardown): + * dismantle one vm's actor world with real drops — container backings are + * malloc'd, so wholesale arena death does NOT cover them (LSan, chat's + * registry map). Cross-shard payloads route home through wo_route_free + * (still live here); the routed kind-2 envelopes are settled by the + * caller's inbox passes. */ +static void vm_drop_actor_world(wo_vm *vm) { + wo_actor *a = vm->actors; + vm->actors = NULL; + while (a) { + wo_actor *nx = a->next_all; + if (a->instance) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)a->instance); + for (uint32_t i = 0; i < a->mlen; i++) { + uint64_t m = a->msgs[(a->mhead + i) % a->mcap].payload; + if (m) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)m); + } + wo_monitor *mo = a->monitors; + while (mo) { + wo_monitor *mnx = mo->next; + if (mo->msg) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)mo->msg); + free(mo); + mo = mnx; + } + free(a->msgs); + free(a); + a = nx; + } + wo_timer *tt = vm->timers; + vm->timers = NULL; + while (tt) { + wo_timer *tnx = tt->next; + if (tt->msg) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)tt->msg); + free(tt); + tt = tnx; + } +} + +/* Settle every inbox after phase 1: home-routed frees execute on their + * owner vm; payload-carrying strays drop (possibly routing again — the + * outer loop runs until everything is quiet). Node memory always freed. */ +static int eng_settle_inboxes(void) { + int moved = 0; + for (uint32_t i = 0; i < wo_eng.nshards && i < WO_ENG_MAX_SHARDS; i++) { + if (!INBOX_READY[i]) continue; + wo_vm *vm = &wo_eng.shards[i]; + wo_inbox *ib = &INBOX[i]; + wo_envelope *e = ib->head; + ib->head = ib->tail = NULL; + while (e) { + wo_envelope *nx = e->next; + switch (e->kind) { + case 2: /* WE are home: the direct drop is the settlement */ + wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)e->payload); + break; + case 0: + case 5: + case 7: /* in-flight payloads: drop (may route -> next pass) */ + if (e->payload) + wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)e->payload); + break; + case 1: /* an unadopted actor shell */ + if (e->actor) { + if (e->actor->instance) + wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)e->actor->instance); + free(e->actor->msgs); + free(e->actor); + } + break; + default: /* 3/4/6: scalar or engine-side payloads, node-only */ + break; + } + free(e); + moved++; + e = nx; + } + } + return moved; +} + int wo_engine_start(const wo_module *mod, size_t heap_cap, uint32_t nshards) { wo_eng.nshards = nshards; eng_heap_cap = heap_cap; @@ -511,6 +590,13 @@ void wo_engine_stop(void) { (void)n; } for (uint32_t i = 1; i < wo_eng.nshards; i++) pthread_join(ts[i - 1], NULL); + /* single-threaded from here: PHASE 1 — real drops while every arena + * and the routing fabric are still alive (malloc'd container backings + * inside actor state need them; iteration 24's registry map). Settle + * passes run until routed frees stop appearing. */ + for (uint32_t i = 0; i < wo_eng.nshards && i < WO_ENG_MAX_SHARDS; i++) + if (wo_eng.shards[i].rt.arena.base) vm_drop_actor_world(&wo_eng.shards[i]); + while (eng_settle_inboxes() > 0) {} /* single-threaded from here. Every arena dies wholesale, so routed * frees and queued payloads need no per-object drops — DISCARD the * envelopes (freeing the malloc'd nodes/actors) and let the arenas @@ -579,19 +665,26 @@ void wo_vm_destroy(wo_vm *vm) { free(fb); } /* actors first — dropping their state and queued messages needs the - * runtime alive */ + * runtime alive. BUT: once the engine is in teardown, arenas die + * WHOLESALE (the standing doctrine) — a moved-in message's home arena + * may belong to an ALREADY-destroyed shard, and even reading its + * header is a use-after-free (ASan, chat's drain). Structures are + * still freed; payload drops are skipped. */ + int drops_ok = !eng_teardown; wo_actor *a = vm->actors; while (a) { wo_actor *nx = a->next_all; - if (a->instance) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)a->instance); - for (uint32_t i = 0; i < a->mlen; i++) { + if (drops_ok && a->instance) + wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)a->instance); + for (uint32_t i = 0; drops_ok && i < a->mlen; i++) { uint64_t m = a->msgs[(a->mhead + i) % a->mcap].payload; if (m) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)m); } wo_monitor *mo = a->monitors; while (mo) { /* undelivered notices are the runtime's to drop */ wo_monitor *mnx = mo->next; - if (mo->msg) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)mo->msg); + if (drops_ok && mo->msg) + wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)mo->msg); free(mo); mo = mnx; } @@ -603,7 +696,8 @@ void wo_vm_destroy(wo_vm *vm) { vm->timers = NULL; while (tt) { /* unfired timers likewise */ wo_timer *tnx = tt->next; - if (tt->msg) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)tt->msg); + if (drops_ok && tt->msg) + wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)tt->msg); free(tt); tt = tnx; } @@ -1415,6 +1509,15 @@ static int vm_run(wo_vm *vm, uint64_t *ret, wo_err *err) { } \ int iorc_ = wo_io_wait(vm); \ if (iorc_ == WO_IO_STOP) { \ + /* iteration 24: a WORKER on stop keeps DRAINING — its \ + * serve loop spins adopting the inbox until the primary \ + * finishes the drain window and sets eng_shutdown, so \ + * queued shutdown messages (close frames!) still run. \ + * Only the PRIMARY's stop ends the program. */ \ + if (!vm->is_primary) { \ + vm->cur = &vm->f0; \ + return 2; \ + } \ fib_reap_all(vm); \ vm->cur = &vm->f0; \ return 1; \ @@ -1917,8 +2020,31 @@ dispatch: vm->cur->frames[vm->cur->depth - 1].pc = pc - 1; vm->cur->ncatch = 0; vm_unwind(vm, 0); - /* a stop ends the PROGRAM: every fiber — the stopped one, - * queued ones, main wherever it is — unwinds clean */ + /* iteration 24 (the drain): a STOPPED wait on a NON-main fiber + * unwinds that fiber ALONE — the rest of the program (main's + * drain code, actors flushing close frames) keeps running. + * Main's own STOPPED still ends the program, as ever. */ + if (vm->cur != &vm->f0) { + wo_fiber *dead = vm->cur; + if (dead->actor) { + wo_actor *da = dead->actor; + if (dead->cur_msg) { + wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)dead->cur_msg); + dead->cur_msg = 0; + } + call_reply_to(vm, dead->msg_caller, dead->msg_caller_shard, + 0, WO_T_ACTOR); + dead->msg_caller = NULL; + da->active = NULL; + } + vm->nfibers--; + fib_retire(vm, dead); + NEXT_RUNNABLE(); + RELOAD(); + NEXT(); + } + /* main: a stop ends the PROGRAM — every remaining fiber + * unwinds clean */ if (vm->cur != &vm->f0) { wo_fiber *dead = vm->cur; vm->cur = &vm->f0; diff --git a/scripts/chat-accept.sh b/scripts/chat-accept.sh new file mode 100755 index 0000000..f471471 --- /dev/null +++ b/scripts/chat-accept.sh @@ -0,0 +1,290 @@ +#!/usr/bin/env bash +# scripts/chat-accept.sh — iteration 24's gate. The chat sample serves +# WebSocket rooms through the framework ([deps], file:// remote); a raw +# RFC 6455 python client (stdlib only, INDEPENDENT accept-key check) +# proves: the handshake, broadcast + presence + isolation across rooms, +# the 1k-clients-one-hot-room soak (fds/RSS accounted), and the SIGTERM +# drain (close frames, exit 0) — functional legs on BOTH WO_IO backends +# plus an ASan run. +set -uo pipefail + +ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +WOC="$ROOT/compiler/_build/default/bin/woc" +WOVM="$ROOT/runtime/wovm" +ASAN="$ROOT/runtime/build/wovm_asan" +PORT0="${CHAT_PORT:-18901}" +PORT="$PORT0" +SOAK_N="${CHAT_SOAK:-1000}" + +pass=0; fail=0 +ok() { echo "ok $1"; pass=$((pass + 1)); } +bad() { echo "FAIL $1 -- $2"; fail=$((fail + 1)); } + +if [[ ! -x "$WOC" || ! -x "$WOVM" ]]; then + echo "chat-accept: build woc and wovm first" >&2; exit 1 +fi +ulimit -n 8192 2>/dev/null || true + +W="$(mktemp -d "${TMPDIR:-/tmp}/chat-accept.XXXXXX")" +SRV="" +cleanup() { + [[ -n "$SRV" ]] && kill -9 "$SRV" 2>/dev/null + rm -rf "$W" +} +trap cleanup EXIT + +cp -r "$ROOT/docs/examples/writeonce-framework" "$W/fw" +git -C "$W/fw" init -q && git -C "$W/fw" add -A +git -C "$W/fw" -c user.email=t@t -c user.name=t commit -qm v01 && git -C "$W/fw" tag v0.1.0 +cp -r "$ROOT/docs/examples/chat" "$W/app" +sed -i "s|https://github.com/shoneyj/writeonce-framework|file://$W/fw|" "$W/app/wo.toml" +printf '[build]\nruntime = "%s"\n' "$WOVM" >> "$W/app/wo.toml" + +if "$WOC" "$W/app" >"$W/build.out" 2>&1 && [[ -x "$W/app/target/chat" ]]; then + ok "deps chain + build" +else + bad "build" "$(grep -m1 error "$W/build.out" || head -1 "$W/build.out")" + echo "chat-accept: 1 checks, 1 failures"; exit 1 +fi + +# the raw client, shared by every leg +CLIENT="$W/wsc.py" +cat > "$CLIENT" <<'PYEOF' +import socket, base64, hashlib, os, time +GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11" +BUF = {} +def connect(port, room, name, timeout=8): + s = socket.create_connection(("127.0.0.1", port), timeout=timeout) + key = base64.b64encode(os.urandom(16)).decode() + s.sendall((f"GET /ws?room={room}&name={name} HTTP/1.1\r\nhost: a\r\n" + f"upgrade: websocket\r\nconnection: Upgrade\r\n" + f"sec-websocket-key: {key}\r\nsec-websocket-version: 13\r\n\r\n").encode()) + d = b"" + while b"\r\n\r\n" not in d: d += s.recv(2000) + head, _, rest = d.partition(b"\r\n\r\n") + BUF[s] = rest # a frame may already ride the same segment + head = head.decode() + assert " 101 " in head.splitlines()[0], head.splitlines()[0] + want = base64.b64encode(hashlib.sha1((key + GUID).encode()).digest()).decode() + assert want in head, "accept-key mismatch (independent check)" + return s +def _take(s, n, timeout): + s.settimeout(timeout) + b = BUF.get(s, b"") + while len(b) < n: + c = s.recv(4096) + if not c: + BUF[s] = b + return None + b += c + BUF[s] = b[n:] + return b[:n] +def send(s, text): + p = text.encode(); mask = os.urandom(4) + if len(p) < 126: hdr = bytes([0x81, 0x80 | len(p)]) + else: hdr = bytes([0x81, 0x80 | 126, len(p) >> 8, len(p) & 255]) + s.sendall(hdr + mask + bytes(b ^ mask[i % 4] for i, b in enumerate(p))) +def recv(s, timeout=5): + h = _take(s, 2, timeout) + if h is None: return (-2, "") # EOF + b0, b1 = h[0], h[1] + ln = b1 & 0x7F + if ln == 126: + e = _take(s, 2, timeout); ln = (e[0] << 8) | e[1] + d = _take(s, ln, timeout) if ln else b"" + return (b0 & 0x0F), (d or b"").decode(errors="replace") +PYEOF + +serve() { # serve PORT [env...] — start + wait for THIS server's listener line + PORT="$1"; shift + : > "$W/srv.out" # stale 'listening' lines from an earlier leg lie + "$@" "$W/app/target/chat" "$PORT" >>"$W/srv.out" 2>&1 & + SRV=$! + for _ in $(seq 1 80); do grep -q listening "$W/srv.out" 2>/dev/null && return 0; sleep 0.1; done + return 1 +} + +functional() { # $1 = leg name + timeout 30 python3 - "$PORT" <<'PYEOF' +import sys; sys.path.insert(0, sys.argv[0].rsplit("/",1)[0]) +port = int(sys.argv[1]) +import importlib.util, os +spec = importlib.util.spec_from_file_location("wsc", os.environ["WSC"]) +wsc = importlib.util.module_from_spec(spec); spec.loader.exec_module(wsc) +a = wsc.connect(port, "lobby", "alice") +assert wsc.recv(a) == (1, "* alice joined") +b = wsc.connect(port, "lobby", "bob") +assert wsc.recv(a) == (1, "* bob joined") +assert wsc.recv(b) == (1, "* bob joined") +c = wsc.connect(port, "other", "carol") +assert wsc.recv(c) == (1, "* carol joined") +wsc.send(a, "hello room") +assert wsc.recv(a) == (1, "alice: hello room") +assert wsc.recv(b) == (1, "alice: hello room") +import socket +try: + k, t = wsc.recv(c, timeout=0.8); assert False, f"leak into other room: {t}" +except socket.timeout: pass +b.close() +k, t = wsc.recv(a) +assert (k, t) == (1, "* bob left"), (k, t) +a.close(); c.close() +print("functional-ok") +PYEOF +} + +# ---- 2. functional on both backends ---- +export WSC="$CLIENT" +serve "$((PORT0 + 0))" env WO_IO=uring || bad "serve-uring" "no listener" +r="$(functional uring)"; [[ "$r" == *functional-ok* ]] \ + && ok "uring: handshake(key verified) + presence + broadcast + isolation + leave" \ + || bad "uring-functional" "$r" +kill -TERM "$SRV" 2>/dev/null; wait "$SRV" 2>/dev/null; SRV="" + +serve "$((PORT0 + 1))" env WO_IO=epoll || bad "serve-epoll" "no listener" +r="$(functional epoll)"; [[ "$r" == *functional-ok* ]] \ + && ok "epoll: the same matrix" || bad "epoll-functional" "$r" +kill -TERM "$SRV" 2>/dev/null; wait "$SRV" 2>/dev/null; SRV="" + +# ---- 3. the soak: N clients, ONE hot room ---- +serve "$((PORT0 + 2))" || bad "serve-soak" "no listener" +fds_before="$(ls /proc/$SRV/fd 2>/dev/null | wc -l)" +r="$(timeout 180 python3 - "$PORT" "$SOAK_N" <<'PYEOF' +import asyncio, sys, os, time, base64, hashlib +port, N = int(sys.argv[1]), int(sys.argv[2]) +GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11" +MARK = "the-hot-room-marker" +sem = asyncio.Semaphore(100) +async def client(i, results): + async with sem: + r, w = await asyncio.open_connection("127.0.0.1", port) + key = base64.b64encode(os.urandom(16)).decode() + w.write((f"GET /ws?room=hot&name=c{i} HTTP/1.1\r\nhost: a\r\n" + f"upgrade: websocket\r\nconnection: Upgrade\r\n" + f"sec-websocket-key: {key}\r\nsec-websocket-version: 13\r\n\r\n").encode()) + await w.drain() + d = b"" + while b"\r\n\r\n" not in d: d += await r.read(2000) + if i == 0: + # the sender: wait for the herd, then one marker line + await asyncio.sleep(0) + results["sender_ready"].set() + try: + buf = b"" + deadline = time.time() + 150 + while time.time() < deadline: + try: + c = await asyncio.wait_for(r.read(8192), timeout=5) + except asyncio.TimeoutError: + if results["sent"].is_set(): break + continue + if not c: break + buf += c + # scan frames for the marker (server frames are unmasked, small) + if MARK.encode() in buf: + results["got"] += 1 + return + finally: + w.close() +async def main(): + results = {"got": 0, "sender_ready": asyncio.Event(), "sent": asyncio.Event()} + conns = [] + # keep the sender's socket outside the tasks: join first + sr, sw = None, None + async def sender(): + nonlocal sr, sw + async with sem: + sr, sw = await asyncio.open_connection("127.0.0.1", port) + key = base64.b64encode(os.urandom(16)).decode() + sw.write((f"GET /ws?room=hot&name=sender HTTP/1.1\r\nhost: a\r\n" + f"upgrade: websocket\r\nconnection: Upgrade\r\n" + f"sec-websocket-key: {key}\r\nsec-websocket-version: 13\r\n\r\n").encode()) + await sw.drain() + d = b"" + while b"\r\n\r\n" not in d: d += await sr.read(2000) + await sender() + tasks = [asyncio.create_task(client(i, results)) for i in range(N)] + await asyncio.sleep(max(2.0, N / 250)) # let the herd join + drain presence + p = MARK.encode(); mask = os.urandom(4) + hdr = bytes([0x81, 0x80 | len(p)]) + sw.write(hdr + mask + bytes(b ^ mask[i % 4] for i, b in enumerate(p))) + await sw.drain() + results["sent"].set() + t0 = time.time() + await asyncio.gather(*tasks, return_exceptions=True) + el = int((time.time() - t0) * 1000) + sw.close() + print(f"{results['got']}|{N}|{el}") +asyncio.run(main()) +PYEOF +)" +got="${r%%|*}"; rest="${r#*|}"; n="${rest%%|*}"; el="${rest#*|}" +[[ "$got" == "$n" ]] \ + && ok "soak: the marker reached all $got/$n hot-room clients (${el}ms after send)" \ + || bad "soak" "$r" +# leave-broadcast storms take a moment to settle after 1k closes +fds_after=99999 +for _ in $(seq 1 20); do + fds_after="$(ls /proc/$SRV/fd 2>/dev/null | wc -l)" + [[ "$fds_after" -le $((fds_before + 8)) ]] && break + sleep 0.5 +done +rss_kb="$(awk '/VmRSS/{print $2}' /proc/$SRV/status 2>/dev/null)" +[[ "$fds_after" -le $((fds_before + 8)) ]] \ + && ok "soak fds came home ($fds_before -> $fds_after)" \ + || bad "soak-fds" "$fds_before -> $fds_after" +[[ -n "$rss_kb" && "$rss_kb" -lt 819200 ]] \ + && ok "soak RSS bounded (${rss_kb}KB < 800MB)" || bad "soak-rss" "${rss_kb}KB" + +# ---- 4. drain: SIGTERM with clients connected -> close frames, exit 0 ---- +r="$(timeout 30 python3 - "$PORT" "$SRV" <<'PYEOF' +import sys, os, time, signal, socket +import importlib.util +spec = importlib.util.spec_from_file_location("wsc", os.environ["WSC"]) +wsc = importlib.util.module_from_spec(spec); spec.loader.exec_module(wsc) +port, srv = int(sys.argv[1]), int(sys.argv[2]) +a = wsc.connect(port, "lobby", "alice"); wsc.recv(a) +b = wsc.connect(port, "lobby", "bob"); wsc.recv(a); wsc.recv(b) +os.kill(srv, signal.SIGTERM) +def drained(s): + try: + while True: + k, _ = wsc.recv(s, timeout=5) + if k == 8: return "close-frame" + if k == -2: return "eof" + except socket.timeout: + return "stuck" + except (ConnectionResetError, BrokenPipeError): + return "reset" +print(drained(a) + "|" + drained(b)) +PYEOF +)" +[[ "$r" == "close-frame|close-frame" ]] \ + && ok "drain: both clients got the close frame" || bad "drain" "$r" +stopped=1 +for _ in $(seq 1 40); do kill -0 "$SRV" 2>/dev/null || { stopped=0; break; }; sleep 0.1; done +[[ $stopped -eq 0 ]] && ok "SIGTERM exits 0" || bad "stop" "still running" +SRV="" + +# ---- 5. the ASan leg: functional matrix, zero leaks ---- +if [[ -x "$ASAN" ]]; then + sed -i "s|runtime = \".*\"|runtime = \"$ASAN\"|" "$W/app/wo.toml" + rm -rf "$W/app/target" + "$WOC" "$W/app" >/dev/null 2>&1 + serve "$((PORT0 + 3))" || bad "serve-asan" "no listener" + r="$(functional asan)" + kill -TERM "$SRV" 2>/dev/null + for _ in $(seq 1 60); do kill -0 "$SRV" 2>/dev/null || break; sleep 0.1; done + SRV="" + if [[ "$r" == *functional-ok* ]] && ! grep -q "AddressSanitizer\|LeakSanitizer" "$W/srv.out"; then + ok "ASan run clean (functional + drain, zero leaks)" + else + bad "asan" "$(grep -m1 -E 'ERROR|SUMMARY' "$W/srv.out" || echo "$r")" + fi +else + bad "asan" "runtime/build/wovm_asan missing — make -C runtime wovm-asan" +fi + +echo +printf 'chat-accept: %d checks, %d failures\n' "$((pass + fail))" "$fail" +[[ $fail -eq 0 ]]