From 800e862adef1c926d17216f0d002f6c8a085fda9 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 22:18:17 +0200 Subject: [PATCH 01/10] test(engine): a connection accepted or adopted after a long park must not be timed out on the parked clock (celeris#713) Failing-first on io_uring: park every worker for longer than ReadTimeout, wake one with an accept or an AdoptConn, and keep that connection busy. Controls: the same after a short park, and the epoll twin of the long-park arms. --- engine/epoll/park_stale_clock_test.go | 169 ++++++++++++++++++++++ engine/iouring/park_stale_clock_test.go | 177 ++++++++++++++++++++++++ 2 files changed, 346 insertions(+) create mode 100644 engine/epoll/park_stale_clock_test.go create mode 100644 engine/iouring/park_stale_clock_test.go diff --git a/engine/epoll/park_stale_clock_test.go b/engine/epoll/park_stale_clock_test.go new file mode 100644 index 00000000..7ef93958 --- /dev/null +++ b/engine/epoll/park_stale_clock_test.go @@ -0,0 +1,169 @@ +//go:build linux + +package epoll + +// celeris#713, the epoll twin. The io_uring worker left its DRAINING→SUSPENDED +// park with the clock it parked with, and stamped the connections it accepted +// or adopted on waking from it; the first checkTimeouts then closed them on +// ReadTimeout. An epoll loop parks the same way, on a Go channel, and it +// stamps from l.cachedNow too, but it refreshes that clock on every +// events-bearing epoll_wait return, and a wake by ResumeAccept or AdoptConn +// comes back with an event (the new listener's accept, the adopt queue's +// eventfd). These are the io_uring tests' arms on this engine: they pin that +// behaviour, so a change to the refresh cannot bring the defect here. + +import ( + "bufio" + "context" + "net" + "net/http" + "testing" + "time" + + "golang.org/x/sys/unix" + + "github.com/goceleris/celeris/engine" + "github.com/goceleris/celeris/resource" +) + +const ( + staleClockReadTimeout = 2 * time.Second + staleClockRequests = 20 + staleClockGap = 100 * time.Millisecond +) + +func staleClockAfterParkEpoll(t *testing.T, park time.Duration, adopt bool) { + 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, err := New(resource.Config{ + Addr: addr, + Protocol: engine.HTTP1, + Resources: resource.Resources{Workers: 2}, + // No TCP_DEFER_ACCEPT: no pause linger (celeris#662), so the + // listeners close as PauseAccept is called and the loops park. + DisableDeferAccept: true, + ReadTimeout: staleClockReadTimeout, + IdleTimeout: 10 * time.Minute, + }, okHandler658{}) + if err != nil { + t.Skipf("epoll engine unavailable: %v", err) + } + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- e.Listen(ctx) }() + t.Cleanup(func() { + cancel() + select { + case <-done: + case <-time.After(5 * time.Second): + t.Error("engine did not stop within 5s") + } + }) + waitFor := func(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) + } + } + if !waitFor(8*time.Second, func() bool { return e.Addr() != nil }) { + t.Skip("epoll engine did not bind") + } + e.mu.Lock() + loops := append([]*Loop(nil), e.loops...) + e.mu.Unlock() + allParked := func() bool { + for _, l := range loops { + if !l.suspended.Load() { + return false + } + } + return true + } + if err := e.PauseAccept(); err != nil { + t.Fatalf("pause: %v", err) + } + if !waitFor(3*time.Second, allParked) { + t.Fatalf("celeris713 PREMISE: the paused loops did not park") + } + parkedAt := time.Now() + time.Sleep(park) + if !allParked() { + t.Fatalf("celeris713 PREMISE: a loop left the park on its own") + } + + m0 := e.Metrics() + var c net.Conn + if adopt { + client, fd := adoptPair658(t) + if err := e.AdoptConn(fd, engine.Carryover{RemoteAddr: client.LocalAddr().String()}); err != nil { + _ = unix.Close(fd) + t.Fatalf("AdoptConn: %v", err) + } + if !waitFor(2*time.Second, func() bool { return e.Metrics().TransplantAdopted == m0.TransplantAdopted+1 }) { + t.Fatalf("celeris713 PREMISE: the adoption did not complete") + } + c = client + } else { + if err := e.ResumeAccept(); err != nil { + t.Fatalf("resume: %v", err) + } + for dl := time.Now().Add(3 * time.Second); ; { + if c, err = net.DialTimeout("tcp", addr, 200*time.Millisecond); err == nil { + break + } + if time.Now().After(dl) { + t.Fatalf("celeris713 PREMISE: no listener after ResumeAccept: %v", err) + } + time.Sleep(5 * time.Millisecond) + } + defer func() { _ = c.Close() }() + } + woke := time.Since(parkedAt) + + br := bufio.NewReader(c) + served, failed, why := 0, -1, "" + for i := range staleClockRequests { + if i > 0 { + time.Sleep(staleClockGap) + } + _ = c.SetDeadline(time.Now().Add(2 * time.Second)) + if _, err := c.Write([]byte("GET / HTTP/1.1\r\nHost: x\r\n\r\n")); err != nil { + failed, why = i, "write: "+err.Error() + break + } + resp, err := http.ReadResponse(br, nil) + if err != nil { + failed, why = i, "read: "+err.Error() + break + } + _ = resp.Body.Close() + served++ + } + m := e.Metrics() + t.Logf("celeris713 epoll adopt=%v park=%v woke_after=%v read_timeout=%v loops=%d served=%d/%d failed_at=%d (%s) closes=%d", + adopt, park, woke.Round(time.Millisecond), staleClockReadTimeout, len(loops), served, staleClockRequests, + failed, why, m.CloseCount-m0.CloseCount) + if failed >= 0 { + t.Errorf("celeris713 STALE: request %d of a connection %s %v after the park began, %v apart, failed "+ + "(%s): the loop closed it on ReadTimeout %v, measured from the clock it parked with", + failed, map[bool]string{true: "adopted", false: "accepted"}[adopt], woke.Round(time.Millisecond), + staleClockGap, why, staleClockReadTimeout) + } +} + +func TestAcceptAfterALongParkIsNotTimedOutEpoll(t *testing.T) { + staleClockAfterParkEpoll(t, staleClockReadTimeout+time.Second, false) +} + +func TestAdoptAfterALongParkIsNotTimedOutEpoll(t *testing.T) { + staleClockAfterParkEpoll(t, staleClockReadTimeout+time.Second, true) +} diff --git a/engine/iouring/park_stale_clock_test.go b/engine/iouring/park_stale_clock_test.go new file mode 100644 index 00000000..f567fd37 --- /dev/null +++ b/engine/iouring/park_stale_clock_test.go @@ -0,0 +1,177 @@ +//go:build linux + +package iouring + +// celeris#713. A worker stamps lastActivity from w.cachedNow, a clock it +// refreshes only on its own iterations (every 64th CQE-bearing one, and in +// checkTimeouts). The DRAINING→SUSPENDED park stops those iterations, so a +// worker leaves the park with the clock it parked with. The connections it +// accepts or adopts on waking were stamped from that clock, and the first +// checkTimeouts after the wake, which compares against a fresh time.Now(), +// saw them idle for the whole park: past ReadTimeout, it closed them, with a +// request already written. Measured through the adaptive engine as 21-156 +// closes at 6 of 6 promotes that came 31 s after a demote (ReadTimeout 30 s), +// and none at 18 s. +// +// Here that is one engine and no adaptive controller: park the workers for +// longer than ReadTimeout, wake one with a new connection (an accept after +// ResumeAccept, or an AdoptConn), and keep that connection busy well inside +// ReadTimeout. It must be served throughout. The short-park arms are the +// controls: the same wake and the same traffic after a park well inside +// ReadTimeout. + +import ( + "bufio" + "net" + "net/http" + "testing" + "time" + + "golang.org/x/sys/unix" + + "github.com/goceleris/celeris/engine" + "github.com/goceleris/celeris/resource" +) + +const ( + staleClockReadTimeout = 2 * time.Second + staleClockRequests = 20 + staleClockGap = 100 * time.Millisecond +) + +// staleClockAfterPark parks every worker for park, wakes one with a new +// connection (adopt: AdoptConn on the still-paused engine; otherwise +// ResumeAccept and a dial), and serves staleClockRequests requests on it, +// staleClockGap apart: two seconds of traffic that no timeout may interrupt. +func staleClockAfterPark(t *testing.T, park time.Duration, adopt bool) { + e, addr := startFDLEngine(t, fdlHandler{}, func(c *resource.Config) { + // No TCP_DEFER_ACCEPT: no pause linger (celeris#662), so the + // listeners close as PauseAccept is called and the workers park at + // once, and startFDLEngine's silent probe dial is accepted at once. + c.DisableDeferAccept = true + c.ReadTimeout = staleClockReadTimeout + c.IdleTimeout = 10 * time.Minute + }) + e.mu.Lock() + ws := append([]*Worker(nil), e.workers...) + e.mu.Unlock() + allParked := func() bool { + for _, w := range ws { + if !w.suspended.Load() { + return false + } + } + return true + } + waitFor := func(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) + } + } + if !waitFor(3*time.Second, func() bool { + m := e.Metrics() + return m.ActiveConnections == 0 && m.AcceptCount >= 1 && m.AcceptCount == m.CloseCount + }) { + t.Fatalf("celeris713 PREMISE: startFDLEngine's probe connection is still live") + } + if err := e.PauseAccept(); err != nil { + t.Fatalf("pause: %v", err) + } + if !waitFor(3*time.Second, allParked) { + t.Fatalf("celeris713 PREMISE: the paused workers did not park") + } + parkedAt := time.Now() + time.Sleep(park) + if !allParked() { + t.Fatalf("celeris713 PREMISE: a worker left the park on its own") + } + + m0 := e.Metrics() + var c net.Conn + if adopt { + client, fd := adoptPair658(t) + if err := e.AdoptConn(fd, engine.Carryover{RemoteAddr: client.LocalAddr().String()}); err != nil { + _ = unix.Close(fd) + t.Fatalf("AdoptConn: %v", err) + } + if !waitFor(2*time.Second, func() bool { return e.Metrics().TransplantAdopted == m0.TransplantAdopted+1 }) { + t.Fatalf("celeris713 PREMISE: the adoption did not complete") + } + c = client + } else { + if err := e.ResumeAccept(); err != nil { + t.Fatalf("resume: %v", err) + } + var err error + for dl := time.Now().Add(3 * time.Second); ; { + if c, err = net.DialTimeout("tcp", addr, 200*time.Millisecond); err == nil { + break + } + if time.Now().After(dl) { + t.Fatalf("celeris713 PREMISE: no listener after ResumeAccept: %v", err) + } + time.Sleep(5 * time.Millisecond) + } + defer func() { _ = c.Close() }() + } + woke := time.Since(parkedAt) + + br := bufio.NewReader(c) + served, failed, why := 0, -1, "" + for i := range staleClockRequests { + if i > 0 { + time.Sleep(staleClockGap) + } + _ = c.SetDeadline(time.Now().Add(2 * time.Second)) + if _, err := c.Write([]byte("GET / HTTP/1.1\r\nHost: x\r\n\r\n")); err != nil { + failed, why = i, "write: "+err.Error() + break + } + resp, err := http.ReadResponse(br, nil) + if err != nil { + failed, why = i, "read: "+err.Error() + break + } + _ = resp.Body.Close() + served++ + } + m := e.Metrics() + t.Logf("celeris713 adopt=%v park=%v woke_after=%v read_timeout=%v workers=%d served=%d/%d failed_at=%d (%s) closes=%d", + adopt, park, woke.Round(time.Millisecond), staleClockReadTimeout, len(ws), served, staleClockRequests, + failed, why, m.CloseCount-m0.CloseCount) + if failed >= 0 { + t.Errorf("celeris713 STALE: request %d of a connection %s %v after the park began, %v apart, failed "+ + "(%s): the worker closed it on ReadTimeout %v, measured from the clock it parked with", + failed, map[bool]string{true: "adopted", false: "accepted"}[adopt], woke.Round(time.Millisecond), + staleClockGap, why, staleClockReadTimeout) + } +} + +// TestAcceptAfterALongParkIsNotTimedOut: a ResumeAccept after a park longer +// than ReadTimeout, then a new connection kept busy. The defect. +func TestAcceptAfterALongParkIsNotTimedOut(t *testing.T) { + staleClockAfterPark(t, staleClockReadTimeout+time.Second, false) +} + +// TestAdoptAfterALongParkIsNotTimedOut: an AdoptConn that wakes a worker +// parked for longer than ReadTimeout. The adaptive promote's shape. +func TestAdoptAfterALongParkIsNotTimedOut(t *testing.T) { + staleClockAfterPark(t, staleClockReadTimeout+time.Second, true) +} + +// TestAcceptAfterAShortParkIsNotTimedOut is the control: the same wake and +// the same traffic after a park well inside ReadTimeout. +func TestAcceptAfterAShortParkIsNotTimedOut(t *testing.T) { + staleClockAfterPark(t, 100*time.Millisecond, false) +} + +// TestAdoptAfterAShortParkIsNotTimedOut is the adoption's control. +func TestAdoptAfterAShortParkIsNotTimedOut(t *testing.T) { + staleClockAfterPark(t, 100*time.Millisecond, true) +} From 45e5405ac1a35b9564fab80b499abf1b2ec57c3b Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 22:22:19 +0200 Subject: [PATCH 02/10] test(iouring): an adoption is stamped with the time it was adopted, not the worker's cached clock (celeris#713) Fault injection: a worker clock an hour stale, as a draining worker's can be for tens of seconds (1 s ring waits, a refresh every 64th CQE-bearing iteration and in checkTimeouts). --- engine/iouring/park_stale_clock_test.go | 32 +++++++++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/engine/iouring/park_stale_clock_test.go b/engine/iouring/park_stale_clock_test.go index f567fd37..5a9e3768 100644 --- a/engine/iouring/park_stale_clock_test.go +++ b/engine/iouring/park_stale_clock_test.go @@ -175,3 +175,35 @@ func TestAcceptAfterAShortParkIsNotTimedOut(t *testing.T) { func TestAdoptAfterAShortParkIsNotTimedOut(t *testing.T) { staleClockAfterPark(t, 100*time.Millisecond, true) } + +// TestAdoptionIsStampedWithTheTimeItWasAdopted is the same rule without the +// park, on a worker whose clock is stale for any other reason: cachedNow is +// refreshed only every 64th CQE-bearing iteration and by checkTimeouts, so a +// draining worker waiting out 1 s ring waits can hold a clock tens of seconds +// old. Injected here directly: an hour. The adopted connection's lastActivity +// must be the time of the adoption, or checkTimeouts reads the hour as idle +// time. +func TestAdoptionIsStampedWithTheTimeItWasAdopted(t *testing.T) { + f := newFDLFixture(t, false) + w := f.w + fd, peer := socketPairFDs(t) + t.Cleanup(func() { _ = unix.Close(peer) }) + _ = unix.SetNonblock(fd, true) + if fd >= len(w.conns) { + t.Fatalf("celeris713 PREMISE: fd %d is past the fixture's table (%d)", fd, len(w.conns)) + } + t0 := time.Now().UnixNano() + w.cachedNow = t0 - int64(time.Hour) + w.attachAdoptedFD(fd, engine.Carryover{RemoteAddr: "127.0.0.1:1"}) + cs := w.conns[fd] + if cs == nil { + t.Fatalf("celeris713 PREMISE: the adoption was refused") + } + age := time.Duration(t0 - cs.lastActivity) + t.Logf("celeris713 adopt stamp: lastActivity is %v before the adoption began (worker clock %v stale)", + age, time.Duration(t0-w.cachedNow)) + if cs.lastActivity < t0 { + t.Errorf("celeris713 STALE: an adopted connection's lastActivity is %v before its adoption: it was "+ + "stamped from the worker's cached clock, and checkTimeouts reads that as idle time", age) + } +} From 2f24d7880d0638cdc112c94bc85f328c04ebc210 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 22:22:19 +0200 Subject: [PATCH 03/10] fix(iouring): read the clock again when a worker leaves its park, and stamp an adoption with a fresh one (celeris#713) The DRAINING->SUSPENDED park stops the iterations that refresh w.cachedNow, so a worker woke with the clock it parked with and stamped the connections it accepted or adopted next from it. The first checkTimeouts after the wake, which compares against a fresh time.Now(), read the whole park as idle time: past ReadTimeout it closed every one of them with a request already written (adaptive: 21-156 closes at 6 of 6 promotes 31 s after a demote, ReadTimeout 30 s). The worker now refreshes cachedNow on leaving the park, and attachAdoptedFD stamps an adoption with time.Now(), so a worker whose clock is stale for any other reason cannot hand an adopted connection its age. --- engine/iouring/transplant.go | 12 +++++++++++- engine/iouring/worker.go | 11 ++++++++++- 2 files changed, 21 insertions(+), 2 deletions(-) diff --git a/engine/iouring/transplant.go b/engine/iouring/transplant.go index e21fc74c..7f6b1c47 100644 --- a/engine/iouring/transplant.go +++ b/engine/iouring/transplant.go @@ -6,6 +6,7 @@ import ( "context" "errors" "fmt" + "time" "golang.org/x/sys/unix" @@ -178,7 +179,16 @@ func (w *Worker) attachAdoptedFD(newFD int, carry engine.Carryover) { if w.transplantCount != nil { w.transplantCount.Add(1) } - cs.lastActivity = w.cachedNow + // A fresh clock, not w.cachedNow: an adoption is drained after the CQE + // batch, often in an iteration that carried none, by a worker that may + // have waited out whole seconds since its last refresh (a draining + // worker's ring wait is up to 1 s, and cachedNow is refreshed only every + // 64th CQE-bearing iteration and by checkTimeouts). A stamp that old is + // read by the next checkTimeouts as idle time, and past ReadTimeout it + // closed the connection it had just adopted (celeris#713). Adoptions + // come one per handed-off connection, so the vDSO call is off the + // request path. + cs.lastActivity = time.Now().UnixNano() // #383 adopts HTTP/1 keep-alive conns only; lock the protocol and install a // fresh parser at the request boundary. diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index 3e4357e4..b649a7df 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -420,7 +420,7 @@ type Worker struct { linkArmBatch uint64 tickCounter uint32 - cachedNow int64 // cached time.Now().UnixNano(), refreshed every 64 iterations + cachedNow int64 // cached time.Now().UnixNano(), refreshed every 64 CQE-bearing iterations, by checkTimeouts, and on leaving the park iterCount uint64 // monotonic event-loop iteration counter (for pendingRelease) // pendingRelease defers returning connState structs to the pool @@ -1513,6 +1513,15 @@ func (w *Worker) run(ctx context.Context) { select { case <-wake: + // The park stopped the iterations that refresh cachedNow, + // so it still reads the time the worker parked at. The + // connections this worker accepts or adopts next are + // stamped from it, and the first checkTimeouts compares + // those stamps with a fresh time.Now(): after a park + // longer than ReadTimeout it closed every one of them, + // with a request already written (celeris#713). Read the + // clock again before anything is stamped. + w.cachedNow = time.Now().UnixNano() case <-ctx.Done(): w.shutdown() return From 239a2d7d5058eb23a306565864f35a7c84e6d829 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 23:14:19 +0200 Subject: [PATCH 04/10] test(iouring): start the celeris#713 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_stale_clock_test.go | 48 ++++++++++++++++++++++--- 1 file changed, 44 insertions(+), 4 deletions(-) diff --git a/engine/iouring/park_stale_clock_test.go b/engine/iouring/park_stale_clock_test.go index 5a9e3768..0e7a35e2 100644 --- a/engine/iouring/park_stale_clock_test.go +++ b/engine/iouring/park_stale_clock_test.go @@ -22,6 +22,8 @@ package iouring import ( "bufio" + "io" + "log/slog" "net" "net/http" "testing" @@ -30,9 +32,47 @@ import ( "golang.org/x/sys/unix" "github.com/goceleris/celeris/engine" + "github.com/goceleris/celeris/protocol/h2/stream" "github.com/goceleris/celeris/resource" ) +// startParkEngine713 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 startParkEngine713(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 +} + const ( staleClockReadTimeout = 2 * time.Second staleClockRequests = 20 @@ -44,10 +84,10 @@ const ( // ResumeAccept and a dial), and serves staleClockRequests requests on it, // staleClockGap apart: two seconds of traffic that no timeout may interrupt. func staleClockAfterPark(t *testing.T, park time.Duration, adopt bool) { - e, addr := startFDLEngine(t, fdlHandler{}, func(c *resource.Config) { + e, addr := startParkEngine713(t, fdlHandler{}, func(c *resource.Config) { // No TCP_DEFER_ACCEPT: no pause linger (celeris#662), so the // listeners close as PauseAccept is called and the workers park at - // once, and startFDLEngine's silent probe dial is accepted at once. + // once. c.DisableDeferAccept = true c.ReadTimeout = staleClockReadTimeout c.IdleTimeout = 10 * time.Minute @@ -76,9 +116,9 @@ func staleClockAfterPark(t *testing.T, park time.Duration, adopt bool) { } if !waitFor(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("celeris713 PREMISE: startFDLEngine's probe connection is still live") + t.Fatalf("celeris713 PREMISE: the engine is not idle") } if err := e.PauseAccept(); err != nil { t.Fatalf("pause: %v", err) From 013b2ab6faae9f7532ce4af4dd53543ee6c3dd49 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 23:17:51 +0200 Subject: [PATCH 05/10] test(engine): bound the celeris#713 arms' WriteTimeout as their ReadTimeout checkTimeouts judges a connection with a SEND in flight against WriteTimeout (60 s by default), and the first check after the wake refreshes the worker's clock, so a stale stamp was forgiven whenever that check caught the response on its way out. On a paused engine, whose worker iterates only on the adopted connection's own completions, that is about half the time: main failed the adoption arm 8/10 (m8) and 4/10 (unl) against 10/10 for the accept arm. --- engine/epoll/park_stale_clock_test.go | 1 + engine/iouring/park_stale_clock_test.go | 8 ++++++++ 2 files changed, 9 insertions(+) diff --git a/engine/epoll/park_stale_clock_test.go b/engine/epoll/park_stale_clock_test.go index 7ef93958..322b412d 100644 --- a/engine/epoll/park_stale_clock_test.go +++ b/engine/epoll/park_stale_clock_test.go @@ -47,6 +47,7 @@ func staleClockAfterParkEpoll(t *testing.T, park time.Duration, adopt bool) { // listeners close as PauseAccept is called and the loops park. DisableDeferAccept: true, ReadTimeout: staleClockReadTimeout, + WriteTimeout: staleClockReadTimeout, // as the io_uring arms: see there IdleTimeout: 10 * time.Minute, }, okHandler658{}) if err != nil { diff --git a/engine/iouring/park_stale_clock_test.go b/engine/iouring/park_stale_clock_test.go index 0e7a35e2..4464afd6 100644 --- a/engine/iouring/park_stale_clock_test.go +++ b/engine/iouring/park_stale_clock_test.go @@ -90,6 +90,14 @@ func staleClockAfterPark(t *testing.T, park time.Duration, adopt bool) { // once. c.DisableDeferAccept = true c.ReadTimeout = staleClockReadTimeout + // checkTimeouts judges a connection with a SEND in flight against + // WriteTimeout instead (60 s by default), so with the default a + // stale stamp was forgiven whenever the first check after the wake + // caught the response on its way out: on a paused engine, whose + // worker iterates only on this connection's own completions, about + // half the time. The same bound on both keeps the arm from depending + // on where in a request that check lands. + c.WriteTimeout = staleClockReadTimeout c.IdleTimeout = 10 * time.Minute }) e.mu.Lock() From ebbeda64af3cb1251689707c1fadaa2d22b96d75 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 02:18:59 +0200 Subject: [PATCH 06/10] test(iouring): say that the celeris#713 short-park adoption arm is not a control With WriteTimeout bounded like ReadTimeout it failed on main too, 10/10 in both shapes: a paused worker's clock stays at its last pre-park refresh until its first checkTimeouts after the wake, 32 iterations into the adopted connection's own traffic. Comments only. --- engine/iouring/park_stale_clock_test.go | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/engine/iouring/park_stale_clock_test.go b/engine/iouring/park_stale_clock_test.go index 4464afd6..c5c80570 100644 --- a/engine/iouring/park_stale_clock_test.go +++ b/engine/iouring/park_stale_clock_test.go @@ -16,9 +16,9 @@ package iouring // Here that is one engine and no adaptive controller: park the workers for // longer than ReadTimeout, wake one with a new connection (an accept after // ResumeAccept, or an AdoptConn), and keep that connection busy well inside -// ReadTimeout. It must be served throughout. The short-park arms are the -// controls: the same wake and the same traffic after a park well inside -// ReadTimeout. +// ReadTimeout. It must be served throughout. The short-park accept arm is the +// control: the same wake and the same traffic after a park well inside +// ReadTimeout (the short-park adoption arm is not one; see there). import ( "bufio" @@ -219,7 +219,15 @@ func TestAcceptAfterAShortParkIsNotTimedOut(t *testing.T) { staleClockAfterPark(t, 100*time.Millisecond, false) } -// TestAdoptAfterAShortParkIsNotTimedOut is the adoption's control. +// TestAdoptAfterAShortParkIsNotTimedOut is the adoption arm after a 100 ms +// park. It is NOT a control: on main it failed too, 10 of 10 in both +// shapes. A paused worker's clock stays at its last refresh before the park +// until its first checkTimeouts after the wake, and on a paused engine the +// worker iterates only on this connection's own completions, so that check +// comes 32 iterations, about 13 requests, into the traffic. The idle time +// before the park, the park and that stretch together pass ReadTimeout. The +// accept arm above is the control: a resumed worker's short listener waits +// bring its first check within a request or two. func TestAdoptAfterAShortParkIsNotTimedOut(t *testing.T) { staleClockAfterPark(t, 100*time.Millisecond, true) } From 05a5cfe77a2d2b956457ccd2910bc055f28b39bf Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 03:31:24 +0200 Subject: [PATCH 07/10] test(iouring): the celeris#713 class without a park: an idle worker, a paused worker that holds a connection, an adoption onto it The round-1 fix read the clock again when a worker left its park and stamped an adoption with a fresh clock. The stale clock is not the park's alone: any worker that waits long on its ring and stamps from a clock read before the wait reads its own wait as the connection's idle time. Three arms with no park, one connection kept busy well inside ReadTimeout (2 s), the first two with a control: - a running worker, idle 5 s with ReadHeaderTimeout off (checkTimeouts every 1024th iteration), then a request every 500 ms; the control keeps the default ReadHeaderTimeout; - a worker paused while it holds the connection, requests 300 ms apart; the control has them 100 ms apart; - an adoption onto a paused worker that has not parked, 3 s after the pause: sixteen idle keep-alive holders keep every worker from parking. The traffic loop is now serveSteadily, shared by every arm. The adoption arm's comment no longer calls it the adaptive promote's shape: a promote resumes the new active engine before it adopts. --- engine/iouring/park_stale_clock_test.go | 219 ++++++++++++++++++++---- 1 file changed, 189 insertions(+), 30 deletions(-) diff --git a/engine/iouring/park_stale_clock_test.go b/engine/iouring/park_stale_clock_test.go index c5c80570..5ae41427 100644 --- a/engine/iouring/park_stale_clock_test.go +++ b/engine/iouring/park_stale_clock_test.go @@ -19,6 +19,14 @@ package iouring // ReadTimeout. It must be served throughout. The short-park accept arm is the // control: the same wake and the same traffic after a park well inside // ReadTimeout (the short-park adoption arm is not one; see there). +// +// The park is one way to wait out whole seconds between refreshes; the class +// is any worker that waits long on its ring and stamps from a clock read +// before the wait. The live-connection arms below keep a worker from parking +// and still stretch its clock past ReadTimeout: a running worker with +// ReadHeaderTimeout off (checkTimeouts every 1024th iteration, idle waits of +// up to 100 ms), a paused worker that still holds a connection (waits of up to +// 1 s, checkTimeouts every 32nd), and an adoption onto such a worker. import ( "bufio" @@ -79,6 +87,29 @@ const ( staleClockGap = 100 * time.Millisecond ) +// serveSteadily sends n requests on c, gap apart, and reads each response. It +// returns how many were served and, if one failed, which and why (-1, ""). +func serveSteadily(c net.Conn, n int, gap time.Duration) (served, failed int, why string) { + br := bufio.NewReader(c) + failed = -1 + for i := range n { + if i > 0 { + time.Sleep(gap) + } + _ = c.SetDeadline(time.Now().Add(2 * time.Second)) + if _, err := c.Write([]byte("GET / HTTP/1.1\r\nHost: x\r\n\r\n")); err != nil { + return served, i, "write: " + err.Error() + } + resp, err := http.ReadResponse(br, nil) + if err != nil { + return served, i, "read: " + err.Error() + } + _ = resp.Body.Close() + served++ + } + return served, failed, "" +} + // staleClockAfterPark parks every worker for park, wakes one with a new // connection (adopt: AdoptConn on the still-paused engine; otherwise // ResumeAccept and a dial), and serves staleClockRequests requests on it, @@ -170,25 +201,7 @@ func staleClockAfterPark(t *testing.T, park time.Duration, adopt bool) { } woke := time.Since(parkedAt) - br := bufio.NewReader(c) - served, failed, why := 0, -1, "" - for i := range staleClockRequests { - if i > 0 { - time.Sleep(staleClockGap) - } - _ = c.SetDeadline(time.Now().Add(2 * time.Second)) - if _, err := c.Write([]byte("GET / HTTP/1.1\r\nHost: x\r\n\r\n")); err != nil { - failed, why = i, "write: "+err.Error() - break - } - resp, err := http.ReadResponse(br, nil) - if err != nil { - failed, why = i, "read: "+err.Error() - break - } - _ = resp.Body.Close() - served++ - } + served, failed, why := serveSteadily(c, staleClockRequests, staleClockGap) m := e.Metrics() t.Logf("celeris713 adopt=%v park=%v woke_after=%v read_timeout=%v workers=%d served=%d/%d failed_at=%d (%s) closes=%d", adopt, park, woke.Round(time.Millisecond), staleClockReadTimeout, len(ws), served, staleClockRequests, @@ -208,7 +221,10 @@ func TestAcceptAfterALongParkIsNotTimedOut(t *testing.T) { } // TestAdoptAfterALongParkIsNotTimedOut: an AdoptConn that wakes a worker -// parked for longer than ReadTimeout. The adaptive promote's shape. +// parked for longer than ReadTimeout. In the engine that is the reclaim onto +// a draining source (reclaimTransplant); an adaptive promote resumes the new +// active engine before it adopts (adaptive/engine.go), so there io_uring +// adopts on a resumed worker, the accept arm's shape. func TestAdoptAfterALongParkIsNotTimedOut(t *testing.T) { staleClockAfterPark(t, staleClockReadTimeout+time.Second, true) } @@ -225,20 +241,163 @@ func TestAcceptAfterAShortParkIsNotTimedOut(t *testing.T) { // until its first checkTimeouts after the wake, and on a paused engine the // worker iterates only on this connection's own completions, so that check // comes 32 iterations, about 13 requests, into the traffic. The idle time -// before the park, the park and that stretch together pass ReadTimeout. The -// accept arm above is the control: a resumed worker's short listener waits -// bring its first check within a request or two. +// before the park, the park and that stretch together pass ReadTimeout. It is +// the paused live-connection class below, on an adopted connection. The accept +// arm above is the control: a resumed worker's short listener waits bring its +// first check within a request or two. func TestAdoptAfterAShortParkIsNotTimedOut(t *testing.T) { staleClockAfterPark(t, 100*time.Millisecond, true) } -// TestAdoptionIsStampedWithTheTimeItWasAdopted is the same rule without the -// park, on a worker whose clock is stale for any other reason: cachedNow is -// refreshed only every 64th CQE-bearing iteration and by checkTimeouts, so a -// draining worker waiting out 1 s ring waits can hold a clock tens of seconds -// old. Injected here directly: an hour. The adopted connection's lastActivity -// must be the time of the adoption, or checkTimeouts reads the hour as idle -// time. +// staleClockLiveConn serves n requests, gap apart, on one connection of a +// worker that never parks: running, after idle, with ReadHeaderTimeout off +// (rhtOff), or paused after the first request (pause), still holding it. +// ReadTimeout and WriteTimeout are 2 s, and every gap is well inside them. +func staleClockLiveConn(t *testing.T, rhtOff bool, idle time.Duration, pause bool, gap time.Duration, n int) { + e, addr := startParkEngine713(t, fdlHandler{}, func(c *resource.Config) { + c.DisableDeferAccept = true + c.ReadTimeout = staleClockReadTimeout + c.WriteTimeout = staleClockReadTimeout + c.IdleTimeout = 10 * time.Minute + if rhtOff { + c.ReadHeaderTimeout = -1 + } + }) + time.Sleep(idle) + c, err := net.DialTimeout("tcp", addr, 2*time.Second) + if err != nil { + t.Fatalf("dial: %v", err) + } + defer func() { _ = c.Close() }() + m0 := e.Metrics() + served, failed, why := serveSteadily(c, 1, 0) + if failed < 0 && pause { + if err := e.PauseAccept(); err != nil { + t.Fatalf("pause: %v", err) + } + } + if failed < 0 { + time.Sleep(gap) + var s int + s, failed, why = serveSteadily(c, n-1, gap) + served += s + if failed >= 0 { + failed++ + } + } + m := e.Metrics() + t.Logf("celeris713 live rht_off=%v idle=%v pause=%v gap=%v read_timeout=%v workers=%d served=%d/%d "+ + "failed_at=%d (%s) closes=%d active_after=%d", + rhtOff, idle, pause, gap, staleClockReadTimeout, e.NumWorkers(), served, n, failed, why, + m.CloseCount-m0.CloseCount, m.ActiveConnections) + if failed >= 0 { + t.Errorf("celeris713 STALE: request %d of a connection with requests %v apart failed (%s): the "+ + "worker closed it on ReadTimeout %v, measured from a stamp taken with a clock older than the "+ + "wait before it", failed, gap, why, staleClockReadTimeout) + } +} + +// TestIdleWorkerWithoutAHeaderTimeoutIsNotTimedOut: a running worker, idle +// for 5 s with ReadHeaderTimeout off, then one connection with a request +// every 500 ms. That worker's clock was refreshed only by checkTimeouts, +// every 1024th iteration, and at a CQE batch whose iteration count was a +// multiple of 64, so the connection's stamps were seconds old. +func TestIdleWorkerWithoutAHeaderTimeoutIsNotTimedOut(t *testing.T) { + staleClockLiveConn(t, true, 5*time.Second, false, 500*time.Millisecond, 20) +} + +// TestIdleWorkerWithAHeaderTimeoutIsNotTimedOut is its control: the default +// ReadHeaderTimeout runs checkTimeouts every 32nd iteration, and its waits +// of at most 25 ms keep that within a request. +func TestIdleWorkerWithAHeaderTimeoutIsNotTimedOut(t *testing.T) { + staleClockLiveConn(t, false, 5*time.Second, false, 500*time.Millisecond, 20) +} + +// TestPausedWorkerWithALiveConnectionIsNotTimedOut: a worker paused while it +// holds a connection does not park, and waits up to 1 s on its ring between +// that connection's requests. Its clock was refreshed by checkTimeouts every +// 32nd iteration: with a request every 300 ms, past ReadTimeout. +func TestPausedWorkerWithALiveConnectionIsNotTimedOut(t *testing.T) { + staleClockLiveConn(t, false, 0, true, 300*time.Millisecond, 30) +} + +// TestPausedWorkerWithABusyConnectionIsNotTimedOut is its control: requests +// 100 ms apart bring the 32nd iteration within ReadTimeout. +func TestPausedWorkerWithABusyConnectionIsNotTimedOut(t *testing.T) { + staleClockLiveConn(t, false, 0, true, 100*time.Millisecond, 30) +} + +// TestAdoptOntoADrainingWorkerIsNotTimedOut: an adoption onto a paused worker +// that has not parked, because it still holds idle keep-alive connections, +// 3 s after the pause. The adoption's own stamp is fresh, and the +// connection's first request overwrote it with the worker's clock. Sixteen +// holders, so that every worker holds one (SO_REUSEPORT spreads them) and +// none parks: the adoption lands on a draining worker, whichever it is. +func TestAdoptOntoADrainingWorkerIsNotTimedOut(t *testing.T) { + e, addr := startParkEngine713(t, fdlHandler{}, func(c *resource.Config) { + c.DisableDeferAccept = true + c.ReadTimeout = staleClockReadTimeout + c.WriteTimeout = staleClockReadTimeout + c.IdleTimeout = 10 * time.Minute + }) + e.mu.Lock() + ws := append([]*Worker(nil), e.workers...) + e.mu.Unlock() + const holders = 16 + for i := range holders { + h, err := net.DialTimeout("tcp", addr, 2*time.Second) + if err != nil { + t.Fatalf("dial holder %d: %v", i, err) + } + defer func() { _ = h.Close() }() + if _, failed, why := serveSteadily(h, 1, 0); failed >= 0 { + t.Fatalf("celeris713 PREMISE: holder %d's request failed: %s", i, why) + } + } + if err := e.PauseAccept(); err != nil { + t.Fatalf("pause: %v", err) + } + pausedAt := time.Now() + time.Sleep(3 * time.Second) + alive := e.Metrics().ActiveConnections + parked := 0 + for _, w := range ws { + if w.suspended.Load() { + parked++ + } + } + if alive != holders || parked != 0 { + t.Fatalf("celeris713 PREMISE: at the adoption %d of %d holders are open and %d of %d workers parked", + alive, holders, parked, len(ws)) + } + m0 := e.Metrics() + client, fd := adoptPair658(t) + if err := e.AdoptConn(fd, engine.Carryover{RemoteAddr: client.LocalAddr().String()}); err != nil { + _ = unix.Close(fd) + t.Fatalf("AdoptConn: %v", err) + } + for dl := time.Now().Add(2 * time.Second); e.Metrics().TransplantAdopted != m0.TransplantAdopted+1; { + if time.Now().After(dl) { + t.Fatalf("celeris713 PREMISE: the adoption did not complete") + } + time.Sleep(2 * time.Millisecond) + } + adoptedAfter := time.Since(pausedAt) + served, failed, why := serveSteadily(client, staleClockRequests, staleClockGap) + m := e.Metrics() + t.Logf("celeris713 drain-adopt adopted_after=%v read_timeout=%v served=%d/%d failed_at=%d (%s) closes=%d", + adoptedAfter.Round(time.Millisecond), staleClockReadTimeout, served, staleClockRequests, failed, why, + m.CloseCount-m0.CloseCount) + if failed >= 0 { + t.Errorf("celeris713 STALE: request %d of a connection adopted by a draining worker, %v apart, "+ + "failed (%s)", failed, staleClockGap, why) + } +} + +// TestAdoptionIsStampedWithTheTimeItWasAdopted is the adoption stamp alone, +// on a worker whose clock is stale for any reason: an hour, injected. The +// adopted connection's lastActivity must be the time of the adoption, or +// checkTimeouts reads the hour as idle time. func TestAdoptionIsStampedWithTheTimeItWasAdopted(t *testing.T) { f := newFDLFixture(t, false) w := f.w From 6267e5608c2a7c08fd527b6900945c06839dee9b Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 03:32:06 +0200 Subject: [PATCH 08/10] fix(iouring): read a worker's clock on every CQE batch, as epoll does, and drop the park-exit read (celeris#713) A worker stamps lastActivity from cachedNow, which it read on every 64th iteration only. An idle or paused worker waits up to 100 ms or 1 s per iteration, so its stamps could be seconds old, and checkTimeouts, which compares them with a fresh time.Now(), read that age as idle time. Round 1 read the clock again on leaving the park, which covered the park and nothing else: a running worker with ReadHeaderTimeout off, a paused worker still holding a connection, and an adoption onto it closed a connection in the middle of steady traffic all the same. The clock is now read once per CQE batch, after the wait that delivered it, as epoll reads it on every events-bearing epoll_wait return. The first stamp after a park comes from such a batch too (an accept's CQE, or the eventfd wake of an adoption), so the park-exit read is gone. The adoption keeps its own fresh stamp, for an adoption drained in an iteration that carried no CQE. One vDSO call per batch, not per request. --- engine/iouring/park_stale_clock_test.go | 7 ++++-- engine/iouring/transplant.go | 14 +++++------- engine/iouring/worker.go | 30 ++++++++++++------------- 3 files changed, 25 insertions(+), 26 deletions(-) diff --git a/engine/iouring/park_stale_clock_test.go b/engine/iouring/park_stale_clock_test.go index 5ae41427..886216f8 100644 --- a/engine/iouring/park_stale_clock_test.go +++ b/engine/iouring/park_stale_clock_test.go @@ -3,7 +3,7 @@ package iouring // celeris#713. A worker stamps lastActivity from w.cachedNow, a clock it -// refreshes only on its own iterations (every 64th CQE-bearing one, and in +// refreshed only on its own iterations (every 64th CQE-bearing one, and in // checkTimeouts). The DRAINING→SUSPENDED park stops those iterations, so a // worker leaves the park with the clock it parked with. The connections it // accepts or adopts on waking were stamped from that clock, and the first @@ -397,7 +397,10 @@ func TestAdoptOntoADrainingWorkerIsNotTimedOut(t *testing.T) { // TestAdoptionIsStampedWithTheTimeItWasAdopted is the adoption stamp alone, // on a worker whose clock is stale for any reason: an hour, injected. The // adopted connection's lastActivity must be the time of the adoption, or -// checkTimeouts reads the hour as idle time. +// checkTimeouts reads the hour as idle time. The engine-level arms do not +// need this stamp: the adoption's eventfd wake is a CQE, and its batch reads +// the clock before the adoption is drained. It covers an adoption drained in +// an iteration that carried no CQE. func TestAdoptionIsStampedWithTheTimeItWasAdopted(t *testing.T) { f := newFDLFixture(t, false) w := f.w diff --git a/engine/iouring/transplant.go b/engine/iouring/transplant.go index 7f6b1c47..21422acf 100644 --- a/engine/iouring/transplant.go +++ b/engine/iouring/transplant.go @@ -180,14 +180,12 @@ func (w *Worker) attachAdoptedFD(newFD int, carry engine.Carryover) { w.transplantCount.Add(1) } // A fresh clock, not w.cachedNow: an adoption is drained after the CQE - // batch, often in an iteration that carried none, by a worker that may - // have waited out whole seconds since its last refresh (a draining - // worker's ring wait is up to 1 s, and cachedNow is refreshed only every - // 64th CQE-bearing iteration and by checkTimeouts). A stamp that old is - // read by the next checkTimeouts as idle time, and past ReadTimeout it - // closed the connection it had just adopted (celeris#713). Adoptions - // come one per handed-off connection, so the vDSO call is off the - // request path. + // batch, which is what reads cachedNow, and in an iteration that carried + // none cachedNow is as old as the wait before it, up to 1 s on a draining + // worker, or the whole park (celeris#713). checkTimeouts reads a stamp's + // age as idle time. The adoption's eventfd wake is normally a CQE of the + // same iteration; this covers one drained without it. Adoptions come one + // per handed-off connection, so the vDSO call is off the request path. cs.lastActivity = time.Now().UnixNano() // #383 adopts HTTP/1 keep-alive conns only; lock the protocol and install a diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index b649a7df..64bb4d4e 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -420,7 +420,7 @@ type Worker struct { linkArmBatch uint64 tickCounter uint32 - cachedNow int64 // cached time.Now().UnixNano(), refreshed every 64 CQE-bearing iterations, by checkTimeouts, and on leaving the park + cachedNow int64 // cached time.Now().UnixNano(), refreshed on every CQE-bearing iteration and by checkTimeouts iterCount uint64 // monotonic event-loop iteration counter (for pendingRelease) // pendingRelease defers returning connState structs to the pool @@ -1189,12 +1189,19 @@ func (w *Worker) run(ctx context.Context) { if cqHead != cqTail { w.emptyIters = 0 // Reset adaptive timeout on activity. - // Refresh cached timestamp every 64 iterations to amortize - // time.Now() vDSO cost (~50ns on ARM64). Timeout detection - // uses multi-second windows so ~1ms resolution is sufficient. - if w.tickCounter&0x3F == 0 { - w.cachedNow = time.Now().UnixNano() - } + // Read the clock once per CQE batch, after the wait that + // delivered it: every stamp this batch takes (an accept's, a + // recv's, a send completion's lastActivity) is then no older + // than the batch, as on epoll, which reads it on every + // events-bearing epoll_wait return. It was read every 64th + // iteration, and an idle or paused worker waits up to 100 ms + // or 1 s per iteration, so its stamps could be seconds old, + // older still after a park, and checkTimeouts, which compares + // them with a fresh time.Now(), read that age as idle time: + // past ReadTimeout it closed a connection in the middle of + // steady traffic (celeris#713). One vDSO call per batch, not + // per request. + w.cachedNow = time.Now().UnixNano() now := w.cachedNow for cqHead != cqTail { entry := w.ring.cqeAt(cqHead) @@ -1513,15 +1520,6 @@ func (w *Worker) run(ctx context.Context) { select { case <-wake: - // The park stopped the iterations that refresh cachedNow, - // so it still reads the time the worker parked at. The - // connections this worker accepts or adopts next are - // stamped from it, and the first checkTimeouts compares - // those stamps with a fresh time.Now(): after a park - // longer than ReadTimeout it closed every one of them, - // with a request already written (celeris#713). Read the - // clock again before anything is stamped. - w.cachedNow = time.Now().UnixNano() case <-ctx.Done(): w.shutdown() return From d3e46c3aa3686ca3e16ec341c97c7ddc0db88ef8 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 05:16:35 +0200 Subject: [PATCH 09/10] test(iouring): keep the celeris#713 draining workers from parking with driver connections, not idle holders TestAdoptOntoADrainingWorkerIsNotTimedOut kept every worker from parking with sixteen idle keep-alive holders. ReadTimeout (2 s) closed a worker's holders during the 3 s wait whenever that worker's checkTimeouts came due, and the worker then parked: a PREMISE failure in 3 of 13 runs at unlimited memlock (two workers), on the fix and on main's engine alike, and none in 50 runs at 8 MiB. A registered driver connection on each worker keeps it from parking instead: the park waits for hasDriverConns to clear, no timeout touches a driver connection, and an idle one completes nothing, so the worker still waits out 1 s ring waits on a clock it does not refresh. --- engine/iouring/park_stale_clock_test.go | 54 ++++++++++++++++--------- 1 file changed, 36 insertions(+), 18 deletions(-) diff --git a/engine/iouring/park_stale_clock_test.go b/engine/iouring/park_stale_clock_test.go index 886216f8..ab51b507 100644 --- a/engine/iouring/park_stale_clock_test.go +++ b/engine/iouring/park_stale_clock_test.go @@ -328,13 +328,17 @@ func TestPausedWorkerWithABusyConnectionIsNotTimedOut(t *testing.T) { } // TestAdoptOntoADrainingWorkerIsNotTimedOut: an adoption onto a paused worker -// that has not parked, because it still holds idle keep-alive connections, -// 3 s after the pause. The adoption's own stamp is fresh, and the -// connection's first request overwrote it with the worker's clock. Sixteen -// holders, so that every worker holds one (SO_REUSEPORT spreads them) and -// none parks: the adoption lands on a draining worker, whichever it is. +// that has not parked, 3 s after the pause. The adoption's own stamp is fresh, +// and the connection's first request overwrote it with the worker's clock. +// What keeps every worker from parking is a registered driver connection on +// each (the park waits for hasDriverConns to clear), not idle keep-alive +// holders: ReadTimeout closed a worker's holders during the wait whenever its +// checkTimeouts came due, and that worker parked (3 of 13 runs at unl). A driver +// connection is not in liveConns, so no timeout touches it, and while idle it +// completes nothing: the worker waits out 1 s ring waits on a clock it does +// not refresh, as it did with the holders. func TestAdoptOntoADrainingWorkerIsNotTimedOut(t *testing.T) { - e, addr := startParkEngine713(t, fdlHandler{}, func(c *resource.Config) { + e, _ := startParkEngine713(t, fdlHandler{}, func(c *resource.Config) { c.DisableDeferAccept = true c.ReadTimeout = staleClockReadTimeout c.WriteTimeout = staleClockReadTimeout @@ -343,32 +347,46 @@ func TestAdoptOntoADrainingWorkerIsNotTimedOut(t *testing.T) { e.mu.Lock() ws := append([]*Worker(nil), e.workers...) e.mu.Unlock() - const holders = 16 - for i := range holders { - h, err := net.DialTimeout("tcp", addr, 2*time.Second) - if err != nil { - t.Fatalf("dial holder %d: %v", i, err) + for i := range ws { + a, b := nonblockSocketPair(t) + wl := e.WorkerLoop(i) + if err := wl.RegisterConn(a, func([]byte) {}, func(error) {}); err != nil { + _ = unix.Close(a) + _ = unix.Close(b) + t.Fatalf("RegisterConn on worker %d: %v", i, err) + } + t.Cleanup(func() { + _ = wl.UnregisterConn(a) + _ = unix.Close(a) + _ = unix.Close(b) + }) + } + for dl := time.Now().Add(2 * time.Second); ; { + all := true + for _, w := range ws { + all = all && w.hasDriverConns.Load() + } + if all { + break } - defer func() { _ = h.Close() }() - if _, failed, why := serveSteadily(h, 1, 0); failed >= 0 { - t.Fatalf("celeris713 PREMISE: holder %d's request failed: %s", i, why) + if time.Now().After(dl) { + t.Fatalf("celeris713 PREMISE: the driver connections were not registered on every worker") } + time.Sleep(2 * time.Millisecond) } if err := e.PauseAccept(); err != nil { t.Fatalf("pause: %v", err) } pausedAt := time.Now() time.Sleep(3 * time.Second) - alive := e.Metrics().ActiveConnections parked := 0 for _, w := range ws { if w.suspended.Load() { parked++ } } - if alive != holders || parked != 0 { - t.Fatalf("celeris713 PREMISE: at the adoption %d of %d holders are open and %d of %d workers parked", - alive, holders, parked, len(ws)) + if parked != 0 { + t.Fatalf("celeris713 PREMISE: at the adoption %d of %d workers parked", parked, len(ws)) } m0 := e.Metrics() client, fd := adoptPair658(t) From 9fb68fe9301dcd2c9df06beec8d821e9d6da41b4 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 05:29:05 +0200 Subject: [PATCH 10/10] test(iouring): keep the run count out of the draining-adoption arm's comment (celeris#713) The comment quoted 3 of 13 unl runs; with the merged tree's package suite the holders' premise failed in 4 of 15 unl runs and in none of 50 at 8 MiB (celeris#768's body, from the lane's logs). The count belongs there, not in the code. Comment only. --- engine/iouring/park_stale_clock_test.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/engine/iouring/park_stale_clock_test.go b/engine/iouring/park_stale_clock_test.go index ab51b507..116aea00 100644 --- a/engine/iouring/park_stale_clock_test.go +++ b/engine/iouring/park_stale_clock_test.go @@ -333,10 +333,10 @@ func TestPausedWorkerWithABusyConnectionIsNotTimedOut(t *testing.T) { // What keeps every worker from parking is a registered driver connection on // each (the park waits for hasDriverConns to clear), not idle keep-alive // holders: ReadTimeout closed a worker's holders during the wait whenever its -// checkTimeouts came due, and that worker parked (3 of 13 runs at unl). A driver -// connection is not in liveConns, so no timeout touches it, and while idle it -// completes nothing: the worker waits out 1 s ring waits on a clock it does -// not refresh, as it did with the holders. +// checkTimeouts came due, and that worker parked, a failed premise seen with +// two workers. A driver connection is not in liveConns, so no timeout touches +// it, and while idle it completes nothing: the worker waits out 1 s ring waits +// on a clock it does not refresh, as it did with the holders. func TestAdoptOntoADrainingWorkerIsNotTimedOut(t *testing.T) { e, _ := startParkEngine713(t, fdlHandler{}, func(c *resource.Config) { c.DisableDeferAccept = true