From 19106c1d1b66ccd1446bd557818efec9a0066523 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 09:00:34 +0200 Subject: [PATCH] test(iouring): void a linger attempt whose close did not linger, and retry it (celeris#763) TestDriverLingeringCloseDoesNotStallTheWorker/linger failed in 3 of 30 linger runs on main's CI with "L's onClose fired while its close still lingered" when the close had not lingered at all: the kernel ends a close's linger wait early on a pending signal, the fixed engine then fires onClose at once, and the rig could not tell that from onClose firing before a lingering close returned. Each attempt now reads /proc/net/tcp, after V was served and after onClose was looked for, for a row that still carries the socket's inode: tcp_close orphans the socket only once its linger wait has ended. An attempt whose close has returned by then is void and is retried on a new engine and socket, up to 5 times; the arm fails as apparatus if none lingers. The no-linger control now also checks the probe reads a close with SO_LINGER off as returned. --- .../iouring/driver_linger_close_linux_test.go | 244 +++++++++++++----- 1 file changed, 174 insertions(+), 70 deletions(-) diff --git a/engine/iouring/driver_linger_close_linux_test.go b/engine/iouring/driver_linger_close_linux_test.go index bc4c3a67..01e4ddb3 100644 --- a/engine/iouring/driver_linger_close_linux_test.go +++ b/engine/iouring/driver_linger_close_linux_test.go @@ -6,6 +6,9 @@ import ( "errors" "io" "net" + "os" + "strconv" + "strings" "testing" "time" @@ -187,84 +190,185 @@ func waitFinalized(t *testing.T, w *Worker, fd int) { } } +// waitCloseStarted waits until the engine's close(2) of its duplicate has +// begun: close(2) takes the number out of the descriptor table before the +// socket's release, so opFD no longer names the socket once it has started. +func waitCloseStarted(t *testing.T, opFD int, id socketIdentity) { + t.Helper() + for deadline := time.Now().Add(5 * time.Second); stillNames(opFD, id); { + if time.Now().After(deadline) { + t.Fatalf("the engine's duplicate (fd %d) of the socket was never closed", opFD) + } + time.Sleep(time.Millisecond) + } +} + +// closeInLingerWait reports whether the last close of the socket id names is +// still in its linger wait, and the socket's TCP state (hex, as +// /proc/net/tcp prints it) when it is. +// +// /proc/net/tcp lists the IPv4 TCP sockets of the network namespace, each +// with the inode of the socket file that owns it. tcp_close orphans the +// socket (sock_orphan, after which that inode reads 0) only once its linger +// wait has ended, however it ends: the FIN acknowledged, the linger time up, +// or a signal. So once the last close has begun, a row that still carries +// the socket's inode says that close has not returned. +func closeInLingerWait(t *testing.T, id socketIdentity) (waiting bool, state string) { + t.Helper() + b, err := os.ReadFile("/proc/net/tcp") + if err != nil { + t.Fatalf("apparatus: read /proc/net/tcp: %v", err) + } + ino := strconv.FormatUint(id.ino, 10) + for i, line := range strings.Split(string(b), "\n") { + // sl local rem st tx:rx tr:when retrnsmt uid timeout inode ... + if f := strings.Fields(line); i > 0 && len(f) > 9 && f[9] == ino { + return true, f[3] + } + } + return false, "" +} + +// lingerAttempts is how many sockets the linger arm tries for one whose +// close lingers before it fails as apparatus. +const lingerAttempts = 5 + // TestDriverLingeringCloseDoesNotStallTheWorker is the issue's measurement as // a test. The linger arm fails on the unfixed engine; the no-linger arm is // the control that shows the rig measures a worker that is not blocked. +// +// A close that should linger does not always do so. The kernel leaves the +// linger wait early on a pending signal (the review of #744 measured that in +// this package's other celeris#735 test), and then the close returns at +// once, the fixed engine fires onClose at once, and the attempt looks like +// the no-linger control: V's byte in microseconds and onClose before it +// (3 of 30 linger runs on main's CI, celeris#763 item 5). Such an attempt +// proves nothing about either check, so each attempt reads whether the close +// is still in its linger wait (closeInLingerWait) after V was served and after +// onClose was looked for: an attempt whose close has returned by then is +// void, and the arm tries again on a new engine and socket. It fails as +// apparatus when no attempt's close lingers. A V held past the budget fails +// at once, lingering or not. func TestDriverLingeringCloseDoesNotStallTheWorker(t *testing.T) { - for _, tc := range []struct { - name string - linger int - }{{"no_linger_control", 0}, {"linger", lingerStallSecs}} { - t.Run(tc.name, func(t *testing.T) { - e, stop := startTestEngine(t) - defer stop() - w, wl := e.workers[0], e.WorkerLoop(0) - v := registerSettledDriver(t, w, wl) - defer v.unregisterAndWait(t, wl) - p := newWorkerPark(t, w, wl) - defer p.d.unregisterAndWait(t, wl) - - fd, drain, closed, opFD, id := lingeringDriver(t, w, tc.linger) - drained := false - defer func() { - if !drained { - drain() - } - }() - unregisterAndCloseInOrder(t, p, wl, fd) - waitFinalized(t, w, fd) - finalized := time.Now() - // The issue's probe: give the worker the cancel and the start of - // the close before V's peer writes. - time.Sleep(100 * time.Millisecond) - - start := time.Now() - if _, err := unix.Write(v.peer, []byte{'v'}); err != nil { - t.Fatalf("write V's peer: %v", err) - } - var lat time.Duration - select { - case <-v.recv: - lat = time.Since(start) - case <-time.After(10 * time.Second): - lat = -1 - } - var closedEarly bool - select { - case err := <-closed: - closed <- err // put it back for the check below - closedEarly = true - default: - } - t.Logf("celeris735 LINGER arm=%s linger=%ds v_onrecv_ms=%.3f onclose_before_v=%v", - tc.name, tc.linger, float64(lat)/1e6, closedEarly) - if lat < 0 || lat > lingerStallBudget { - t.Errorf("celeris#735: V's byte reached onRecv after %v (budget %v) while another conn's "+ - "close lingered (SO_LINGER %d s): the worker was blocked in close(2)", lat, lingerStallBudget, tc.linger) - } - if tc.linger > 0 && closedEarly { - t.Errorf("L's onClose fired while its close still lingered: a driver must see the " + - "socket closed before onClose (engine/provider.go)") + t.Run("no_linger_control", func(t *testing.T) { + // With SO_LINGER off the close never waits, so the probe must read + // it as returned: a probe that reads it lingering is broken. + if lingeringCloseAttempt(t, "no_linger_control", 0, 1) { + t.Fatal("apparatus: /proc/net/tcp still shows the socket owned by a file after a close with " + + "SO_LINGER off: closeInLingerWait cannot tell a returned close from a lingering one") + } + }) + t.Run("linger", func(t *testing.T) { + void := 0 + for attempt := 1; attempt <= lingerAttempts; attempt++ { + lingered := lingeringCloseAttempt(t, "linger", lingerStallSecs, attempt) + switch { + case t.Failed(): + t.Logf("celeris735 LINGER arm=linger attempts=%d void=%d result=fail", attempt, void) + return + case lingered: + t.Logf("celeris735 LINGER arm=linger attempts=%d void=%d result=lingered", attempt, void) + return } + void++ + } + t.Fatalf("apparatus: L's close did not linger in any of %d attempts (%d void), so neither "+ + "check was made", lingerAttempts, void) + }) +} - // The peer drains, the close returns, and then onClose fires, - // once, with the unregister's nil error, the socket closed. +// lingeringCloseAttempt is one attempt of +// TestDriverLingeringCloseDoesNotStallTheWorker on a new engine: L with +// SO_LINGER {1, linger}, or SO_LINGER off when linger is 0. It reports +// whether L's close was still in its linger wait once V had been served and +// onClose looked for; only then are the two checks evidence. V's latency +// fails at once either way: nothing but a blocked worker holds V's byte for +// a second. +func lingeringCloseAttempt(t *testing.T, arm string, linger, attempt int) (lingered bool) { + t.Helper() + e, stop := startTestEngine(t) + defer stop() + w, wl := e.workers[0], e.WorkerLoop(0) + v := registerSettledDriver(t, w, wl) + defer v.unregisterAndWait(t, wl) + p := newWorkerPark(t, w, wl) + defer p.d.unregisterAndWait(t, wl) + + fd, drain, closed, opFD, id := lingeringDriver(t, w, linger) + drained := false + defer func() { + if !drained { drain() - drained = true - select { - case err := <-closed: - if err != nil { - t.Errorf("onClose(%v), want nil after UnregisterConn", err) - } - case <-time.After(10 * time.Second): - t.Fatal("L's onClose never fired") - } - t.Logf("celeris735 LINGER arm=%s onclose_after_finalize_ms=%.1f", tc.name, float64(time.Since(finalized))/1e6) - if stillNames(opFD, id) { - t.Errorf("onClose fired with the engine's duplicate (fd %d) still open on the socket", opFD) - } - }) + } + }() + unregisterAndCloseInOrder(t, p, wl, fd) + waitFinalized(t, w, fd) + finalized := time.Now() + waitCloseStarted(t, opFD, id) + // The issue's probe: give the worker the cancel and the start of the + // close before V's peer writes. + time.Sleep(100 * time.Millisecond) + + start := time.Now() + if _, err := unix.Write(v.peer, []byte{'v'}); err != nil { + t.Fatalf("write V's peer: %v", err) + } + var lat time.Duration + select { + case <-v.recv: + lat = time.Since(start) + case <-time.After(10 * time.Second): + lat = -1 + } + var closedEarly bool + select { + case err := <-closed: + closed <- err // put it back for the check below + closedEarly = true + default: + } + // After onClose was looked for: a close still waiting now was waiting + // when V was served, and when onClose was looked for. + lingered, state := closeInLingerWait(t, id) + stalled := lat < 0 || lat > lingerStallBudget + early := linger > 0 && lingered && closedEarly + verdict := "void" // the close did not linger: neither check is evidence + switch { + case stalled || early: + verdict = "fail" + case linger == 0: + verdict = "control" + case lingered: + verdict = "valid" + } + t.Logf("celeris735 LINGER arm=%s attempt=%d linger=%ds v_onrecv_ms=%.3f onclose_before_v=%v close_lingering=%v tcp_state=%q verdict=%s", + arm, attempt, linger, float64(lat)/1e6, closedEarly, lingered, state, verdict) + if stalled { + t.Errorf("celeris#735: V's byte reached onRecv after %v (budget %v) while another conn's "+ + "close lingered (SO_LINGER %d s): the worker was blocked in close(2)", lat, lingerStallBudget, linger) + } + if early { + t.Errorf("L's onClose fired while its close still lingered: a driver must see the " + + "socket closed before onClose (engine/provider.go)") + } + + // The peer drains, the close returns, and then onClose fires, + // once, with the unregister's nil error, the socket closed. + drain() + drained = true + select { + case err := <-closed: + if err != nil { + t.Errorf("onClose(%v), want nil after UnregisterConn", err) + } + case <-time.After(10 * time.Second): + t.Fatal("L's onClose never fired") + } + t.Logf("celeris735 LINGER arm=%s attempt=%d onclose_after_finalize_ms=%.1f", arm, attempt, float64(time.Since(finalized))/1e6) + if stillNames(opFD, id) { + t.Errorf("onClose fired with the engine's duplicate (fd %d) still open on the socket", opFD) } + return lingered } // TestDriverShutdownWaitsForAHandedOffClose: a conn finalized before the