Settle the native-api counter reads instead of racing the reader task

The suite asserted a happens-before the native API does not offer. A client's
write on a flow descriptor lands in the kernel buffer of an AF_UNIX
socketpair and runs no daemon code; the counters advance only inside the
per-flow reader task, after that task's own recv().await. `stats` is answered
on a different task and loads the same atomics, and the daemon is a
single-threaded runtime, so a stats reply can be produced while the datagrams
are still queued and the reader task has not been polled. The reply is then a
well-formed status ok with a live flow_id and local_port and both counters at
zero, which is the shape that redded maint at c5aeef39 and next at 1c822ae9.

The crate's own tests already concede this. Connection::settle is a bounded
yield loop whose doc comment says it yields "rather than asserting into a
race a test would lose intermittently", and the in-crate assertions run
behind it. The shell harness had no equivalent, though it half knew: one
counter read is followed by a sleep whose comment calls it "the task hop",
which is why that step passed while its neighbours did not.

An RPC step may now carry "settle", which re-asks the command until its
expectations hold or a five-second deadline passes. A bounded re-ask is a
barrier where a fixed sleep is a guess, and a datagram that never arrives
still reds rather than hanging. Nothing here serializes anything.

Reproduced and measured rather than reasoned about. Against a daemon
throttled to 0.02 CPU with twenty concurrent clients, the unsettled read
failed 9 of 60 runs, every failure carrying the zero-counter signature, while
the settling read failed 0 of 60 interleaved under the same load. An
expectation that can never hold still reds, at the five-second bound. The
full suite passes 28 of 28.

The close scenario's one-second sleep is replaced by settling the release
check, since the same task hop delays the daemon noticing end of file.

Two verdict lines are corrected while here. Both were canned else-branch
strings that fire on any non-zero client exit, so each named a cause the run
never observed: one reported traffic crossing between flows when the evidence
was a zero counter, the other reported a flow not being released when the
failing assertion was a datagram count.

