Skip to content
37 changes: 37 additions & 0 deletions .github/scripts/mutant-812-backstop-releases-zc.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
#!/usr/bin/env python3
"""celeris#812 detector control: the MUTANT, applied in CI and never committed.

Disables the hold in drainPendingRelease (engine/iouring/worker.go): past the
release backstop, an entry that still owes a SEND_ZC is released like any
other, to connStatePool with its send buffer (a detached one to the GC), as
before the fix. Everything else of the fix stays: the accounting, the
counters and the shutdown retention.

Against the mutated tree every run of every case of
TestBackstopHoldsASendBufferAZCNotificationStillReads MUST detect the release:
its result line reports held_past_backstop=false and corrupt_bytes above 0
(the peer read the array's next owner's bytes), and the case FAILs. If one
does not, the test is not watching what it claims to on this runner (a kernel
that copied the send's pages before the release would pass it with or without
the fix), and its green run on the tree as committed proves nothing.

The statement 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/worker.go"

HOLD = "\t\t\tif w.closedZCOwed(cs) {\n\t\t\t\tw.holdZCPastBackstop(entry)\n"

src = open(PATH).read()
if src.count(HOLD) != 1:
print(f"mutant-812: expected exactly one backstop hold, found {src.count(HOLD)}", file=sys.stderr)
sys.exit(2)
mutated = src.replace(
HOLD,
"\t\t\t// MUTANT celeris#812: the backstop releases a SEND_ZC's send buffer\n"
"\t\t\tif false && w.closedZCOwed(cs) {\n\t\t\t\tw.holdZCPastBackstop(entry)\n",
)
open(PATH, "w").write(mutated)
print(f"mutant-812: backstop hold disabled ({PATH})")
70 changes: 70 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -1123,6 +1123,76 @@ jobs:
echo "SEND_ZC notification hold a descriptor or stall shutdown, or was a test renamed or skipped?"
exit 1
fi
# celeris#812: the release backstop gave up a closed connection 5 s after
# the close whatever it still owed, to connStatePool with its send buffer.
# With a SEND_ZC's notification owed, the kernel was still sending the
# unsent tail from that array to a peer that had stopped reading, and the
# array's next owner overwrote it: the peer got another connection's
# bytes. The backstop now holds such a send buffer until the notification,
# off the release walk, keeping the array alone, in gauges, and with no
# new SEND_ZC armed while a worker holds 16 MiB; worker shutdown keeps the
# buffers the ring closes on. Both arches, by name, at the runner's own
# 8 MiB memlock (the pages SEND_ZC pins count against it), -race, three
# runs, with a tally: every run of each of the ten tests and their
# twenty-four cases must PASS, with no FAIL and no SKIP line. A SKIP
# would mean the runner gave the engine no working SEND_ZC.
- name: celeris#812 a SEND_ZC's send buffer outlives the release backstop (both arches, skipping forbidden)
if: ${{ !cancelled() }}
shell: bash
env:
CELERIS_REQUIRE_IOURING_WORKERS: "1"
run: |
set -o pipefail
echo "memlock (KiB): $(ulimit -l)"
names='TestBackstopHoldsASendBufferAZCNotificationStillReads|TestBackstopCountsItsZCHolds|TestPendingReleaseBackstopHoldsAZCSendPastItsDeadline|TestClosedOpsCountsTheZCSendApart|TestDroppingAnIdentityWithAZCOwedIsCounted|TestShutdownRetainsSendBuffersAZCMayStillRead|TestZCHoldIsOffTheReleaseWalk|TestZCHoldKeepsOnlyTheSendBuffer|TestZCHoldCoversACollidingSibling|TestSendZCStopsWhileHeldBuffersReachTheCap'
runs=3
want=$(( $(tr '|' '\n' <<<"$names" | wc -l) * runs ))
wantsub=$(( 24 * runs ))
go test -race -count="$runs" -timeout=300s -v -run "^(${names})\$" \
./engine/iouring/ 2>&1 | tee /tmp/zc812.log || true
grep -E 'celeris812 ' /tmp/zc812.log || true
passed=$(grep -cE "^--- PASS: (${names}) \(" /tmp/zc812.log || true)
subpassed=$(grep -cE "^ --- PASS: (${names})/" /tmp/zc812.log || true)
failed=$(grep -cE '^[[:space:]]*--- FAIL' /tmp/zc812.log || true)
skipped=$(grep -cE '^[[:space:]]*--- SKIP' /tmp/zc812.log || true)
echo "celeris#812 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 the backstop"
echo "give up a send buffer a SEND_ZC still reads, or was a test renamed or skipped?"
exit 1
fi
# Its detector control. .github/scripts/mutant-812-backstop-releases-zc.py
# disables the hold, so the backstop releases the connState with its send
# buffer as before the fix. Every run of each of the four wire cases must
# then report the release (held_past_backstop=false) and bytes on the wire
# that were not sent (corrupt_bytes above 0), and FAIL. A runner whose
# kernel copied the send before the release would let the wire test pass
# with or without the fix; this is what shows it does not. The mutation is
# undone before the step ends.
- name: celeris#812 detector control (the hold removed; every run of the four wire cases must see foreign bytes)
if: ${{ !cancelled() }}
shell: bash
env:
CELERIS_REQUIRE_IOURING_WORKERS: "1"
run: |
set -o pipefail
python3 .github/scripts/mutant-812-backstop-releases-zc.py
runs=3
log=/tmp/zc812-mutant.log
go test -race -count="$runs" -timeout=300s -v -run '^TestBackstopHoldsASendBufferAZCNotificationStillReads$' \
./engine/iouring/ 2>&1 | tee "$log" || true
git checkout -- engine/iouring/worker.go
bad=0
for c in notif-only recv-and-notif send-done-after-close detached-notif-only; do
lines=$(grep -cE "celeris812 backstop case=$c " "$log" || true)
det=$(grep -E "celeris812 backstop case=$c " "$log" | grep -cE ' held_past_backstop=false .* corrupt_bytes=[1-9]' || true)
f=$(grep -cE "^ --- FAIL: TestBackstopHoldsASendBufferAZCNotificationStillReads/$c \(" "$log" || true)
echo "mutant, case $c: result lines $lines, detected $det, FAIL $f (want $runs, $runs, $runs)"
if [ "$lines" -ne "$runs" ] || [ "$det" -ne "$runs" ] || [ "$f" -ne "$runs" ]; 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"
# 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
Expand Down
10 changes: 10 additions & 0 deletions adaptive/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -1177,6 +1177,16 @@ func (e *Engine) Metrics() engine.EngineMetrics {
// close of the standby's residue lands on the standby.
CloseFDDeferred: pm.CloseFDDeferred + sm.CloseFDDeferred,
CloseFDForced: pm.CloseFDForced + sm.CloseFDForced,
// The send buffers a SEND_ZC may still read (celeris#812): a hold,
// a forced release and a retention at shutdown are each an event
// on the one sub-engine whose connection it was.
CloseZCNotifHeld: pm.CloseZCNotifHeld + sm.CloseZCNotifHeld,
CloseZCNotifForced: pm.CloseZCNotifForced + sm.CloseZCNotifForced,
ShutdownZCBufRetained: pm.ShutdownZCBufRetained + sm.ShutdownZCBufRetained,
// Gauges of what each sub-engine's workers hold right now: a held
// buffer is held by one of them, so the engine holds the sum.
CloseZCNotifHeldNow: pm.CloseZCNotifHeldNow + sm.CloseZCNotifHeldNow,
CloseZCNotifHeldBytes: pm.CloseZCNotifHeldBytes + sm.CloseZCNotifHeldBytes,
// 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
9 changes: 9 additions & 0 deletions adaptive/handoff_loss_metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,8 @@ func TestMetricsSumsTheHandoffLossWitnesses(t *testing.T) {
TransplantHoldRescued: 8, TransplantDoubleClaim: 9,
TransplantClaimDeferred: 12, TransplantReapFailed: 13, TransplantReapUnsupported: 14,
CloseFDDeferred: 15, CloseFDForced: 16,
CloseZCNotifHeld: 17, CloseZCNotifForced: 18, ShutdownZCBufRetained: 19,
CloseZCNotifHeldNow: 21, CloseZCNotifHeldBytes: 22,
})
e.secondary.(*mockEngine).SetMetrics(engine.EngineMetrics{
StaleRecvDataClosed: 10, StaleRecvDataTransplanted: 20,
Expand All @@ -34,6 +36,8 @@ func TestMetricsSumsTheHandoffLossWitnesses(t *testing.T) {
TransplantHoldRescued: 80, TransplantDoubleClaim: 90,
TransplantClaimDeferred: 120, TransplantReapFailed: 130, TransplantReapUnsupported: 140,
CloseFDDeferred: 150, CloseFDForced: 160,
CloseZCNotifHeld: 170, CloseZCNotifForced: 180, ShutdownZCBufRetained: 190,
CloseZCNotifHeldNow: 210, CloseZCNotifHeldBytes: 220,
})

m := e.Metrics()
Expand All @@ -55,6 +59,11 @@ func TestMetricsSumsTheHandoffLossWitnesses(t *testing.T) {
{"TransplantReapUnsupported", m.TransplantReapUnsupported, 154},
{"CloseFDDeferred", m.CloseFDDeferred, 165},
{"CloseFDForced", m.CloseFDForced, 176},
{"CloseZCNotifHeld", m.CloseZCNotifHeld, 187},
{"CloseZCNotifForced", m.CloseZCNotifForced, 198},
{"ShutdownZCBufRetained", m.ShutdownZCBufRetained, 209},
{"CloseZCNotifHeldNow", m.CloseZCNotifHeldNow, 231},
{"CloseZCNotifHeldBytes", m.CloseZCNotifHeldBytes, 242},
} {
if c.got != c.want {
t.Errorf("Metrics().%s = %d, want %d (sum of both sub-engines)",
Expand Down
40 changes: 40 additions & 0 deletions engine/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -609,6 +609,46 @@ type EngineMetrics struct { //nolint:revive // user-approved name
// engine each is the sum over both sub-engines.
CloseFDDeferred uint64
CloseFDForced uint64
// CloseZCNotifHeld, CloseZCNotifHeldNow, CloseZCNotifHeldBytes,
// CloseZCNotifForced and ShutdownZCBufRetained show how the io_uring
// engine keeps a SEND_ZC's send buffer for as long as the kernel may read
// it (celeris#812). A zero-copy send leaves its unsent part queued on the
// socket as references to the buffer's pages, and a peer that stops
// reading keeps it there, after a close too, until the peer reads or the
// kernel gives up on the socket; the kernel then sends it from whatever
// the buffer holds. Released to be reused before that, the buffer
// delivered another connection's bytes to that peer. So the release
// backstop, 5 s after a close, holds such a buffer until the kernel says
// it is done, for as long as the peer keeps the socket alive: a peer that
// keeps reading, however slowly, can keep it for as long as it likes. A
// hold keeps the buffer alone (its connection's other state is released)
// and costs no descriptor, no connection slot and no work per event-loop
// pass; a worker holding 16 MiB of them stops using SEND_ZC, and copies,
// until some are released.
//
// - CloseZCNotifHeld: holds started. A rate: a connection the server
// closed while its peer had stopped reading mid-send.
// - CloseZCNotifHeldNow and CloseZCNotifHeldBytes are GAUGES: the send
// buffers held right now, and their capacity in bytes. A worker that
// shuts down takes its share out (its buffers then count in
// ShutdownZCBufRetained).
// - CloseZCNotifForced: send buffers given up while a SEND_ZC was still
// owed on them. Must stay 0. The hold is decided so that the backstop
// never gives one up, so this is a tripwire for a change that breaks
// that decision, not a measure of what peers or the kernel do: it
// cannot move on the shipped code.
// - ShutdownZCBufRetained: send buffers engine shutdown kept for the
// life of the process, because a SEND_ZC may still read them when the
// io_uring ring closes and nothing can say when it stops.
//
// io_uring-only; zero on other engines. All but the two gauges are
// cumulative. On the adaptive engine each is the sum over both
// sub-engines.
CloseZCNotifHeld uint64
CloseZCNotifHeldNow uint64
CloseZCNotifHeldBytes uint64
CloseZCNotifForced uint64
ShutdownZCBufRetained 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
14 changes: 9 additions & 5 deletions engine/iouring/conn.go
Original file line number Diff line number Diff line change
Expand Up @@ -610,14 +610,18 @@ func releaseConnState(cs *connState) {
cs.sendBody = nil
// bodyRecvPin is cleared here, after every kernel-held op delivered its
// terminal CQE (releaseConnState is only called from drainPendingRelease
// once cs.kernelInflight drained — or its backstop fired), so the kernel
// can no longer be writing into the pinned bodyBuf array (#256
// body-buffer UAF guard).
// once cs.kernelInflight drained — or its backstop fired — and from the
// backstop's SEND_ZC hold once a SEND_ZC, which writes nothing, is all
// that is owed; releaseHeldConnState), so the kernel can no longer be
// writing into the pinned bodyBuf array (#256 body-buffer UAF guard).
cs.bodyRecvPin = nil
cs.detectAccum = cs.detectAccum[:0]
// kernelInflight is zero on every normal release (drainPendingRelease
// gates on it); reset defensively for the wall-clock-backstop path,
// where the worker gave up waiting on a CQE the kernel never produced.
// gates on it); reset for the wall-clock-backstop path, where the worker
// gave up waiting on a CQE the kernel never produced, and for the
// backstop's SEND_ZC hold, which keeps the send buffer's array and lets
// the connState go with the SEND_ZC still owed (releaseHeldConnState;
// the identity's count, not this one, waits for its notification).
cs.kernelInflight = 0
cs.recvArmed = false
cs.recvOutstanding = 0
Expand Down
8 changes: 8 additions & 0 deletions engine/iouring/engine.go
Original file line number Diff line number Diff line change
Expand Up @@ -584,6 +584,14 @@ func (e *Engine) Metrics() engine.EngineMetrics {
TransplantReapUnsupported: e.metrics.handoffLoss.reapUnsupported.Load(),
CloseFDDeferred: e.metrics.handoffLoss.closeFDDeferred.Load(),
CloseFDForced: e.metrics.handoffLoss.closeFDForced.Load(),
CloseZCNotifHeld: e.metrics.handoffLoss.zcNotifHeld.Load(),
CloseZCNotifForced: e.metrics.handoffLoss.zcNotifForced.Load(),
ShutdownZCBufRetained: e.metrics.handoffLoss.zcBufRetained.Load(),
// Gauges: each worker takes off only what it added, after it, so
// neither goes below 0; the clamp only keeps a bug from showing
// as 2^64.
CloseZCNotifHeldNow: uint64(max(0, e.metrics.handoffLoss.zcHeldNow.Load())),
CloseZCNotifHeldBytes: uint64(max(0, e.metrics.handoffLoss.zcHeldBytes.Load())),

TransplantSweepPasses: e.metrics.sweep.passes.Load(),
TransplantResidualDetached: e.metrics.sweep.residual[resDetached].Load(),
Expand Down
64 changes: 60 additions & 4 deletions engine/iouring/handoff_loss.go
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,24 @@ import (
// batched per loop iteration (Worker.closeFDDeferredBatch).
// - closeFDForced: such descriptors the pendingRelease backstop closed
// with an op still owed. Must stay 0.
//
// And what a SEND_ZC's send buffer is kept for (celeris#812; see
// zc_send_buffer.go), three more and two gauges:
//
// - zcNotifHeld: holds the pendingRelease backstop started past its 5 s
// because a SEND_ZC was still owed on the send buffer, each until the
// notification. A rate: a connection closed on a peer that stopped
// reading mid-send.
// - zcHeldNow, zcHeldBytes: GAUGES, the send buffers held right now and
// their capacity in bytes, summed over the workers: each adds at a hold
// and takes off at its end, and a worker that shuts down takes off what
// it still held (retractZCHolds).
// - zcNotifForced: closed identities whose accounting was dropped with a
// SEND_ZC still owed, i.e. a send buffer given up to the pool or the GC
// while the kernel may still send from it. Must stay 0. The hold is
// decided on the identity, so only a change that breaks it can move this.
// - zcBufRetained: send buffers worker shutdown kept for the life of the
// process, because a SEND_ZC may still read them when the ring closes.
type handoffLossStats struct {
staleRecvDataClosed atomic.Uint64
staleRecvDataTransplanted atomic.Uint64
Expand All @@ -127,6 +145,11 @@ type handoffLossStats struct {
reapUnsupported atomic.Uint64
closeFDDeferred atomic.Uint64
closeFDForced atomic.Uint64
zcNotifHeld atomic.Uint64
zcNotifForced atomic.Uint64
zcBufRetained atomic.Uint64
zcHeldNow atomic.Int64
zcHeldBytes atomic.Int64
}

func (s *handoffLossStats) noteCloseFDForced() {
Expand All @@ -135,6 +158,32 @@ func (s *handoffLossStats) noteCloseFDForced() {
}
}

func (s *handoffLossStats) noteCloseZCNotifHeld() {
if s != nil {
s.zcNotifHeld.Add(1)
}
}

func (s *handoffLossStats) noteCloseZCNotifForced() {
if s != nil {
s.zcNotifForced.Add(1)
}
}

func (s *handoffLossStats) noteShutdownZCBufRetained(n uint64) {
if s != nil {
s.zcBufRetained.Add(n)
}
}

// addZCHeld moves the held-now gauges by n buffers of bytes in all.
func (s *handoffLossStats) addZCHeld(n, bytes int64) {
if s != nil && (n != 0 || bytes != 0) {
s.zcHeldNow.Add(n)
s.zcHeldBytes.Add(bytes)
}
}

// The fd-lifetime counters are nil-safe: a hand-built test Worker has none.

func (s *handoffLossStats) noteHeld() {
Expand Down Expand Up @@ -214,10 +263,17 @@ func (w *Worker) noteStaleRecvData(ud uint64) {
// closedOps entry. Worker thread only.
func (w *Worker) noteStaleRecvExemplar(c *completionEntry, fd int, ud uint64) {
var head []byte
if e := w.closedOps[connOpKey(ud)]; e != nil && len(e.conns) > 0 && w.bufRing == nil {
buf := e.conns[0].buf
n := min(int(c.Res), len(buf), 64)
head = buf[:n]
if e := w.closedOps[connOpKey(ud)]; e != nil && w.bufRing == nil {
// The first conn still registered: a SEND_ZC hold may have released
// one and left its slot nil (releaseHeldConnState).
for _, cs := range e.conns {
if cs != nil {
buf := cs.buf
n := min(int(c.Res), len(buf), 64)
head = buf[:n]
break
}
}
}
recvtheft.NoteStaleRecvData(w.id, fd, decodeGen(ud), c.Res, head)
}
Expand Down
10 changes: 10 additions & 0 deletions engine/iouring/handoff_loss_metrics_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,11 @@ func TestMetricsCarriesTheHandoffLossWitnesses(t *testing.T) {
e.metrics.handoffLoss.reapUnsupported.Store(41)
e.metrics.handoffLoss.closeFDDeferred.Store(43)
e.metrics.handoffLoss.closeFDForced.Store(47)
e.metrics.handoffLoss.zcNotifHeld.Store(53)
e.metrics.handoffLoss.zcNotifForced.Store(59)
e.metrics.handoffLoss.zcBufRetained.Store(61)
e.metrics.handoffLoss.zcHeldNow.Store(67)
e.metrics.handoffLoss.zcHeldBytes.Store(71)

m := e.Metrics()
for _, c := range []struct {
Expand All @@ -48,6 +53,11 @@ func TestMetricsCarriesTheHandoffLossWitnesses(t *testing.T) {
{"TransplantReapUnsupported", m.TransplantReapUnsupported, 41},
{"CloseFDDeferred", m.CloseFDDeferred, 43},
{"CloseFDForced", m.CloseFDForced, 47},
{"CloseZCNotifHeld", m.CloseZCNotifHeld, 53},
{"CloseZCNotifForced", m.CloseZCNotifForced, 59},
{"ShutdownZCBufRetained", m.ShutdownZCBufRetained, 61},
{"CloseZCNotifHeldNow", m.CloseZCNotifHeldNow, 67},
{"CloseZCNotifHeldBytes", m.CloseZCNotifHeldBytes, 71},
} {
if c.got != c.want {
t.Errorf("Metrics().%s = %d, want %d — the witness exists but cannot "+
Expand Down
Loading
Loading