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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions engine/epoll/adopt.go
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,9 @@ func (l *Loop) drainAdoptQueue(ctx context.Context, now int64) {
for _, it := range l.adoptQSpare {
l.attachAdoptedFD(ctx, it.fd, it.carry, now)
}
// Same guard as the detach queue: adoptItem carries a Carryover with
// the peer address string, so stale slots pin it after the handoff.
clear(l.adoptQSpare)
l.adoptQSpare = l.adoptQSpare[:0]
}

Expand Down
81 changes: 81 additions & 0 deletions engine/epoll/detach_queue_retention_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
//go:build linux

package epoll

import (
"runtime"
"testing"
)

// TestDrainDetachQueueDropsSlotRefs is the epoll half of the same guard the
// io_uring worker carries. Truncating the hand-off queue with [:0] reuses the
// backing array but leaves every *connState in it reachable, so the queue
// pins its own high-water mark worth of conns — buffers, H1State and, on the
// WebSocket path, whatever the middleware hung off that state — until some
// later drain overwrites each slot.
//
// detachClosed short-circuits the drain body, so the queue can be exercised
// end to end without an epoll fd or a live connection.
func TestDrainDetachQueueDropsSlotRefs(t *testing.T) {
const depth = 64

l := &Loop{}
for range depth {
l.detachQueue = append(l.detachQueue, &connState{fd: -1, detachClosed: true})
}
l.detachQPending.Store(1)

l.drainDetachQueue()

if got := len(l.detachQSpare); got != 0 {
t.Fatalf("spare queue not truncated: len = %d, want 0", got)
}
full := l.detachQSpare[:cap(l.detachQSpare)]
if len(full) < depth {
t.Fatalf("backing array shrank to %d, expected at least the %d enqueued slots", len(full), depth)
}
held := 0
for _, cs := range full[:depth] {
if cs != nil {
held++
}
}
if held != 0 {
t.Fatalf("drainDetachQueue left %d of %d connStates reachable in the reused backing array", held, depth)
}
}

// TestDrainDetachQueueLetsDrainedConnsBeCollected states the same guarantee
// as an observable outcome: a conn the queue has already drained must become
// unreachable once the caller drops its own reference.
func TestDrainDetachQueueLetsDrainedConnsBeCollected(t *testing.T) {
l := &Loop{}
collected := make(chan struct{})

func() {
cs := &connState{fd: -1, detachClosed: true}
runtime.AddCleanup(cs, func(ch chan struct{}) { close(ch) }, collected)
l.detachQueue = append(l.detachQueue, cs)
for range 15 {
l.detachQueue = append(l.detachQueue, &connState{fd: -1, detachClosed: true})
}
}()
l.detachQPending.Store(1)
l.drainDetachQueue()

for range 5 {
runtime.GC()
select {
case <-collected:
runtime.KeepAlive(l)
return
default:
}
}
// KeepAlive is load-bearing: without it the compiler may treat l as dead
// after its last use, the Loop and its queues become collectable, and the
// cleanup runs regardless of what the queue did — passing against the
// unfixed code.
runtime.KeepAlive(l)
t.Fatal("a connState the detach queue already drained is still reachable after 5 GC cycles")
}
5 changes: 5 additions & 0 deletions engine/epoll/loop.go
Original file line number Diff line number Diff line change
Expand Up @@ -2180,6 +2180,11 @@ func (l *Loop) drainDetachQueue() {
}
l.markDirty(cs)
}
// Drop the strong refs before reusing the array (see the io_uring
// worker's drainDetachQueue and drainPendingRelease for the same
// guard): [:0] alone leaves every *connState reachable in the
// backing array until some later drain overwrites its slot.
clear(l.detachQSpare)
l.detachQSpare = l.detachQSpare[:0]
}

Expand Down
96 changes: 96 additions & 0 deletions engine/iouring/detach_queue_retention_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
//go:build linux

package iouring

import (
"runtime"
"testing"
)

