Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
244 changes: 174 additions & 70 deletions engine/iouring/driver_linger_close_linux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,9 @@ import (
"errors"
"io"
"net"
"os"
"strconv"
"strings"
"testing"
"time"

Expand Down Expand Up @@ -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
Expand Down
Loading