Skip to content

io_uring: the pendingRelease backstop returns a connection's send buffer to the pool while a SEND_ZC notification is still owed, so a peer that stalls >10 s then reads receives another connection's bytes #812

Description

@FumingPower3925

On io_uring, the 5 s pendingRelease backstop releases a closed connection's connState while the kernel still owes that connection's SEND_ZC notification. A plain connection's connState goes back to connStatePool with its sendBuf array. 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 connState on, or any heap data once the GC has reused the memory.

Measured on main 00d985c with a probe that stands in for the next owner by writing 0xAA over the released array:

  • Peer stalled 5.5 s (the backstop fires): every byte past what the peer's receive buffer already held came off the wire as 0xAA. That is 61,440 of the 65,536 sent, in 5 of 5 runs of each backstop case.
  • Peer stalled 300 ms (the notification arrives before any release): 0 corrupt bytes, 5 of 5.
  • Release held until the notification (a one-line fix sketch): 0 corrupt bytes in every case, 3 of 3.

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)

  1. The send is zero-copy. An unlinked send of sendZCMinBytes (4096, engine/iouring/consts.go:114) or more goes out as IORING_OP_SEND_ZC (useSendZC, engine/iouring/worker.go:5251-5253; prepSendSQE, worker.go:5217), pointing at cs.sendBuf. flushSend swaps writeBuf into sendBuf (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.
  2. The peer stops reading. The send's first CQE (F_MORE) arrives. zcNotifPending is set, and kernelInflight stays at 1 for the notification.
  3. The connection is closed. Examples: a timeout, a WS idle deadline, or a handler that ends a stream. closeConn defers the close because zcNotifPending is set (worker.go:3760-3773). The closing-drain sweep in checkTimeouts then tears it down after closingDrainTimeoutNanos, which is 5 s (worker.go:268, worker.go:5535-5545), through finishCloseAny.
  4. finishClose queues the connState and closes the descriptor (worker.go:4032-4034, then worker.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.
    • noteClosedInflight registers the owed op in closedOps.
    • queuePendingRelease (worker.go:3903) sets releaseAtNanos = now + pendingReleaseHoldNanos, which is 5 s (worker.go:231).
    • The descriptor is closed: close() alone on the H1 fast path, and shutdown(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.
    • Detached connections take finishCloseDetached and then queuePendingReleaseDetached instead (worker.go:4148-4150, worker.go:3924).
  5. The backstop releases the connState 5 s later. drainPendingRelease (worker.go:3953) finds kernelInflight == 1 past releaseAtNanos (worker.go:3963-3973). It logs the WARN releasing connState with kernel ops unaccounted for after backstop hold, calls dropClosedOps, and calls releaseConnState(cs), which ends in connStatePool.Put(cs) (engine/iouring/conn.go:567).
    • The array stays attached to the pooled connState. acquireConnState only truncates sendBuf (conn.go:471), so the next occupant's handler output lands in the array after its next flushSend swap.
    • A detached entry is not pooled. The worker drops its strong reference, and the array is free for any allocation once the goroutine's closures let go of it.
  6. The peer reads again. Its window opens, and the orphaned socket sends the queued tail from the pinned pages. The peer receives whatever the new owner wrote there.

The comment on pendingReleaseHoldNanos (worker.go:206-231) says the wall clock only matters "if the kernel never delivers a terminal CQE at all", and drainPendingRelease calls 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's newZCCloseWorker / startZCSend / sweepClose, ported to main.

Setup:

  • A synthetic worker over loopback TCP.
  • SEND_ZC on under the default auto policy. The startup probe reports copy fallback on loopback, and auto enables it.
  • A peer with SO_RCVBUF=4096 that reads nothing.
  • One 64 KiB flushSend of byte(i).

Steps:

  1. The probe closes the connection the way the sweep does: closeConn, then finishCloseAny. 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.
  2. It runs the ring and drainPendingRelease.
  3. The moment pendingRelease empties, it writes 0xAA over the array the SEND_ZC was given.
  4. After the stall, the peer reads to EOF, and every byte is compared with what was sent.

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.

go test -race -c -o /tmp/pkg.test ./engine/iouring/
cd engine/iouring && /tmp/pkg.test -test.v -test.count=3 -test.run '^TestZZBackstopZCWireBytes$'
tree case (peer stall) runs connState released after the close ops owed at release corrupt bytes per run
main 00d985c backstop-notif-only (5.5 s) 5 FAIL of 5 5.002-5.005 s 1 61,200
main 00d985c backstop-send-done-after-close (5.5 s) 5 FAIL of 5 5.001-5.013 s 1 61,200
main 00d985c notif-only (300 ms) 5 PASS of 5 552-563 ms 0 0
main 00d985c send-done-after-close (300 ms) 5 PASS of 5 549-562 ms 0 0
main + hold sketch both backstop cases (5.5 s) 3 PASS of 3 each 5.746-5.762 s (after the peer read) 0 0
main + hold sketch both 300 ms cases 3 PASS of 3 each 547-558 ms 0 0

Main 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's aa_bytes_from_4096=61440 in every failing run. Of those positions, 240 hold 0xAA in byte(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.
  • bf129c8 with 0f36096's worker.go and fd_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_ZC unset or auto enables it whenever the startup functional probe passes, copy fallback included (resolveSendZCPolicy, engine/iouring/probe.go:301). on enables it too.
    • Only off avoids it.
  • An unlinked send of 4096 bytes or more. That covers:

    • H2: the write-queue drain uses flushSend (worker.go:1347).
    • Detached WS, SSE and async connections, which flushSendLink never links (worker.go:5296-5298).
    • The remainder of a partial send (worker.go:5301-5303). A large H1 response to a slow reader gets here.
    • Every send on a worker with a provided-buffer ring (multishot recv).
    • A response held for a transplant.
    • A send that found only one free SQE.

    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.

    • On a worker that keeps accepting, the pool hands the connState to a new connection, and that connection's second flush writes into the array.
    • A detached connection's array goes to the GC.

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:

#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 == 0 gates cannot catch it.

Fix direction

  • Hold instead of releasing. Never release a 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 in pendingRelease past the backstop, detached entries included.
    • The hold costs memory only, the connState and its buffer, for as long as the kernel keeps the pages.
    • The kernel bounds that itself. An orphaned socket that cannot deliver is eventually ended (orphan retries, tcp_max_orphans), its write queue is purged, and the notification is posted. How long that takes is unmeasured here.
    • Keep the identity in closedOps (do not dropClosedOps it), so that the notification still finds the entry and releases it the usual way.
  • Count such holds, with a counter of holds and a gauge of what is held now, so a population of stalled peers is visible. The backstop's release of anything else stays the anomaly it is today.
  • Test: the probe above, as a failing-first regression test, on both arches. The backstop cases must read back byte(i).

A sketch of the hold in drainPendingRelease makes the probe pass on main: 3 of 3 runs, every case. The connState is 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 sending or zcNotifPending, for example the SEND_ZC EINVAL/ENOMEM fallback.

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)
//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)
			}
		})
	}
}

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    area/engineEngine interface or implementationbugSomething isn't workingengine/iouringio_uring engine specificssecuritySecurity hardening

    Type

    No type

    Projects

    No projects

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions