//go:build linux
package iouring
// Issue probe (the pendingRelease backstop releases a connState while its
// SEND_ZC notification is owed). Not for commit. Runs on main as is: the
// fixture is #793's newZCCloseWorker / startZCSend / sweepClose, ported
// verbatim minus #793-only fields (closeFDForced), under zcbs names.
//
// The kernel reads a SEND_ZC's send buffer (cs.sendBuf) until its
// notification. This probe models what the pool does with that buffer once
// the connState is released: the next owner writes into it. On every
// drainPendingRelease pass, once cs has left pendingRelease, the array the
// SEND_ZC pinned is overwritten with 0xAA (once). Then the peer reads to EOF
// and every byte is compared with what was sent (byte(i)).
//
// Cases: peer stalled 300 ms after the close (the notification arrives
// before any release: the probe's negative arm), and peer stalled 5.5 s
// (past the 5 s pendingRelease backstop).
import (
"os"
"runtime"
"strconv"
"sync/atomic"
"testing"
"time"
"golang.org/x/sys/unix"
"github.com/goceleris/celeris/engine"
"github.com/goceleris/celeris/engine/internal/errclass"
"github.com/goceleris/celeris/internal/conn"
)
const zcbsPayload = 64 << 10
func zcbsFDTarget(fd int) string {
s, err := os.Readlink("/proc/self/fd/" + strconv.Itoa(fd))
if err != nil {
return ""
}
return s
}
func zcbsRunRing(t *testing.T, w *Worker, d time.Duration) {
t.Helper()
if err := w.ring.SubmitAndWaitTimeout(d); err != nil {
t.Fatalf("submit: %v", err)
}
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)
}
}
w.ring.EndCQ(head)
}
func zcbsWorker(t *testing.T) (*Worker, *connState, int) {
t.Helper()
res, reason := probeSendZCCached()
if on, _ := resolveSendZCPolicy(res == SendZCTrueZeroCopy || res == SendZCCopyFallback, "auto"); !on {
t.Skipf("the engine does not turn SEND_ZC on here under the default policy (probe: %v %s)", res, reason)
}
t.Logf("zcbs send_zc_probe=%v policy=auto", res)
ring := newTestRing(t)
lfd, err := unix.Socket(unix.AF_INET, unix.SOCK_STREAM|unix.SOCK_CLOEXEC, 0)
if err != nil {
t.Fatalf("listen socket: %v", err)
}
defer func() { _ = unix.Close(lfd) }()
if err := unix.Bind(lfd, &unix.SockaddrInet4{Addr: [4]byte{127, 0, 0, 1}}); err != nil {
t.Fatalf("bind: %v", err)
}
if err := unix.Listen(lfd, 1); err != nil {
t.Fatalf("listen: %v", err)
}
sa, err := unix.Getsockname(lfd)
if err != nil {
t.Fatalf("getsockname: %v", err)
}
peer, err := unix.Socket(unix.AF_INET, unix.SOCK_STREAM|unix.SOCK_CLOEXEC, 0)
if err != nil {
t.Fatalf("peer socket: %v", err)
}
if err := unix.SetsockoptInt(peer, unix.SOL_SOCKET, unix.SO_RCVBUF, 4096); err != nil {
_ = unix.Close(peer)
t.Fatalf("SO_RCVBUF: %v", err)
}
if err := unix.Connect(peer, sa); err != nil {
_ = unix.Close(peer)
t.Fatalf("connect: %v", err)
}
local, _, err := unix.Accept4(lfd, unix.SOCK_NONBLOCK|unix.SOCK_CLOEXEC)
if err != nil {
_ = unix.Close(peer)
t.Fatalf("accept: %v", err)
}
localTarget := zcbsFDTarget(local)
t.Cleanup(func() {
_ = unix.Close(peer)
if zcbsFDTarget(local) == localTarget {
_ = unix.Close(local)
}
})
w := &Worker{
ring: ring,
conns: make([]*connState, local+1),
liveConns: make([]int, 0, 4),
errs: &errclass.Counters{},
activeConns: &atomic.Int64{},
closeCount: &atomic.Uint64{},
recvArm: &recvArmStats{},
handoffLoss: &handoffLossStats{},
sendZC: true,
}
w.cachedNow = time.Now().UnixNano()
cs := &connState{
fd: local,
liveIdx: -1,
generation: 13,
buf: make([]byte, 4096),
h1State: conn.NewH1State(),
detected: true,
}
cs.protocol.Store(int32(engine.HTTP1))
w.conns[local] = cs
w.addLiveConn(cs)
w.connCount = 1
w.activeConns.Add(1)
return w, cs, peer
}
func zcbsStartSend(t *testing.T, w *Worker, cs *connState, reap bool) int32 {
t.Helper()
payload := make([]byte, zcbsPayload)
for i := range payload {
payload[i] = byte(i)
}
cs.writeBuf = payload
if w.flushSend(cs) || !cs.sendIsZC || cs.kernelInflight != 1 {
t.Fatalf("no SEND_ZC placed: sendIsZC=%v kernelInflight=%d", cs.sendIsZC, cs.kernelInflight)
}
if _, err := w.ring.Submit(); err != nil {
t.Fatalf("submit: %v", err)
}
var c *completionEntry
for end := time.Now().Add(2 * time.Second); c == nil; {
if head, tail := w.ring.BeginCQ(); head != tail {
c = w.ring.cqeAt(head)
break
}
if time.Now().After(end) {
t.Fatal("no completion for the SEND_ZC within 2s")
}
time.Sleep(time.Millisecond)
}
if c.UserData&udMask != udSend || !cqeHasMore(c.Flags) || c.Res <= 0 {
t.Fatalf("first completion ud=%#x flags=%#x res=%d, want the SEND_ZC's result (F_MORE, res > 0)", c.UserData, c.Flags, c.Res)
}
sent := c.Res
if reap {
head, _ := w.ring.BeginCQ()
if !w.staleConnCQE(c, cs.fd, c.UserData) {
w.handleSend(c, cs.fd, time.Now().UnixNano())
}
w.ring.EndCQ(head + 1)
if !cs.zcNotifPending || !cs.sending || cs.kernelInflight != 1 {
t.Fatalf("after the first CQE: zcNotifPending=%v sending=%v kernelInflight=%d, want true true 1", cs.zcNotifPending, cs.sending, cs.kernelInflight)
}
}
time.Sleep(100 * time.Millisecond)
want := uint32(1)
if reap {
want = 0
}
if head, tail := w.ring.BeginCQ(); tail-head != want {
t.Fatalf("%d completions in the ring 100 ms after the send, want %d: the notification arrived with the peer "+
"reading nothing (sent %d of %d), so the case did not form", tail-head, want, sent, zcbsPayload)
}
return sent
}
func zcbsSweepClose(t *testing.T, w *Worker, cs *connState) {
t.Helper()
fd := cs.fd
w.closeConn(fd)
if w.conns[fd] != cs || !cs.closing {
t.Fatalf("closeConn did not defer the close of a conn with a send owed: registered=%v closing=%v", w.conns[fd] == cs, cs.closing)
}
w.removeDirty(cs)
w.finishCloseAny(fd, cs)
if w.conns[fd] != nil {
t.Fatal("the conn is still registered after the sweep's close")
}
}
func TestZZBackstopZCWireBytes(t *testing.T) {
for _, tc := range []struct {
name string
reap bool
stall time.Duration
}{
{"notif-only", true, 300 * time.Millisecond},
{"send-done-after-close", false, 300 * time.Millisecond},
{"backstop-notif-only", true, 5500 * time.Millisecond},
{"backstop-send-done-after-close", false, 5500 * time.Millisecond},
} {
t.Run(tc.name, func(t *testing.T) {
runtime.LockOSThread()
defer runtime.UnlockOSThread()
w, cs, peer := zcbsWorker(t)
fd := cs.fd
target := zcbsFDTarget(fd)
sent := zcbsStartSend(t, w, cs, tc.reap)
pinned := cs.sendBuf[:cap(cs.sendBuf)]
if len(cs.sendBuf) != zcbsPayload || pinned[1] != 1 || pinned[255] != 255 {
t.Fatalf("sendBuf is not the payload: len=%d", len(cs.sendBuf))
}
closedAt := time.Now()
zcbsSweepClose(t, w, cs)
fdAfter := time.Duration(-1)
csAfter := time.Duration(-1)
owedAtRelease := int32(-1)
reused := false
pass := func() {
if fdAfter < 0 && zcbsFDTarget(fd) != target {
fdAfter = time.Since(closedAt)
}
inflight := cs.kernelInflight
w.cachedNow = time.Now().UnixNano()
w.drainPendingRelease()
if !reused && len(w.pendingRelease) == 0 {
csAfter = time.Since(closedAt)
owedAtRelease = inflight
for i := range pinned {
pinned[i] = 0xAA
}
reused = true
}
}
for time.Since(closedAt) < tc.stall {
zcbsRunRing(t, w, 10*time.Millisecond)
pass()
}
if err := unix.SetNonblock(peer, true); err != nil {
t.Fatal(err)
}
got := make([]byte, 0, zcbsPayload)
buf := make([]byte, 64<<10)
eof := false
for end := time.Now().Add(5 * time.Second); time.Now().Before(end) && !(eof && len(w.pendingRelease) == 0); {
for !eof {
n, err := unix.Read(peer, buf)
if n > 0 {
got = append(got, buf[:n]...)
continue
}
if n == 0 && err == nil {
eof = true
}
break
}
zcbsRunRing(t, w, 5*time.Millisecond)
pass()
}
bad, first, aa := 0, -1, 0
for i, b := range got {
if b != byte(i) {
bad++
if first < 0 {
first = i
}
}
if b == 0xAA && i >= 4096 {
aa++
}
}
t.Logf("zcbswire case=%s sent=%d fd_released_after=%v cs_released_after=%v inflight_at_cs_release=%d peer_got=%d eof=%v corrupt_bytes=%d first_corrupt=%d aa_bytes_from_4096=%d",
tc.name, sent, fdAfter.Round(time.Microsecond), csAfter.Round(time.Microsecond), owedAtRelease, len(got), eof, bad, first, aa)
if bad != 0 {
t.Errorf("CORRUPT: %d of %d bytes the peer read are not what was sent (first at %d): the send buffer was reused before the kernel let go of it", bad, len(got), first)
}
if len(got) != int(sent) || !eof {
t.Errorf("peer got %d of %d bytes, eof=%v", len(got), sent, eof)
}
})
}
}
On io_uring, the 5 s
pendingReleasebackstop releases a closed connection'sconnStatewhile the kernel still owes that connection's SEND_ZC notification. A plain connection'sconnStategoes back toconnStatePoolwith itssendBufarray. A detached one (WS, SSE, async) is dropped for the GC. In both cases the kernel is not done with that array. The unsent tail of the zero-copy send is still queued on the closed, orphaned socket, and the kernel reads it from those pages when the peer opens its window again.So a peer can stall for more than about 10 s (the 5 s closing-drain sweep plus the 5 s backstop) and then read. It receives whatever the array holds by then: another connection's response bytes once the pool has handed the
connStateon, or any heap data once the GC has reused the memory.Measured on main
00d985cwith a probe that stands in for the next owner by writing0xAAover the released array:0xAA. That is 61,440 of the 65,536 sent, in 5 of 5 runs of each backstop case.This is pre-existing and not introduced by #793. The round-3 review of #793 found it. The probe gives the same bytes on #793's heads.
Mechanism (line numbers at main
00d985c)sendZCMinBytes(4096,engine/iouring/consts.go:114) or more goes out asIORING_OP_SEND_ZC(useSendZC,engine/iouring/worker.go:5251-5253;prepSendSQE,worker.go:5217), pointing atcs.sendBuf.flushSendswapswriteBufintosendBuf(worker.go:5164). The kernel reads those pages until every skb that references them is freed. For TCP, that is once the bytes have been sent and ACKed, and only then does it post the notification CQE.F_MORE) arrives.zcNotifPendingis set, andkernelInflightstays at 1 for the notification.closeConndefers the close becausezcNotifPendingis set (worker.go:3760-3773). The closing-drain sweep incheckTimeoutsthen tears it down afterclosingDrainTimeoutNanos, which is 5 s (worker.go:268,worker.go:5535-5545), throughfinishCloseAny.finishClosequeues theconnStateand closes the descriptor (worker.go:4032-4034, thenworker.go:4079-4093):cancelConnOps(worker.go:3812) submits an ASYNC_CANCEL for the send. It cannot end a notification, because the request itself already completed.noteClosedInflightregisters the owed op inclosedOps.queuePendingRelease(worker.go:3903) setsreleaseAtNanos = now + pendingReleaseHoldNanos, which is 5 s (worker.go:231).close()alone on the H1 fast path, andshutdown(SHUT_WR)plus a receive-queue drain on the graceful path. With nothing unread in the receive queue, no RST is sent. The socket is now orphaned but keeps the unsent tail, which references the pinned pages, and keeps probing the peer's zero window.finishCloseDetachedand thenqueuePendingReleaseDetachedinstead (worker.go:4148-4150,worker.go:3924).connState5 s later.drainPendingRelease(worker.go:3953) findskernelInflight == 1pastreleaseAtNanos(worker.go:3963-3973). It logs the WARNreleasing connState with kernel ops unaccounted for after backstop hold, callsdropClosedOps, and callsreleaseConnState(cs), which ends inconnStatePool.Put(cs)(engine/iouring/conn.go:567).connState.acquireConnStateonly truncatessendBuf(conn.go:471), so the next occupant's handler output lands in the array after its nextflushSendswap.The comment on
pendingReleaseHoldNanos(worker.go:206-231) says the wall clock only matters "if the kernel never delivers a terminal CQE at all", anddrainPendingReleasecalls its firing a kernel anomaly. A SEND_ZC notification behind a stalled peer is not an anomaly. The kernel is still using the buffer legitimately, for as long as the peer stays stalled. The comment calls releasing at the backstop "a potential use-after-free". For this op it is an actual one: the kernel reads the buffer after the release whenever the peer comes back.Reproduction
The probe is
engine/iouring/zz_zcbs_main_test.go, full source below. Its fixture is #793'snewZCCloseWorker/startZCSend/sweepClose, ported to main.Setup:
autopolicy. The startup probe reportscopy fallbackon loopback, andautoenables it.SO_RCVBUF=4096that reads nothing.flushSendofbyte(i).Steps:
closeConn, thenfinishCloseAny. It skips the sweep's 5 s wait, so the backstop fires 5 s after the close here, where production takes 10 s from the moment the close starts.drainPendingRelease.pendingReleaseempties, it writes0xAAover the array the SEND_ZC was given.Each run used one Docker container:
golang:1.27,--security-opt seccomp=unconfined,--cpus 4,-race, kernel 7.0.12-linuxkit, aarch64. Shapes were memlock 8 MiB (CI's unit job) and unlimited.connStatereleased after the close00d985c00d985c00d985c00d985cMain was run 3 times at memlock 8 MiB and 2 times unlimited. There were 0 data races.
Where the 61,200 comes from: 61,200 = 65,536 - 4,096 - 240. The peer's 4 KiB receive buffer took the first 4,096 bytes before the stall. Every one of the remaining 61,440 bytes arrived as
0xAA: the probe'saa_bytes_from_4096=61440in every failing run. Of those positions, 240 hold0xAAinbyte(i)anyway, so they do not count as corrupt.Copy fallback does not protect. The probe classifies this host's SEND_ZC as copy fallback. Even so, the bytes show the tail was read from the pages after the overwrite. On loopback the kernel's copy happens when a segment is transmitted, not when the SEND_ZC is submitted. On a NIC with true zero-copy, the DMA reads the pages at transmit time.
#793's heads show the same thing (the round-3 review's runs):
bf129c8: 8 runs of each backstop case.0f36096: 2 runs.bf129c8with0f36096'sworker.goandfd_lifetime.go(the control): 1 run.Every backstop run read 61,200 corrupt bytes, and every 300 ms run read 0. A mutant that releases the buffer at the close corrupted all five of that probe's cases, which shows the probe detects an early release in every case it builds.
amd64 was compile-checked only (
GOOS=linux GOARCH=amd64 go vet ./engine/iouring/on both trees). All the runs above are arm64. The mechanism is kernel TCP plus the Go pool, with nothing arch-specific in it, but it has not been measured on amd64.Production path and preconditions
Engine and SEND_ZC. The io_uring engine (or adaptive while on io_uring) with SEND_ZC enabled.
CELERIS_IOURING_SEND_ZCunset orautoenables it whenever the startup functional probe passes, copy fallback included (resolveSendZCPolicy,engine/iouring/probe.go:301).onenables it too.offavoids it.An unlinked send of 4096 bytes or more. That covers:
flushSend(worker.go:1347).flushSendLinknever links (worker.go:5296-5298).worker.go:5301-5303). A large H1 response to a slow reader gets here.The plain H1 keep-alive response is linked and goes out as a plain SEND.
A stalled peer. Part of that send is still queued behind the peer's closed window, because the peer stopped reading.
The server closes the connection during the stall. Examples: a write or idle timeout, a WS or SSE close, a handler that ends the stream.
The peer resumes reading more than about 10 s after that close. That is 5 s of closing drain plus 5 s of backstop.
The memory is reused before the peer reads.
connStateto a new connection, and that connection's second flush writes into the array.Who receives what. A client that deliberately stops reading, waits for the server to give up on it, and then reads can receive up to the unsent tail of its own response. For example, it can open an SSE or WS stream or an H2 request for 4 KiB or more. The tail's bytes are overwritten with another client's response or with other heap memory.
Signal on main. The only signal today is the backstop WARN log.
What #793 changes
#793 (head
bf129c8, approved, not merged) does not change this release: the probe reads the same 61,200 bytes on its heads. It changes the signal:0f36096), the close path kept the descriptor while the notification was owed. The backstop force-closed it and countedCloseFDForced, which was 1 in every backstop run.bf129c8), the fix for Follow-ups from #793: a SEND_ZC notification counted as an owed op, the unpinned park gate and shutdown drain, and two counter/exit nits #798 item 1 releases the descriptor as soon as only notifications are owed.CloseFDForcedstays 0 in this case: 0 in every backstop run onbf129c8. The release of theconnStateis then visible only through the WARN log, which now also carriesholds_fd.#798 item 1 says "The connState, whose send buffer the notification guards, still waits for them". That wait is bounded by the backstop. So once #793 lands, no counter sees this case, and the
CloseFDForced == 0gates cannot catch it.Fix direction
connState, or drop the last reference to its send buffer, while a SEND_ZC on it is unfinished, meaning its first CQE or its notification is still owed. Keep such an entry inpendingReleasepast the backstop, detached entries included.connStateand its buffer, for as long as the kernel keeps the pages.tcp_max_orphans), its write queue is purged, and the notification is posted. How long that takes is unmeasured here.closedOps(do notdropClosedOpsit), so that the notification still finds the entry and releases it the usual way.byte(i).A sketch of the hold in
drainPendingReleasemakes the probe pass on main: 3 of 3 runs, every case. TheconnStateis released at 5.75 s, after the peer has read and the notification has arrived, with nothing owed.if cs.kernelInflight > 0 { - if entry.releaseAtNanos > w.cachedNow { + if entry.releaseAtNanos > w.cachedNow || (cs.sendIsZC && (cs.sending || cs.zcNotifPending)) { kept = append(kept, *entry) continue }This is only a sketch. It does not count the hold, and its condition still has to be checked against every path that clears
sendingorzcNotifPending, for example the SEND_ZCEINVAL/ENOMEMfallback.Refs: #793, #798, #685, #256 and #498 (the UAF class this queue exists for, and the closing-drain bound), #591, #585.
The probe,
engine/iouring/zz_zcbs_main_test.go(runs on main as is)