diff --git a/engine/epoll/park_stale_clock_test.go b/engine/epoll/park_stale_clock_test.go new file mode 100644 index 00000000..322b412d --- /dev/null +++ b/engine/epoll/park_stale_clock_test.go @@ -0,0 +1,170 @@ +//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, + WriteTimeout: staleClockReadTimeout, // as the io_uring arms: see there + 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..116aea00 --- /dev/null +++ b/engine/iouring/park_stale_clock_test.go @@ -0,0 +1,445 @@ +//go:build linux + +package iouring + +// celeris#713. A worker stamps lastActivity from w.cachedNow, a clock it +// 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 +// 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 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" + "io" + "log/slog" + "net" + "net/http" + "testing" + "time" + + "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 + 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, +// staleClockGap apart: two seconds of traffic that no timeout may interrupt. +func staleClockAfterPark(t *testing.T, park time.Duration, adopt bool) { + 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. + 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() + 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 == m.CloseCount + }) { + t.Fatalf("celeris713 PREMISE: the engine is not idle") + } + 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) + + 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, + 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. 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) +} + +// 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 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. 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) +} + +// 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, 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, 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 + c.ReadTimeout = staleClockReadTimeout + c.WriteTimeout = staleClockReadTimeout + c.IdleTimeout = 10 * time.Minute + }) + e.mu.Lock() + ws := append([]*Worker(nil), e.workers...) + e.mu.Unlock() + 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 + } + 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) + parked := 0 + for _, w := range ws { + if w.suspended.Load() { + parked++ + } + } + 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) + 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. 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 + 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) + } +} diff --git a/engine/iouring/transplant.go b/engine/iouring/transplant.go index e21fc74c..21422acf 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,14 @@ 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, 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 // fresh parser at the request boundary. diff --git a/engine/iouring/worker.go b/engine/iouring/worker.go index be138d70..d6a9ac3d 100644 --- a/engine/iouring/worker.go +++ b/engine/iouring/worker.go @@ -442,7 +442,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 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 @@ -1287,12 +1287,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)