writeonce/docs/examples/chat/main.wo
shoney.arickathil 4af1e8bcdd fix(chat gate): every leg starts its own server — and it found a real bug
Gate defects, all measured:

- fd check was core-count dependent: `fds_before + 8` read LAZY per-shard
  init as a leak. Shards init on first fiber, each taking one io_uring +
  one eventfd, capped at nproc; on 20 cores the first wave legitimately
  adds 18. Measured 26 -> 44 after 20 clients, still 44 after 40 more.
  Replaced with the invariant the check is for: a second wave must not
  raise the count. Core-count independent, and catches a slow leak that
  any fixed slack would hide
- a failed leg ORPHANED its server: drain inherited $SRV from the soak
  leg, so its python died on int("") and the soak server was never
  killed — its listener then broke the next run's soak on the same port.
  drain now starts its own server; cleanup kills every server a run
  started, matched on the run's unique temp dir
- two legs the plan requires were missing: WO_SHARDS=1 (the single-shard
  control) and WO_MAILBOX=8 (drop-slow-member backpressure). Both added,
  both green. The mailbox leg shrinks the slow client's SO_RCVBUF so it
  needs no sleeps
- chat adopted the porch naming (use porch/..., [deps] key) after the
  rename landed on master

Decoupling the legs exposed a REAL drain bug, traced and documented in
docs/2026-08-27-chat-drain-finding.md, NOT fixed here:

- on a FRESH server the SIGTERM drain is flaky: 5 of 16 runs left a
  client at EOF with no close frame and no diagnostic
- traced: main -> Registry -> Room -> Writer. Registry runs (diag
  confirms), the Room NEVER processes its shutdown message, so the
  Writer's close branch never runs. Clients that do get a frame are
  saved by their own Reader seeing env.stopping()
- ruled out: the spin budget (a 1s wall-clock deadline still failed 2 of
  12 — reverted, it fixed nothing and cost 1s per shutdown),
  dummy_writer() spawning during shutdown, and write failure
- the fix is an engine guarantee — a send issued before the stop flag is
  delivered — which belongs to the actor lifecycle, not a spin count

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-27 23:27:46 +02:00

335 lines
9.6 KiB
Text

-- 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 porch
use porch/http
use porch/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<Text, RoomRef>
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>&name=<who>"); }
let who = req.query["name"];
if who == nil { return bad_request("expected ?room=<name>&name=<who>"); }
-- 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>&name=<who>\"}");
}
}
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 <port>");
return 2;
}
let port = parse_int(args[0]);
if port == nil {
print_err("chat: <port> 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 });
}
}
}