-- 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 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 }); } } }