From 8011f40d5fd77747e274109a1fea0063aaeab68a Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 07:14:17 +0200 Subject: [PATCH 1/5] fix(iouring): stop a ring SEND completion parking the worker on a running async handler's detachMu (celeris#750) handleSend applied every SEND completion of a conn with a detachMu under a blocking Lock (the notification, F_MORE, SEND_ZC-fallback and error branches, and completeSend). runAsyncHandler holds that mutex across the user handler, so a SEND completing while the conn's next handler ran (a pipelining client reading a large response) parked the worker, and every connection of its ring, until the handler returned: the fifth site of the celeris#704 class. handleSend now takes the lock once, with #704's TryLock/dispatchBusy pattern, and every branch releases it (completeSend runs under the caller's lock). When the dispatch goroutine holds it across a handler, the completion is held on the conn (heldSends, in arrival order; a later one is held behind an earlier one whatever the lock) and the goroutine owes the conn back (relinkOwed). replayHeldSends applies them through handleSend at the hand-back in drainDetachQueue, before anything else acts on the entry, and at the top of closeConn, so a close deferred behind cs.sending never waits for a completion that has already arrived. cs.sending stays set while a completion is held, so no other SEND starts. --- engine/iouring/conn.go | 22 +- ...send_completion_stall_fields_linux_test.go | 9 + .../send_completion_stall_linux_test.go | 386 ++++++++++++++++++ engine/iouring/worker.go | 168 +++++--- 4 files changed, 536 insertions(+), 49 deletions(-) create mode 100644 engine/iouring/send_completion_stall_fields_linux_test.go create mode 100644 engine/iouring/send_completion_stall_linux_test.go diff --git a/engine/iouring/conn.go b/engine/iouring/conn.go index 732d3142..f6aa5ed1 100644 --- a/engine/iouring/conn.go +++ b/engine/iouring/conn.go @@ -272,10 +272,25 @@ 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). 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 @@ -531,6 +546,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..c35720fe --- /dev/null +++ b/engine/iouring/send_completion_stall_linux_test.go @@ -0,0 +1,386 @@ +//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. +func TestIouringCloseAppliesAHeldSendCompletion(t *testing.T) { + const tail = "the rest of response 1" + rig := newStallRig704(t) + w, cs := rig.w, rig.cs + inFlight750(rig, tail, false) + release := holdAsHandler704(t, cs, true) + 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) + } + 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 held=%d closing=%v slot_owned=%v", heldSends750(cs), cs.closing, w.conns[rig.local] == cs) + rig.expectClosedOnce(t) + if heldSends750(cs) != 0 { + t.Errorf("%d completion(s) still held after the close", heldSends750(cs)) + } + cs.asyncInMu.Lock() + cs.asyncRun = false + cs.asyncInMu.Unlock() +} + +// ---- 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 b734d332..caea1f06 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -3258,6 +3258,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. @@ -3267,28 +3276,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 @@ -3312,9 +3316,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 { @@ -3353,12 +3357,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) @@ -3377,12 +3377,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] @@ -3400,11 +3397,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 @@ -3458,16 +3519,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 @@ -3627,6 +3685,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) @@ -4809,6 +4879,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 } From 85314d046f9ece4fab5f53430128bf9c96dfbaa2 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 07:25:54 +0200 Subject: [PATCH 2/5] docs(iouring): completeSend runs under its caller's detachMu (celeris#750) --- engine/iouring/worker.go | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index caea1f06..45b02975 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -3500,15 +3500,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 @@ -3637,8 +3638,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 @@ -3661,7 +3662,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) } From 2423242b2869ccb6c0828fa811e68b0d449a9eea Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 07:35:24 +0200 Subject: [PATCH 3/5] ci: the celeris#587 detachMu mutant releases the lock handleSend now takes at its top (celeris#750) handleSend takes detachMu once, before its branches, and the CQE_F_MORE branch defers the release, so the mutant's anchor (the branch's own Lock) is gone and the script exited 2. The mutant now releases the lock at the top of that branch, so the SEND_ZC first completion again writes cs.sending / cs.zcNotifPending / cs.zcSentBytes with no lock. --- .../scripts/mutant-587-unlock-zc-first-cqe.py | 23 ++++++++++++------- .github/workflows/ci.yml | 3 ++- 2 files changed, 17 insertions(+), 9 deletions(-) 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 393bca9a..6eac9fc0 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -912,7 +912,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 From 2e575b404ae6ae469d235511c54911cc309c1e82 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 08:22:10 +0200 Subject: [PATCH 4/5] test(iouring): a held send completion that fails closes once at closeConn, and a bounded holder is still waited out (celeris#750) TestIouringCloseAppliesAHeldSendCompletion gains an error arm: the held completion is a failure, so applying it at closeConn closes the conn (OnError once) and the outer close must not tear it down again. TestIouringSendCompletionStillWaitsForABoundedHolder is the negative control of the unit arms, #704's for the send completion: with the goroutine parked or past a Detach, the holder is a guarded writeFn in one write, and the completion must wait for it and be applied, not be held for a hand-back nothing owes. --- .../send_completion_stall_linux_test.go | 109 ++++++++++++++---- 1 file changed, 88 insertions(+), 21 deletions(-) diff --git a/engine/iouring/send_completion_stall_linux_test.go b/engine/iouring/send_completion_stall_linux_test.go index c35720fe..2150512f 100644 --- a/engine/iouring/send_completion_stall_linux_test.go +++ b/engine/iouring/send_completion_stall_linux_test.go @@ -224,32 +224,99 @@ func TestIouringSendCompletionsAreAppliedInOrder(t *testing.T) { // 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. +// 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" - rig := newStallRig704(t) - w, cs := rig.w, rig.cs - inFlight750(rig, tail, false) - release := holdAsHandler704(t, cs, true) - 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) + 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() + }) } - 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 held=%d closing=%v slot_owned=%v", heldSends750(cs), cs.closing, w.conns[rig.local] == cs) - rig.expectClosedOnce(t) - if heldSends750(cs) != 0 { - t.Errorf("%d completion(s) still held after the close", heldSends750(cs)) +// 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) + } + }) } - cs.asyncInMu.Lock() - cs.asyncRun = false - cs.asyncInMu.Unlock() } // ---- end to end: the issue's measurement -------------------------------- From f8f51fdf604423f6b52006414b5b59e893cc549e Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 10:40:52 +0200 Subject: [PATCH 5/5] fix(iouring): the dirty pass gives up a conn whose send completion is held, so the worker parks while the handler runs (celeris#750) A held completion keeps cs.sending set until the dispatch goroutine's hand-back, and flushDirty kept a sending conn on the dirty list, so baseTimeout returned 0 and the worker waited with a zero timeout, a spin, for as long as the handler ran; the blocking Lock this replaces had parked it. flushDirty now unlinks a conn with held completions, as the #704 give-up does; the hand-back's drain entry lists it again. The check is in the pass, not where the completion is held, because that entry lists the conn after holding the completion again when the goroutine is already in its next handler. TestIouringHeldSendCompletionDoesNotKeepTheRingPolling: arms held and held_again (found in review of #801). --- engine/iouring/conn.go | 5 +- .../send_completion_stall_linux_test.go | 88 +++++++++++++++++++ engine/iouring/worker.go | 17 ++++ 3 files changed, 108 insertions(+), 2 deletions(-) diff --git a/engine/iouring/conn.go b/engine/iouring/conn.go index f6aa5ed1..14efeb9a 100644 --- a/engine/iouring/conn.go +++ b/engine/iouring/conn.go @@ -288,8 +288,9 @@ type connState struct { // 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). The kernel's side of each is done: - // kernelInflight was settled when it was dispatched. + // 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 diff --git a/engine/iouring/send_completion_stall_linux_test.go b/engine/iouring/send_completion_stall_linux_test.go index 2150512f..d9ae70eb 100644 --- a/engine/iouring/send_completion_stall_linux_test.go +++ b/engine/iouring/send_completion_stall_linux_test.go @@ -271,6 +271,94 @@ func TestIouringCloseAppliesAHeldSendCompletion(t *testing.T) { } } +// 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 diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index 45b02975..3ed25186 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -5137,6 +5137,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