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>
335 lines
9.6 KiB
Text
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 });
|
|
}
|
|
}
|
|
}
|