Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 53 additions & 0 deletions .github/scripts/mutant-685-release-owed-fd.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
#!/usr/bin/env python3
"""celeris#685 detector control: the MUTANT, applied in CI and never committed.

Makes fdOwed (engine/iouring/fd_lifetime.go) report false for every
connection, which takes the fd-lifetime rule off every path it guards at
once: the close paths close the descriptor at once again (no kept number,
no read-side shutdown), hijackConn no longer submits its cancels before
handing the socket over, and worker shutdown no longer ends the owed ops
before closing descriptors.

Against the mutated tree every run of the five trials that judge the rule
MUST detect the theft, as its result line reports it (the CI step reads
that line, not the --- line, so a failure for another reason does not
count): TestRecvTheft715ArmA and its hole twin TestRecvTheft715ArmAHole (a
recv still in the SQ ring at the close), TestRecvTheft685Linked and
TestRecvTheft685LinkedHole (a recv linked behind a SEND, not issued at the
close), each with hit=true reused=true stolen=true, and
TestRecvTheft685HijackMultishotCoop (a multishot recv owed at a Hijack on a
ring without DEFER_TASKRUN) with hijacker_read=false. If one does not, it
is not watching the path it claims to judge, and its green run on the tree
as committed proves nothing.

The mutant also takes the rule off worker shutdown (endOwedOpsAtShutdown
ends only the ops fdOwed reports), but no trial judges that half: it is
covered by construction, not by a detector.

The function is matched EXACTLY; any drift makes this script exit 2 instead
of silently mutating nothing.
"""
import sys

PATH = sys.argv[1] if len(sys.argv) > 1 else "engine/iouring/fd_lifetime.go"

BODY = (
"func fdOwed(cs *connState) bool {\n"
"\treturn cs != nil && fdOps(cs) > 0 && !cs.fixedFile\n"
"}\n"
)

src = open(PATH).read()
if src.count(BODY) != 1:
print(f"mutant-685: expected exactly one fdOwed body, found {src.count(BODY)}", file=sys.stderr)
sys.exit(2)
mutated = src.replace(
BODY,
"func fdOwed(cs *connState) bool {\n"
"\t// MUTANT celeris#685: the fd-lifetime rule is off everywhere\n"
"\t_ = cs\n"
"\treturn false\n"
"}\n",
)
open(PATH, "w").write(mutated)
print(f"mutant-685: fdOwed always false ({PATH})")
164 changes: 164 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -1008,6 +1008,170 @@ jobs:
tally mutant kill || ok=1
exit "$ok"

