diff --git a/engine/iouring/park_close_fin_test.go b/engine/iouring/park_close_fin_test.go new file mode 100644 index 00000000..80e5a827 --- /dev/null +++ b/engine/iouring/park_close_fin_test.go @@ -0,0 +1,445 @@ +//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. +// +// 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 a fix that makes one +// enter before the park reaches only the first of them. + +import ( + "errors" + "io" + "log/slog" + "net" + "os" + "strconv" + "strings" + "sync" + "sync/atomic" + "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 +} + +// 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. +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 := 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 + // 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...) + 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") + } + 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 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) + } + 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 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") + } + 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) + } +} + +// 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. +func TestParkedWorkerSendsFINForAHeaderTimeoutClose(t *testing.T) { + finAfterClose(t, true, false, false) +} + +// TestParkedWorkerSendsFINForAReadTimeoutClose: the same, with checkTimeouts' +// 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 +// 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) +} + +// 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) +}