// TestDrainDetachQueueDropsSlotRefs guards the hand-off queue against the
// retention hazard drainPendingRelease already guards against: truncating a
// []*connState with [:0] reuses the backing array but leaves every pointer in
// it reachable, so the queue pins its own high-water mark worth of
// connStates until some later drain happens to overwrite each slot.
//
// That matters here more than it would for a plain buffer. A detached
// connState carries its read and write buffers, its H1State, and — on the
// WebSocket path — whatever the middleware hung off that state, so a queue
// that once peaked deep holds all of it for as long as the worker lives.
//
// The test enqueues a deep batch of already-closed conns (detachClosed short-
// circuits the drain body, so no ring is needed), drains, and then reads the
// spare slice out to its capacity. Every slot past the new length must be nil.
func TestDrainDetachQueueDropsSlotRefs(t *testing.T) {
const depth = 64

w := &Worker{}
for range depth {
w.detachQueue = append(w.detachQueue, &connState{fd: -1, detachClosed: true})
}
w.detachQPending.Store(1)

w.drainDetachQueue()

if got := len(w.detachQSpare); got != 0 {
t.Fatalf("spare queue not truncated: len = %d, want 0", got)
}
full := w.detachQSpare[:cap(w.detachQSpare)]
if len(full) < depth {
t.Fatalf("backing array shrank to %d, expected it to keep at least the %d enqueued slots", len(full), depth)
}
held := 0
for _, cs := range full[:depth] {
if cs != nil {
held++
}
}
if held != 0 {
t.Fatalf("drainDetachQueue left %d of %d connStates reachable in the reused backing array; "+
"the queue pins its high-water mark (celeris#571 sibling)", held, depth)
}
}

// TestDrainDetachQueueLetsDrainedConnsBeCollected is the same guarantee stated
// as an observable outcome rather than as a property of the slice: once the
// queue has drained a conn and the caller drops its own reference, the
// connState must become unreachable. A finalizer on one of the enqueued
// conns proves it, where reading the backing array only proves the slot was
// cleared.
func TestDrainDetachQueueLetsDrainedConnsBeCollected(t *testing.T) {
w := &Worker{}
collected := make(chan struct{})

func() {
// Enqueued inside a function literal so the local goes out of scope
// before the GC below; a live stack slot would keep it reachable
// regardless of what the queue does.
cs := &connState{fd: -1, detachClosed: true}
runtime.AddCleanup(cs, func(ch chan struct{}) { close(ch) }, collected)
w.detachQueue = append(w.detachQueue, cs)
// Pad the queue so the drained conn is not the only slot; a
// single-element array is the easiest case to get right by accident.
for range 15 {
w.detachQueue = append(w.detachQueue, &connState{fd: -1, detachClosed: true})
}
}()
w.detachQPending.Store(1)
w.drainDetachQueue()

for range 5 {
runtime.GC()
select {
case <-collected:
runtime.KeepAlive(w)
return
default:
}
}
// KeepAlive is load-bearing, not defensive. Without it the compiler is
// free to treat w as dead after its last use above, the whole Worker
// (queues included) becomes collectable, and the cleanup runs no matter
// what the queue did -- the test then passes against the unfixed code.
runtime.KeepAlive(w)
t.Fatal("a connState the detach queue already drained is still reachable after 5 GC cycles: " +
"the queue is holding it in its reused backing array")
}
3 changes: 3 additions & 0 deletions engine/iouring/driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -232,6 +232,9 @@ func (w *Worker) drainDriverActions() {
w.attachAdoptedFD(a.adoptFD, a.adoptCarry)
}
}
// Same guard as the detach queue: a driverAction holds pointers, so
// stale slots keep them reachable after the action has been applied.
clear(w.driverActionSpare)
w.driverActionSpare = w.driverActionSpare[:0]
}

Expand Down
7 changes: 7 additions & 0 deletions engine/iouring/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -3808,6 +3808,13 @@ func (w *Worker) drainDetachQueue() {
}
w.markDirty(cs)
}
// Drop the strong refs before reusing the array. Truncating to [:0]
// leaves every *connState in the backing array reachable until some
// later drain overwrites that slot, so the queue pins its own
// high-water mark worth of connStates -- each one holding its buffers
// and, for a detached conn, its H1State and whatever the middleware
// hung off it. drainPendingRelease already guards the same hazard.
clear(w.detachQSpare)
w.detachQSpare = w.detachQSpare[:0]
}

Expand Down