recv-theft:
name: io_uring close paths never release a descriptor an op can still resolve (celeris#685, ${{ matrix.os }})
# Both arches: the theft was measured on both (celeris#715).
strategy:
fail-fast: false
matrix:
os: [ubuntu-latest, ubuntu-24.04-arm]
runs-on: ${{ matrix.os }}
timeout-minutes: 20
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
with:
persist-credentials: false
- uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0
with:
go-version: "1.27.0"
# The celeris#715 / celeris#685 trials (-tags=validation only; the
# holds are internal/recvtheft). Each parks the worker that closed a
# connection A with a recv still able to resolve A's descriptor number
# (TestRecvTheft715ArmA: a recv SQE not yet submitted;
# TestRecvTheft685Linked: a recv linked behind a SEND, not yet
# issued), lets the sibling worker accept a fresh connection B, and
# asserts that B is answered and that no stale recv read B's bytes.
# On a tree that frees the number at the close, the sibling is given
# it and the old recv reads B's request: both fail in every trial.
# Their *Hole twins open a free number under A's once the closer is
# parked: the sibling's accept of the hole tests nothing and is
# skipped, so they must fail the same way (a verdict that counted it
# passed a tree without the fix). TestRecvTheft715Control (the ring
# submitted before the close) and TestRecvTheft715ArmC (celeris#715
# hypothesis (c)) must pass either way.
# The three TestRecvTheft685Hijack* arms hold the worker right after a
# Hijack and send the hijacker's first bytes meanwhile: with multishot
# recv (owed across its request) on a ring without DEFER_TASKRUN, the
# recv reads them unless hijackConn submits its cancel first (Coop);
# the DEFER_TASKRUN and single-shot arms pass either way.
#
# The trials need TWO io_uring workers (12 MiB of RLIMIT_MEMLOCK each;
# the runner defaults to 8 MiB), so memlock is raised as in the
# `iouring` job, and CELERIS_REQUIRE_IOURING_WORKERS=1 turns a missing
# sibling into a failure. -v, -count=5, and a tally that requires
# every run of every test to PASS, with no FAIL and no SKIP line: a
# renamed test selects nothing and go test would still exit 0.
- name: celeris#685 recv-theft trials (two workers, skipping forbidden)
shell: bash
env:
CELERIS_REQUIRE_IOURING_WORKERS: "1"
CELERIS_RECV_THEFT_715: "1"
run: |
set -o pipefail
sudo prlimit --pid "$$" --memlock=unlimited:unlimited
echo "memlock (KiB): $(ulimit -l)"
names='TestRecvTheft715ArmA|TestRecvTheft715ArmAHole|TestRecvTheft715Control|TestRecvTheft715ArmC|TestRecvTheft685Linked|TestRecvTheft685LinkedHole|TestRecvTheft685HijackSingleShot|TestRecvTheft685HijackMultishotDefer|TestRecvTheft685HijackMultishotCoop'
runs=5
want=$(( $(tr '|' '\n' <<<"$names" | wc -l) * runs ))
go test -tags=validation -race -count="$runs" -timeout=600s -v -run "^(${names})\$" \
./engine/iouring/ 2>&1 | tee /tmp/recvtheft685.log || true
grep -E 'RECVTHEFT(715|685) .*result' /tmp/recvtheft685.log || true
passed=$(grep -cE "^--- PASS: (${names}) \(" /tmp/recvtheft685.log || true)
failed=$(grep -cE '^[[:space:]]*--- FAIL' /tmp/recvtheft685.log || true)
skipped=$(grep -cE '^[[:space:]]*--- SKIP' /tmp/recvtheft685.log || true)
for n in ${names//|/ }; do
echo "$n: PASS $(grep -cE "^--- PASS: $n \(" /tmp/recvtheft685.log || true) FAIL $(grep -cE "^--- FAIL: $n \(" /tmp/recvtheft685.log || true)"
done
echo "celeris#685 recv-theft trials: want $want PASS, passed $passed, FAIL lines $failed, SKIP lines $skipped"
if [ "$passed" -ne "$want" ] || [ "$failed" -ne 0 ] || [ "$skipped" -ne 0 ]; then
echo "expected exactly $want trial runs to PASS with no FAIL or SKIP line -- did a close path"
echo "release a descriptor number a recv could still resolve, or was a test renamed or skipped?"
exit 1
fi
# celeris#798: the rule counts only the ops that name the descriptor. A
# SEND_ZC's notification names none, yet a peer that stops reading holds
# it for as long as the socket is open. So it must not keep a closed
# connection's descriptor (the release backstop forced it after 5 s and
# counted CloseFDForced), nor stall worker shutdown's drain, while the
# connState, whose send buffer the notification guards, still waits for
# it. The unit job runs these two tests in its package step, on x86 only;
# here they run on both arches, by name, at the runner's own 8 MiB
# memlock (the pages SEND_ZC pins count against it), with a tally: every
# run of each of the five cases must PASS, with no FAIL and no SKIP line.
# A SKIP would mean the runner's kernel gave the engine no working
# SEND_ZC, so nothing was tested.
- name: celeris#798 a SEND_ZC notification holds no descriptor (both arches, skipping forbidden)
if: ${{ !cancelled() }}
shell: bash
env:
CELERIS_REQUIRE_IOURING_WORKERS: "1"
run: |
set -o pipefail
echo "memlock (KiB): $(ulimit -l)"
names='TestCloseReleasesDescriptorWithOnlyAZCNotificationOwed|TestShutdownDoesNotWaitForAZCNotification'
runs=3
want=$(( $(tr '|' '\n' <<<"$names" | wc -l) * runs ))
wantsub=$(( 5 * runs ))
go test -race -count="$runs" -timeout=300s -v -run "^(${names})\$" \
./engine/iouring/ 2>&1 | tee /tmp/zc798.log || true
grep -E 'celeris798 ' /tmp/zc798.log || true
passed=$(grep -cE "^--- PASS: (${names}) \(" /tmp/zc798.log || true)
subpassed=$(grep -cE "^ --- PASS: (${names})/" /tmp/zc798.log || true)
failed=$(grep -cE '^[[:space:]]*--- FAIL' /tmp/zc798.log || true)
skipped=$(grep -cE '^[[:space:]]*--- SKIP' /tmp/zc798.log || true)
echo "celeris#798 SEND_ZC tests: want $want PASS ($wantsub cases), passed $passed ($subpassed cases), FAIL lines $failed, SKIP lines $skipped"
if [ "$passed" -ne "$want" ] || [ "$subpassed" -ne "$wantsub" ] || [ "$failed" -ne 0 ] || [ "$skipped" -ne 0 ]; then
echo "expected exactly $want runs ($wantsub cases) to PASS with no FAIL or SKIP line -- did a"
echo "SEND_ZC notification hold a descriptor or stall shutdown, or was a test renamed or skipped?"
exit 1
fi
# The detector control. .github/scripts/mutant-685-release-owed-fd.py
# makes fdOwed report false, which takes the fd-lifetime rule off the
# close paths, Hijack and shutdown at once. Against that tree every run
# of the five trials that judge the rule must DETECT the theft, read
# from its result line, not from its --- line: a FAIL for any other
# reason (INCONCLUSIVE, a setup error, a race report) is not a
# detection. The close-path trials (ArmA, ArmAHole, Linked, LinkedHole)
# must report hit=true, reused=true and stolen=true: the sibling was
# given A's number and A's stale recv read B's request. The hijack
# trial (MultishotCoop) must report op_owed=+1 and hijacker_read=false.
# The trials that pass on any tree must still PASS. No trial judges the
# shutdown half of the rule (see the script). The mutation is undone
# before the step ends.
- name: celeris#685 detector control (the rule removed; every run of the five judging trials must detect the theft)
if: ${{ !cancelled() }}
shell: bash
env:
CELERIS_REQUIRE_IOURING_WORKERS: "1"
CELERIS_RECV_THEFT_715: "1"
run: |
set -o pipefail
sudo prlimit --pid "$$" --memlock=unlimited:unlimited
python3 .github/scripts/mutant-685-release-owed-fd.py
others='TestRecvTheft715Control|TestRecvTheft685HijackSingleShot|TestRecvTheft685HijackMultishotDefer'
runs=3
log=/tmp/recvtheft685-mutant.log
go test -tags=validation -race -count="$runs" -timeout=600s -v \
-run "^(TestRecvTheft715ArmA|TestRecvTheft715ArmAHole|TestRecvTheft685Linked|TestRecvTheft685LinkedHole|TestRecvTheft685HijackMultishotCoop|${others})\$" \
./engine/iouring/ 2>&1 | tee "$log" || true
git checkout -- engine/iouring/fd_lifetime.go
grep -E 'RECVTHEFT(715|685) .*result' "$log" || true
bad=0
# judged <test> <its result line> <the detection that line must report>
judged() {
local n=$1 line=$2 want=$3 lines det p
lines=$(grep -cE "$line" "$log" || true)
det=$(grep -E "$line" "$log" | grep -cE "$want" || true)
p=$(grep -cE "^--- PASS: $n \(" "$log" || true)
echo "mutant, judged $n: result lines $lines, detected $det, PASS $p (want $runs, $runs, 0)"
if [ "$lines" -ne "$runs" ] || [ "$det" -ne "$runs" ] || [ "$p" -ne 0 ]; then bad=1; fi
}
stolen=' hit=true .* reused=true .* stolen=true '
judged TestRecvTheft715ArmA 'RECVTHEFT715 arm=A result ' "$stolen"
judged TestRecvTheft715ArmAHole 'RECVTHEFT715 arm=A-hole result ' "$stolen"
judged TestRecvTheft685Linked 'RECVTHEFT685 linked result ' "$stolen"
judged TestRecvTheft685LinkedHole 'RECVTHEFT685 linked-hole result ' "$stolen"
judged TestRecvTheft685HijackMultishotCoop 'RECVTHEFT685 hijack result tier=2 ' ' op_owed=\+1 .* hijacker_read=false '
for n in ${others//|/ }; do
p=$(grep -cE "^--- PASS: $n \(" "$log" || true)
f=$(grep -cE "^--- FAIL: $n \(" "$log" || true)
echo "mutant, control $n: PASS $p FAIL $f (want PASS $runs)"
if [ "$p" -ne "$runs" ] || [ "$f" -ne 0 ]; then bad=1; fi
done
skipped=$(grep -cE '^[[:space:]]*--- SKIP' "$log" || true)
if [ "$skipped" -ne 0 ]; then echo "a SKIP is not a result ($skipped SKIP lines)"; bad=1; fi
exit "$bad"

conformance:
name: Conformance
runs-on: ubuntu-latest
Expand Down
5 changes: 5 additions & 0 deletions adaptive/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -1158,6 +1158,11 @@ func (e *Engine) Metrics() engine.EngineMetrics {
TransplantClaimDeferred: pm.TransplantClaimDeferred + sm.TransplantClaimDeferred,
TransplantReapFailed: pm.TransplantReapFailed + sm.TransplantReapFailed,
TransplantReapUnsupported: pm.TransplantReapUnsupported + sm.TransplantReapUnsupported,
// The same rule on the close paths (celeris#685): each close is an
// event on the one sub-engine that owned the connection, and a
// close of the standby's residue lands on the standby.
CloseFDDeferred: pm.CloseFDDeferred + sm.CloseFDDeferred,
CloseFDForced: pm.CloseFDForced + sm.CloseFDForced,
// The post-switch sweep (celeris#657 PR-3). Both sub-engines sweep,
// in opposite directions, and only the one draining runs passes at
// all, so the pass count sums as a rate. The residual entries are
Expand Down
4 changes: 4 additions & 0 deletions adaptive/handoff_loss_metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,13 +25,15 @@ func TestMetricsSumsTheHandoffLossWitnesses(t *testing.T) {
TransplantHeld: 5, TransplantReaps: 6, TransplantReapMisses: 7,
TransplantHoldRescued: 8, TransplantDoubleClaim: 9,
TransplantClaimDeferred: 12, TransplantReapFailed: 13, TransplantReapUnsupported: 14,
CloseFDDeferred: 15, CloseFDForced: 16,
})
e.secondary.(*mockEngine).SetMetrics(engine.EngineMetrics{
StaleRecvDataClosed: 10, StaleRecvDataTransplanted: 20,
StaleRecvDataUnattributed: 30, TransplantHandoffInFlight: 40,
TransplantHeld: 50, TransplantReaps: 60, TransplantReapMisses: 70,
TransplantHoldRescued: 80, TransplantDoubleClaim: 90,
TransplantClaimDeferred: 120, TransplantReapFailed: 130, TransplantReapUnsupported: 140,
CloseFDDeferred: 150, CloseFDForced: 160,
})

m := e.Metrics()
Expand All @@ -51,6 +53,8 @@ func TestMetricsSumsTheHandoffLossWitnesses(t *testing.T) {
{"TransplantClaimDeferred", m.TransplantClaimDeferred, 132},
{"TransplantReapFailed", m.TransplantReapFailed, 143},
{"TransplantReapUnsupported", m.TransplantReapUnsupported, 154},
{"CloseFDDeferred", m.CloseFDDeferred, 165},
{"CloseFDForced", m.CloseFDForced, 176},
} {
if c.got != c.want {
t.Errorf("Metrics().%s = %d, want %d (sum of both sub-engines)",
Expand Down
39 changes: 30 additions & 9 deletions engine/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -469,15 +469,15 @@ type EngineMetrics struct { //nolint:revive // user-approved name
// - Closed: a connection this engine closed or hijacked (the close
// paths and Hijack register the identity the same way). Usually
// the client's bytes raced a server-side close, and that client
// sees its connection end. It is not always benign, so a non-zero
// Closed does not prove that no live client lost a request. After
// a Hijack the socket lives on under the hijacker's net.Conn, so a
// recv that completes with data before its cancel lands has taken
// the first bytes of a connection its client is still using. And a
// recv that had not reached the kernel when the descriptor was
// closed resolves the descriptor NUMBER when it does; if a new
// connection holds that number by then, the recv reads that
// client's request.
// sees its connection end. It was not always benign: after a
// Hijack, a recv that completed with data before its cancel landed
// took the first bytes of a connection its client was still using,
// and a recv that had not been issued when the descriptor was
// closed resolved the descriptor NUMBER when it was, reading the
// request of whatever new connection held that number by then
// (celeris#715). Since celeris#685 neither can happen (see
// CloseFDDeferred), so what Closed counts is a closed connection's
// recv reading its own client's late bytes.
// - Transplanted: a connection this engine handed to the other
// sub-engine. A recv armed before the hand-off outlived it and
// read a request meant for the new owner — or, through a reused
Expand Down Expand Up @@ -580,6 +580,27 @@ type EngineMetrics struct { //nolint:revive // user-approved name
// probe finds the flags. io_uring-only; on the adaptive engine the sum
// over both sub-engines.
TransplantReapUnsupported uint64
// CloseFDDeferred and CloseFDForced count how the io_uring close paths
// keep the same fd-lifetime rule (celeris#685): a connection's descriptor
// NUMBER is released only when no operation that names it can still be
// issued. Closing it earlier let a receive the kernel had not issued yet
// read the request of a new connection that another thread had been
// given the freed number for; that request was dropped as
// StaleRecvDataClosed, and its client was never answered.
//
// - CloseFDDeferred: closes whose descriptor stayed open until the
// last operation the kernel owed on it had completed, and was closed
// then (normally one loop iteration later). A rate: on an
// async-handler engine it is close to one per connection the server
// closes, and on a sync-mode engine it is the server-side closes of
// connections with a receive armed (timeouts).
// - CloseFDForced: such descriptors closed by the release backstop with
// an operation still owed. Must stay 0.
//
// io_uring-only and cumulative; zero on other engines. On the adaptive
// engine each is the sum over both sub-engines.
CloseFDDeferred uint64
CloseFDForced uint64
// TransplantSweepPasses counts passes of the post-switch sweep, the
// re-examination that moves a connection the drain would otherwise
// reach only at that connection's own next event — which, for a
Expand Down
8 changes: 8 additions & 0 deletions engine/iouring/conn.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"sync/atomic"

"github.com/goceleris/celeris/internal/conn"
"github.com/goceleris/celeris/internal/recvtheft"
)

// maxSendQueueBytes is the per-connection back-pressure limit for
Expand Down Expand Up @@ -397,6 +398,13 @@ type connState struct {
// has repurposed. drainPendingRelease only releases a connState once
// this counter reaches zero (with a wall-clock backstop for kernel
// anomalies). Mirrors driverConn.inflightOps.
//
// recvArmSeq (validation builds only; zero-size in production, and not
// the last field, so it adds no padding) is the SQ ring sequence number
// of this conn's latest recv SQE: the close paths compare it with the
// kernel's SQ head to tell a recv the kernel has not consumed yet
// (celeris#715, Worker.recvUnsubmitted). Set by noteRecvPlaced.
recvArmSeq recvtheft.ArmSeq
kernelInflight int32
// recvArmed is true while a recv SQE (single-shot or multishot) is
// kernel-held for this conn. Set by prepareRecv / flushSendLink's
Expand Down
2 changes: 2 additions & 0 deletions engine/iouring/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -570,6 +570,8 @@ func (e *Engine) Metrics() engine.EngineMetrics {
TransplantClaimDeferred: e.metrics.handoffLoss.claimDeferred.Load(),
TransplantReapFailed: e.metrics.handoffLoss.reapFailed.Load(),
TransplantReapUnsupported: e.metrics.handoffLoss.reapUnsupported.Load(),
CloseFDDeferred: e.metrics.handoffLoss.closeFDDeferred.Load(),
CloseFDForced: e.metrics.handoffLoss.closeFDForced.Load(),

TransplantSweepPasses: e.metrics.sweep.passes.Load(),
TransplantResidualDetached: e.metrics.sweep.residual[resDetached].Load(),
Expand Down
Loading
Loading