From 7786a50ab00dc9c0128f0e10801f75e8c022a2d5 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 22:18:17 +0200 Subject: [PATCH 1/7] test(iouring): the client of a connection closed in the parking iteration must see the close (celeris#712) Measure-first: a sync-mode HTTP/1 connection whose header timer or ReadTimeout closes it on a paused worker, which parks in that iteration. Controls: the same close on a worker that keeps its listener, and on an AsyncHandlers engine. --- engine/iouring/park_close_fin_test.go | 232 ++++++++++++++++++++++++++ 1 file changed, 232 insertions(+) create mode 100644 engine/iouring/park_close_fin_test.go diff --git a/engine/iouring/park_close_fin_test.go b/engine/iouring/park_close_fin_test.go new file mode 100644 index 00000000..48a6f3b3 --- /dev/null +++ b/engine/iouring/park_close_fin_test.go @@ -0,0 +1,232 @@ +//go:build linux + +package iouring + +// celeris#712. finishClose's HTTP/1 fast path is a plain close(fd), and the +// recv the connection still has armed holds its own reference to the file, so +// close(fd) alone sends nothing: the socket goes, and the FIN with it, when +// the cancelled recv completes and puts that reference. On a DEFER_TASKRUN +// ring that completion is task work, and task work runs only inside an +// io_uring_enter with GETEVENTS. A running worker makes one on its next wait. +// A worker that parks in the iteration of the close made none: the pre-park +// submit (celeris#657 A5) is an enter without GETEVENTS, and the park itself +// waits on a Go channel. So the client of a connection closed in the parking +// iteration saw no FIN, and both ends stayed ESTABLISHED, until something +// woke the worker. +// +// Every arm below closes one sync-mode connection the engine gave up on and +// asks one thing of the client's side of it: that the close arrives. + +import ( + "errors" + "net" + "os" + "strconv" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/goceleris/celeris/protocol/h2/stream" + "github.com/goceleris/celeris/resource" +) + +// asyncFdlHandler is fdlHandler with every route async: the engine then runs +// AsyncHandlers, every HTTP/1 connection gets a detachMu, and closes go +// through finishCloseDetached, whose HTTP/1 branch shuts down SHUT_WR first. +type asyncFdlHandler struct{ fdlHandler } + +func (asyncFdlHandler) RouteAsync(_, _ string) bool { return true } +func (asyncFdlHandler) HasAsyncRoutes() bool { return true } + +var _ stream.AsyncRouteResolver = asyncFdlHandler{} + +// tcpState712 returns the /proc/net/tcp state column of the socket with this +// local and remote port ("01" ESTABLISHED, "04" FIN_WAIT1, "05" FIN_WAIT2, +// "08" CLOSE_WAIT), or "none". Diagnostic only. +func tcpState712(localPort, remotePort int) string { + for _, f := range []string{"/proc/net/tcp", "/proc/net/tcp6"} { + b, err := os.ReadFile(f) + if err != nil { + continue + } + for _, ln := range strings.Split(string(b), "\n")[1:] { + fs := strings.Fields(ln) + if len(fs) < 4 { + continue + } + l, _ := strconv.ParseInt(fs[1][strings.LastIndex(fs[1], ":")+1:], 16, 32) + r, _ := strconv.ParseInt(fs[2][strings.LastIndex(fs[2], ":")+1:], 16, 32) + if int(l) == localPort && int(r) == remotePort { + return fs[3] + } + } + } + return "none" +} + +// clientSeesClose polls c until its peer's close arrives (EOF or reset) or +// bound runs out; it returns how long that took, or -1. +func clientSeesClose(c net.Conn, bound time.Duration) (time.Duration, string) { + t0 := time.Now() + for time.Since(t0) < bound { + _ = c.SetReadDeadline(time.Now().Add(10 * time.Millisecond)) + _, err := c.Read(make([]byte, 64)) + if err == nil { + continue + } + var ne net.Error + if errors.As(err, &ne) && ne.Timeout() { + continue + } + return time.Since(t0), err.Error() + } + return -1, "no close seen" +} + +// finAfterClose runs one arm. park: PauseAccept first, so the close takes the +// worker's last connection and the worker parks in that iteration. readTimeout: +// the close is checkTimeouts' ReadTimeout branch instead of the header timer's +// CQE. async: the engine runs AsyncHandlers. +func finAfterClose(t *testing.T, park, readTimeout, async bool) { + var h stream.Handler = fdlHandler{} + if async { + h = asyncFdlHandler{} + } + e, addr := startFDLEngine(t, h, func(c *resource.Config) { + // No TCP_DEFER_ACCEPT, so no pause linger (celeris#662): the + // listeners close as PauseAccept is called, well inside the + // connection's deadline, and the close lands on a worker that has + // no listener, the only kind that parks. + c.DisableDeferAccept = true + c.AsyncHandlers = async + if readTimeout { + c.ReadHeaderTimeout = 60 * time.Second + c.ReadTimeout = 300 * time.Millisecond + c.IdleTimeout = 10 * time.Minute + } else { + c.ReadHeaderTimeout = 400 * time.Millisecond + } + }) + e.mu.Lock() + ws := append([]*Worker(nil), e.workers...) + tier := e.tier.Tier().String() + e.mu.Unlock() + allParked := func() bool { + for _, w := range ws { + if !w.suspended.Load() { + return false + } + } + return true + } + if !parkWait712(3*time.Second, func() bool { + m := e.Metrics() + return m.ActiveConnections == 0 && m.AcceptCount >= 1 && m.AcceptCount == m.CloseCount + }) { + t.Fatalf("celeris712 PREMISE: startFDLEngine's probe connection is still live") + } + c, err := net.DialTimeout("tcp", addr, 2*time.Second) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer func() { _ = c.Close() }() + srvPort := c.RemoteAddr().(*net.TCPAddr).Port + cliPort := c.LocalAddr().(*net.TCPAddr).Port + t0 := time.Now() + const partial = "GET /slow HTTP/1.1\r\nHost: x\r\n" // no blank line: mid-headers until a deadline + if _, err := c.Write([]byte(partial)); err != nil { + t.Fatalf("write: %v", err) + } + if !parkWait712(2*time.Second, func() bool { return e.Metrics().ActiveConnections == 1 }) { + t.Fatalf("celeris712 PREMISE: the connection was never accepted") + } + if park { + if err := e.PauseAccept(); err != nil { + t.Fatalf("pause: %v", err) + } + if n := e.Metrics().ActiveConnections; n != 1 { + t.Fatalf("celeris712 PREMISE: the connection closed (active=%d) before the listeners did", n) + } + } + // The engine's close, timed from the client's write. + var closedAt atomic.Int64 + go func() { + if parkWait712(10*time.Second, func() bool { return e.Metrics().ActiveConnections == 0 }) { + closedAt.Store(int64(time.Since(t0))) + } + }() + seen, why := clientSeesClose(c, 3*time.Second) + seenMs := int64(-1) + if seen >= 0 { + seenMs = time.Since(t0).Milliseconds() + } + if !parkWait712(10*time.Second, func() bool { return closedAt.Load() != 0 }) { + t.Fatalf("celeris712 PREMISE: the engine never closed the connection") + } + engineMs := time.Duration(closedAt.Load()).Milliseconds() + parked := allParked() + if park && !parked { + // The client's wait can end before the worker sets suspended. + parked = parkWait712(time.Second, allParked) + } + stSrv, stCli := tcpState712(srvPort, cliPort), tcpState712(cliPort, srvPort) + wakeMs := int64(-1) + if seen < 0 && park { + tw := time.Now() + _ = e.ResumeAccept() // ends the park + if d, _ := clientSeesClose(c, 2*time.Second); d >= 0 { + wakeMs = time.Since(tw).Milliseconds() + } + } + t.Logf("celeris712 park=%v read_timeout=%v async=%v tier=%s workers=%d engine_close_ms=%d parked=%v "+ + "client_close_ms=%d (%s) tcp_state_srv=%s tcp_state_cli=%s wake_to_close_ms=%d", + park, readTimeout, async, tier, len(ws), engineMs, parked, seenMs, why, stSrv, stCli, wakeMs) + if park && !parked { + t.Fatalf("celeris712 PREMISE: the worker did not park after the close") + } + if !park && parked { + t.Fatalf("celeris712 PREMISE: a worker with its listener open parked") + } + if seen < 0 { + t.Errorf("celeris712 NOFIN: the engine counted the connection closed %d ms after the client's write, "+ + "but the client saw no close for 3 s (server socket state %s, client %s; the close arrived %d ms "+ + "after a wake)", engineMs, stSrv, stCli, wakeMs) + } +} + +func parkWait712(d time.Duration, f func() bool) bool { + for dl := time.Now().Add(d); ; { + if f() { + return true + } + if time.Now().After(dl) { + return false + } + time.Sleep(2 * time.Millisecond) + } +} + +// TestParkedWorkerSendsFINForAHeaderTimeoutClose: the header timer's CQE +// closes the last connection of a paused worker, which parks in that +// iteration. The measured case of the issue. +func TestParkedWorkerSendsFINForAHeaderTimeoutClose(t *testing.T) { + finAfterClose(t, true, false, false) +} + +// TestParkedWorkerSendsFINForAReadTimeoutClose: the same, with checkTimeouts' +// ReadTimeout branch as the close. +func TestParkedWorkerSendsFINForAReadTimeoutClose(t *testing.T) { finAfterClose(t, true, true, false) } + +// TestRunningWorkerSendsFINForAHeaderTimeoutClose is the control: the same +// close on a worker that keeps its listener and keeps waiting on its ring. +func TestRunningWorkerSendsFINForAHeaderTimeoutClose(t *testing.T) { + finAfterClose(t, false, false, false) +} + +// TestParkedAsyncWorkerSendsFINForAHeaderTimeoutClose is the second control: +// an AsyncHandlers engine closes through finishCloseDetached, whose +// shutdown(SHUT_WR) acts on the socket and needs no recv completion. +func TestParkedAsyncWorkerSendsFINForAHeaderTimeoutClose(t *testing.T) { + finAfterClose(t, true, false, true) +} From d655435c7aaa1772a933a7e96de6f00cc70d127d Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 22:27:32 +0200 Subject: [PATCH 2/7] test(iouring): close the ReadTimeout arm on a draining worker, whose sweep caps the ring wait (celeris#712) A paused worker with no drain waits up to 1 s per iteration and runs checkTimeouts every 32nd, so the ReadTimeout close came past the arm's 10 s bound (measured: never within 13 s, 10/10). With a transplant set, as on the adaptive standby, the sweep caps the wait. --- engine/iouring/park_close_fin_test.go | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/engine/iouring/park_close_fin_test.go b/engine/iouring/park_close_fin_test.go index 48a6f3b3..7f754e17 100644 --- a/engine/iouring/park_close_fin_test.go +++ b/engine/iouring/park_close_fin_test.go @@ -142,6 +142,16 @@ func finAfterClose(t *testing.T, park, readTimeout, async bool) { t.Fatalf("celeris712 PREMISE: the connection was never accepted") } if park { + if readTimeout { + // A paused worker with no drain set waits up to 1 s on its + // ring and runs checkTimeouts every 32nd iteration, so its + // ReadTimeout close can come half a minute late. A drain caps + // the wait at the sweep's cadence: the adaptive standby's + // shape, and a close within a few seconds. + tgt := &fdlTarget{} + e.StartTransplant(tgt) + defer e.StopTransplant() + } if err := e.PauseAccept(); err != nil { t.Fatalf("pause: %v", err) } @@ -215,7 +225,8 @@ func TestParkedWorkerSendsFINForAHeaderTimeoutClose(t *testing.T) { } // TestParkedWorkerSendsFINForAReadTimeoutClose: the same, with checkTimeouts' -// ReadTimeout branch as the close. +// ReadTimeout branch as the close, on a draining worker (a transplant set, as +// on the adaptive standby). func TestParkedWorkerSendsFINForAReadTimeoutClose(t *testing.T) { finAfterClose(t, true, true, false) } // TestRunningWorkerSendsFINForAHeaderTimeoutClose is the control: the same From a77611bb63d52f65ad7550d8b6d30cd25562cf2c Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 22:23:37 +0200 Subject: [PATCH 3/7] fix(iouring): run the ring's deferred completions before a worker parks (celeris#712) On a DEFER_TASKRUN ring a cancelled recv completes as task work, which runs only inside an io_uring_enter with GETEVENTS, and until then the recv keeps its reference to the file. finishClose's HTTP/1 fast path is a plain close(fd), so a sync-mode connection closed in the iteration that parks a paused worker kept its socket ESTABLISHED, and its client got no FIN, until something woke the worker: the pre-park submit (celeris#657 A5) was an enter without GETEVENTS, and the park waits on a Go channel. The pre-park enter is now Ring.SubmitAndFlush, io_uring_enter(fd, pending, 0, GETEVENTS): it submits what is queued and runs the deferred work without waiting, and it is made even when nothing is left to submit. One syscall per park; parks are rare. --- engine/iouring/ring.go | 31 +++++++++++++++++++++++++++++++ engine/iouring/worker.go | 27 +++++++++++++++++++-------- 2 files changed, 50 insertions(+), 8 deletions(-) diff --git a/engine/iouring/ring.go b/engine/iouring/ring.go index 63ab1576..fee07533 100644 --- a/engine/iouring/ring.go +++ b/engine/iouring/ring.go @@ -308,6 +308,37 @@ func (r *Ring) Submit() (int, error) { return int(ret), nil } +// SubmitAndFlush submits pending SQEs and runs the ring's deferred completion +// work without waiting for a CQE: io_uring_enter(fd, pending, 0, GETEVENTS). +// +// On a DEFER_TASKRUN ring the completion of a request the kernel has already +// finished with, or cancelled, is task work that runs only inside an enter +// with GETEVENTS (celeris#79). Until it runs, the request still holds what it +// held, including its reference to the file: a socket closed with a recv +// still armed does not go, and sends no FIN, until then. The worker's own +// waits all pass GETEVENTS; this is for a caller that is about to stop +// entering the ring (the DRAINING→SUSPENDED park, celeris#712). With +// min_complete 0 the kernel runs that work and returns at once; the CQEs it +// posts wait in the ring for the next BeginCQ. On a ring without +// DEFER_TASKRUN it is Submit. +func (r *Ring) SubmitAndFlush() (int, error) { + n := r.pending + ret, _, errno := unix.Syscall6( + uintptr(sysIOUringEnter), + uintptr(r.fd), + uintptr(n), + 0, + uintptr(enterGetEvents), + 0, 0, + ) + if errno != 0 { + r.pending = retryPending(n, errno) + return 0, fmt.Errorf("io_uring_enter submit+flush: %w", errno) + } + r.pending = n - uint32(ret) + return int(ret), nil +} + // WaitCQE waits for at least one CQE to become available. func (r *Ring) WaitCQE() error { _, _, errno := unix.Syscall6( diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index 3e4357e4..74a361fe 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -1492,14 +1492,25 @@ func (w *Worker) run(ctx context.Context) { // queues its close-path cancels (the header timer's, a send's) // in the same pass that finds the worker idle, and the park is // indefinite: those SQEs used to sit unsubmitted until something - // woke the worker, measured 1-24 pending at parks. Outside - // wakeMu, which is a leaf. Not under SQPOLL, where the kernel's - // SQ thread submits; an SQ thread that has gone idle would need - // the NEED_WAKEUP kick the submit branch of this loop gives it, - // but no tier enables SQPOLL today (SQPollIdle is 0 in all - // three), so that case is not handled here. - if !w.sqpoll && w.ring.Pending() > 0 { - _, _ = w.ring.Submit() + // woke the worker, measured 1-24 pending at parks. + // + // Submitting is not enough on a DEFER_TASKRUN ring, and the + // enter runs the deferred completion work too, even with + // nothing left to submit (celeris#712). A cancelled recv + // completes as task work that only an enter with GETEVENTS + // runs, and until it does the recv keeps its reference to the + // file: finishClose's HTTP/1 fast path is a plain close(fd), + // so the socket of a connection closed in this iteration stayed + // ESTABLISHED, and its client got no FIN, for as long as the + // park lasted. The park waits on a Go channel, not in the ring. + // + // Outside wakeMu, which is a leaf. Not under SQPOLL, where the + // kernel's SQ thread submits; an SQ thread that has gone idle + // would need the NEED_WAKEUP kick the submit branch of this + // loop gives it, but no tier enables SQPOLL today (SQPollIdle is + // 0 in all three), so that case is not handled here. + if !w.sqpoll { + _, _ = w.ring.SubmitAndFlush() } w.wakeMu.Lock() if !w.acceptPaused.Load() || From 76917bd21d9a3261c4142e28ad3443181774c4ac Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 23:14:19 +0200 Subject: [PATCH 4/7] test(iouring): start the celeris#712 engines with the ring ENOMEM retried At CI's 8 MiB memlock the kernel gives a closed ring's pages back 12-23 ms after the close, and startFDLEngine does not wait that out: an engine started right after another ring closed failed with ENOMEM (measured on the celeris#713 branch, 3 of 10 runs per shape right after a fixture ring), and its cleanup then waited 5 s on a Listen result it had already consumed. The tests now start through startRingRetried662, as the other back-to-back engine tests do. --- engine/iouring/park_close_fin_test.go | 46 +++++++++++++++++++++++++-- 1 file changed, 43 insertions(+), 3 deletions(-) diff --git a/engine/iouring/park_close_fin_test.go b/engine/iouring/park_close_fin_test.go index 7f754e17..0b5cdbdd 100644 --- a/engine/iouring/park_close_fin_test.go +++ b/engine/iouring/park_close_fin_test.go @@ -19,6 +19,8 @@ package iouring import ( "errors" + "io" + "log/slog" "net" "os" "strconv" @@ -27,10 +29,48 @@ import ( "testing" "time" + "github.com/goceleris/celeris/engine" "github.com/goceleris/celeris/protocol/h2/stream" "github.com/goceleris/celeris/resource" ) +// startParkEngine712 is startFDLEngine with its ring ENOMEM retried +// (startRingRetried662): at CI's 8 MiB memlock the kernel gives a closed +// ring's pages back 12-23 ms after the close, so an engine started right +// after another ring closed can fail on memory nothing holds any more. No +// probe dial: the engine is idle when it returns. +func startParkEngine712(t *testing.T, h stream.Handler, mut func(*resource.Config)) (*Engine, string) { + t.Helper() + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("pick port: %v", err) + } + addr := ln.Addr().String() + _ = ln.Close() + e, cancel, done := startRingRetried662(t, func() (*Engine, error) { + cfg := resource.Config{ + Addr: addr, + Protocol: engine.HTTP1, + Resources: resource.Resources{Workers: 2}, + Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + } + if mut != nil { + mut(&cfg) + } + return New(cfg, h) + }) + t.Cleanup(func() { + cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("engine did not stop within 5s") + } + }) + t.Logf("celeris657 engine workers=%d", e.NumWorkers()) + return e, addr +} + // asyncFdlHandler is fdlHandler with every route async: the engine then runs // AsyncHandlers, every HTTP/1 connection gets a detachMu, and closes go // through finishCloseDetached, whose HTTP/1 branch shuts down SHUT_WR first. @@ -93,7 +133,7 @@ func finAfterClose(t *testing.T, park, readTimeout, async bool) { if async { h = asyncFdlHandler{} } - e, addr := startFDLEngine(t, h, func(c *resource.Config) { + e, addr := startParkEngine712(t, h, func(c *resource.Config) { // No TCP_DEFER_ACCEPT, so no pause linger (celeris#662): the // listeners close as PauseAccept is called, well inside the // connection's deadline, and the close lands on a worker that has @@ -122,9 +162,9 @@ func finAfterClose(t *testing.T, park, readTimeout, async bool) { } if !parkWait712(3*time.Second, func() bool { m := e.Metrics() - return m.ActiveConnections == 0 && m.AcceptCount >= 1 && m.AcceptCount == m.CloseCount + return m.ActiveConnections == 0 && m.AcceptCount == m.CloseCount }) { - t.Fatalf("celeris712 PREMISE: startFDLEngine's probe connection is still live") + t.Fatalf("celeris712 PREMISE: the engine is not idle") } c, err := net.DialTimeout("tcp", addr, 2*time.Second) if err != nil { From 80131894e7fd37a2e2cc26e9a2326a06b1e88c23 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 03:27:59 +0200 Subject: [PATCH 5/7] test(iouring): close 64 connections in the parking iteration, and log whether the rings defer their task work (celeris#712) One io_uring_enter runs at most IO_LOCAL_TW_DEFAULT_MAX (20) deferred completions per local-work pass on kernel 6.13 and later, and an enter with min_complete 0 makes at most two such passes. The single-connection arms cannot see that limit. The two new arms close 64 mid-header connections on a paused worker together, by one checkTimeouts pass (ReadTimeout, on a draining worker) or by their header timers, and require every client to see its close while the worker stays parked. The arms now log defer_taskrun as well as the tier: a COOP_TASKRUN ring logs tier=high too, and the defect needs a DEFER_TASKRUN ring. Under CELERIS_REQUIRE_IOURING_WORKERS=1 a ring without it fails the premise. --- engine/iouring/park_close_fin_test.go | 168 +++++++++++++++++++++++++- 1 file changed, 162 insertions(+), 6 deletions(-) diff --git a/engine/iouring/park_close_fin_test.go b/engine/iouring/park_close_fin_test.go index 0b5cdbdd..97d0a479 100644 --- a/engine/iouring/park_close_fin_test.go +++ b/engine/iouring/park_close_fin_test.go @@ -14,8 +14,12 @@ package iouring // iteration saw no FIN, and both ends stayed ESTABLISHED, until something // woke the worker. // -// Every arm below closes one sync-mode connection the engine gave up on and -// asks one thing of the client's side of it: that the close arrives. +// Every arm below closes sync-mode connections the engine gave up on and asks +// one thing of the client's side of them: that the close arrives. The first +// arms close one connection; the many-connection arms close 64 together. +// One enter runs at most 20 deferred completions per local-work pass +// (IO_LOCAL_TW_DEFAULT_MAX, kernel 6.13 and later), so one enter before the +// park reached only the first of them. import ( "errors" @@ -25,6 +29,7 @@ import ( "os" "strconv" "strings" + "sync" "sync/atomic" "testing" "time" @@ -71,6 +76,23 @@ func startParkEngine712(t *testing.T, h stream.Handler, mut func(*resource.Confi return e, addr } +// ringKind712 reports the engine's tier and whether its rings run completions +// as deferred task work (IORING_SETUP_DEFER_TASKRUN), the ring kind the defect +// needs: a COOP_TASKRUN ring logs tier=high as well. With +// CELERIS_REQUIRE_IOURING_WORKERS=1 (CI, the cluster) a ring without it is a +// failed premise, not a pass. +func ringKind712(t *testing.T, e *Engine) (tier string, deferTaskrun bool) { + t.Helper() + e.mu.Lock() + tier = e.tier.Tier().String() + deferTaskrun = e.tier.SetupFlags()&setupDeferTaskrun != 0 + e.mu.Unlock() + if !deferTaskrun && os.Getenv(envRequireIOUring656) == "1" { + t.Fatalf("celeris712 PREMISE: tier %s rings do not use DEFER_TASKRUN", tier) + } + return tier, deferTaskrun +} + // asyncFdlHandler is fdlHandler with every route async: the engine then runs // AsyncHandlers, every HTTP/1 connection gets a detachMu, and closes go // through finishCloseDetached, whose HTTP/1 branch shuts down SHUT_WR first. @@ -150,8 +172,8 @@ func finAfterClose(t *testing.T, park, readTimeout, async bool) { }) e.mu.Lock() ws := append([]*Worker(nil), e.workers...) - tier := e.tier.Tier().String() e.mu.Unlock() + tier, deferTaskrun := ringKind712(t, e) allParked := func() bool { for _, w := range ws { if !w.suspended.Load() { @@ -229,9 +251,9 @@ func finAfterClose(t *testing.T, park, readTimeout, async bool) { wakeMs = time.Since(tw).Milliseconds() } } - t.Logf("celeris712 park=%v read_timeout=%v async=%v tier=%s workers=%d engine_close_ms=%d parked=%v "+ - "client_close_ms=%d (%s) tcp_state_srv=%s tcp_state_cli=%s wake_to_close_ms=%d", - park, readTimeout, async, tier, len(ws), engineMs, parked, seenMs, why, stSrv, stCli, wakeMs) + t.Logf("celeris712 park=%v read_timeout=%v async=%v tier=%s defer_taskrun=%v workers=%d engine_close_ms=%d "+ + "parked=%v client_close_ms=%d (%s) tcp_state_srv=%s tcp_state_cli=%s wake_to_close_ms=%d", + park, readTimeout, async, tier, deferTaskrun, len(ws), engineMs, parked, seenMs, why, stSrv, stCli, wakeMs) if park && !parked { t.Fatalf("celeris712 PREMISE: the worker did not park after the close") } @@ -257,6 +279,128 @@ func parkWait712(d time.Duration, f func() bool) bool { } } +// finToManyAfterClose runs a many-connection arm: n sync connections, each +// mid-headers, on a paused engine, closed together by one deadline (the +// header timers' CQEs, or one checkTimeouts pass's ReadTimeout on a draining +// worker) so that the last of them go in the iteration that parks the worker. +// Every client must see its close while the worker stays parked. +func finToManyAfterClose(t *testing.T, n int, readTimeout bool) { + e, addr := startParkEngine712(t, fdlHandler{}, func(c *resource.Config) { + c.DisableDeferAccept = true // no pause linger (celeris#662), as above + if readTimeout { + c.ReadHeaderTimeout = 60 * time.Second + // Long enough for n dials and the pause to land before it. + c.ReadTimeout = time.Second + c.IdleTimeout = 10 * time.Minute + } else { + c.ReadHeaderTimeout = time.Second + } + }) + e.mu.Lock() + ws := append([]*Worker(nil), e.workers...) + e.mu.Unlock() + tier, deferTaskrun := ringKind712(t, e) + allParked := func() bool { + for _, w := range ws { + if !w.suspended.Load() { + return false + } + } + return true + } + if !parkWait712(3*time.Second, func() bool { + m := e.Metrics() + return m.ActiveConnections == 0 && m.AcceptCount == m.CloseCount + }) { + t.Fatalf("celeris712 PREMISE: the engine is not idle") + } + const partial = "GET /slow HTTP/1.1\r\nHost: x\r\n" // no blank line: mid-headers until a deadline + read0 := e.Metrics().BytesRead + conns := make([]net.Conn, n) + for i := range conns { + c, err := net.DialTimeout("tcp", addr, 2*time.Second) + if err != nil { + t.Fatalf("dial %d: %v", i, err) + } + conns[i] = c + defer func() { _ = c.Close() }() + } + for i, c := range conns { + if _, err := c.Write([]byte(partial)); err != nil { + t.Fatalf("write %d: %v", i, err) + } + } + if !parkWait712(3*time.Second, func() bool { return e.Metrics().ActiveConnections == int64(n) }) { + t.Fatalf("celeris712 PREMISE: %d of %d connections accepted", e.Metrics().ActiveConnections, n) + } + // Every partial header read: a connection whose bytes the engine has not + // read yet is at a request boundary, and the drain below hands it off. + want := uint64(n * len(partial)) + if !parkWait712(3*time.Second, func() bool { return e.Metrics().BytesRead-read0 >= want }) { + t.Fatalf("celeris712 PREMISE: the engine read %d of %d header bytes", e.Metrics().BytesRead-read0, want) + } + if readTimeout { + // A drain caps the paused worker's ring wait at the sweep's + // cadence, as in the single-connection ReadTimeout arm. + tgt := &fdlTarget{} + e.StartTransplant(tgt) + defer e.StopTransplant() + } + if err := e.PauseAccept(); err != nil { + t.Fatalf("pause: %v", err) + } + if a := e.Metrics().ActiveConnections; a != int64(n) { + t.Fatalf("celeris712 PREMISE: closes began before the listeners closed (active=%d of %d)", a, n) + } + if !parkWait712(15*time.Second, func() bool { return e.Metrics().ActiveConnections == 0 }) { + t.Fatalf("celeris712 PREMISE: the engine never closed the connections (active=%d)", e.Metrics().ActiveConnections) + } + closedAt := time.Now() + if !parkWait712(2*time.Second, allParked) { + t.Fatalf("celeris712 PREMISE: the worker did not park after the last close") + } + parkMs := time.Since(closedAt).Milliseconds() + var seen atomic.Int64 + var wg sync.WaitGroup + for _, c := range conns { + wg.Add(1) + go func(c net.Conn) { + defer wg.Done() + if d, _ := clientSeesClose(c, 2*time.Second); d >= 0 { + seen.Add(1) + } + }(c) + } + wg.Wait() + parked := allParked() + // How many of the rest arrive once something wakes the worker. + var seenAfterWake atomic.Int64 + if seen.Load() < int64(n) { + _ = e.ResumeAccept() + var wg2 sync.WaitGroup + for _, c := range conns { + wg2.Add(1) + go func(c net.Conn) { + defer wg2.Done() + if d, _ := clientSeesClose(c, 2*time.Second); d >= 0 { + seenAfterWake.Add(1) + } + }(c) + } + wg2.Wait() + } + t.Logf("celeris712many read_timeout=%v tier=%s defer_taskrun=%v workers=%d n=%d park_after_last_close_ms=%d "+ + "parked_through_poll=%v fin_seen_while_parked=%d/%d seen_after_wake=%d", + readTimeout, tier, deferTaskrun, len(ws), n, parkMs, parked, seen.Load(), n, seenAfterWake.Load()) + if !parked { + t.Fatalf("celeris712 PREMISE: the worker left the park while the clients polled") + } + if seen.Load() != int64(n) { + t.Errorf("celeris712 NOFIN: %d of %d clients saw no close in 2 s while the worker was parked "+ + "(%d of them saw it after a wake)", int64(n)-seen.Load(), n, seenAfterWake.Load()) + } +} + // TestParkedWorkerSendsFINForAHeaderTimeoutClose: the header timer's CQE // closes the last connection of a paused worker, which parks in that // iteration. The measured case of the issue. @@ -281,3 +425,15 @@ func TestRunningWorkerSendsFINForAHeaderTimeoutClose(t *testing.T) { func TestParkedAsyncWorkerSendsFINForAHeaderTimeoutClose(t *testing.T) { finAfterClose(t, true, false, true) } + +// TestParkedWorkerSendsFINToManyReadTimeoutCloses: one checkTimeouts pass on +// a draining, paused worker closes 64 connections, and the worker parks in +// that iteration. Every client must see its close. +func TestParkedWorkerSendsFINToManyReadTimeoutCloses(t *testing.T) { finToManyAfterClose(t, 64, true) } + +// TestParkedWorkerSendsFINToManyHeaderTimeoutCloses: the same with 64 header +// timers, whose CQEs arrive together or in a few batches: only the batch +// that empties the worker is closed in the parking iteration. +func TestParkedWorkerSendsFINToManyHeaderTimeoutCloses(t *testing.T) { + finToManyAfterClose(t, 64, false) +} From d9307b557cad552dc30163a3d00a520807ce6cfb Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 03:28:39 +0200 Subject: [PATCH 6/7] fix(iouring): do not park a worker while a closed connection's kernel ops are in flight (celeris#712) Round 1 made one io_uring_enter with GETEVENTS before the park, to run the deferred completion of the recv each connection closed in the parking iteration still had armed: until that completion runs, the recv holds its reference to the file, and the fast-path close(fd) sends no FIN. One enter runs at most IO_LOCAL_TW_DEFAULT_MAX (20) deferred completions per local-work pass on kernel 6.13 and later, and 64 connections closed by one checkTimeouts pass got 20 FINs. pendingRelease already holds each closed connState until its kernelInflight reports every op's terminal CQE. So the park now drains it first, with a fresh clock for its wall-clock backstop, and goes round again while an entry is left: the ring wait runs the deferred work and returns with the CQEs, as it does on a running worker, however many there are. An op the kernel never completes holds the park back by pendingReleaseHoldNanos (5 s) and one ring wait at most. The pre-park flush and Ring.SubmitAndFlush are gone; the A5 submit is main's again. --- engine/iouring/ring.go | 31 ----------------------- engine/iouring/worker.go | 54 ++++++++++++++++++++++++++-------------- 2 files changed, 35 insertions(+), 50 deletions(-) diff --git a/engine/iouring/ring.go b/engine/iouring/ring.go index fee07533..63ab1576 100644 --- a/engine/iouring/ring.go +++ b/engine/iouring/ring.go @@ -308,37 +308,6 @@ func (r *Ring) Submit() (int, error) { return int(ret), nil } -// SubmitAndFlush submits pending SQEs and runs the ring's deferred completion -// work without waiting for a CQE: io_uring_enter(fd, pending, 0, GETEVENTS). -// -// On a DEFER_TASKRUN ring the completion of a request the kernel has already -// finished with, or cancelled, is task work that runs only inside an enter -// with GETEVENTS (celeris#79). Until it runs, the request still holds what it -// held, including its reference to the file: a socket closed with a recv -// still armed does not go, and sends no FIN, until then. The worker's own -// waits all pass GETEVENTS; this is for a caller that is about to stop -// entering the ring (the DRAINING→SUSPENDED park, celeris#712). With -// min_complete 0 the kernel runs that work and returns at once; the CQEs it -// posts wait in the ring for the next BeginCQ. On a ring without -// DEFER_TASKRUN it is Submit. -func (r *Ring) SubmitAndFlush() (int, error) { - n := r.pending - ret, _, errno := unix.Syscall6( - uintptr(sysIOUringEnter), - uintptr(r.fd), - uintptr(n), - 0, - uintptr(enterGetEvents), - 0, 0, - ) - if errno != 0 { - r.pending = retryPending(n, errno) - return 0, fmt.Errorf("io_uring_enter submit+flush: %w", errno) - } - r.pending = n - uint32(ret) - return int(ret), nil -} - // WaitCQE waits for at least one CQE to become available. func (r *Ring) WaitCQE() error { _, _, errno := unix.Syscall6( diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index 74a361fe..0fb7e105 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -1487,30 +1487,46 @@ func (w *Worker) run(ctx context.Context) { if w.listenFD < 0 && w.connCount == 0 && !w.hasDriverConns.Load() && w.driverActionPending.Load() == 0 && w.detachQPending.Load() == 0 && w.acceptPaused.Load() { + // Do not park while a closed connection's kernel ops are still + // in flight (celeris#712). finishClose's HTTP/1 fast path is a + // plain close(fd), and the recv the connection still has armed + // holds its own reference to the file, so the socket goes, and + // its client gets a FIN, only once that recv's cancelled + // completion has run. On a DEFER_TASKRUN ring the completion is + // task work, which runs only inside an io_uring_enter with + // GETEVENTS, and the park waits on a Go channel: a connection + // closed in the parking iteration stayed ESTABLISHED for as long + // as the park lasted. One more enter is not enough either, as it + // runs at most 20 deferred completions per local-work pass + // (IO_LOCAL_TW_DEFAULT_MAX, kernel 6.13 and later). + // + // pendingRelease already holds each closed connState until its + // kernelInflight reports every op's terminal CQE, so an entry a + // drain leaves is such an op: go round again, and the ring wait + // runs the work and returns with those CQEs. The clock is read + // for the drain's wall-clock backstop, so an op the kernel never + // completes holds the park back by pendingReleaseHoldNanos and + // one ring wait at most. + if len(w.pendingRelease) > 0 { + w.cachedNow = time.Now().UnixNano() + w.drainPendingRelease() + if len(w.pendingRelease) > 0 { + continue + } + } // Submit what this iteration queued before parking (celeris#657, // A5). The iteration that closes or hands off the last conn // queues its close-path cancels (the header timer's, a send's) // in the same pass that finds the worker idle, and the park is // indefinite: those SQEs used to sit unsubmitted until something - // woke the worker, measured 1-24 pending at parks. - // - // Submitting is not enough on a DEFER_TASKRUN ring, and the - // enter runs the deferred completion work too, even with - // nothing left to submit (celeris#712). A cancelled recv - // completes as task work that only an enter with GETEVENTS - // runs, and until it does the recv keeps its reference to the - // file: finishClose's HTTP/1 fast path is a plain close(fd), - // so the socket of a connection closed in this iteration stayed - // ESTABLISHED, and its client got no FIN, for as long as the - // park lasted. The park waits on a Go channel, not in the ring. - // - // Outside wakeMu, which is a leaf. Not under SQPOLL, where the - // kernel's SQ thread submits; an SQ thread that has gone idle - // would need the NEED_WAKEUP kick the submit branch of this - // loop gives it, but no tier enables SQPOLL today (SQPollIdle is - // 0 in all three), so that case is not handled here. - if !w.sqpoll { - _, _ = w.ring.SubmitAndFlush() + // woke the worker, measured 1-24 pending at parks. Outside + // wakeMu, which is a leaf. Not under SQPOLL, where the kernel's + // SQ thread submits; an SQ thread that has gone idle would need + // the NEED_WAKEUP kick the submit branch of this loop gives it, + // but no tier enables SQPOLL today (SQPollIdle is 0 in all + // three), so that case is not handled here. + if !w.sqpoll && w.ring.Pending() > 0 { + _, _ = w.ring.Submit() } w.wakeMu.Lock() if !w.acceptPaused.Load() || From b37d11b2c68ff49fc6cfa3e6f955b0888810b36f Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 14:16:51 +0200 Subject: [PATCH 7/7] test(iouring): keep only the celeris#712 regression arms; drop the park change, which #793 superseded On main after #793 (3e7abba) the six park-FIN tests pass 18/18 without this branch's worker.go block, and the four park arms fail on main before #793 (1fdfcd4). #793 keeps the descriptor while a recv is owed and does not park while closeFDOwed is non-zero, so the park change is dropped and the tests stay as the regression guard. --- engine/iouring/park_close_fin_test.go | 10 ++++++++-- engine/iouring/worker.go | 27 --------------------------- 2 files changed, 8 insertions(+), 29 deletions(-) diff --git a/engine/iouring/park_close_fin_test.go b/engine/iouring/park_close_fin_test.go index 97d0a479..80e5a827 100644 --- a/engine/iouring/park_close_fin_test.go +++ b/engine/iouring/park_close_fin_test.go @@ -14,12 +14,18 @@ package iouring // iteration saw no FIN, and both ends stayed ESTABLISHED, until something // woke the worker. // +// celeris#685 (#793) fixed it: a close path that still owes the kernel an op +// on the descriptor leaves the descriptor to its pendingRelease entry, which +// closes it at the op's terminal CQE, and the worker does not park while such +// a close is outstanding (closeFDOwed). The loop's ring waits run the deferred +// completions, however many there are, before the park. +// // Every arm below closes sync-mode connections the engine gave up on and asks // one thing of the client's side of them: that the close arrives. The first // arms close one connection; the many-connection arms close 64 together. // One enter runs at most 20 deferred completions per local-work pass -// (IO_LOCAL_TW_DEFAULT_MAX, kernel 6.13 and later), so one enter before the -// park reached only the first of them. +// (IO_LOCAL_TW_DEFAULT_MAX, kernel 6.13 and later), so a fix that makes one +// enter before the park reaches only the first of them. import ( "errors" diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index 0fb7e105..3e4357e4 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -1487,33 +1487,6 @@ func (w *Worker) run(ctx context.Context) { if w.listenFD < 0 && w.connCount == 0 && !w.hasDriverConns.Load() && w.driverActionPending.Load() == 0 && w.detachQPending.Load() == 0 && w.acceptPaused.Load() { - // Do not park while a closed connection's kernel ops are still - // in flight (celeris#712). finishClose's HTTP/1 fast path is a - // plain close(fd), and the recv the connection still has armed - // holds its own reference to the file, so the socket goes, and - // its client gets a FIN, only once that recv's cancelled - // completion has run. On a DEFER_TASKRUN ring the completion is - // task work, which runs only inside an io_uring_enter with - // GETEVENTS, and the park waits on a Go channel: a connection - // closed in the parking iteration stayed ESTABLISHED for as long - // as the park lasted. One more enter is not enough either, as it - // runs at most 20 deferred completions per local-work pass - // (IO_LOCAL_TW_DEFAULT_MAX, kernel 6.13 and later). - // - // pendingRelease already holds each closed connState until its - // kernelInflight reports every op's terminal CQE, so an entry a - // drain leaves is such an op: go round again, and the ring wait - // runs the work and returns with those CQEs. The clock is read - // for the drain's wall-clock backstop, so an op the kernel never - // completes holds the park back by pendingReleaseHoldNanos and - // one ring wait at most. - if len(w.pendingRelease) > 0 { - w.cachedNow = time.Now().UnixNano() - w.drainPendingRelease() - if len(w.pendingRelease) > 0 { - continue - } - } // Submit what this iteration queued before parking (celeris#657, // A5). The iteration that closes or hands off the last conn // queues its close-path cancels (the header timer's, a send's)