writeonce/docs/examples/chat/main.wo
shoney.arickathil 735fd270db feat: chat sample + gate (T8/T9, IN PROGRESS) + stop-drain semantics
- 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 <noreply@anthropic.com>
2026-08-23 09:53:21 +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 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<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 });
}
}
}