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>
This commit is contained in:
parent
4092074201
commit
735fd270db
8 changed files with 809 additions and 12 deletions
335
docs/examples/chat/main.wo
Normal file
335
docs/examples/chat/main.wo
Normal file
|
|
@ -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<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 });
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
9
docs/examples/chat/wo.toml
Normal file
9
docs/examples/chat/wo.toml
Normal file
|
|
@ -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" }
|
||||||
7
justfile
7
justfile
|
|
@ -60,6 +60,13 @@ site:
|
||||||
fibers:
|
fibers:
|
||||||
./scripts/fibers-accept.sh
|
./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
|
# db-actor: arc stage 3's gate (docs/examples/db-actor) — worker-shard
|
||||||
# actors read/write the database through the transparent DB actor; WAL
|
# actors read/write the database through the transparent DB actor; WAL
|
||||||
# replay pair included. `just db-actor` runs it.
|
# replay pair included. `just db-actor` runs it.
|
||||||
|
|
|
||||||
|
|
@ -4,6 +4,7 @@
|
||||||
* 2 = usage or load failure (loader's message on stderr)
|
* 2 = usage or load failure (loader's message on stderr)
|
||||||
* Heap cap defaults to 64 MiB, overridable via WO_HEAP_MB. */
|
* Heap cap defaults to 64 MiB, overridable via WO_HEAP_MB. */
|
||||||
#include <fcntl.h>
|
#include <fcntl.h>
|
||||||
|
#include <signal.h>
|
||||||
#include <stdio.h>
|
#include <stdio.h>
|
||||||
#include <stdlib.h>
|
#include <stdlib.h>
|
||||||
#include <string.h>
|
#include <string.h>
|
||||||
|
|
@ -178,6 +179,10 @@ int main(int argc, char **argv) {
|
||||||
return 2;
|
return 2;
|
||||||
}
|
}
|
||||||
wo_tls_set(&VM);
|
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.
|
/* The database engine boots with the VM: every class IS a table.
|
||||||
* Durability is opt-in — WO_DATA=<dir> opens <dir>/shard-0.wal,
|
* Durability is opt-in — WO_DATA=<dir> opens <dir>/shard-0.wal,
|
||||||
* replays it before the entry runs (boot-before-listeners doctrine),
|
* replays it before the entry runs (boot-before-listeners doctrine),
|
||||||
|
|
|
||||||
|
|
@ -325,7 +325,27 @@ static void efd_drain(wo_vm *vm) {
|
||||||
|
|
||||||
int wo_io_wait(wo_vm *vm) {
|
int wo_io_wait(wo_vm *vm) {
|
||||||
for (;;) {
|
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) {
|
if (vm->io_kind == 0) {
|
||||||
/* keep the wake eventfd armed (oneshot POLL_ADD, re-armed
|
/* keep the wake eventfd armed (oneshot POLL_ADD, re-armed
|
||||||
* after each firing) so inbox pushes interrupt the wait */
|
* after each firing) so inbox pushes interrupt the wait */
|
||||||
|
|
|
||||||
|
|
@ -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 (stop_pending()) return WO_SYS_STOPPED;
|
||||||
}
|
}
|
||||||
if (fd < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) {
|
if (fd < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) {
|
||||||
|
if (stop_pending()) return WO_SYS_STOPPED;
|
||||||
/* arc T4: park until the listener is readable, then retry */
|
/* arc T4: park until the listener is readable, then retry */
|
||||||
vm->cur->park_fd = (int)R[B];
|
vm->cur->park_fd = (int)R[B];
|
||||||
vm->cur->park_deadline = 0;
|
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
|
/* arc T4: nothing readable yet — free the buffer (the retry
|
||||||
* re-allocates) and park until the fd is readable */
|
* re-allocates) and park until the fd is readable */
|
||||||
wo_str_free(rt, s);
|
wo_str_free(rt, s);
|
||||||
|
if (stop_pending()) return WO_SYS_STOPPED;
|
||||||
vm->cur->park_fd = (int)R[B];
|
vm->cur->park_fd = (int)R[B];
|
||||||
vm->cur->park_deadline = 0;
|
vm->cur->park_deadline = 0;
|
||||||
vm->cur->park_events = POLLIN;
|
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;
|
continue;
|
||||||
}
|
}
|
||||||
if (errno == EAGAIN || errno == EWOULDBLOCK) {
|
if (errno == EAGAIN || errno == EWOULDBLOCK) {
|
||||||
|
if (stop_pending()) return WO_SYS_STOPPED;
|
||||||
vm->cur->park_wr_at = at;
|
vm->cur->park_wr_at = at;
|
||||||
vm->cur->park_fd = (int)R[B];
|
vm->cur->park_fd = (int)R[B];
|
||||||
vm->cur->park_deadline = 0;
|
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)) {
|
if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) {
|
||||||
wo_str_free(rt, s);
|
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;
|
fb->dl_active = 0;
|
||||||
R[A] = 0; /* ?Text nil: the deadline expired */
|
R[A] = 0; /* ?Text nil: the deadline expired */
|
||||||
return 0;
|
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 (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;
|
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;
|
return 0;
|
||||||
}
|
}
|
||||||
fb->park_fd = (int)R[B];
|
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;
|
continue;
|
||||||
}
|
}
|
||||||
if (errno == EAGAIN || errno == EWOULDBLOCK) {
|
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;
|
fb->dl_active = 0;
|
||||||
R[A] = 0; /* false: torn mid-write — close the fd */
|
R[A] = 0; /* false: torn mid-write — close the fd */
|
||||||
return 0;
|
return 0;
|
||||||
|
|
|
||||||
140
runtime/src/vm.c
140
runtime/src/vm.c
|
|
@ -467,6 +467,85 @@ int wo_engine_primary_inbox(int wake_efd) {
|
||||||
return 0;
|
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) {
|
int wo_engine_start(const wo_module *mod, size_t heap_cap, uint32_t nshards) {
|
||||||
wo_eng.nshards = nshards;
|
wo_eng.nshards = nshards;
|
||||||
eng_heap_cap = heap_cap;
|
eng_heap_cap = heap_cap;
|
||||||
|
|
@ -511,6 +590,13 @@ void wo_engine_stop(void) {
|
||||||
(void)n;
|
(void)n;
|
||||||
}
|
}
|
||||||
for (uint32_t i = 1; i < wo_eng.nshards; i++) pthread_join(ts[i - 1], NULL);
|
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
|
/* single-threaded from here. Every arena dies wholesale, so routed
|
||||||
* frees and queued payloads need no per-object drops — DISCARD the
|
* frees and queued payloads need no per-object drops — DISCARD the
|
||||||
* envelopes (freeing the malloc'd nodes/actors) and let the arenas
|
* envelopes (freeing the malloc'd nodes/actors) and let the arenas
|
||||||
|
|
@ -579,19 +665,26 @@ void wo_vm_destroy(wo_vm *vm) {
|
||||||
free(fb);
|
free(fb);
|
||||||
}
|
}
|
||||||
/* actors first — dropping their state and queued messages needs the
|
/* 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;
|
wo_actor *a = vm->actors;
|
||||||
while (a) {
|
while (a) {
|
||||||
wo_actor *nx = a->next_all;
|
wo_actor *nx = a->next_all;
|
||||||
if (a->instance) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)a->instance);
|
if (drops_ok && a->instance)
|
||||||
for (uint32_t i = 0; i < a->mlen; i++) {
|
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;
|
uint64_t m = a->msgs[(a->mhead + i) % a->mcap].payload;
|
||||||
if (m) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)m);
|
if (m) wo_drop_obj(&vm->rt, (wo_hdr *)(uintptr_t)m);
|
||||||
}
|
}
|
||||||
wo_monitor *mo = a->monitors;
|
wo_monitor *mo = a->monitors;
|
||||||
while (mo) { /* undelivered notices are the runtime's to drop */
|
while (mo) { /* undelivered notices are the runtime's to drop */
|
||||||
wo_monitor *mnx = mo->next;
|
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);
|
free(mo);
|
||||||
mo = mnx;
|
mo = mnx;
|
||||||
}
|
}
|
||||||
|
|
@ -603,7 +696,8 @@ void wo_vm_destroy(wo_vm *vm) {
|
||||||
vm->timers = NULL;
|
vm->timers = NULL;
|
||||||
while (tt) { /* unfired timers likewise */
|
while (tt) { /* unfired timers likewise */
|
||||||
wo_timer *tnx = tt->next;
|
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);
|
free(tt);
|
||||||
tt = tnx;
|
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); \
|
int iorc_ = wo_io_wait(vm); \
|
||||||
if (iorc_ == WO_IO_STOP) { \
|
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); \
|
fib_reap_all(vm); \
|
||||||
vm->cur = &vm->f0; \
|
vm->cur = &vm->f0; \
|
||||||
return 1; \
|
return 1; \
|
||||||
|
|
@ -1917,8 +2020,31 @@ dispatch:
|
||||||
vm->cur->frames[vm->cur->depth - 1].pc = pc - 1;
|
vm->cur->frames[vm->cur->depth - 1].pc = pc - 1;
|
||||||
vm->cur->ncatch = 0;
|
vm->cur->ncatch = 0;
|
||||||
vm_unwind(vm, 0);
|
vm_unwind(vm, 0);
|
||||||
/* a stop ends the PROGRAM: every fiber — the stopped one,
|
/* iteration 24 (the drain): a STOPPED wait on a NON-main fiber
|
||||||
* queued ones, main wherever it is — unwinds clean */
|
* 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) {
|
if (vm->cur != &vm->f0) {
|
||||||
wo_fiber *dead = vm->cur;
|
wo_fiber *dead = vm->cur;
|
||||||
vm->cur = &vm->f0;
|
vm->cur = &vm->f0;
|
||||||
|
|
|
||||||
290
scripts/chat-accept.sh
Executable file
290
scripts/chat-accept.sh
Executable file
|
|
@ -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 ]]
|
||||||
Loading…
Reference in a new issue