diff --git a/docs/2026-08-27-chat-drain-finding.md b/docs/2026-08-27-chat-drain-finding.md new file mode 100644 index 0000000..f1268ee --- /dev/null +++ b/docs/2026-08-27-chat-drain-finding.md @@ -0,0 +1,93 @@ +# Iteration 24 T9 — the drain bug the gate was hiding + +**Found 2026-08-27** while finishing T8/T9 on branch `chat-ws-lifecycle`. +Not fixed: the fix is an engine-level decision, recorded here so it is not +rediscovered. + +## The symptom + +`just chat`'s drain leg asserts both connected clients receive a WebSocket +close frame on `SIGTERM`. Against a **fresh** server it is flaky: + +| Sample | Result | +| --- | --- | +| 5 fresh servers, 2 clients each | 4 × `close\|close`, 1 × `eof\|close` | +| 12 fresh servers | 3 failures, one of them `eof\|eof` | +| 16 fresh servers | 5 failures | + +A failing client's socket reaches EOF with **no close frame and no +diagnostic** — the process exits and the kernel closes the fd. + +## Why the gate never caught it + +The drain leg did not start its own server. It inherited `$SRV` from the soak +leg — a server the soak had already pushed 1000 clients through, so every +shard was warm and every actor already scheduled. Draining a warm server hides +the cold-start race. Fixed in this change: **every leg now starts its own +server**, which is what exposed the bug. + +## Root cause, traced + +Instrumented the sample's actors (diagnostics not committed) and correlated +against failing runs: + +1. `DIAG registry-shutdown rooms=1` — main's `send(reg, kind: 2)` **is** + delivered and the Registry runs. +2. `DIAG room-shutdown` — **never printed on a failing run.** The Room never + processes the `kind: 4` shutdown the Registry sends it. +3. The Writer's close branch never runs for the affected client, so no close + frame is written and the fd is never closed by the Writer. Its + `try net.write_dl(...)` is **not** failing — a diagnostic on that path + printed zero times. +4. A client that *does* get a close frame is usually saved by its own + **Reader** noticing `env.stopping()` and running its tail + (`DIAG reader-tail bob r2=1`), not by the room broadcast. + +So the drain chain is main → Registry → Room → Writer, three hops across +shards, and **the Room's shard does not reliably adopt its inbox before the +engine stops.** + +## What was ruled out + +- **Not the spin budget.** Replacing `spin < 20000000` with a wall-clock + deadline of 1 s (`time.ticks()`) still failed 2 of 12. More time does not + help, which is the strongest evidence the room's shard is not being + scheduled at all rather than being scheduled late. That change was reverted: + it fixed nothing and cost a fixed 1 s on every shutdown. +- **Not `dummy_writer()` spawning during shutdown.** Hoisting it to a + Registry field spawned once at startup left 5 of 16 failing. +- **Not a write failure.** See point 3. + +## The decision this needs + +`main` cannot park after the stop flag (a park unwinds), so it spins — and +spinning is not a barrier. Either: + +- **the engine drains pending inboxes before stopping**, so a `send` issued + before the stop flag is guaranteed delivered; or +- **the sample gets a real barrier** — the drain is acknowledged back to main, + which requires main to observe a reply without parking. + +The first is the honest fix and belongs to the actor lifecycle (iteration 31, +absorbed into 24). It is a semantic guarantee — "a send before shutdown is +delivered" — not a tuning parameter, and it should be stated in the runtime's +lifecycle docs and pinned by a corpus fixture, not left to a spin count. + +## Gate defects fixed alongside (all committed) + +1. **fd check was core-count dependent.** `fds_before + 8` read lazy per-shard + init as a leak: shards initialise on first fiber, each taking one + `io_uring` + one `eventfd`, capped at `nproc`. On a 20-core box the first + wave legitimately adds 18. Measured 26 → 44 after 20 clients, then **still + 44 after 40 more**. Replaced with the invariant the check is actually for: + a second wave must not raise the count. Core-count independent, and it + catches a slow leak that any fixed slack would hide. +2. **A failed leg orphaned its server.** The drain leg's python died on + `int("")` when `$SRV` was empty, so the soak server was never killed and + its listener broke the *next* run's soak on the same port. `cleanup` now + kills every server a run started, matched on the run's unique temp dir. +3. **Two legs the plan requires were missing** — `WO_SHARDS=1` (the + single-shard control that says a failure is placement's fault) and + `WO_MAILBOX=8` (the drop-slow-member backpressure path). Both added, both + green. The mailbox leg manufactures a genuinely slow member by shrinking + its `SO_RCVBUF`, so it needs no sleeps. diff --git a/docs/examples/chat/main.wo b/docs/examples/chat/main.wo index 2721878..23993ce 100644 --- a/docs/examples/chat/main.wo +++ b/docs/examples/chat/main.wo @@ -16,9 +16,9 @@ use env use net use time -use framework -use framework/http -use framework/router +use porch +use porch/http +use porch/router -- ---- message types (one per actor) -------------------------------------- diff --git a/docs/examples/chat/wo.toml b/docs/examples/chat/wo.toml index c58a51a..4b8be1c 100644 --- a/docs/examples/chat/wo.toml +++ b/docs/examples/chat/wo.toml @@ -6,4 +6,4 @@ description = "Iteration 24's acceptance workload: rooms + presence + broadcast wo = ">= 0.1" [deps] -framework = { git = "https://github.com/shoneyj/writeonce-framework", rev = "v0.1.0" } +porch = { git = "https://github.com/shoneyj/porch", rev = "v0.1.0" } diff --git a/scripts/chat-accept.sh b/scripts/chat-accept.sh index f471471..e4a4b48 100755 --- a/scripts/chat-accept.sh +++ b/scripts/chat-accept.sh @@ -28,16 +28,21 @@ ulimit -n 8192 2>/dev/null || true W="$(mktemp -d "${TMPDIR:-/tmp}/chat-accept.XXXXXX")" SRV="" cleanup() { + # kill EVERY server this run started, not merely the most recent $SRV: a leg + # that dies before clearing SRV used to orphan a listener, which then broke + # the next run on the same port. $W is unique per run, so matching on it + # cannot touch another run's processes. [[ -n "$SRV" ]] && kill -9 "$SRV" 2>/dev/null + pkill -9 -f "$W/app/target/chat" 2>/dev/null rm -rf "$W" } trap cleanup EXIT -cp -r "$ROOT/docs/examples/writeonce-framework" "$W/fw" +cp -r "$ROOT/docs/examples/porch" "$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" +sed -i "s|https://github.com/shoneyj/porch|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 @@ -53,8 +58,17 @@ 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) +def connect(port, room, name, timeout=8, rcvbuf=None): + # rcvbuf: shrink THIS client's receive buffer so the server's socket fills + # quickly — how the WO_MAILBOX leg manufactures a genuinely slow member + # without sleeping. Must be set before connect() to take effect. + if rcvbuf is None: + s = socket.create_connection(("127.0.0.1", port), timeout=timeout) + else: + s = socket.socket() + s.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, rcvbuf) + s.settimeout(timeout) + s.connect(("127.0.0.1", port)) 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" @@ -149,6 +163,7 @@ 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)" +fds_prev=99999 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]) @@ -223,20 +238,60 @@ got="${r%%|*}"; rest="${r#*|}"; n="${rest%%|*}"; el="${rest#*|}" && 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_w1="$(ls /proc/$SRV/fd 2>/dev/null | wc -l)" + [[ "$fds_w1" -le "$fds_prev" ]] && break + fds_prev="$fds_w1" + sleep 0.5 +done +# The fd check is for a per-CONNECTION leak, and a fixed tolerance cannot +# express that. Shards initialise LAZILY (runtime/src/vm.c: a worker's vm is +# not paid for until its first fiber arrives), so the first wave legitimately +# adds one io_uring + one eventfd PER SHARD, capped at nproc — on a 20-core +# box that is +18, which the old `fds_before + 8` read as a leak. Measured +# 2026-08-27: 26 -> 44 after 20 clients, then still 44 after 40 more. +# +# So assert the invariant itself: a SECOND wave must not raise the count. +# Core-count independent, and it catches a slow leak that any fixed +# tolerance would hide inside its own slack. +timeout 60 python3 - "$PORT" 20 <<'PYEOF' >/dev/null 2>&1 +import socket, base64, os, sys, time +port, n = int(sys.argv[1]), int(sys.argv[2]) +socks = [] +for i in range(n): + s = socket.create_connection(("127.0.0.1", port), timeout=8) + k = base64.b64encode(os.urandom(16)).decode() + s.sendall((f"GET /ws?room=fdwave&name=w{i} HTTP/1.1\r\nhost: a\r\n" + f"upgrade: websocket\r\nconnection: Upgrade\r\n" + f"sec-websocket-key: {k}\r\nsec-websocket-version: 13\r\n\r\n").encode()) + h = b"" + while b"\r\n\r\n" not in h: + h += s.recv(4096) + socks.append(s) +time.sleep(0.5) +for s in socks: + s.close() +PYEOF +fds_after="$fds_w1" for _ in $(seq 1 20); do fds_after="$(ls /proc/$SRV/fd 2>/dev/null | wc -l)" - [[ "$fds_after" -le $((fds_before + 8)) ]] && break + [[ "$fds_after" -le "$fds_w1" ]] && 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" +[[ "$fds_after" -le "$fds_w1" ]] \ + && ok "no per-connection fd leak (start $fds_before, after $SOAK_N: $fds_w1, after 20 more: $fds_after)" \ + || bad "soak-fds" "second wave grew fds: $fds_w1 -> $fds_after (start $fds_before)" [[ -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 ---- +# Starts its OWN server. It used to inherit the soak leg's $SRV, which meant +# any leg inserted between them silently handed drain an empty pid: its python +# died on int(""), the leg reported a bare failure, AND the soak server was +# never killed — orphaning a listener that then broke the NEXT run's soak on +# the same port. No leg may depend on another leg's server. +serve "$((PORT0 + 6))" || bad "serve-drain" "no listener" r="$(timeout 30 python3 - "$PORT" "$SRV" <<'PYEOF' import sys, os, time, signal, socket import importlib.util @@ -266,6 +321,75 @@ for _ in $(seq 1 40); do kill -0 "$SRV" 2>/dev/null || { stopped=0; break; }; sl [[ $stopped -eq 0 ]] && ok "SIGTERM exits 0" || bad "stop" "still running" SRV="" +# ---- 4b. WO_SHARDS=1: the same matrix on one shard ---- +# The plan requires `just chat` green at default cores AND on a single shard: +# cross-shard placement is where the actor work can hide a bug, so the +# one-shard run is the control that says a failure is placement's fault. +serve "$((PORT0 + 4))" env WO_SHARDS=1 || bad "serve-shards1" "no listener" +r="$(functional shards1)"; [[ "$r" == *functional-ok* ]] \ + && ok "WO_SHARDS=1: the same matrix on a single shard" \ + || bad "shards1-functional" "$r" +kill -TERM "$SRV" 2>/dev/null; wait "$SRV" 2>/dev/null; SRV="" + +# ---- 4c. WO_MAILBOX=8: the drop-slow-member path FIRES and the room lives ---- +# The backpressure policy earning its keep. A member that stops reading makes +# its writer block on write_dl; with the mailbox capped at 8 the room's +# broadcast send traps (WO_T_ACTOR), and the room must CATCH that, drop the +# member, and keep serving everyone else. Asserting the room survives is the +# point — a room that dies with its slowest member is the bug this policy +# exists to prevent. +serve "$((PORT0 + 5))" env WO_MAILBOX=8 || bad "serve-mailbox" "no listener" +r="$(timeout 90 python3 - "$PORT" <<'PYEOF' +import importlib.util, os, socket, sys, time +spec = importlib.util.spec_from_file_location("wsc", os.environ["WSC"]) +wsc = importlib.util.module_from_spec(spec); spec.loader.exec_module(wsc) +port = int(sys.argv[1]) + +fast = wsc.connect(port, "bp", "fast") +wsc.recv(fast) # * fast joined +# the slow member: a tiny receive buffer so the server's socket fills fast, +# and it never reads a single frame +slow = wsc.connect(port, "bp", "slow", rcvbuf=2048) +wsc.recv(fast) # * slow joined + +# storm: big frames the slow member never drains +blob = "x" * 1024 +for i in range(400): + try: + wsc.send(fast, f"{i}-{blob}") + except OSError: + break +# drain what fast owes us so its own mailbox cannot be the thing that fills +deadline = time.time() + 20 +seen = 0 +while time.time() < deadline: + try: + k, t = wsc.recv(fast, timeout=0.5) + seen += 1 + except Exception: + break + +# the room must still be alive and serving the fast member +survivor = wsc.connect(port, "bp", "late") +ok_join = False +deadline = time.time() + 15 +while time.time() < deadline: + try: + k, t = wsc.recv(fast, timeout=1.0) + if "late joined" in t: + ok_join = True + break + except Exception: + break +print("mailbox-ok" if ok_join else f"mailbox-dead seen={seen}") +slow.close(); fast.close(); survivor.close() +PYEOF +)" +[[ "$r" == *mailbox-ok* ]] \ + && ok "WO_MAILBOX=8: slow member dropped, room survived and kept serving" \ + || bad "mailbox-backpressure" "$r" +kill -TERM "$SRV" 2>/dev/null; wait "$SRV" 2>/dev/null; SRV="" + # ---- 5. the ASan leg: functional matrix, zero leaks ---- if [[ -x "$ASAN" ]]; then sed -i "s|runtime = \".*\"|runtime = \"$ASAN\"|" "$W/app/wo.toml"