diff --git a/.github/scripts/mutant-587-unlock-zc-first-cqe.py b/.github/scripts/mutant-587-unlock-zc-first-cqe.py index c155ba3e..6ffc34ac 100644 --- a/.github/scripts/mutant-587-unlock-zc-first-cqe.py +++ b/.github/scripts/mutant-587-unlock-zc-first-cqe.py @@ -1,10 +1,12 @@ #!/usr/bin/env python3 """celeris#587 detector control: the MUTANT, applied in CI and never committed. -Deletes the cs.detachMu acquire in handleSend's CQE_F_MORE branch +Releases cs.detachMu at the top of handleSend's CQE_F_MORE branch (engine/iouring/worker.go), i.e. the SEND_ZC first completion then writes cs.sending / cs.zcNotifPending / cs.zcSentBytes with no lock, while the inline-egress guard on the dispatch goroutine reads them under the lock. +handleSend takes the lock once, at its top (celeris#750), and the branch +defers its release; the mutant releases it at once instead. TestSendZCWindowGuardUnderRace, run under -race against the mutated tree, MUST fail with a "WARNING: DATA RACE" report. If it does not, the race @@ -22,24 +24,29 @@ ANCHOR = "\tif cqeHasMore(c.Flags) {\n" LOCK = ( - "\t\tif mu := cs.detachMu; mu != nil {\n" - "\t\t\tmu.Lock()\n" + "\t\tif mu != nil {\n" "\t\t\tdefer mu.Unlock()\n" "\t\t}\n" ) +MUTANT = ( + "\t\t// MUTANT celeris#587: detachMu released before the F_MORE writes\n" + "\t\tif mu != nil {\n" + "\t\t\tmu.Unlock()\n" + "\t\t}\n" +) src = open(PATH).read() if src.count(ANCHOR) != 1: print(f"mutant-587: expected exactly one F_MORE branch anchor, found {src.count(ANCHOR)}", file=sys.stderr) sys.exit(2) start = src.index(ANCHOR) + len(ANCHOR) -# The lock must be the first statement block of the branch (after its -# comment lines), within a short distance of the anchor. +# The deferred release must be the first statement block of the branch +# (after its comment lines), within a short distance of the anchor. window = src[start:start + 600] if window.count(LOCK) != 1: - print("mutant-587: the detachMu acquire was not found at the top of the F_MORE branch", file=sys.stderr) + print("mutant-587: the deferred detachMu release was not found at the top of the F_MORE branch", file=sys.stderr) sys.exit(2) i = start + window.index(LOCK) -mutated = src[:i] + "\t\t// MUTANT celeris#587: detachMu acquire deleted\n" + src[i + len(LOCK):] +mutated = src[:i] + MUTANT + src[i + len(LOCK):] open(PATH, "w").write(mutated) -print(f"mutant-587: deleted the detachMu acquire in handleSend's CQE_F_MORE branch ({PATH})") +print(f"mutant-587: released detachMu before the writes of handleSend's CQE_F_MORE branch ({PATH})") diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 97b6f6a5..9757180a 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -947,7 +947,8 @@ jobs: # 2. CELERIS_IOURING_SEND_ZC=off: every process PASS with every ZC # witness at 0 (the witnesses track SEND_ZC, not traffic), no race; # 3. the MUTANT (.github/scripts/mutant-587-unlock-zc-first-cqe.py - # deletes the detachMu acquire in handleSend's CQE_F_MORE branch): + # releases detachMu before the writes of handleSend's CQE_F_MORE + # branch): # EVERY process FAIL, and every process with its own # "WARNING: DATA RACE" report and "race detected during execution # of test" line. A process that fails for any other reason does diff --git a/engine/iouring/conn.go b/engine/iouring/conn.go index a450b9c8..5dfa5ac0 100644 --- a/engine/iouring/conn.go +++ b/engine/iouring/conn.go @@ -273,10 +273,26 @@ type connState struct { closeOwed bool // relinkOwed (guarded by asyncInMu) is set by the dirty-list pass when it // gives the conn up because its dispatch goroutine holds detachMu across - // a handler (celeris#704). The goroutine hands cs back through the detach - // queue at the top of its next loop, after the handler's own flush, and - // drainDetachQueue puts it on the dirty list again. + // a handler (celeris#704), and by a send completion held for the same + // reason (heldSends, celeris#750). The goroutine hands cs back through + // the detach queue at the top of its next loop, after the handler's own + // flush, and drainDetachQueue applies the held completions and puts it + // on the dirty list again. relinkOwed bool + // heldSends (worker-thread only) are the ring SEND completions of this + // conn that arrived while its dispatch goroutine held detachMu across a + // handler, in arrival order (a SEND_ZC gives two). handleSend applies a + // completion under detachMu, and waiting for the lock parked the worker, + // and every connection of its ring, until the handler returned + // (celeris#750, a fifth celeris#704 site). A completion cannot be dropped: + // it is held, with relinkOwed set, and replayHeldSends applies it when the + // goroutine hands the conn back, or when the conn is closed, before + // anything else acts on the conn. Until then cs.sending (or + // zcNotifPending) stays set, so no other SEND starts and every raw write + // waits (celeris#751), and the dirty-list pass gives the conn up, which + // it would spin on otherwise (flushDirty). The kernel's side of each is + // done: kernelInflight was settled when it was dispatched. + heldSends []completionEntry // closeErr (worker-thread only) is the error handleRecv's peer-FIN or // recv-error branch owes a detached middleware (OnError) when it met a // running handler holding detachMu. The branch used to deliver it under @@ -539,6 +555,7 @@ func releaseConnState(cs *connState) { cs.asyncParked = false cs.closeOwed = false cs.relinkOwed = false + cs.heldSends = cs.heldSends[:0] cs.closeErr = nil cs.transplantPending.Store(false) cs.sweepKick = nil diff --git a/engine/iouring/send_completion_stall_fields_linux_test.go b/engine/iouring/send_completion_stall_fields_linux_test.go new file mode 100644 index 00000000..5b548ea2 --- /dev/null +++ b/engine/iouring/send_completion_stall_fields_linux_test.go @@ -0,0 +1,9 @@ +//go:build linux + +package iouring + +// The connState field the celeris#750 tests read. + +// heldSends750 reports how many of cs's send completions are held for its +// dispatch goroutine's hand-back. Worker thread (the test's). +func heldSends750(cs *connState) int { return len(cs.heldSends) } diff --git a/engine/iouring/send_completion_stall_linux_test.go b/engine/iouring/send_completion_stall_linux_test.go new file mode 100644 index 00000000..d9ae70eb --- /dev/null +++ b/engine/iouring/send_completion_stall_linux_test.go @@ -0,0 +1,541 @@ +//go:build linux + +package iouring + +import ( + "bufio" + "bytes" + "context" + "io" + "net" + "net/http" + "slices" + "strconv" + "syscall" + "testing" + "time" + + "golang.org/x/sys/unix" + + "github.com/goceleris/celeris/internal/ctxkit" + "github.com/goceleris/celeris/protocol/h2/stream" + "github.com/goceleris/celeris/resource" +) + +// celeris#750, the fifth worker-thread site of the celeris#704 class: with +// AsyncHandlers, runAsyncHandler holds cs.detachMu across ProcessH1, i.e. for +// as long as the user handler runs, and handleSend applied every ring SEND +// completion of a conn that has a detachMu under a blocking Lock (the +// notification, F_MORE, SEND_ZC-fallback and error branches, and all of +// completeSend). A SEND completing while the conn's next handler ran parked +// the LockOSThread'd worker, and every connection of its ring, until the +// handler returned. The common shape is a pipelining client: the rest of a +// large response goes out as a ring SEND, the next request's handler starts, +// and then the client reads. +// +// The unit arms hold detachMu the way a running handler does, deliver the +// completion, and ask handleSend to return while the lock is held. The +// completion must then be applied, in order, when the dispatch goroutine +// hands the conn back, or when the conn is closed. The end-to-end arm is the +// issue's measurement: a fast keep-alive conn on the same worker must not +// carry the handler's duration. + +// sendCQE750 is a SEND completion of rig's conn: res bytes (or -errno), with +// the CQE flags of a plain SEND (0), a SEND_ZC's first completion (F_MORE) or +// its notification (F_NOTIF). +func sendCQE750(rig *stallRig704, res int32, flags uint32) *completionEntry { + return &completionEntry{UserData: encodeUserDataGen(udSend, rig.local, rig.cs.generation), Res: res, Flags: flags} +} + +// inFlight750 puts a ring SEND of tail on rig's conn, as flushSend leaves it. +func inFlight750(rig *stallRig704, tail string, zc bool) { + cs := rig.cs + cs.sendBuf = append(cs.sendBuf[:0], tail...) + cs.sending = true + cs.sendIsZC = zc +} + +// handBack750 runs the REAL dispatch loop after the handler returned: its loop +// top hands the conn back (relinkOwed), and it parks. Then the worker drains +// the detach queue, which is where the held completions are applied. +func handBack750(t *testing.T, rig *stallRig704) { + t.Helper() + w := rig.w + w.asyncWG.Add(1) + go w.runAsyncHandler(rig.cs) + for dl := time.Now().Add(5 * time.Second); w.detachQPending.Load() == 0; { + if time.Now().After(dl) { + t.Fatal("the dispatch goroutine parked without handing the conn back") + } + time.Sleep(time.Millisecond) + } + w.drainDetachQueue() +} + +// endDispatch750 ends the rig's dispatch goroutine, if any, and waits for it. +func endDispatch750(rig *stallRig704) { + if rig.w.conns[rig.local] == rig.cs { + rig.w.closeConn(rig.local) + } + rig.w.asyncWG.Wait() +} + +func sqeOps750(sqes []sqeRec) []uint8 { + var ops []uint8 + for _, s := range sqes { + ops = append(ops, s.op) + } + return ops +} + +// TestIouringSendCompletionDoesNotWaitForARunningAsyncHandler: each kind of +// SEND completion, delivered while the conn's dispatch goroutine holds +// detachMu across a handler, must not park the worker, and must be applied +// exactly as it would have been, in order, at the goroutine's hand-back. +func TestIouringSendCompletionDoesNotWaitForARunningAsyncHandler(t *testing.T) { + const tail = "the rest of response 1" + const next = "response 2, written by the running handler" + type step struct { + res int32 + flags uint32 + } + for _, tc := range []struct { + name string + zc bool + cqes []step + check func(t *testing.T, rig *stallRig704, placed []sqeRec) + }{ + { + // The SEND finished: the handler's response goes out next. + name: "send", cqes: []step{{int32(len(tail)), 0}}, + check: func(t *testing.T, rig *stallRig704, placed []sqeRec) { + cs := rig.cs + if len(placed) != 1 || placed[0].op != opSEND || !cs.sending || string(cs.sendBuf) != next || len(cs.writeBuf) != 0 { + t.Errorf("after the hand-back: placed %v sending=%v sendBuf=%q writeBuf=%q; want one SEND of %q", + sqeOps750(placed), cs.sending, cs.sendBuf, cs.writeBuf, next) + } + }, + }, + { + // A short SEND: its remainder is re-sent before the handler's bytes. + name: "partial", cqes: []step{{5, 0}}, + check: func(t *testing.T, rig *stallRig704, placed []sqeRec) { + cs := rig.cs + if len(placed) != 1 || placed[0].op != opSEND || !cs.sending || string(cs.sendBuf) != tail[5:] || + string(cs.writeBuf) != next { + t.Errorf("after the hand-back: placed %v sending=%v sendBuf=%q writeBuf=%q; want a SEND of the "+ + "remainder %q with %q still queued", sqeOps750(placed), cs.sending, cs.sendBuf, cs.writeBuf, + tail[5:], next) + } + }, + }, + { + // A SEND_ZC: its first completion and its notification both arrive + // while the handler runs; they are applied in that order. + name: "zc", zc: true, cqes: []step{{int32(len(tail)), cqeFMore}, {0, cqeFNotif}}, + check: func(t *testing.T, rig *stallRig704, placed []sqeRec) { + cs := rig.cs + if cs.zcNotifPending || len(placed) != 1 || placed[0].op != opSEND || string(cs.sendBuf) != next { + t.Errorf("after the hand-back: zcNotifPending=%v placed %v sendBuf=%q; want the notification "+ + "applied and one SEND of %q", cs.zcNotifPending, sqeOps750(placed), cs.sendBuf, next) + } + }, + }, + { + // The peer reset: the error is delivered, once, and the conn closed. + name: "error", cqes: []step{{-int32(unix.ECONNRESET), 0}}, + check: func(t *testing.T, rig *stallRig704, _ []sqeRec) { + rig.expectClosedOnce(t) + want := errIORingSend(-int32(unix.ECONNRESET)).Error() + if len(rig.notified) != 1 || rig.notified[0].Error() != want { + t.Errorf("OnError calls = %v, want exactly one %q", rig.notified, want) + } + }, + }, + } { + t.Run(tc.name, func(t *testing.T) { + rig := newStallRig704(t) + w, cs := rig.w, rig.cs + inFlight750(rig, tail, tc.zc) + release := holdAsHandler704(t, cs, true) + cs.writeBuf = append(cs.writeBuf, next...) // the handler writes, under the lock it holds + + for i, s := range tc.cqes { + c := sendCQE750(rig, s.res, s.flags) + if !returnsWhileHeld704(t, release, func() { w.handleSend(c, rig.local, time.Now().UnixNano()) }) { + t.Fatalf("celeris#750: send completion %d of %d (res=%d flags=%#x) waited %v on cs.detachMu held "+ + "by a running async handler; the worker, and every connection of its ring, is parked until "+ + "the handler returns", i+1, len(tc.cqes), s.res, s.flags, stallWait704) + } + } + t.Logf("celeris750 HELD arm=%s held=%d sending=%v relinkOwed=%v", tc.name, heldSends750(cs), cs.sending, + relinkOwed704(cs)) + rig.expectWhole(t) + if !cs.sending || string(cs.sendBuf) != tail || len(rig.notified) != 0 { + t.Fatalf("the completion was applied under a running handler: sending=%v sendBuf=%q OnError=%v", + cs.sending, cs.sendBuf, rig.notified) + } + if !relinkOwed704(cs) { + t.Fatalf("the completion was held without a hand-back owed: it would never be applied") + } + + release() + handBack750(t, rig) + placed := takeSQEs(w.ring) + if heldSends750(cs) != 0 { + t.Errorf("%d completion(s) still held after the hand-back", heldSends750(cs)) + } + tc.check(t, rig, placed) + endDispatch750(rig) + }) + } +} + +// TestIouringSendCompletionsAreAppliedInOrder: a SEND_ZC's first completion +// arrives while the handler runs and is held; its notification arrives after +// the handler returned, with the lock free, before the hand-back is drained. +// Applying the notification first completed the send against a stale byte +// count and left the first completion to set zcNotifPending for a +// notification that had already come: a conn stuck sending (the celeris#519 +// class). The notification waits behind the held completion. +func TestIouringSendCompletionsAreAppliedInOrder(t *testing.T) { + const tail = "the rest of response 1" + rig := newStallRig704(t) + w, cs := rig.w, rig.cs + inFlight750(rig, tail, true) + release := holdAsHandler704(t, cs, true) + first := sendCQE750(rig, int32(len(tail)), cqeFMore) + if !returnsWhileHeld704(t, release, func() { w.handleSend(first, rig.local, time.Now().UnixNano()) }) { + t.Fatalf("celeris#750: the SEND_ZC first completion waited %v on a running handler's detachMu", stallWait704) + } + release() + w.handleSend(sendCQE750(rig, 0, cqeFNotif), rig.local, time.Now().UnixNano()) + t.Logf("celeris750 ORDER held=%d zcNotifPending=%v sending=%v", heldSends750(cs), cs.zcNotifPending, cs.sending) + handBack750(t, rig) + if cs.sending || cs.zcNotifPending || len(cs.sendBuf) != 0 || heldSends750(cs) != 0 { + t.Errorf("after the hand-back: sending=%v zcNotifPending=%v sendBuf=%q held=%d; want the send complete "+ + "(the notification applied after its first completion)", cs.sending, cs.zcNotifPending, cs.sendBuf, + heldSends750(cs)) + } + endDispatch750(rig) +} + +// TestIouringCloseAppliesAHeldSendCompletion: a close that runs after the +// handler returned, before the goroutine's hand-back is drained (the timeout +// sweep, a recv FIN), must not defer itself behind cs.sending for a SEND +// whose completion has already arrived and is held: it applies the +// completion, and closes now, once. Arm error: the held completion is itself +// a failure, which closes the conn as it is applied (OnError once), and the +// close that applied it must not run a second teardown. +func TestIouringCloseAppliesAHeldSendCompletion(t *testing.T) { + const tail = "the rest of response 1" + for _, tc := range []struct { + name string + res int32 + }{ + {"send", int32(len(tail))}, + {"error", -int32(unix.ECONNRESET)}, + } { + t.Run(tc.name, func(t *testing.T) { + rig := newStallRig704(t) + w, cs := rig.w, rig.cs + inFlight750(rig, tail, false) + release := holdAsHandler704(t, cs, true) + c := sendCQE750(rig, tc.res, 0) + if !returnsWhileHeld704(t, release, func() { w.handleSend(c, rig.local, time.Now().UnixNano()) }) { + t.Fatalf("celeris#750: the send completion waited %v on a running handler's detachMu", stallWait704) + } + release() + // The goroutine returns from its handler to its park (no hand-back drained). + cs.asyncInMu.Lock() + setParked704(cs, true) + cs.asyncInMu.Unlock() + + w.closeConn(rig.local) + t.Logf("celeris750 CLOSE arm=%s held=%d closing=%v slot_owned=%v OnError=%d", tc.name, heldSends750(cs), + cs.closing, w.conns[rig.local] == cs, len(rig.notified)) + rig.expectClosedOnce(t) + if heldSends750(cs) != 0 { + t.Errorf("%d completion(s) still held after the close", heldSends750(cs)) + } + if tc.res < 0 { + want := errIORingSend(tc.res).Error() + if len(rig.notified) != 1 || rig.notified[0].Error() != want { + t.Errorf("OnError calls = %v, want exactly one %q", rig.notified, want) + } + } + cs.asyncInMu.Lock() + cs.asyncRun = false + cs.asyncInMu.Unlock() + }) + } +} + +// TestIouringHeldSendCompletionDoesNotKeepTheRingPolling: while a completion +// is held, cs.sending stays set until the goroutine's hand-back, and a conn on +// the dirty list makes the worker wait with a zero timeout. Left listed, the +// held conn kept the ring polling for as long as the handler ran, where the +// blocking Lock had parked it (found in review of PR #801). The pass must +// give the conn up, as the celeris#704 give-up does, and the hand-back must +// list it again, with nothing the handler wrote lost. +// +// Arm held: the pass submitted the SEND of the rest of a response (the conn is +// listed while it is in flight), and it completes under the next handler. +// Arm held_again: the goroutine hands the conn back at the top of its loop and +// goes straight into the next pipelined request's handler, before the worker +// drains the hand-back; the drain's entry holds the completion again, and then +// lists the conn, as every entry does. +func TestIouringHeldSendCompletionDoesNotKeepTheRingPolling(t *testing.T) { + const tail = "the rest of response 1" + const next = "response 2, written by the running handler" + for _, again := range []bool{false, true} { + name := "held" + if again { + name = "held_again" + } + t.Run(name, func(t *testing.T) { + rig := newStallRig704(t) + w, cs := rig.w, rig.cs + // The goroutine's direct write was short: the rest is queued, the + // drain listed the conn, and the pass submits the SEND. + cs.writeBuf = append(cs.writeBuf[:0], tail...) + w.markDirty(cs) + w.flushDirty() + if placed := takeSQEs(w.ring); len(placed) != 1 || placed[0].op != opSEND || !cs.sending || !cs.dirty { + t.Fatalf("apparatus: placed %v sending=%v dirty=%v; want the SEND in flight with the conn listed", + sqeOps750(placed), cs.sending, cs.dirty) + } + release := holdAsHandler704(t, cs, true) + cs.writeBuf = append(cs.writeBuf, next...) // the handler writes, under the lock it holds + c := sendCQE750(rig, int32(len(tail)), 0) + if !returnsWhileHeld704(t, release, func() { w.handleSend(c, rig.local, time.Now().UnixNano()) }) { + t.Fatalf("celeris#750: the send completion waited %v on a running handler's detachMu", stallWait704) + } + if again { + // The handler returns; the loop top hands the conn back, and the + // next handler takes detachMu before the worker drains it. + release() + cs.asyncInMu.Lock() + cs.relinkOwed = false + w.enqueueDetach(cs) + cs.asyncInMu.Unlock() + release = holdAsHandler704(t, cs, true) + } + // The worker's per-iteration passes while the handler runs. + listed := 0 + for range 3 { + w.drainDetachQueue() + w.flushDirty() + if cs.dirty || w.dirtyHead != nil { + listed++ + } + } + t.Logf("celeris750 HELDLIST arm=%s held=%d sending=%v relinkOwed=%v dirty=%v baseTimeout=%v listed_passes=%d/3", + name, heldSends750(cs), cs.sending, relinkOwed704(cs), cs.dirty, w.baseTimeout(), listed) + if heldSends750(cs) != 1 || !cs.sending || !relinkOwed704(cs) { + t.Fatalf("apparatus: held=%d sending=%v relinkOwed=%v; want the completion held, with a hand-back owed", + heldSends750(cs), cs.sending, relinkOwed704(cs)) + } + if listed != 0 { + t.Errorf("with a completion held, the conn stayed on the dirty list after %d of 3 passes: the "+ + "worker waits with a zero timeout, a spin, for as long as the handler runs", listed) + } + + // The handler returns; the REAL dispatch loop hands the conn back. + release() + handBack750(t, rig) + placed := takeSQEs(w.ring) + t.Logf("celeris750 HELDLIST arm=%s after_handback held=%d placed=%v sendBuf=%q dirty=%v", name, + heldSends750(cs), sqeOps750(placed), cs.sendBuf, cs.dirty) + if heldSends750(cs) != 0 || len(placed) != 1 || placed[0].op != opSEND || string(cs.sendBuf) != next { + t.Errorf("after the hand-back: held=%d placed %v sendBuf=%q; want the completion applied and one "+ + "SEND of %q", heldSends750(cs), sqeOps750(placed), cs.sendBuf, next) + } + if !cs.dirty { + t.Errorf("the hand-back did not put the conn back on the dirty list") + } + endDispatch750(rig) + }) + } +} + +// TestIouringSendCompletionStillWaitsForABoundedHolder is the negative control +// for the unit arms, as #704's TestIouringCloseStillWaitsForABoundedHolder is +// for the close. The dispatch goroutine is PARKED, or past a Detach (it no +// longer holds the lock across a handler), so whoever holds detachMu is a +// guarded writeFn in one write, a hold bounded by a syscall. A completion must +// wait for it, as it always has, and be applied then: held instead, it would +// wait for a hand-back that nothing owes. +func TestIouringSendCompletionStillWaitsForABoundedHolder(t *testing.T) { + const tail = "the rest of response 1" + for _, tc := range []struct { + name string + setup func(cs *connState) (running bool) + }{ + {"parked", func(*connState) bool { return false }}, + {"after_detach", func(cs *connState) bool { cs.asyncDetachUnlocked = true; return true }}, + } { + t.Run(tc.name, func(t *testing.T) { + rig := newStallRig704(t) + w, cs := rig.w, rig.cs + inFlight750(rig, tail, false) + release := holdAsHandler704(t, cs, tc.setup(cs)) + c := sendCQE750(rig, int32(len(tail)), 0) + done := make(chan struct{}) + go func() { + defer close(done) + w.handleSend(c, rig.local, time.Now().UnixNano()) + }() + select { + case <-done: + t.Fatal("handleSend returned while a bounded holder held detachMu: the completion was held, and " + + "nothing owes it back") + case <-time.After(200 * time.Millisecond): + } + release() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("handleSend never returned after the holder released detachMu") + } + t.Logf("celeris750 BOUNDED arm=%s held=%d sending=%v", tc.name, heldSends750(cs), cs.sending) + if heldSends750(cs) != 0 || cs.sending || len(cs.sendBuf) != 0 { + t.Errorf("the completion was not applied after the bounded holder let go: held=%d sending=%v "+ + "sendBuf=%q", heldSends750(cs), cs.sending, cs.sendBuf) + } + }) + } +} + +// ---- end to end: the issue's measurement -------------------------------- + +// sendStallBig750 is /big's body: far larger than the client's receive buffer +// and the socket's send buffer, so the dispatch goroutine's direct write is +// short and the rest goes out as a ring SEND that waits for the client. +const sendStallBig750 = 3 << 20 + +type sendStallHandler750 struct{ big []byte } + +func (h sendStallHandler750) HandleStream(ctx context.Context, s *stream.Stream) error { + if s.ResponseWriter == nil { + return nil + } + var body []byte + switch s.Path { + case "/big": + body = h.big + case "/slow": + time.Sleep(stallSlow704) + body = []byte("slow") + default: + id, _ := ctxkit.WorkerIDFrom(ctx) + body = []byte("w=" + strconv.Itoa(id)) + } + return s.ResponseWriter.WriteResponse(s, 200, + [][2]string{{"content-type", "text/plain"}, {"content-length", strconv.Itoa(len(body))}}, body) +} +func (sendStallHandler750) RouteAsync(_, p string) bool { return p == "/big" || p == "/slow" } +func (sendStallHandler750) HasAsyncRoutes() bool { return true } + +// TestIouringSendCompletionDuringASlowAsyncHandlerDoesNotStallItsWorker: a +// pipelining client asks for /big, whose tail goes out as a ring SEND, and +// then for /slow, an 800 ms async handler; 50 ms into it the client reads, so +// the SEND completes while the handler runs. A fast keep-alive conn on the +// same worker is pinged throughout and must stay within the celeris#704 +// budget. +func TestIouringSendCompletionDuringASlowAsyncHandlerDoesNotStallItsWorker(t *testing.T) { + big := bytes.Repeat([]byte("0123456789abcdef"), sendStallBig750/16) + e, addr := startFDLEngine(t, sendStallHandler750{big: big}, func(c *resource.Config) { c.AsyncHandlers = true }) + d := net.Dialer{Timeout: 3 * time.Second, Control: func(_, _ string, rc syscall.RawConn) error { + return rc.Control(func(fd uintptr) { _ = unix.SetsockoptInt(int(fd), unix.SOL_SOCKET, unix.SO_RCVBUF, 4096) }) + }} + sc, err := d.Dial("tcp", addr) + if err != nil { + t.Fatalf("dial: %v", err) + } + t.Cleanup(func() { _ = sc.Close() }) + slow := &stallConn704{c: sc, br: bufio.NewReaderSize(sc, 4096)} + fast, worker := colocate704(t, addr, slow) + + stop := make(chan struct{}) + type result struct { + lat []time.Duration + err error + } + pinged := make(chan result, 1) + go func() { + var r result + for { + select { + case <-stop: + pinged <- r + return + default: + } + t0 := time.Now() + if _, err := fast.get("/fast", 5*time.Second); err != nil { + r.err = err + pinged <- r + return + } + r.lat = append(r.lat, time.Since(t0)) + time.Sleep(2 * time.Millisecond) + } + }() + + _ = sc.SetDeadline(time.Now().Add(15 * time.Second)) + if _, err := sc.Write([]byte("GET /big HTTP/1.1\r\nHost: x\r\n\r\n")); err != nil { + t.Fatalf("send /big: %v", err) + } + time.Sleep(150 * time.Millisecond) // the direct write fills the buffers; the rest is a ring SEND + if _, err := sc.Write([]byte("GET /slow HTTP/1.1\r\nHost: x\r\n\r\n")); err != nil { + t.Fatalf("send /slow: %v", err) + } + time.Sleep(50 * time.Millisecond) // the slow handler has started + // The client reads /big: the SEND completes while /slow's handler runs. + var bigOK bool + var bigErr error + if r1, err := http.ReadResponse(slow.br, nil); err != nil { + bigErr = err + } else { + b1, err := io.ReadAll(r1.Body) + _ = r1.Body.Close() + bigOK, bigErr = bytes.Equal(b1, big), err + } + // /slow's answer, read for the log only: whether it follows /big intact is + // celeris#751, not this test. + var slowBody []byte + r2, slowErr := http.ReadResponse(slow.br, nil) + if slowErr == nil { + slowBody, slowErr = io.ReadAll(r2.Body) + _ = r2.Body.Close() + } + time.Sleep(50 * time.Millisecond) + close(stop) + r := <-pinged + + var worst time.Duration + for _, x := range r.lat { + worst = max(worst, x) + } + sorted := slices.Clone(r.lat) + slices.Sort(sorted) + var p50 time.Duration + if len(sorted) > 0 { + p50 = sorted[len(sorted)/2] + } + t.Logf("celeris750 SENDSTALL worker=%s workers=%d samples=%d max_ms=%.1f p50_ms=%.2f fast_err=%v big_ok=%v "+ + "big_err=%v slow_body=%q slow_err=%v", worker, e.NumWorkers(), len(r.lat), float64(worst)/1e6, + float64(p50)/1e6, r.err, bigOK, bigErr, slowBody, slowErr) + if r.err != nil { + t.Errorf("celeris#750: the fast conn on the same worker broke while the slow handler ran: %v", r.err) + } + if worst > stallBudget704 { + t.Errorf("celeris#750: a fast request on worker %s took %v (budget %v) while a %v async handler ran on "+ + "another conn of the same worker, whose ring SEND completed mid-handler: the worker was parked on "+ + "that conn's detachMu", worker, worst, stallBudget704, stallSlow704) + } + if bigErr != nil { + t.Errorf("reading /big: %v", bigErr) + } +} diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index f9e079f7..a6256725 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -3479,6 +3479,15 @@ func (w *Worker) handleSend(c *completionEntry, fd int, now int64) { if cs == nil { return } + // The completion is applied under detachMu when the conn has one: the + // dispatch goroutine and a detached conn's guarded writeFn read the send + // state under it. Every branch below releases it. When the goroutine + // holds it across a handler, the completion is held instead + // (celeris#750); see holdOrLockSend. + mu := cs.detachMu + if mu != nil && !w.holdOrLockSend(cs, c) { + return + } // SEND_ZC notification CQE: the NIC has finished DMA-reading the buffer. // Now safe to modify/reuse sendBuf. Process the deferred result. @@ -3488,28 +3497,23 @@ func (w *Worker) handleSend(c *completionEntry, fd int, now int64) { w.zc.noteNotif() validation.IouringSendZCNotifs.Add(1) // zcNotifPending is read by the inline-egress guard on the dispatch - // goroutine under detachMu; clear it under the lock (completeSend - // re-acquires detachMu, so release first). - if mu := cs.detachMu; mu != nil { - mu.Lock() - // celeris#591: the NOTIF is the instant the guard reopens. If - // writeBuf already holds queued bytes, the very next inline - // unix.Write is admitted against data the worker has not yet - // flushed — the ordering celeris#587 exercises. Read under - // detachMu, the same lock the dispatch goroutine writes it - // under, so this witness is not itself a race. - if len(cs.writeBuf) > 0 { - validation.IouringZCCompletionWithPendingWrite.Add(1) - } - cs.zcNotifPending = false + // goroutine under detachMu, which is held here. + // + // celeris#591: the NOTIF is the instant the guard reopens. If + // writeBuf already holds queued bytes, the very next inline + // unix.Write is admitted against data the worker has not yet + // flushed — the ordering celeris#587 exercises. Read under + // detachMu, the same lock the dispatch goroutine writes it + // under, so this witness is not itself a race. + if len(cs.writeBuf) > 0 { + validation.IouringZCCompletionWithPendingWrite.Add(1) + } + cs.zcNotifPending = false + closeAfter := w.completeSend(cs, fd, int(cs.zcSentBytes), now, true) + if mu != nil { mu.Unlock() - } else { - if len(cs.writeBuf) > 0 { - validation.IouringZCCompletionWithPendingWrite.Add(1) - } - cs.zcNotifPending = false } - if w.completeSend(cs, fd, int(cs.zcSentBytes), now, true) { + if closeAfter { w.closeConn(fd) } return @@ -3533,9 +3537,9 @@ func (w *Worker) handleSend(c *completionEntry, fd int, now int64) { // kernel only for SEND_ZC, so it is the accurate test. if cqeHasMore(c.Flags) { // cs.sending / cs.zcNotifPending are read by the inline-egress guard on - // the dispatch goroutine under detachMu; mutate them under the lock. - if mu := cs.detachMu; mu != nil { - mu.Lock() + // the dispatch goroutine under detachMu; mutate them under the lock + // (held). + if mu != nil { defer mu.Unlock() } if c.Res < 0 { @@ -3574,12 +3578,8 @@ func (w *Worker) handleSend(c *completionEntry, fd int, now int64) { if cs.sendIsZC && (c.Res == -int32(unix.EINVAL) || c.Res == -int32(unix.ENOMEM)) { w.retireSendZC(-c.Res, "SEND_ZC unavailable, falling back to regular SEND") // cs.sending is read by the inline-egress guard under detachMu; clear it - // and re-flush under the lock (flushSend for a detached conn is always - // called under detachMu, as in the dirty-flush loop). - mu := cs.detachMu - if mu != nil { - mu.Lock() - } + // and re-flush under the lock (held; flushSend for a detached conn is + // always called under detachMu, as in the dirty-flush loop). cs.sending = false if w.flushSend(cs) { w.markDirty(cs) @@ -3598,12 +3598,9 @@ func (w *Worker) handleSend(c *completionEntry, fd int, now int64) { // nothing at all. w.errs.SendFailed(unix.Errno(-c.Res)) // cs.sending / cs.sendBuf are read by the inline-egress guard under - // detachMu; reset them (and writeBuf) inside the lock rather than before - // it, so the dispatch-goroutine read never races this error completion. - mu := cs.detachMu - if mu != nil { - mu.Lock() - } + // detachMu; reset them (and writeBuf) inside the lock (held) rather than + // before it, so the dispatch-goroutine read never races this error + // completion. cs.sending = false cs.sendBuf = cs.sendBuf[:0] cs.writeBuf = cs.writeBuf[:0] @@ -3621,11 +3618,75 @@ func (w *Worker) handleSend(c *completionEntry, fd int, now int64) { return } - if w.completeSend(cs, fd, int(c.Res), now, cs.sendIsZC) { + closeAfter := w.completeSend(cs, fd, int(c.Res), now, cs.sendIsZC) + if mu != nil { + mu.Unlock() + } + if closeAfter { w.closeConn(fd) } } +// holdOrLockSend takes cs.detachMu for send completion c and reports true, +// or holds c on cs (heldSends) and reports false. Worker thread. +// +// The dispatch goroutine holds detachMu across ProcessH1, i.e. for as long as +// the user handler runs, and a blocking Lock here parked the LockOSThread'd +// worker, and every connection of its ring, until the handler returned +// (celeris#750, the fifth site of the celeris#704 class). The common shape is +// a pipelining client: the rest of a response goes out as a ring SEND, and +// the next request's handler is running when the client reads it. The +// completion cannot be skipped, since its result must be applied, so it is +// held, and the goroutine owes the conn back (relinkOwed, set in the same +// asyncInMu section by dispatchBusy); replayHeldSends applies it then. +// +// A completion is also held, whatever the lock, while an earlier one of the +// conn is: they are applied in arrival order (a SEND_ZC's notification after +// its first completion). When the goroutine is parked or gone, whoever holds +// the lock holds it for one write, and waiting for it is bounded, as it +// always was; see dispatchBusy. +func (w *Worker) holdOrLockSend(cs *connState, c *completionEntry) bool { + if len(cs.heldSends) == 0 { + if cs.detachMu.TryLock() { + return true + } + if !dispatchBusy(cs, &cs.relinkOwed) { + cs.detachMu.Lock() + return true + } + } + cs.heldSends = append(cs.heldSends, *c) + return false +} + +// replayHeldSends applies cs's held send completions (celeris#750) through +// handleSend, in arrival order. It runs on the worker thread wherever the +// conn is next acted on: drainDetachQueue, for the dispatch goroutine's +// hand-back (any entry of the conn), and closeConn, so that a close deferred +// behind cs.sending never waits for a completion that has already arrived. +// If the goroutine holds detachMu across a handler again, handleSend holds +// that completion again, and every later one behind it, and the hand-back is +// owed again. The CQE dispatch's post-send hand-off attempt (tryTransplant) +// is not repeated: a conn whose goroutine lives is that goroutine's to claim, +// and a claimed one is finished by the drain entry this runs from. +func (w *Worker) replayHeldSends(cs *connState) { + fd := cs.fd + held := cs.heldSends + // A completion held again is appended to the front of held's array, at + // an index no later than the one being applied: the loop has read it. + cs.heldSends = held[:0] + for i := range held { + c := held[i] + // A conn that left its slot since (its completions were dispatched + // for it) has nothing left to apply them to. + if fd < 0 || fd >= len(w.conns) || w.conns[fd] != cs || decodeGen(c.UserData) != cs.generation { + cs.heldSends = cs.heldSends[:0] + return + } + w.handleSend(&c, fd, w.cachedNow) + } +} + // retireSendZC handles a send completion that failed for a reason that // condemns the SEND_ZC opcode rather than the connection: EINVAL (the kernel // does not support it) or ENOMEM (SEND_ZC pins the user buffer against @@ -3660,15 +3721,16 @@ func (w *Worker) retireSendZC(errno int32, msg string) { // exists, otherwise the goroutine read races the event-loop write — // observed via -race in TestNativeEngineLargePayload/io_uring. // Reports whether the CALLER must close the connection. closeConn takes -// cs.detachMu, and this function holds that same lock for its whole body; -// sync.Mutex is not reentrant, so closing inline wedges the worker thread +// cs.detachMu, and this function runs under that same lock (its caller, +// handleSend, holds it for the whole body); sync.Mutex is not reentrant, so +// closing inline wedges the worker thread // against itself -- and with it the entire event loop: no CQE is processed, // the detach queue is never drained (so every WebSocket recv-pause the // middleware asked to lift stays paused), no timeout sweep runs, and // graceful shutdown never completes. Releasing the lock early instead is // NOT the fix: it opens the window the lock exists to close, and measurably // corrupts streams (protoErr 0 -> 47 on the celeris#519 reproduction). The -// caller closes once the deferred unlock has run. +// caller closes once it has released the lock. // // A multi-worker engine only loses the one worker to this, so its // connections hang while the others keep serving; on a single-worker engine @@ -3679,16 +3741,13 @@ func (w *Worker) retireSendZC(errno int32, msg string) { // and the plain-SEND call site passes cs.sendIsZC, the provenance recorded // when that SQE was armed. Never w.sendZC (celeris#609). func (w *Worker) completeSend(cs *connState, fd int, sent int, now int64, fromZC bool) (closeAfter bool) { - // Take the lock up-front for detached connections so the entire state - // mutation (cs.sending clear / sendBuf truncate / writeBuf reset / OnError - // fire) is serialized against the goroutine writeFn path. The inline-egress - // fast path (the initProtocol guarded closure) reads cs.sending under - // detachMu to decide whether a ring SEND is in-flight, so the clear MUST be - // inside the lock — otherwise that read races this completion. - if mu := cs.detachMu; mu != nil { - mu.Lock() - defer mu.Unlock() - } + // The caller (handleSend) holds detachMu for detached and async + // connections, so the entire state mutation (cs.sending clear / sendBuf + // truncate / writeBuf reset / OnError fire) is serialized against the + // goroutine writeFn path. The inline-egress fast path (the initProtocol + // guarded closure) reads cs.sending under detachMu to decide whether a + // ring SEND is in-flight, so the clear MUST be inside the lock — otherwise + // that read races this completion. cs.sending = false // celeris#607 witness. cs.recvLinked is cleared only when the chained // recv's own CQE is processed, and that recv cannot run before this @@ -3800,8 +3859,8 @@ func (w *Worker) completeSend(cs *connState, fd int, sent int, now int64, fromZC } else { cs.sendBuf = cs.sendBuf[:0] } - // detachMu (if any) is held by the deferred unlock at the top of - // the function — no per-branch Unlock needed below. + // detachMu (if any) is held by the caller for the whole function — no + // per-branch Unlock needed below. if cs.closing && len(cs.sendBuf) == 0 && len(cs.writeBuf) == 0 { w.finishCloseAny(fd, cs) return @@ -3824,7 +3883,7 @@ func (w *Worker) completeSend(cs *connState, fd int, sent int, now int64, fromZC } // Re-send remainder or flush new data. Only markDirty on SQ ring full. - // detachMu (if any) is held by the deferred unlock at the top. + // detachMu (if any) is held by the caller. if w.flushSend(cs) { w.markDirty(cs) } @@ -3848,6 +3907,18 @@ func (w *Worker) closeConn(fd int) { if cs == nil { return } + // A send completion held for the dispatch goroutine (celeris#750) is + // applied first: the close below defers itself behind cs.sending, and + // the completion that would end that wait has already arrived. If the + // goroutine still holds detachMu across a handler, it stays held, and so + // does the close (closeOwed, below). + if len(cs.heldSends) > 0 { + closing := cs.closing + w.replayHeldSends(cs) + if w.conns[fd] != cs || cs.closing != closing { + return // the completion closed it, or deferred its close + } + } detached := cs.detachMu != nil if detached { cs.asyncClosed.Store(true) @@ -5138,6 +5209,12 @@ func (w *Worker) drainDetachQueue() { w.detachQPending.Store(0) w.detachQMu.Unlock() for _, cs := range w.detachQSpare { + // The dispatch goroutine's hand-back of send completions held while + // its handler ran (celeris#750): applied before anything else acts + // on the conn, whichever entry this is. + if len(cs.heldSends) > 0 { + w.replayHeldSends(cs) + } if cs.detachClosed { continue } @@ -5389,6 +5466,23 @@ func (w *Worker) removeDirty(cs *connState) { func (w *Worker) flushDirty() { for cs := w.dirtyHead; cs != nil; { next := cs.dirtyNext + if len(cs.heldSends) > 0 { + // A send completion of this conn is held for its dispatch + // goroutine (celeris#750), so cs.sending (or zcNotifPending) + // stays set, and nothing is sent or armed for the conn until the + // goroutine's hand-back applies the completion; the hand-back + // puts the conn on this list again (drainDetachQueue). Kept + // listed until then, it held the ring at a zero wait, a spin, + // for as long as the handler ran, where a blocking Lock had + // parked the worker. Give it up, as the celeris#704 give-up + // below does. Here rather than where the completion is held: + // the hand-back's own entry lists the conn again when the + // goroutine is already in its next handler and the completion + // is held again. + w.removeDirty(cs) + cs = next + continue + } if cs.sending { // celeris#607 witness. The retry below is gated on the // send, so a connection that is owed a recv arm and has a