fix(chat gate): every leg starts its own server — and it found a real bug
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>
This commit is contained in:
parent
ebc3522c40
commit
4af1e8bcdd
4 changed files with 230 additions and 13 deletions
93
docs/2026-08-27-chat-drain-finding.md
Normal file
93
docs/2026-08-27-chat-drain-finding.md
Normal file
|
|
@ -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.
|
||||
|
|
@ -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) --------------------------------------
|
||||
|
||||
|
|
|
|||
|
|
@ -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" }
|
||||
|
|
|
|||
|
|
@ -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):
|
||||
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"
|
||||
|
|
|
|||
Loading…
Reference in a new issue