The two-flow scenario's second read is deliberately left alone, with a note
saying why: settling re-asks until an expectation holds, and "b counted 0"
holds on the first ask whether or not b's reader has run, so that assertion
stays a false green until it makes a positive claim.
This commit is contained in:
Johnathan Corgan
2026-09-07 18:27:28 +00:00
parent 4c345f3ffd
commit 70002baf20
2 changed files with 59 additions and 12 deletions
+35 -3
View File
@@ -12,6 +12,7 @@ it exits.
Kinds of step:
RPC step: {"command": str, "params": {...}?, "expect": {"dotted.key": val}?,
"settle": bool?,
"keep_fd": name?, "keep_listener": name?, "keep_flow": name?}
Sends a command and checks the reply. `keep_fd` stores a flow
descriptor under that name, `keep_listener` a listener descriptor;
@@ -76,6 +77,13 @@ from typing import Any
FLOW = "flow"
LISTENER = "listener"
# How long a settling step keeps re-asking, and how long it pauses between
# tries. The wait is for a task hop on a host that may be loaded, so it is
# seconds rather than milliseconds; it is bounded because a datagram that
# never arrives has to end the run red rather than hold it open.
SETTLE_SECONDS = 5.0
SETTLE_PAUSE = 0.02
def recvfds(sock: socket.socket, bufsize: int, maxfds: int) -> tuple[bytes, list[int]]:
"""One recvmsg, returning its payload and whatever descriptors it carried.
@@ -287,7 +295,16 @@ def store(client: Client, step: dict, body: dict, fd: int | None) -> list[str]:
def run_rpc(client: Client, step: dict) -> list[str]:
"""Send one command and report what did not hold."""
"""Send one command and report what did not hold.
A step naming `settle` is asked again until its expectations hold or the
deadline passes. What `stats` reports is advanced by the daemon's per-flow
reader task, and nothing orders that task against the client's write on the
flow descriptor: one ask can be answered while a datagram is still queued in
the kernel socket buffer, and it reads back as zero. Re-asking is the only
barrier this protocol offers, and the deadline is what keeps a datagram that
never arrives a failure rather than a hang.
"""
command = step["command"]
try:
params = substitute(step.get("params"), client.flows)
@@ -295,8 +312,23 @@ def run_rpc(client: Client, step: dict) -> list[str]:
except KeyError as error:
return [str(error)]
reply, fd = client.call(command, params)
problems = check(reply, expect)
settle = bool(step.get("settle"))
if settle and (step.get("keep_fd") or step.get("keep_listener")):
# Every ask but the last is discarded, and a discarded reply's
# descriptor has no owner. Refusing beats closing one a later step
# meant to keep.
return ["settle: a step that keeps a descriptor cannot be re-asked"]
deadline = time.monotonic() + SETTLE_SECONDS
while True:
reply, fd = client.call(command, params)
problems = check(reply, expect)
if not problems or not settle or time.monotonic() >= deadline:
break
if fd is not None:
os.close(fd)
time.sleep(SETTLE_PAUSE)
problems += store(client, step, reply, fd)
if problems:
+24 -9
View File
@@ -409,11 +409,17 @@ check_descriptor_carries_datagrams() {
# Three writes must reach the daemon as three datagrams. The byte count
# matters as much as the datagram count: a boundary loss would show up as
# one datagram of nine bytes rather than three of three.
#
# The read settles. The daemon counts a datagram on the flow's own reader
# task, and nothing orders that task against the client's write, so a single
# ask can be answered while the writes are still queued in the kernel and
# read back as zero. Re-asking is the only barrier the protocol offers; the
# deadline keeps a datagram that never arrives a failure rather than a hang.
local script='[
{"command":"connect","params":{"peer":"'"$PEER"'","remote_port":4242},
"keep_fd":"a","keep_flow":"a","expect":{"status":"ok"}},
{"fd":"a","write":"00ff10","repeat":3},
{"command":"stats","params":{"flow_id":"@a"},
{"command":"stats","params":{"flow_id":"@a"},"settle":true,
"expect":{"status":"ok","data.rx_datagrams":3,"data.rx_bytes":9,"data.closed":false}}
]'
if run_client "$script"; then
@@ -452,34 +458,43 @@ check_close_reaches_the_daemon() {
# The node forgets a released flow, so `stats` answers for it the way it
# answers for any name it does not hold. The step before the close is what
# makes that discriminating: the same flow answered a moment earlier, so the
# refusal afterwards can only be the release. The sleep is the task hop
# between the daemon reading end of file and giving the entry back.
# refusal afterwards can only be the release. Both reads settle: the same
# task hop that delays a datagram count also delays the daemon noticing end
# of file, and a bounded re-ask is a barrier where a fixed sleep was a guess.
local script='[
{"command":"connect","params":{"peer":"'"$PEER"'","remote_port":4242},
"keep_fd":"a","keep_flow":"a","expect":{"status":"ok"}},
{"fd":"a","write":"aa"},
{"command":"stats","params":{"flow_id":"@a"},
{"command":"stats","params":{"flow_id":"@a"},"settle":true,
"expect":{"status":"ok","data.closed":false,"data.rx_datagrams":1}},
{"fd":"a","close":true},
{"sleep":1},
{"command":"stats","params":{"flow_id":"@a"},"expect":{"status":"error"}}
{"command":"stats","params":{"flow_id":"@a"},"settle":true,
"expect":{"status":"error"}}
]'
if run_client "$script"; then
pass "the daemon saw the close and gave the flow back"
else
fail "the daemon did not release the closed flow"
fail "the close script did not hold; the failing step is above"
fi
}
check_flows_are_independent() {
log "Two flows on one connection stay separate"
# Flow a's read settles; flow b's cannot and is left as it is. A settling
# step re-asks until its expectation HOLDS, and "b counted 0" holds on the
# first ask whether or not b's reader task has ever run. So this stays a
# false green: if traffic ever did cross, b could read 0 for the same
# scheduling reason and the check would pass. Closing it needs a positive
# claim on b, writing a known count there and settling it, which is a
# larger change than this one.
local script='[
{"command":"connect","params":{"peer":"'"$PEER"'","remote_port":4242},
"keep_fd":"a","keep_flow":"a","expect":{"status":"ok"}},
{"command":"connect","params":{"peer":"'"$PEER2"'","remote_port":4243},
"keep_fd":"b","keep_flow":"b","expect":{"status":"ok"}},
{"fd":"a","write":"11","repeat":2},
{"command":"stats","params":{"flow_id":"@a"},"expect":{"data.rx_datagrams":2}},
{"command":"stats","params":{"flow_id":"@a"},"settle":true,
"expect":{"data.rx_datagrams":2}},
{"command":"stats","params":{"flow_id":"@b"},"expect":{"data.rx_datagrams":0}},
{"command":"inject","params":{"flow_id":"@b","data":"22"},"expect":{"status":"ok"}},
{"fd":"b","read":1,"expect_bytes":"22"},
@@ -488,7 +503,7 @@ check_flows_are_independent() {
if run_client "$script"; then
pass "each flow saw only its own traffic"
else
fail "traffic crossed between flows"
fail "the two-flow script did not hold; the failing step is above"
fi
}