diff --git a/.github/scripts/mutant-812-backstop-releases-zc.py b/.github/scripts/mutant-812-backstop-releases-zc.py new file mode 100755 index 00000000..fe7793ec --- /dev/null +++ b/.github/scripts/mutant-812-backstop-releases-zc.py @@ -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})") diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9757180a..dee23999 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -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 diff --git a/adaptive/engine.go b/adaptive/engine.go index 324904d6..c5656c24 100644 --- a/adaptive/engine.go +++ b/adaptive/engine.go @@ -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 diff --git a/adaptive/handoff_loss_metrics_test.go b/adaptive/handoff_loss_metrics_test.go index 1b256121..9bd02992 100644 --- a/adaptive/handoff_loss_metrics_test.go +++ b/adaptive/handoff_loss_metrics_test.go @@ -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, @@ -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() @@ -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)", diff --git a/engine/engine.go b/engine/engine.go index 3599f075..1d9fa79f 100644 --- a/engine/engine.go +++ b/engine/engine.go @@ -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 diff --git a/engine/iouring/conn.go b/engine/iouring/conn.go index 54d3d040..c167d4ed 100644 --- a/engine/iouring/conn.go +++ b/engine/iouring/conn.go @@ -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 diff --git a/engine/iouring/engine.go b/engine/iouring/engine.go index 0918d325..af0d6c16 100644 --- a/engine/iouring/engine.go +++ b/engine/iouring/engine.go @@ -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(), diff --git a/engine/iouring/handoff_loss.go b/engine/iouring/handoff_loss.go index a36bcd38..67c9cfda 100644 --- a/engine/iouring/handoff_loss.go +++ b/engine/iouring/handoff_loss.go @@ -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 @@ -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() { @@ -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() { @@ -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) } diff --git a/engine/iouring/handoff_loss_metrics_test.go b/engine/iouring/handoff_loss_metrics_test.go index f62a4e96..c490e989 100644 --- a/engine/iouring/handoff_loss_metrics_test.go +++ b/engine/iouring/handoff_loss_metrics_test.go @@ -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 { @@ -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 "+ diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index 31d41ea4..31cb0faf 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -191,11 +191,19 @@ func clampBufRingCount(n int) int { // so that the zero entry holds nothing. Both sit in detached's padding, so the // entry stays 24 bytes on 64-bit platforms // (TestPendingReleaseEntryStaysTwentyFourBytes). +// +// zcHeld marks an entry the backstop found still owing a SEND_ZC and held +// (celeris#812, see closedZCOwed and holdZCPastBackstop): the kernel may still +// read cs.sendBuf, so its array is held, off this queue, until that op's +// notification arrives, however long past releaseAtNanos. It records only +// that the hold was counted (CloseZCNotifHeld), once per entry, for an entry +// that comes back to the queue and is held again; the same padding holds it. type pendingReleaseEntry struct { cs *connState releaseAtNanos int64 detached bool holdsFD bool + zcHeld bool fd int32 } @@ -215,6 +223,17 @@ type closedOpsEntry struct { // (TestClosedOpsEntryStaysThirtyTwoBytes). 32-bit platforms have no // padding there, and the entry grows from 16 to 20 bytes. handoff bool + // zcOwed is the part of inflight that is a SEND_ZC (celeris#812): the + // send, and after its first CQE its notification. The kernel may read + // the conn's sendBuf until that notification. noteClosedInflight adds + // one per conn whose send in flight is a SEND_ZC, and + // noteStaleTerminalOp takes it off at that op's terminal CQE (see + // there). closedZCOwed reads it, so the release backstop holds the + // connState, and the send buffer with it, until the kernel is done with + // the buffer. A conn has one send in flight at most, so a byte is ample + // even under a collision; it takes the byte of padding between handoff + // and fdOps, and the entry keeps its size on every platform. + zcOwed uint8 // fdOps is the part of inflight that still names the descriptor // (celeris#798): all of it but the SEND_ZC notifications whose send has // completed. noteClosedInflight adds each conn's fdOps; staleConnCQE @@ -244,6 +263,24 @@ type closedOpsEntry struct { // last-resort trade of a potential use-after-free against an // unbounded memory leak. // +// One op is exempt, because for it the trade is not a last resort: a +// SEND_ZC, whose terminal CQE is its notification (celeris#812). The +// kernel posts it once no queued segment references the send buffer's +// pages any more, and a peer that stopped reading keeps the unsent tail +// queued on the socket for as long as the socket lives: after a close, +// until the kernel gives up on the orphaned socket, minutes on a peer that +// keeps acknowledging its closed window. Until then the kernel reads +// cs.sendBuf whenever the peer opens its window, so a release there was a +// certain use-after-free, not a potential one: the peer received whatever +// the array's next owner wrote into it. The backstop holds such a +// send buffer until the notification arrives instead (closedZCOwed, +// holdZCPastBackstop). Nothing but the peer decides how long that is: a peer +// that keeps reading, however slowly, keeps the orphaned socket alive. So the +// hold is kept off this queue's walk, keeps the array alone once the SEND_ZC +// is all that is owed, shows in the CloseZCNotifHeldNow/Bytes gauges, and +// stops the worker arming new SEND_ZCs while it holds zcHoldBytesMax (see +// zc_send_buffer.go). It never costs a descriptor. +// // 5 s is comfortably above any plausible straggler window: TCP // retransmits begin at RTO_MIN = 200 ms and a 4 KiB POST tail // straggling through several backoffs stays well under a second — @@ -480,7 +517,9 @@ type Worker struct { // "s.allocCount != s.nelems" span-corruption (v1.4.15/7beebb9, both bench runs). Neither is // race-detectable because the writer is the kernel, not Go code. // releaseAtNanos is only the anomaly backstop — see - // pendingReleaseHoldNanos. + // pendingReleaseHoldNanos, which also says why it never gives up a + // send buffer a SEND_ZC may still read (celeris#812), and holds it off + // this queue instead (zcHolds). pendingRelease []pendingReleaseEntry // closeFDOwed counts the pendingRelease entries that still hold their // descriptor open (holdsFD, celeris#685). The worker does not park @@ -508,6 +547,19 @@ type Worker struct { // releases all colliding conns only after every expected terminal // CQE arrived (errs toward holding longer; see closedOpsEntry). closedOps map[uint64]*closedOpsEntry + // zcHolds are the send buffers the pendingRelease backstop holds past + // its deadline for a SEND_ZC still owed on them (celeris#812), keyed by + // the closed identity (connOpKey) whose CQEs settle them + // (settleZCHolds, from noteStaleTerminalOp). Kept off pendingRelease, + // which the loop walks every pass, because a hold lasts as long as the + // peer keeps its orphaned socket alive (see zc_send_buffer.go). + // zcHoldCount and zcHoldBytes are what they hold: the worker's share of + // the held-now gauges, and zcHoldBytes is what prepSendSQE checks + // against zcHoldBytesMax. Worker thread only; nil/zero whenever nothing + // is held. + zcHolds map[uint64][]zcHold + zcHoldCount int + zcHoldBytes int // shutdownDrainDeadline is the wall-clock (UnixNano) bound on the // send drain the run loop performs once its context is cancelled @@ -1725,7 +1777,11 @@ func (w *Worker) staleConnCQE(c *completionEntry, fd int, ud uint64) bool { // still owed (CloseFDForced, must stay 0). Either way the // gen collision comes ON TOP. When it fires, the closed conn's // closedOps entry is left orphaned and the 5 s backstop WARN - // in drainPendingRelease is the production signal. + // in drainPendingRelease is the production signal. If the CQE + // taken was the closed conn's SEND_ZC notification, its zcOwed + // stays set and the backstop holds that connState for good + // instead (celeris#812, CloseZCNotifHeld): one connState + // leaked, the safe side of the trade. if cs.kernelInflight > 0 { cs.kernelInflight-- } @@ -1768,6 +1824,18 @@ func (w *Worker) staleConnCQE(c *completionEntry, fd int, ud uint64) bool { // (conn closed with zero in-flight ops, or already backstop-released). // namedFD says whether the op still named the descriptor until this CQE // (closedOpsEntry.fdOps): false only for a SEND_ZC notification. +// +// A send's terminal CQE also ends a SEND_ZC's hold on the send buffer +// (closedOpsEntry.zcOwed, celeris#812) when it is the SEND_ZC's: its +// notification (namedFD false), or, for a SEND_ZC whose first CQE came +// without IORING_CQE_F_MORE and so has no notification to follow, that first +// CQE. The latter looks like a plain send's CQE, so it is taken for the +// SEND_ZC's only where the identity holds one conn, whose one send in flight +// it is. Under an (fd, generation) collision only a notification ends a hold: +// holding late costs memory, releasing early sends a peer another +// connection's bytes. Whatever the CQE changed, the send buffers held for the +// identity past the backstop are then settled (settleZCHolds): the lookup +// that costs is paid only while some hold exists. func (w *Worker) noteStaleTerminalOp(ud uint64, namedFD bool) { if len(w.closedOps) == 0 { return @@ -1781,11 +1849,21 @@ func (w *Worker) noteStaleTerminalOp(ud uint64, namedFD bool) { if namedFD && e.fdOps > 0 { e.fdOps-- } + if ud&udMask == udSend && e.zcOwed > 0 && (!namedFD || len(e.conns) == 1) { + e.zcOwed-- + } + if len(w.zcHolds) > 0 { + w.settleZCHolds(key, e) + } if e.inflight > 0 { return } for _, cs := range e.conns { - cs.kernelInflight = 0 + // nil: a conn whose SEND_ZC hold released it and kept its array + // alone (releaseHeldConnState); the pool may have handed it on. + if cs != nil { + cs.kernelInflight = 0 + } } delete(w.closedOps, key) } @@ -4275,6 +4353,9 @@ func (w *Worker) noteClosedInflight(cs *connState) { } e.inflight += cs.kernelInflight e.fdOps += int16(fdOps(cs)) + if zcSendOwed(cs) { + e.zcOwed++ + } e.conns = append(e.conns, cs) } @@ -4285,11 +4366,21 @@ func (w *Worker) noteClosedInflight(cs *connState) { // conn. Any conns colliding on the same identity lose their accounting // too and will be reaped by their own backstop — acceptable for a path // that only fires on kernel anomalies. +// +// It is also where the accounting lets go of a SEND_ZC the kernel may still +// read cs.sendBuf for (celeris#812): the backstop holds an entry that owes +// one (closedZCOwed), so an identity dropped with zcOwed above zero is a send +// buffer given up while the kernel may still send from it. Counted +// (CloseZCNotifForced); must stay 0. func (w *Worker) dropClosedOps(cs *connState) { if len(w.closedOps) == 0 { return } - delete(w.closedOps, encodeConnOpKey(cs.fd, cs.generation)) + key := encodeConnOpKey(cs.fd, cs.generation) + if e := w.closedOps[key]; e != nil && e.zcOwed > 0 { + w.handoffLoss.noteCloseZCNotifForced() + } + delete(w.closedOps, key) } // queuePendingRelease enqueues cs for deferred release: drainPendingRelease @@ -4347,7 +4438,11 @@ func (w *Worker) queuePendingReleaseDetached(cs *connState) { // the kernel never delivered a terminal CQE for an op we believe it // holds — log a WARN, since releasing now trades a potential // use-after-free against an unbounded leak, and scrub the conn from -// closedOps so a later CQE cannot touch the released memory. +// closedOps so a later CQE cannot touch the released memory. The one +// exception is an entry that still owes a SEND_ZC (celeris#812, +// closedZCOwed): its send buffer is held past the backstop until the +// notification arrives, off this walk (Worker.zcHolds), see +// pendingReleaseHoldNanos. // // Detached entries skip the pool recycle (releaseConnState would // reset fields that goroutine closures may still observe via the @@ -4376,6 +4471,15 @@ func (w *Worker) drainPendingRelease() { kept = append(kept, *entry) continue } + // Past the backstop. A SEND_ZC the kernel may still read + // cs.sendBuf for is no anomaly (celeris#812): the send buffer is + // held, off this walk, until that op's notification, and the + // backstop gives up only the descriptor, as it always has + // (holdZCPastBackstop). + if w.closedZCOwed(cs) { + w.holdZCPastBackstop(entry) + continue + } // Backstop: kernel anomaly, not normal flow. if w.logger != nil { w.logger.Warn("releasing connState with kernel ops unaccounted for after backstop hold", @@ -5722,7 +5826,11 @@ func (w *Worker) prepSendSQE(sqe unsafe.Pointer, cs *connState, linked bool) { // classifier reads cs.sendIsZC, never w.sendZC, because the fallbacks // clear w.sendZC while sends armed under it are still in flight // (celeris#609). - cs.sendIsZC = useSendZC(w.sendZC, linked, len(cs.sendBuf)) + // + // A worker holding zcHoldBytesMax of send buffers past the release + // backstop arms no new SEND_ZC until those holds end: that is what bounds + // them (celeris#812, see zc_send_buffer.go). + cs.sendIsZC = useSendZC(w.sendZC && w.zcHoldBytes < zcHoldBytesMax, linked, len(cs.sendBuf)) if cs.sendIsZC { // celeris#591 exposure witnesses. Deliberately inside the ZC arm: // every sub-sendZCMinBytes and every linked send — the per-request @@ -6300,7 +6408,10 @@ func (w *Worker) shutdown() { // still-running sibling worker while this ring's kernel side can // still write into it (same #256-class UAF, shutdown variant). // The conns remain reachable via w.conns until the Worker itself - // is collected, well after the ring teardown cancels its ops. + // is collected, well after the ring teardown cancels its ops. A + // SEND_ZC's send buffer is the exception: no cancel or teardown + // ends the kernel's use of it, so retainZCSendBufsAtShutdown keeps + // the ones still owed past the Worker (celeris#812). } // celeris#657 R2: this worker is gone, so it must not leave its last // cycle's residue standing in the engine-wide gauges. Nothing else @@ -6320,6 +6431,12 @@ func (w *Worker) shutdown() { // goroutines this function only joins below — can write this descriptor // number once it is free to be recycled. w.wakeFD.Close() + // celeris#812: the ring's last word on which SEND_ZC send buffers the + // kernel is done with. Every buffer it cannot clear by now is kept for + // the life of the process: nothing would ever say when the kernel lets + // go of it, and the Worker, which is all that holds it, can be + // collected once the engine is dropped. + w.retainZCSendBufsAtShutdown() if w.bufRing != nil && w.ring != nil { w.bufRing.Close(w.ring) } diff --git a/engine/iouring/zc_send_buffer.go b/engine/iouring/zc_send_buffer.go new file mode 100644 index 00000000..f8706d24 --- /dev/null +++ b/engine/iouring/zc_send_buffer.go @@ -0,0 +1,361 @@ +//go:build linux + +package iouring + +import ( + "sync" + "time" + + "golang.org/x/sys/unix" +) + +// The send buffer of a SEND_ZC (celeris#812). +// +// A SEND_ZC does not copy cs.sendBuf into the socket: the queued segments +// reference its pages, and the kernel reads them whenever it transmits one, +// a retransmission included. Its notification CQE, the op's terminal CQE, says +// that no segment references them any more. Until then the array must not be +// written, which on a live connection flushSend and flushSendLink see to (they +// wait out cs.zcNotifPending), and must not be given to anyone else, which is +// what the release gate on a closed connection is for (drainPendingRelease, +// on cs.kernelInflight, which counts the SEND_ZC until its notification). +// +// That gate had a wall-clock backstop that released anyway 5 s after the +// close, on the reasoning that an op still owed by then is a kernel anomaly. +// A notification is not one. A peer that stops reading keeps the unsent tail +// of the send queued behind its closed window, and a close leaves the socket +// orphaned with that tail still queued: the kernel keeps offering it for as +// long as it keeps the orphan, and posts the notification when it gives up on +// the socket or the peer takes the bytes. The backstop returned such a +// connState to connStatePool with its sendBuf (a detached one to the GC); the +// array's next owner wrote into it; and when the peer read again, it received +// the tail from that array, another connection's bytes. Measured by +// TestBackstopHoldsASendBufferAZCNotificationStillReads: 61,200 of the 65,536 +// bytes of the send, in every case, on the tree before this fix. +// +// So a closed connection's closedOps entry counts its SEND_ZC apart +// (closedOpsEntry.zcOwed), and past the backstop the array is held until the +// notification (closedZCOwed, holdZCPastBackstop). Nothing but the peer +// decides how long that is. A peer that keeps reading, however slowly, keeps +// the orphaned socket alive; one that stays at a zero window was measured +// keeping it 5 min 36 s. Neither a descriptor (celeris#798) nor a connection +// slot counts a hold. So the hold is built to cost as little as it can for as +// long as it lasts, and to stop growing (celeris#813 round 2): +// +// - Off the release walk. A held array waits in Worker.zcHolds, keyed by +// its closed identity, not in pendingRelease, which the loop walks every +// pass. Only a CQE of that identity reaches it (settleZCHolds, from +// noteStaleTerminalOp), so a hold costs O(1) when it starts, at each of +// its identity's CQEs and when it ends, and nothing per loop pass. +// - The array alone. The kernel reads nothing of the connState but the +// send buffer's array, so once the SEND_ZC is all the identity owes, the +// connState is released as usual (to the pool without the array, or, +// detached, to the GC) and the hold keeps the array only +// (releaseHeldConnState). While another op is owed too (a recv the close's +// cancel has not ended yet, which may still write cs.buf), the whole entry +// is held, and handed back to the walk at that op's CQE. +// - Shown. CloseZCNotifHeldNow and CloseZCNotifHeldBytes are gauges of +// what is held right now. +// - Bounded. A worker whose held arrays reach zcHoldBytesMax arms no new +// SEND_ZC (prepSendSQE): its sends copy, and a copied send's unsent tail +// is the kernel's own socket memory, which the kernel bounds itself. The +// held bytes can then grow only by the SEND_ZCs already in flight on live +// connections, which the connection limit bounds. +// +// Worker shutdown, whose ring teardown ends no such use of the pages either, +// keeps every send buffer still owed for the life of the process +// (retainZCSendBufsAtShutdown). +// +// Counted, in handoffLossStats and EngineMetrics: +// +// - CloseZCNotifHeld: holds the backstop started. A rate of connections +// closed on a stalled peer mid-send. +// - CloseZCNotifHeldNow, CloseZCNotifHeldBytes: gauges, the arrays held +// right now and their capacity. +// - CloseZCNotifForced: identities dropped while they still owed a SEND_ZC +// (dropClosedOps), i.e. a send buffer given up while the kernel may still +// send from it. Must stay 0. The hold is decided on the identity, so the +// backstop never drops one that owes a SEND_ZC: this is a tripwire for a +// change that breaks that, not a measure of the kernel. +// - ShutdownZCBufRetained: send buffers worker shutdown kept for the life +// of the process. + +// zcSendOwed reports whether cs's send in flight is a SEND_ZC: armed as one +// (sendIsZC, set per send by prepSendSQE) and not finished, which for a +// SEND_ZC is its notification (handleSend clears zcNotifPending there, and +// completeSend then cs.sending; a first CQE that failed clears sending and +// leaves zcNotifPending set). On a closed connection the flags are the ones +// it had at the close: nothing touches them after it. Worker thread. +func zcSendOwed(cs *connState) bool { + return cs.sendIsZC && (cs.sending || cs.zcNotifPending) +} + +// closedZCOwed reports whether closed cs's identity still owes a SEND_ZC, +// the send or its notification, so the kernel may still read a send buffer +// closed under it: its closedOps entry counts one (closedOpsEntry.zcOwed). +// It is decided on the identity, not on cs's own send: under an (fd, +// generation) collision the accounting cannot tell the conns' CQEs apart, so +// a conn with no SEND_ZC of its own is held with its sibling's, and none of +// them reaches the backstop's release and drops the identity while the +// sibling's array is still read. No entry, or kernelInflight at zero, means +// the accounting has retired the identity: every op delivered its terminal +// CQE. Worker thread only. +func (w *Worker) closedZCOwed(cs *connState) bool { + if cs.kernelInflight <= 0 { + return false + } + e := w.closedOps[encodeConnOpKey(cs.fd, cs.generation)] + return e != nil && e.zcOwed > 0 +} + +// zcHold is a send buffer the release backstop holds for a SEND_ZC past its +// deadline, in Worker.zcHolds, off the release walk (celeris#812, #813). +type zcHold struct { + // sendBuf is the array the kernel may read, at its full capacity: what + // the hold keeps, and what it counts in bytes. + sendBuf []byte + // entry is the whole pendingRelease entry while its identity owes an op + // other than a SEND_ZC too, which may still write the connState's other + // buffers: settleZCHolds hands it back to the walk at that identity's + // next CQE. Zero (cs nil) once the SEND_ZC is all that is owed: the + // connState has been released, and the hold keeps the array alone. + entry pendingReleaseEntry +} + +// zcHoldBytesMax is how many bytes of send buffers a worker holds past the +// backstop before it stops arming SEND_ZC (prepSendSQE): a hold lasts as long +// as the peer keeps its orphaned socket alive, and no descriptor or connection +// slot counts it, so this is what bounds it. At or above it the worker's sends +// copy, and a copied send's unsent tail is the kernel's own socket memory, +// which its orphan and socket-memory limits bound. 16 MiB is some 2,000 +// held responses of 8 KiB per worker, far above the few a server closes on +// peers that stopped reading mid-send, and a small share of memory. +const zcHoldBytesMax = 16 << 20 + +// holdZCPastBackstop takes e, which the backstop found still owing a +// SEND_ZC (closedZCOwed), off the release walk and holds its send buffer in +// Worker.zcHolds until the SEND_ZC's notification (settleZCHolds). The +// caller does not keep e. The descriptor is not held with it. A +// notification names none, and the entry let go of it already (celeris#798); +// if it still keeps it, an op that does name it is owed past the backstop +// (the SEND_ZC's send itself, or a recv), and it is closed and counted as the +// backstop has always done (CloseFDForced, must stay 0). If the SEND_ZC is +// all the identity owes, the connState goes now (releaseHeldConnState) and +// the hold keeps the array alone; otherwise the whole entry is held. The +// hold is counted once per entry (CloseZCNotifHeld), however often its entry +// comes back to the walk. Worker thread only. +func (w *Worker) holdZCPastBackstop(e *pendingReleaseEntry) { + cs := e.cs + if e.holdsFD { + if w.logger != nil { + w.logger.Warn("closing a kept descriptor with kernel ops unaccounted for after backstop hold", + "worker", w.id, "fd", cs.fd, "generation", cs.generation, + "inflight", cs.kernelInflight, "detached", e.detached) + } + w.handoffLoss.noteCloseFDForced() + w.releaseKeptFD(e) + } + if !e.zcHeld { + e.zcHeld = true + w.handoffLoss.noteCloseZCNotifHeld() + if w.logger != nil { + w.logger.Debug("holding a closed connection's send buffer past the release backstop until its SEND_ZC notification", + "worker", w.id, "fd", cs.fd, "generation", cs.generation, + "inflight", cs.kernelInflight, "detached", e.detached) + } + } + key := encodeConnOpKey(cs.fd, cs.generation) + h := zcHold{sendBuf: cs.sendBuf[:cap(cs.sendBuf)], entry: *e} + if ce := w.closedOps[key]; ce != nil && ce.inflight == int32(ce.zcOwed) { + w.releaseHeldConnState(ce, cs, e.detached) + h.entry = pendingReleaseEntry{} + } + if w.zcHolds == nil { + w.zcHolds = make(map[uint64][]zcHold) + } + w.zcHolds[key] = append(w.zcHolds[key], h) + n := cap(h.sendBuf) + w.zcHoldCount++ + w.zcHoldBytes += n + w.handoffLoss.addZCHeld(1, int64(n)) +} + +// releaseHeldConnState releases held cs, whose identity ce owes nothing but +// SEND_ZC ops, while the hold keeps its send buffer's array: the kernel reads +// nothing else of it (a SEND_ZC's SQE carries the array's address and +// length; it is never a WRITEV, and nothing owed writes cs.buf). cs leaves +// the identity first, its slot set to nil so the number of conns the +// collision rule reads is kept (noteStaleTerminalOp), and so no later CQE of +// the identity writes into a connState the pool has handed on. A plain one +// then goes to the pool without the array, which its next owner would write +// into; a detached one to the GC as it is, since a dispatch goroutine may +// still read it (queuePendingReleaseDetached). Worker thread only. +func (w *Worker) releaseHeldConnState(ce *closedOpsEntry, cs *connState, detached bool) { + for i, c := range ce.conns { + if c == cs { + ce.conns[i] = nil + } + } + if detached { + return + } + cs.sendBuf = nil + releaseConnState(cs) +} + +// settleZCHolds is what a CQE of closed identity key does to the send +// buffers held for it, once noteStaleTerminalOp has taken the CQE off e's +// counts. An array held alone is let go once the identity owes no SEND_ZC +// (zcOwed 0) or nothing at all: the kernel is done with it, and the GC takes +// it. A whole entry goes back to the release walk at any such CQE, where the +// next pass releases it (the identity retired), holds it again (the array +// alone, if the SEND_ZC is all that is left) or, if its SEND_ZC has ended +// with another op still owed, gives up on it as the backstop always has. +// Worker thread only. +func (w *Worker) settleZCHolds(key uint64, e *closedOpsEntry) { + hs, ok := w.zcHolds[key] + if !ok { + return + } + done := e.inflight <= 0 || e.zcOwed == 0 + kept := hs[:0] + for _, h := range hs { + if h.entry.cs == nil && !done { + kept = append(kept, h) + continue + } + n := cap(h.sendBuf) + w.zcHoldCount-- + w.zcHoldBytes -= n + w.handoffLoss.addZCHeld(-1, -int64(n)) + if h.entry.cs != nil { + w.pendingRelease = append(w.pendingRelease, h.entry) + } + } + clear(hs[len(kept):]) + if len(kept) == 0 { + delete(w.zcHolds, key) + } else { + w.zcHolds[key] = kept + } +} + +// liveZCOwed is zcSendOwed for a live connection at worker shutdown, whose +// flags lag its count: the shutdown drains (endOwedOpsAtShutdown, and +// retainZCSendBufsAtShutdown's reap) retire its CQEs through staleConnCQE, +// which takes each terminal one off kernelInflight (and clears recvArmed at a +// recv's) but leaves the send flags as they were. kernelInflight counts the +// armed recv, if any, and the send, so the send is still owed while it counts +// more than the recv. +func liveZCOwed(cs *connState) bool { + n := cs.kernelInflight + if cs.recvArmed { + n-- + } + return zcSendOwed(cs) && n > 0 +} + +// zcRetained holds, for the life of the process, the send buffers worker +// shutdown found a SEND_ZC may still read (retainZCSendBufsAtShutdown). +var zcRetained struct { + mu sync.Mutex + bufs [][]byte +} + +// retainZCSendBufsAtShutdown is worker shutdown's half of celeris#812, run +// just before the ring is closed. Closing the ring ends no SEND_ZC's use of +// its send buffer (the queued segments hold the pages, and the orphaned socket +// keeps sending from them), and once the ring is closed no notification can +// say when that use ends. The connStates are not pooled at shutdown, but the +// Worker is the last thing that holds them, and the engine may be dropped +// and collected long before the kernel lets go: the GC would then hand the +// array to new allocations while a peer can still receive it. +// +// So every send buffer a SEND_ZC may still read, of a live connection +// (liveZCOwed), of one queued for release (closedZCOwed) or held past the +// backstop (zcHolds), goes into zcRetained, which is never freed, and is +// counted (ShutdownZCBufRetained). No wait is added for them: the run loop's +// send drain (shutdownSendDrainNanos) already gave the live connections' +// sends their time, and a notification owed to a stalled peer can take +// minutes. Only when there is something to keep does it first take what the +// kernel has already completed: an enter that submits nothing +// (shutdownDrivers closed the driver descriptors on the promise that nothing +// is submitted after it) and waits at most zcShutdownFlushWait, which also +// runs the task work a DEFER_TASKRUN ring holds its completions in until it +// is entered; then it retires the recv and send CQEs through staleConnCQE as +// the loop does, and closes any accepted descriptor, as endOwedOpsAtShutdown +// does. The cost is the kept arrays' memory, for a shutdown that finds peers +// stalled mid-send. Either way the worker then takes its share out of the +// held-now gauges: it holds nothing any more. Worker thread only. +func (w *Worker) retainZCSendBufsAtShutdown() { + defer w.retractZCHolds() + bufs := w.zcSendBufsOwed() + if len(bufs) == 0 { + return + } + if w.ring != nil && !w.sqpoll { + _ = w.ring.WaitCQETimeout(zcShutdownFlushWait) + head, tail := w.ring.BeginCQ() + for ; head != tail; head++ { + c := w.ring.cqeAt(head) + ud := c.UserData + switch ud & udMask { + case udRecv, udSend: + w.staleConnCQE(c, int(ud&fdMask), ud) + case udAccept: + if c.Res >= 0 && !w.fixedFiles { + _ = unix.Close(int(c.Res)) + } + } + } + w.ring.EndCQ(head) + bufs = w.zcSendBufsOwed() + if len(bufs) == 0 { + return + } + } + zcRetained.mu.Lock() + zcRetained.bufs = append(zcRetained.bufs, bufs...) + zcRetained.mu.Unlock() + w.handoffLoss.noteShutdownZCBufRetained(uint64(len(bufs))) + if w.logger != nil { + w.logger.Info("keeping send buffers a SEND_ZC may still read for the life of the process", + "worker", w.id, "buffers", len(bufs)) + } +} + +// retractZCHolds takes a shut-down worker's holds out of the held-now gauges +// and drops them: retainZCSendBufsAtShutdown has kept every array still owed. +// Worker thread only. +func (w *Worker) retractZCHolds() { + w.handoffLoss.addZCHeld(-int64(w.zcHoldCount), -int64(w.zcHoldBytes)) + w.zcHolds, w.zcHoldCount, w.zcHoldBytes = nil, 0, 0 +} + +// zcShutdownFlushWait bounds the one enter retainZCSendBufsAtShutdown makes +// to collect the completions the kernel already has. It waits only while +// nothing at all has completed. +const zcShutdownFlushWait = time.Millisecond + +// zcSendBufsOwed lists the send buffers a SEND_ZC may still read: of the live +// connections, of those queued for release, and those held past the backstop, +// every one of which is owed. Worker thread only. +func (w *Worker) zcSendBufsOwed() [][]byte { + var bufs [][]byte + for _, fd := range w.liveConns { + if cs := w.conns[fd]; cs != nil && liveZCOwed(cs) { + bufs = append(bufs, cs.sendBuf) + } + } + for i := range w.pendingRelease { + if cs := w.pendingRelease[i].cs; cs != nil && w.closedZCOwed(cs) { + bufs = append(bufs, cs.sendBuf) + } + } + for _, hs := range w.zcHolds { + for _, h := range hs { + bufs = append(bufs, h.sendBuf) + } + } + return bufs +} diff --git a/engine/iouring/zc_send_buffer_bound_test.go b/engine/iouring/zc_send_buffer_bound_test.go new file mode 100644 index 00000000..bc327c52 --- /dev/null +++ b/engine/iouring/zc_send_buffer_bound_test.go @@ -0,0 +1,363 @@ +//go:build linux + +package iouring + +import ( + "fmt" + "sync" + "testing" + "time" + "unsafe" + + "golang.org/x/sys/unix" + + "github.com/goceleris/celeris/internal/conn" +) + +// What the celeris#812 hold costs while it lasts, and what bounds it +// (celeris#813 round 2). The notification a held send buffer waits for comes +// when the peer takes the unsent tail or the kernel ends the orphaned socket, +// and a peer that keeps reading, however slowly, keeps the socket alive: the +// hold lasts as long as the peer wants. So the hold must cost nothing per loop +// pass, keep no more than the array the kernel reads, show what it keeps, and +// stop growing on its own. + +// zcHeldWorker is a Worker with n closed connections, each owing only its +// SEND_ZC's notification, taken past the release backstop, which holds them. +func zcHeldWorker(tb testing.TB, n int) *Worker { + tb.Helper() + w := &Worker{conns: make([]*connState, 16), handoffLoss: &handoffLossStats{}} + for i := range n { + cs := &connState{fd: 100 + i, generation: uint32(i + 1), sendIsZC: true, sending: true, zcNotifPending: true, + kernelInflight: 1, sendBuf: make([]byte, 0, 64)} + w.noteClosedInflight(cs) + w.queuePendingReleaseFD(cs, false, -1) + } + w.cachedNow = time.Now().UnixNano() + 2*pendingReleaseHoldNanos + w.drainPendingRelease() + if got := w.handoffLoss.zcNotifHeld.Load(); got != uint64(n) { + tb.Fatalf("held %d of %d", got, n) + } + return w +} + +// zcNotifCQE delivers the notification of the SEND_ZC of the closed identity +// (fd, gen) as the loop does, through staleConnCQE. +func zcNotifCQE(t *testing.T, w *Worker, fd int, gen uint32) { + t.Helper() + c := &completionEntry{UserData: encodeUserDataGen(udSend, fd, gen), Flags: cqeFNotif} + if !w.staleConnCQE(c, fd, c.UserData) { + t.Fatalf("notification of (%d, %d) not taken as stale", fd, gen) + } +} + +// zcHeldGauges reads the held-now gauges. +func zcHeldGauges(w *Worker) (int64, int64) { + return w.handoffLoss.zcHeldNow.Load(), w.handoffLoss.zcHeldBytes.Load() +} + +// TestZCHoldIsOffTheReleaseWalk: the loop walks pendingRelease on every pass +// while it is not empty, and in round 1 a held entry stayed on it, one +// closedOps lookup a pass each, for the whole hold: 94 us a pass at 10,000 +// held (the round-2 review's BenchmarkRVBDrainPass). Past the backstop a held +// send buffer must leave the walk at once and wait in zcHolds, where only a +// CQE of its own identity reaches it: the notification releases that hold and +// no other. +func TestZCHoldIsOffTheReleaseWalk(t *testing.T) { + const n = 3 + w := &Worker{conns: make([]*connState, 16), handoffLoss: &handoffLossStats{}} + bufs := make([][]byte, n) + for i := range n { + bufs[i] = make([]byte, 5000, 8192) + cs := &connState{fd: 3 + i, generation: uint32(40 + i), sendIsZC: true, sending: true, zcNotifPending: true, + kernelInflight: 1, sendBuf: bufs[i]} + w.noteClosedInflight(cs) + w.queuePendingReleaseFD(cs, false, -1) + } + w.cachedNow = w.pendingRelease[n-1].releaseAtNanos + 1 + for pass := range 4 { + if len(w.pendingRelease) > 0 { + w.drainPendingRelease() + } + if len(w.pendingRelease) != 0 { + t.Fatalf("pass %d past the backstop: %d held entries are still on the release walk, which every loop pass makes", + pass, len(w.pendingRelease)) + } + } + for i := range n { + hs := w.zcHolds[encodeConnOpKey(3+i, uint32(40+i))] + if len(hs) != 1 || unsafe.SliceData(hs[0].sendBuf) != unsafe.SliceData(bufs[i]) || cap(hs[0].sendBuf) != cap(bufs[i]) { + t.Fatalf("conn %d: its send buffer is not held: %+v", i, hs) + } + } + if held := w.handoffLoss.zcNotifHeld.Load(); held != n { + t.Fatalf("CloseZCNotifHeld = %d, want %d: one per hold, however many passes", held, n) + } + if now, b := zcHeldGauges(w); now != n || b != n*8192 || w.zcHoldBytes != n*8192 || w.zcHoldCount != n { + t.Fatalf("held-now gauges %d buffers, %d bytes (worker %d, %d), want %d and %d", now, b, w.zcHoldCount, w.zcHoldBytes, n, n*8192) + } + // One notification releases its own hold, and only that one. + zcNotifCQE(t, w, 4, 41) + if len(w.zcHolds) != n-1 || len(w.zcHolds[encodeConnOpKey(4, 41)]) != 0 { + t.Fatalf("after conn 1's notification: holds %+v, want conns 0 and 2 only", w.zcHolds) + } + if now, b := zcHeldGauges(w); now != n-1 || b != (n-1)*8192 { + t.Fatalf("held-now gauges %d buffers, %d bytes after one release, want %d and %d", now, b, n-1, (n-1)*8192) + } + zcNotifCQE(t, w, 3, 40) + zcNotifCQE(t, w, 5, 42) + if len(w.zcHolds) != 0 || len(w.pendingRelease) != 0 || len(w.closedOps) != 0 || w.zcHoldBytes != 0 || w.zcHoldCount != 0 { + t.Fatalf("after every notification: holds=%d walk=%d identities=%d worker bytes=%d count=%d, want all 0", + len(w.zcHolds), len(w.pendingRelease), len(w.closedOps), w.zcHoldBytes, w.zcHoldCount) + } + if now, b := zcHeldGauges(w); now != 0 || b != 0 { + t.Fatalf("held-now gauges %d buffers, %d bytes after every release, want 0 and 0", now, b) + } + if forced := w.handoffLoss.zcNotifForced.Load(); forced != 0 { + t.Fatalf("CloseZCNotifForced = %d, want 0", forced) + } +} + +// TestZCHoldKeepsOnlyTheSendBuffer: what the kernel may still read of a closed +// connection that owes only a SEND_ZC is its send buffer's array, and that is +// all the hold may keep; round 1 kept the whole connState (23.9 KB a held +// connection in the round-2 review's load probe, for a 7 KB response). Its +// connState is released at the hold: to the pool without the array (the pool's +// next owner would write into it), or, detached, to the GC untouched (a +// dispatch goroutine may still read it). While another op is owed too (a recv +// the close's cancel has not ended yet, which may still write cs.buf), the +// whole entry is held, and the array alone from that op's CQE on. +func TestZCHoldKeepsOnlyTheSendBuffer(t *testing.T) { + for _, tc := range []struct { + name string + detached, recv bool + }{ + {"notif-only", false, false}, + {"detached-notif-only", true, false}, + {"recv-and-notif", false, true}, + } { + t.Run(tc.name, func(t *testing.T) { + const fd, gen = 9, 31 + key := encodeConnOpKey(fd, gen) + w := &Worker{conns: make([]*connState, 16), handoffLoss: &handoffLossStats{}} + arr := make([]byte, 6000, 8192) + h1 := conn.NewH1State() + cs := &connState{fd: fd, generation: gen, sendIsZC: true, sending: true, zcNotifPending: true, kernelInflight: 1, + sendBuf: arr, buf: make([]byte, 4096), writeBuf: make([]byte, 0, 4096), h1State: h1} + if tc.recv { + cs.recvArmed = true + cs.kernelInflight++ + } + if tc.detached { + cs.detachMu = new(sync.Mutex) + } + w.noteClosedInflight(cs) + w.queuePendingReleaseFD(cs, tc.detached, -1) + w.cachedNow = w.pendingRelease[0].releaseAtNanos + 1 + w.drainPendingRelease() + heldArray := func(when string) zcHold { + t.Helper() + hs := w.zcHolds[key] + if len(w.pendingRelease) != 0 || len(hs) != 1 || unsafe.SliceData(hs[0].sendBuf) != unsafe.SliceData(arr) || cap(hs[0].sendBuf) != cap(arr) { + t.Fatalf("%s: the send buffer's array is not held off the release walk: walk=%+v holds=%+v", when, w.pendingRelease, hs) + } + return hs[0] + } + h := heldArray("past the backstop") + if tc.recv { + if h.entry.cs != cs || cs.h1State != h1 || w.closedOps[key].conns[0] != cs { + t.Fatalf("a recv is owed with the notification, and may still write cs.buf: the whole entry must be held") + } + c := &completionEntry{UserData: encodeUserDataGen(udRecv, fd, gen), Res: -int32(unix.ECANCELED)} + if !w.staleConnCQE(c, fd, c.UserData) { + t.Fatal("recv CQE not taken as stale") + } + // The recv's end hands the entry back to the walk, which holds the array alone this time. + if len(w.pendingRelease) != 1 || len(w.zcHolds[key]) != 0 { + t.Fatalf("the recv's CQE did not hand the held entry back to the walk: walk=%+v holds=%+v", w.pendingRelease, w.zcHolds) + } + w.drainPendingRelease() + h = heldArray("after the recv's CQE") + } + if h.entry.cs != nil { + t.Fatal("the hold keeps the whole connState although only the SEND_ZC is owed, which reads nothing of it but the array") + } + if e := w.closedOps[key]; e == nil || len(e.conns) != 1 || e.conns[0] != nil { + t.Fatalf("the identity still refers to the released connState: %+v", e) + } + if tc.detached { + if cs.h1State != h1 || cs.fd != fd { + t.Fatal("a detached connState was reset: it goes to the GC as it is, since a dispatch goroutine may still read it") + } + } else if cs.h1State != nil || unsafe.SliceData(cs.sendBuf) == unsafe.SliceData(arr) { + t.Fatalf("the connState was not released to the pool without its array: h1State=%v sendBuf shares the array=%v", + cs.h1State != nil, unsafe.SliceData(cs.sendBuf) == unsafe.SliceData(arr)) + } + if now, b := zcHeldGauges(w); now != 1 || b != int64(cap(arr)) { + t.Fatalf("held-now gauges %d buffers, %d bytes, want 1 and %d", now, b, cap(arr)) + } + if held := w.handoffLoss.zcNotifHeld.Load(); held != 1 { + t.Fatalf("CloseZCNotifHeld = %d, want 1: one hold, whatever it keeps", held) + } + zcNotifCQE(t, w, fd, gen) + if len(w.zcHolds) != 0 || len(w.pendingRelease) != 0 || len(w.closedOps) != 0 { + t.Fatalf("the notification did not end the hold: holds=%+v walk=%+v identities=%d", w.zcHolds, w.pendingRelease, len(w.closedOps)) + } + if now, b := zcHeldGauges(w); now != 0 || b != 0 || w.handoffLoss.zcNotifForced.Load() != 0 { + t.Fatalf("after the notification: held-now %d buffers, %d bytes, CloseZCNotifForced %d; want 0 0 0", + now, b, w.handoffLoss.zcNotifForced.Load()) + } + }) + } +} + +// TestZCHoldCoversACollidingSibling: two connections closed under one (fd, +// generation) identity, which needs 2^32 accepts during one hold, are +// indistinguishable to the accounting. One owes a SEND_ZC's notification, the +// other a recv and no SEND_ZC. Round 1 decided the hold on each connection's +// own send (cs.sendIsZC): the sibling reached the backstop's release, dropped +// the shared identity (CloseZCNotifForced), and the next pass found no identity +// for the SEND_ZC's connection and pooled it with the array the kernel still +// sent from. The hold is decided on the identity: while it owes a SEND_ZC, +// nothing under it is released, and the notification releases both. +func TestZCHoldCoversACollidingSibling(t *testing.T) { + const fd, gen = 6, 77 + key := encodeConnOpKey(fd, gen) + w := &Worker{conns: make([]*connState, 16), handoffLoss: &handoffLossStats{}} + arr := make([]byte, 6000, 8192) + sibling := &connState{fd: fd, generation: gen, recvArmed: true, kernelInflight: 1, sendBuf: make([]byte, 0, 4096)} + zc := &connState{fd: fd, generation: gen, sendIsZC: true, sending: true, zcNotifPending: true, kernelInflight: 1, sendBuf: arr} + w.noteClosedInflight(sibling) + w.queuePendingReleaseFD(sibling, false, -1) + w.noteClosedInflight(zc) + w.queuePendingReleaseFD(zc, false, -1) + w.cachedNow = w.pendingRelease[1].releaseAtNanos + 1 + held := func(when string) { + t.Helper() + if forced := w.handoffLoss.zcNotifForced.Load(); forced != 0 { + t.Fatalf("%s: CloseZCNotifForced = %d: the identity was dropped with its SEND_ZC owed", when, forced) + } + for _, h := range w.zcHolds[key] { + if unsafe.SliceData(h.sendBuf) == unsafe.SliceData(arr) { + return + } + } + t.Fatalf("%s: the array the SEND_ZC still reads is not held: walk=%+v holds=%+v", when, w.pendingRelease, w.zcHolds[key]) + } + for range 3 { + if len(w.pendingRelease) > 0 { + w.drainPendingRelease() + } + held("past the backstop") + } + c := &completionEntry{UserData: encodeUserDataGen(udRecv, fd, gen), Res: -int32(unix.ECANCELED)} + if !w.staleConnCQE(c, fd, c.UserData) { + t.Fatal("recv CQE not taken as stale") + } + for range 3 { + if len(w.pendingRelease) > 0 { + w.drainPendingRelease() + } + held("after the sibling's recv ended") + } + zcNotifCQE(t, w, fd, gen) + if len(w.pendingRelease) > 0 { + w.drainPendingRelease() + } + if len(w.zcHolds) != 0 || len(w.pendingRelease) != 0 || len(w.closedOps) != 0 { + t.Fatalf("the notification did not release both: holds=%+v walk=%+v identities=%d", w.zcHolds, w.pendingRelease, len(w.closedOps)) + } + if now, b := zcHeldGauges(w); now != 0 || b != 0 || w.handoffLoss.zcNotifForced.Load() != 0 { + t.Fatalf("after the notification: held-now %d buffers, %d bytes, CloseZCNotifForced %d; want 0 0 0", + now, b, w.handoffLoss.zcNotifForced.Load()) + } +} + +// TestSendZCStopsWhileHeldBuffersReachTheCap: what bounds the hold. A held +// array lasts as long as the peer keeps its orphaned socket alive, and no +// descriptor or connection slot counts it, so a worker whose held arrays reach +// zcHoldBytesMax arms no new SEND_ZC: its sends copy, and a copied send's +// unsent tail is the kernel's own memory, which its orphan and socket-memory +// limits bound. The held bytes can then grow only by the SEND_ZCs already in +// flight on live connections. SEND_ZC comes back as the holds end. +func TestSendZCStopsWhileHeldBuffersReachTheCap(t *testing.T) { + opcode := func(w *Worker) (byte, bool) { + var sqe [sqeSize]byte + cs := &connState{fd: 7, sendBuf: make([]byte, sendZCMinBytes)} + w.prepSendSQE(unsafe.Pointer(&sqe[0]), cs, false) + return sqe[0], cs.sendIsZC + } + w := &Worker{sendZC: true, conns: make([]*connState, 16), handoffLoss: &handoffLossStats{}} + holdOne := func(i int) { + cs := &connState{fd: 3 + i, generation: uint32(50 + i), sendIsZC: true, sending: true, zcNotifPending: true, kernelInflight: 1, + sendBuf: make([]byte, sendZCMinBytes, zcHoldBytesMax/2)} + w.noteClosedInflight(cs) + w.queuePendingReleaseFD(cs, false, -1) + w.cachedNow = w.pendingRelease[len(w.pendingRelease)-1].releaseAtNanos + 1 + w.drainPendingRelease() + } + holdOne(0) + if op, zc := opcode(w); op != opSENDZC || !zc { + t.Fatalf("with %d bytes held, below the cap of %d, an unlinked %d-byte send was armed as op %d, want SEND_ZC", + w.zcHoldBytes, zcHoldBytesMax, sendZCMinBytes, op) + } + holdOne(1) + if w.zcHoldBytes != zcHoldBytesMax { + t.Fatalf("the worker counts %d bytes held, want %d (two arrays of %d)", w.zcHoldBytes, zcHoldBytesMax, zcHoldBytesMax/2) + } + if op, zc := opcode(w); op != opSEND || zc { + t.Fatalf("with %d bytes held, at the cap, an unlinked %d-byte send was armed as op %d (sendIsZC %v), want a plain SEND", + w.zcHoldBytes, sendZCMinBytes, op, zc) + } + zcNotifCQE(t, w, 3, 50) + if op, zc := opcode(w); op != opSENDZC || !zc { + t.Fatalf("with %d bytes held after one release, an unlinked send was armed as op %d, want SEND_ZC again", w.zcHoldBytes, op) + } +} + +// BenchmarkZCHeldLoopPass is what the release walk costs one loop pass with n +// send buffers held past the backstop: the loop walks pendingRelease only while +// it is not empty, and a held one is not on it. Round 1 walked every held +// entry every pass. +func BenchmarkZCHeldLoopPass(b *testing.B) { + for _, n := range []int{100, 1000, 10000, 100000} { + b.Run(fmt.Sprintf("held=%d", n), func(b *testing.B) { + w := zcHeldWorker(b, n) + b.ResetTimer() + for range b.N { + if len(w.pendingRelease) > 0 { + w.drainPendingRelease() + } + } + b.StopTimer() + if got := w.handoffLoss.zcNotifHeld.Load(); got != uint64(n) || w.handoffLoss.zcNotifForced.Load() != 0 { + b.Fatalf("held %d of %d, forced %d", got, n, w.handoffLoss.zcNotifForced.Load()) + } + }) + } +} + +// BenchmarkClosedIdentityRetireWithZCHolds is what n held send buffers add to +// the CQE path: a closed connection's recv, cancelled at its close, retiring +// its identity through staleConnCQE, as every server-side close with a recv +// armed does. The hold adds one lookup there while any hold exists. +func BenchmarkClosedIdentityRetireWithZCHolds(b *testing.B) { + for _, n := range []int{0, 10000} { + b.Run(fmt.Sprintf("held=%d", n), func(b *testing.B) { + w := zcHeldWorker(b, n) + const fd, gen = 2, 1 << 31 + cs := &connState{fd: fd, generation: gen, recvArmed: true} + c := &completionEntry{UserData: encodeUserDataGen(udRecv, fd, gen), Res: -int32(unix.ECANCELED)} + b.ResetTimer() + for range b.N { + cs.kernelInflight = 1 + w.noteClosedInflight(cs) + w.staleConnCQE(c, fd, c.UserData) + } + b.StopTimer() + if _, ok := w.closedOps[encodeConnOpKey(fd, gen)]; ok { + b.Fatal("the identity was not retired") + } + }) + } +} diff --git a/engine/iouring/zc_send_buffer_counters_test.go b/engine/iouring/zc_send_buffer_counters_test.go new file mode 100644 index 00000000..7297b05c --- /dev/null +++ b/engine/iouring/zc_send_buffer_counters_test.go @@ -0,0 +1,367 @@ +//go:build linux + +package iouring + +import ( + "runtime" + "testing" + "time" + "unsafe" + + "golang.org/x/sys/unix" +) + +// The accounting and the counters behind the celeris#812 hold (see +// zc_send_buffer.go). TestBackstopHoldsASendBufferAZCNotificationStillReads +// judges the hold by what reaches the wire; these judge what the engine +// records about it, and the shutdown half, which the wire cannot show. + +// TestBackstopCountsItsZCHolds runs every wire case again and reads the +// counters: one hold (CloseZCNotifHeld) per case, however many passes the +// backstop ran past its deadline, no send buffer given up with its SEND_ZC +// owed (CloseZCNotifForced, which must stay 0), and no descriptor forced +// (CloseFDForced): a notification alone names none (celeris#798). +// +// And what the hold costs while it lasts (celeris#813 round 2), at the end of +// the stall, with the notification owed: nothing on the release walk the loop +// makes every pass, one hold for the identity keeping the send buffer's array +// alone, and the held-now gauges at one buffer of its capacity; both back to +// 0 once the notification released it. +func TestBackstopCountsItsZCHolds(t *testing.T) { + for _, tc := range zcBackstopCases { + t.Run(tc.name, func(t *testing.T) { + runtime.LockOSThread() + defer runtime.UnlockOSThread() + r := runZCBackstopCase(t, tc) + held := r.w.handoffLoss.zcNotifHeld.Load() + forced := r.w.handoffLoss.zcNotifForced.Load() + fdForced := r.w.handoffLoss.closeFDForced.Load() + t.Logf("celeris812 counters case=%s close_zc_notif_held=%d close_zc_notif_forced=%d close_fd_forced=%d", tc.name, held, forced, fdForced) + if held != 1 || forced != 0 || fdForced != 0 { + t.Errorf("CloseZCNotifHeld=%d CloseZCNotifForced=%d CloseFDForced=%d, want 1 0 0", held, forced, fdForced) + } + if r.corrupt != 0 || !r.heldPastBackstop || r.owedAtRelease != 0 { + t.Errorf("the hold itself failed: corrupt=%d held_past_backstop=%v owed_at_release=%d", r.corrupt, r.heldPastBackstop, r.owedAtRelease) + } + if r.walkAtStall != 0 || r.holdsAtStall != 1 || !r.arrayAloneAtStall { + t.Errorf("held past the backstop: %d entries left on the release walk, %d holds for the identity, array alone %v; "+ + "want 0, 1, true: every loop pass walks what is left, and the kernel reads nothing of the connState but the array", + r.walkAtStall, r.holdsAtStall, r.arrayAloneAtStall) + } + if r.heldNowAtStall != 1 || r.heldBytesAtStall != int64(r.bufCap) { + t.Errorf("held-now gauges %d buffers, %d bytes while held, want 1 and %d", r.heldNowAtStall, r.heldBytesAtStall, r.bufCap) + } + if r.heldNowAfterRelease != 0 || r.heldBytesAfterRelease != 0 { + t.Errorf("held-now gauges %d buffers, %d bytes after the release, want 0 and 0", r.heldNowAfterRelease, r.heldBytesAfterRelease) + } + }) + } +} + +// TestPendingReleaseBackstopHoldsAZCSendPastItsDeadline is the hold's own +// decision on hand-built state, with no kernel, so it runs where SEND_ZC does +// not (the wire test skips there). Past its deadline, on every pass, an entry +// that still owes a SEND_ZC stays held, off the release walk (zcHolds, see +// TestZCHoldIsOffTheReleaseWalk), counted once (CloseZCNotifHeld), and the +// notification releases it; nothing is forced (CloseZCNotifForced). +// +// - notif-owed: the send's first CQE came before the close; the descriptor +// went at the close (only the notification is owed, celeris#798). +// - send-owed-fd-kept: the send itself was owed at the close, so the entry +// kept the descriptor for it (celeris#685). Past the backstop the +// descriptor goes, forced and counted as the backstop always did +// (CloseFDForced), while the connState stays for the SEND_ZC; the send's +// first CQE and then its notification end the hold. +func TestPendingReleaseBackstopHoldsAZCSendPastItsDeadline(t *testing.T) { + for _, tc := range []struct { + name string + keepFD bool + cqes []uint32 // flags of the send's CQEs after the backstop, in order + fdForced uint64 + }{ + {"notif-owed", false, []uint32{cqeFNotif}, 0}, + {"send-owed-fd-kept", true, []uint32{cqeFMore, cqeFNotif}, 1}, + } { + t.Run(tc.name, func(t *testing.T) { + const fd, gen = 5, 7 + w := &Worker{conns: make([]*connState, 16), handoffLoss: &handoffLossStats{}} + buf := make([]byte, 5000, 8192) + cs := &connState{fd: fd, generation: gen, sendIsZC: true, sending: true, zcNotifPending: !tc.keepFD, kernelInflight: 1, sendBuf: buf} + w.noteClosedInflight(cs) + kept := -1 + if tc.keepFD { + efd, err := unix.Eventfd(0, unix.EFD_CLOEXEC) + if err != nil { + t.Fatalf("eventfd: %v", err) + } + kept = efd + } + w.queuePendingReleaseFD(cs, false, kept) + key := encodeConnOpKey(fd, gen) + deadline := w.pendingRelease[0].releaseAtNanos + for i := range 5 { + w.cachedNow = deadline + int64(i+1)*int64(time.Second) + if len(w.pendingRelease) > 0 { + w.drainPendingRelease() + } + if hs := w.zcHolds[key]; len(w.pendingRelease) != 0 || len(hs) != 1 || unsafe.SliceData(hs[0].sendBuf) != unsafe.SliceData(buf) { + t.Fatalf("pass %d past the backstop: the send buffer was not held for its SEND_ZC off the release walk: walk=%+v holds=%+v", + i, w.pendingRelease, hs) + } + if w.closeFDOwed != 0 { + t.Fatalf("pass %d past the backstop: the descriptor is still kept (closeFDOwed=%d)", i, w.closeFDOwed) + } + } + if held, forced, fdForced := w.handoffLoss.zcNotifHeld.Load(), w.handoffLoss.zcNotifForced.Load(), w.handoffLoss.closeFDForced.Load(); held != 1 || forced != 0 || fdForced != tc.fdForced { + t.Fatalf("CloseZCNotifHeld=%d CloseZCNotifForced=%d CloseFDForced=%d, want 1 0 %d", held, forced, fdForced, tc.fdForced) + } + for i, flags := range tc.cqes { + if len(w.zcHolds[key]) != 1 { + t.Fatalf("released before CQE %d (flags %#x)", i, flags) + } + c := &completionEntry{UserData: encodeUserDataGen(udSend, fd, gen), Res: 100, Flags: flags} + if !w.staleConnCQE(c, fd, c.UserData) { + t.Fatalf("CQE %d not taken as stale", i) + } + if len(w.pendingRelease) > 0 { + w.drainPendingRelease() + } + } + if len(w.pendingRelease) != 0 || len(w.zcHolds) != 0 { + t.Fatalf("the notification did not release the entry: walk=%+v holds=%+v", w.pendingRelease, w.zcHolds) + } + if forced := w.handoffLoss.zcNotifForced.Load(); forced != 0 { + t.Fatalf("CloseZCNotifForced = %d, want 0", forced) + } + }) + } +} + +// TestClosedOpsCountsTheZCSendApart drives the closedOps accounting with +// completions built by hand, so each way a SEND_ZC can end is checked without +// a kernel: zcOwed must be 1 from the close until the op's terminal CQE, and +// only that CQE may clear it. +// +// - notif-owed: closed after the send's first CQE; the notification ends it. +// - send-then-notif: closed with the send itself owed; its first CQE +// (IORING_CQE_F_MORE) leaves the notification owed, which ends it. +// - send-without-notif: a SEND_ZC whose first CQE came without F_MORE has +// no notification to follow; that CQE ends it. +// - recv-and-notif: a recv's terminal CQE must not end it. +// - plain-send: a plain SEND owes no SEND_ZC at all. +// - collision: two conns under one (fd, generation), one owing a SEND_ZC's +// notification and one a plain send: the plain send's CQE cannot be told +// from a SEND_ZC's first CQE without F_MORE, so it must not end the hold; +// the notification does. +func TestClosedOpsCountsTheZCSendApart(t *testing.T) { + const fd, gen = 7, 21 + sendUD := encodeUserDataGen(udSend, fd, gen) + recvUD := encodeUserDataGen(udRecv, fd, gen) + first := &completionEntry{UserData: sendUD, Res: 100, Flags: cqeFMore} + notif := &completionEntry{UserData: sendUD, Flags: cqeFNotif} + plain := &completionEntry{UserData: sendUD, Res: 100} + recvEnd := &completionEntry{UserData: recvUD, Res: -int32(unix.ECANCELED)} + type conn struct{ zc, sending, notifPending, recv bool } + for _, tc := range []struct { + name string + conns []conn + cqes []*completionEntry + owed []uint8 // zcOwed after the close, then after each CQE + closed bool // the identity retired after the last CQE + }{ + {"notif-owed", []conn{{zc: true, sending: true, notifPending: true}}, []*completionEntry{notif}, []uint8{1, 0}, true}, + {"send-then-notif", []conn{{zc: true, sending: true}}, []*completionEntry{first, notif}, []uint8{1, 1, 0}, true}, + {"send-without-notif", []conn{{zc: true, sending: true}}, []*completionEntry{plain}, []uint8{1, 0}, true}, + {"recv-and-notif", []conn{{zc: true, sending: true, notifPending: true, recv: true}}, []*completionEntry{recvEnd, notif}, []uint8{1, 1, 0}, true}, + {"plain-send", []conn{{sending: true}}, []*completionEntry{plain}, []uint8{0, 0}, true}, + {"collision", []conn{{zc: true, sending: true, notifPending: true}, {sending: true}}, []*completionEntry{plain, notif}, []uint8{1, 1, 0}, true}, + } { + t.Run(tc.name, func(t *testing.T) { + w := &Worker{conns: make([]*connState, fd+1), handoffLoss: &handoffLossStats{}} + var css []*connState + for _, c := range tc.conns { + cs := &connState{fd: fd, generation: gen, sendIsZC: c.zc, sending: c.sending, zcNotifPending: c.notifPending, recvArmed: c.recv} + if c.sending || c.notifPending { + cs.kernelInflight++ + } + if c.recv { + cs.kernelInflight++ + } + w.noteClosedInflight(cs) + css = append(css, cs) + } + key := encodeConnOpKey(fd, gen) + owed := func() uint8 { + if e := w.closedOps[key]; e != nil { + return e.zcOwed + } + return 0 + } + if got := owed(); got != tc.owed[0] { + t.Fatalf("zcOwed = %d after the close, want %d", got, tc.owed[0]) + } + if got, want := w.closedZCOwed(css[0]), tc.owed[0] > 0; got != want { + t.Fatalf("closedZCOwed = %v after the close, want %v", got, want) + } + for i, c := range tc.cqes { + if !w.staleConnCQE(c, fd, c.UserData) { + t.Fatalf("CQE %d was not taken as stale", i) + } + if got := owed(); got != tc.owed[i+1] { + t.Fatalf("zcOwed = %d after CQE %d (ud=%#x flags=%#x), want %d", got, i, c.UserData, c.Flags, tc.owed[i+1]) + } + } + if _, ok := w.closedOps[key]; ok == tc.closed { + t.Fatalf("identity still registered = %v after the last CQE, want %v", ok, !tc.closed) + } + for _, cs := range css { + if w.closedZCOwed(cs) { + t.Fatalf("closedZCOwed still true once every op was retired") + } + } + }) + } +} + +// TestDroppingAnIdentityWithAZCOwedIsCounted is the positive control of the +// must-stay-0 counter: dropping the accounting of a closed conn that still +// owes a SEND_ZC, as the backstop does for any entry it gives up on, is a send +// buffer released with the kernel still able to send from it, and must be +// counted (CloseZCNotifForced). Dropping one that owes only a recv must not. +func TestDroppingAnIdentityWithAZCOwedIsCounted(t *testing.T) { + for _, tc := range []struct { + name string + zc bool + want uint64 + }{ + {"zc-owed", true, 1}, + {"recv-only", false, 0}, + } { + t.Run(tc.name, func(t *testing.T) { + w := &Worker{handoffLoss: &handoffLossStats{}} + cs := &connState{fd: 9, generation: 5, recvArmed: true, kernelInflight: 1} + if tc.zc { + cs.sendIsZC, cs.zcNotifPending = true, true + cs.kernelInflight++ + } + w.noteClosedInflight(cs) + w.dropClosedOps(cs) + if got := w.handoffLoss.zcNotifForced.Load(); got != tc.want { + t.Fatalf("CloseZCNotifForced = %d, want %d", got, tc.want) + } + if len(w.closedOps) != 0 { + t.Fatalf("the identity is still registered after dropClosedOps") + } + }) + } +} + +// zcRetainedSince reports whether the array behind buf went into zcRetained +// after its first n entries, and how many entries were added. +func zcRetainedSince(n int, buf []byte) (found bool, added int) { + zcRetained.mu.Lock() + defer zcRetained.mu.Unlock() + for _, b := range zcRetained.bufs[n:] { + if unsafe.SliceData(b) == unsafe.SliceData(buf) { + found = true + } + } + return found, len(zcRetained.bufs) - n +} + +func zcRetainedLen() int { + zcRetained.mu.Lock() + defer zcRetained.mu.Unlock() + return len(zcRetained.bufs) +} + +// TestShutdownRetainsSendBuffersAZCMayStillRead is the shutdown half of +// celeris#812. Closing the ring ends no SEND_ZC's use of its send buffer, and +// after it no notification can say when that use ends, so worker shutdown +// must keep, for the life of the process, every send buffer a SEND_ZC may +// still read, and only those. Run on a synthetic worker right where shutdown +// runs it (retainZCSendBufsAtShutdown, before the ring closes): +// +// - live-notif-owed: a live connection whose peer stopped reading owes its +// notification: its sendBuf array must be retained and counted. +// - live-notif-in-ring: the peer read everything and the notification is in +// the completion ring, not yet read: the retention reads it first, and +// nothing is retained. +// - held-after-close: a closed connection the backstop holds for its +// notification (the #812 hold): retained and counted. +// +// In every case the worker takes its share out of the held-now gauges: a +// worker that has shut down holds nothing, and what it retained is counted +// apart (ShutdownZCBufRetained). +func TestShutdownRetainsSendBuffersAZCMayStillRead(t *testing.T) { + for _, tc := range []struct { + name string + peerReads bool + closeFirst bool + wantRetain bool + wantCounted uint64 + }{ + {"live-notif-owed", false, false, true, 1}, + {"live-notif-in-ring", true, false, false, 0}, + {"held-after-close", false, true, true, 1}, + } { + t.Run(tc.name, func(t *testing.T) { + runtime.LockOSThread() + defer runtime.UnlockOSThread() + w, cs, peer := newZCCloseWorker(t) + sent := startZCSend(t, w, cs, true) + buf := cs.sendBuf + if tc.peerReads { + if err := unix.SetNonblock(peer, true); err != nil { + t.Fatal(err) + } + got, rb := 0, make([]byte, 64<<10) + for end := time.Now().Add(3 * time.Second); got < int(sent) && time.Now().Before(end); { + if n, _ := unix.Read(peer, rb); n > 0 { + got += n + continue + } + time.Sleep(time.Millisecond) + } + if got != int(sent) { + t.Fatalf("the peer read %d of %d", got, sent) + } + for end := time.Now().Add(3 * time.Second); ; { + if head, tail := w.ring.BeginCQ(); tail != head { + break + } + if time.Now().After(end) { + t.Fatal("no notification in the ring 3s after the peer read everything") + } + time.Sleep(time.Millisecond) + } + } + if tc.closeFirst { + sweepClose(t, w, cs) + w.pendingRelease[0].releaseAtNanos = time.Now().UnixNano() + for range 3 { + runRingOnce(t, w, 10*time.Millisecond) + w.cachedNow = time.Now().UnixNano() + w.drainPendingRelease() + } + if len(w.pendingRelease) != 0 || len(w.zcHolds) != 1 || w.handoffLoss.zcHeldNow.Load() != 1 { + t.Fatalf("the backstop did not hold the send buffer off the release walk: walk=%+v holds=%+v held_now=%d", + w.pendingRelease, w.zcHolds, w.handoffLoss.zcHeldNow.Load()) + } + } + n := zcRetainedLen() + w.retainZCSendBufsAtShutdown() + found, added := zcRetainedSince(n, buf) + counted := w.handoffLoss.zcBufRetained.Load() + t.Logf("celeris812 shutdown case=%s sent=%d retained=%v added=%d shutdown_zc_buf_retained=%d kernelInflight=%d", + tc.name, sent, found, added, counted, cs.kernelInflight) + if found != tc.wantRetain || added != int(tc.wantCounted) || counted != tc.wantCounted { + t.Fatalf("retained=%v added=%d ShutdownZCBufRetained=%d, want %v %d %d", found, added, counted, + tc.wantRetain, tc.wantCounted, tc.wantCounted) + } + if now, b := w.handoffLoss.zcHeldNow.Load(), w.handoffLoss.zcHeldBytes.Load(); now != 0 || b != 0 { + t.Fatalf("held-now gauges %d buffers, %d bytes after the worker's shutdown, want 0 and 0", now, b) + } + }) + } +} diff --git a/engine/iouring/zc_send_buffer_hold_test.go b/engine/iouring/zc_send_buffer_hold_test.go new file mode 100644 index 00000000..d01c5700 --- /dev/null +++ b/engine/iouring/zc_send_buffer_hold_test.go @@ -0,0 +1,246 @@ +//go:build linux + +package iouring + +import ( + "runtime" + "sync" + "testing" + "time" + "unsafe" + + "golang.org/x/sys/unix" +) + +// celeris#812: the kernel reads a SEND_ZC's send buffer (cs.sendBuf) until the +// send's notification CQE. A peer that stops reading keeps the unsent tail of +// the send queued on the socket, still referencing those pages, and the +// notification with it: for as long as the socket lives, which after a close +// is as long as the kernel keeps the orphaned socket trying to deliver that +// tail. The pendingRelease backstop released such a connState 5 s after the +// close, to connStatePool with its sendBuf (a detached one to the GC). The +// next occupant's response went into that array, and when the peer read +// again, the orphaned socket sent the tail from it: another connection's +// bytes. + +// zcStallPastBackstop is how long the peer stays stalled with the backstop +// already due: about 30 loop passes, every one of which would release a +// connState the backstop does not hold. +const zcStallPastBackstop = 300 * time.Millisecond + +// zcBackstopCase is one case of the backstop tests: what closes, and how. +type zcBackstopCase struct { + name string + reap, recv, detached bool +} + +// zcBackstopCases are the ways a connection is closed with a SEND_ZC +// notification owed: +// +// - notif-only: the send's first CQE was read before the close; only the +// notification is owed. +// - recv-and-notif: a recv was armed after the send too; the close's +// cancel ends it, and the notification is left. +// - send-done-after-close: the send's first CQE was still in the ring, +// unread, at the close. +// - detached-notif-only: notif-only on a detached connection (detachMu), +// which the release drops for the GC instead of pooling. +var zcBackstopCases = []zcBackstopCase{ + {"notif-only", true, false, false}, + {"recv-and-notif", true, true, false}, + {"send-done-after-close", false, false, false}, + {"detached-notif-only", true, false, true}, +} + +// zcBackstopResult is what one case measured. +type zcBackstopResult struct { + w *Worker + sent int32 + got int // bytes the peer read + eof bool // the peer read to EOF + corrupt, first int // bytes that are not what was sent, and the first one's offset (-1: none) + heldPastBackstop bool + releasedAfter time.Duration // from the close to the array's release; -1: not released + owedAtRelease int32 // ops the closed identity still owed at the release: 0 once every op retired + + // What the hold looked like at the end of the stall, with the backstop + // past and the notification owed (celeris#813 round 2): the entries the + // release walk still visits, the holds kept for the identity off that + // walk, whether the one hold keeps the array alone (its connState + // released), and the held-now gauges. Then the gauges after the + // release. + walkAtStall int + holdsAtStall int + arrayAloneAtStall bool + heldNowAtStall int64 + heldBytesAtStall int64 + heldNowAfterRelease int64 + heldBytesAfterRelease int64 + bufCap int // the array's capacity: what a hold of it counts in bytes +} + +// zcIdentityOwed is what closed identity key still owes the kernel, by its +// closedOps entry: 0 once every op retired (the entry is gone). +func zcIdentityOwed(w *Worker, key uint64) int32 { + if e := w.closedOps[key]; e != nil { + return e.inflight + } + return 0 +} + +// runZCBackstopCase builds the case with the #798 fixture: a 64 KiB SEND_ZC to +// a loopback peer with a 4 KiB receive buffer that reads nothing, closed as +// the closing-drain sweep closes it. The backstop is then made due at once +// (the entry's releaseAtNanos, set to the close: the 5 s wait itself is not +// what is tested), and the loop runs for zcStallPastBackstop with the peer +// still stalled. The moment the array the SEND_ZC was given stops being held +// (the connState leaves pendingRelease, or stays there without the array), +// the test does what the array's next owner does and writes over it (0xAA). +// Then the peer reads to EOF while the loop runs, and every byte it read is +// compared with the byte that was sent (byte(i)). +func runZCBackstopCase(t *testing.T, tc zcBackstopCase) zcBackstopResult { + t.Helper() + w, cs, peer := newZCCloseWorker(t) + if tc.detached { + cs.detachMu = new(sync.Mutex) + } + r := zcBackstopResult{w: w, first: -1, releasedAfter: -1, owedAtRelease: -1} + r.sent = startZCSend(t, w, cs, tc.reap) + pinned := cs.sendBuf[:cap(cs.sendBuf)] + r.bufCap = cap(pinned) + // The identity, taken before the close: a connState released to the pool + // keeps neither its descriptor number nor its count. + key := encodeConnOpKey(cs.fd, cs.generation) + if len(cs.sendBuf) != zcPayload || pinned[1] != 1 || pinned[255] != 255 { + t.Fatalf("sendBuf is not the payload: len=%d", len(cs.sendBuf)) + } + if tc.recv { + if !w.prepareRecv(cs, cs.buf) || cs.kernelInflight != 2 { + t.Fatalf("recv not placed: kernelInflight=%d", cs.kernelInflight) + } + if res := runRingOnce(t, w, 10*time.Millisecond); len(res) != 0 { + t.Fatalf("the recv completed before the close (%v): the peer sent nothing", res) + } + } + sweepClose(t, w, cs) + if len(w.pendingRelease) != 1 || w.pendingRelease[0].cs != cs || w.pendingRelease[0].detached != tc.detached { + t.Fatalf("the close did not queue the connState for release as expected: %+v", w.pendingRelease) + } + // The backstop is due from here on. + closedAt := time.Now() + w.pendingRelease[0].releaseAtNanos = closedAt.UnixNano() + + // held: the array is still the send buffer of a connState queued for + // release, or held past the backstop off the release walk (zcHolds): + // with its whole entry, or alone. Held alone, the connState it came + // from must not still own it, unless that connState went to the GC + // (detached) rather than to the pool. Anything else gives the array to a + // next owner. + same := func(b []byte) bool { return unsafe.SliceData(b) == unsafe.SliceData(pinned) } + held := func() bool { + for i := range w.pendingRelease { + if w.pendingRelease[i].cs == cs { + return same(cs.sendBuf) + } + } + for _, h := range w.zcHolds[key] { + if !same(h.sendBuf) { + continue + } + if h.entry.cs == cs { + return same(cs.sendBuf) + } + return tc.detached || !same(cs.sendBuf) + } + return false + } + pass := func(d time.Duration) { + runRingOnce(t, w, d) + owed := zcIdentityOwed(w, key) + w.cachedNow = time.Now().UnixNano() + w.drainPendingRelease() + if r.releasedAfter < 0 && !held() { + r.releasedAfter = time.Since(closedAt) + r.owedAtRelease = owed + // The next owner of the array writes into it. + for i := range pinned { + pinned[i] = 0xAA + } + } + } + for time.Since(closedAt) < zcStallPastBackstop { + pass(10 * time.Millisecond) + } + r.heldPastBackstop = r.releasedAfter < 0 + r.walkAtStall = len(w.pendingRelease) + r.holdsAtStall = len(w.zcHolds[key]) + r.arrayAloneAtStall = r.holdsAtStall == 1 && w.zcHolds[key][0].entry.cs == nil && same(w.zcHolds[key][0].sendBuf) + r.heldNowAtStall, r.heldBytesAtStall = w.handoffLoss.zcHeldNow.Load(), w.handoffLoss.zcHeldBytes.Load() + + if err := unix.SetNonblock(peer, true); err != nil { + t.Fatalf("peer nonblock: %v", err) + } + got := make([]byte, 0, zcPayload) + buf := make([]byte, 64<<10) + for end := time.Now().Add(5 * time.Second); time.Now().Before(end) && (!r.eof || r.releasedAfter < 0); { + for !r.eof { + n, err := unix.Read(peer, buf) + if n > 0 { + got = append(got, buf[:n]...) + continue + } + if n == 0 && err == nil { + r.eof = true + } + break + } + pass(5 * time.Millisecond) + } + r.got = len(got) + r.heldNowAfterRelease, r.heldBytesAfterRelease = w.handoffLoss.zcHeldNow.Load(), w.handoffLoss.zcHeldBytes.Load() + for i, b := range got { + if b != byte(i) { + r.corrupt++ + if r.first < 0 { + r.first = i + } + } + } + t.Logf("celeris812 backstop case=%s sent=%d held_past_backstop=%v released_after=%v owed_at_release=%d peer_got=%d eof=%v corrupt_bytes=%d first_corrupt=%d "+ + "walk_at_stall=%d holds_at_stall=%d array_alone=%v held_now=%d held_bytes=%d held_now_after=%d held_bytes_after=%d", + tc.name, r.sent, r.heldPastBackstop, r.releasedAfter.Round(time.Microsecond), r.owedAtRelease, r.got, r.eof, r.corrupt, r.first, + r.walkAtStall, r.holdsAtStall, r.arrayAloneAtStall, r.heldNowAtStall, r.heldBytesAtStall, r.heldNowAfterRelease, r.heldBytesAfterRelease) + return r +} + +// TestBackstopHoldsASendBufferAZCNotificationStillReads is celeris#812: in +// every case of zcBackstopCases the backstop must hold the connState, and so +// its send buffer, while the notification is owed, however long past the +// backstop, and release it once the notification arrives. A release before +// the notification puts the 0xAA on the wire: the unsent tail, some 61 KB, in +// every case. Held for good, the connState would be a leak. +func TestBackstopHoldsASendBufferAZCNotificationStillReads(t *testing.T) { + for _, tc := range zcBackstopCases { + t.Run(tc.name, func(t *testing.T) { + runtime.LockOSThread() + defer runtime.UnlockOSThread() + r := runZCBackstopCase(t, tc) + if r.corrupt != 0 { + t.Errorf("CORRUPT: %d of the %d bytes the peer read are not what was sent (first at offset %d): "+ + "the connState, and the send buffer the kernel was still sending from, left pendingRelease "+ + "%v after the close with %d op(s) owed", r.corrupt, r.got, r.first, r.releasedAfter, r.owedAtRelease) + } + if !r.heldPastBackstop { + t.Errorf("the backstop released the connState %v after the close with %d op(s) owed: a SEND_ZC "+ + "notification was still owed, so the kernel could still read its send buffer", r.releasedAfter, r.owedAtRelease) + } + if r.got != int(r.sent) || !r.eof { + t.Errorf("the peer read %d bytes (EOF %v), want the %d the send completed with and then EOF", r.got, r.eof, r.sent) + } + if r.releasedAfter < 0 || r.owedAtRelease != 0 { + t.Errorf("the send buffer was not released at the notification: released_after=%v owed_at_release=%d "+ + "pendingRelease=%d holds=%d", r.releasedAfter, r.owedAtRelease, len(r.w.pendingRelease), len(r.w.zcHolds)) + } + }) + } +}