diff --git a/middleware/websocket/backpressure_oracle_linux_test.go b/middleware/websocket/backpressure_oracle_linux_test.go new file mode 100644 index 00000000..7043723c --- /dev/null +++ b/middleware/websocket/backpressure_oracle_linux_test.go @@ -0,0 +1,1304 @@ +//go:build linux + +package websocket + +// The client and the diagnostics shared by the two real-socket backpressure +// oracles, TestBackpressurePauseDoesNotCancelInflightSend (celeris#482) and +// TestBackpressureInboundSequenceIntegrity (celeris#484). +// +// # Why the client waits on progress, not on a clock +// +// Both tests flood a server whose handler echoes everything, from clients +// that never read during the flood and hold a 32 KiB receive buffer, so the +// server's echo SEND blocks and the engine pauses inbound delivery. That is +// the point of both tests and it is kept. +// +// After the flood the client used to finish its last frame and its Close +// frame with fixed deadlines (15 s and 10 s) and wait a fixed 10 s for the +// server to close. celeris#633 measured what those deadlines judge instead of +// the engine, on GitHub's runners (kernel 6.17): +// +// - The server is often simply behind. Under coverage the echo runs at +// 0.05-0.19x, and 457 of 459 close-timeouts were still receiving bytes +// when the 10 s ran out. With waits that give up only when nothing moves, +// the same connections finished: frame writes after up to 48 s, close +// handshakes after up to 150 s. +// - A client that never reads while it waits can stall itself for good. +// Its receive queue stays full, so since Linux 6.17 (tcp_sequence(), +// 9ca48d616ed7) its kernel drops the server's end-of-window segments +// BEFORE processing the ACK they carry, and it never learns the server +// reopened its window. In 70 of 70 such give-ups the handler sat idle +// in ReadMessage, nothing buffered, not paused, server receive queue +// empty. Emptying the client's receive queue between write slices took +// those give-ups from 70 to 0. +// - The server's write and read errors the tests counted were the echo of +// the client's own give-up: a socket closed with unread data sends an +// RST. None of 1,165 came before its own client's close. +// +// So the client writes in wsoSlice slices, empties its receive queue +// (without blocking) after every slice that times out, and gives up only +// after wsoWriteIdle with no progress (its send queue did not shrink and +// nothing was written), or at wsoWaitCap in total. The close wait is re-armed +// by every byte received and gives up after wsoCloseIdle of silence, or at +// wsoWaitCap. A real wedge still fails, and now it fails with its state: each +// give-up prints the connection's timeline from both ends (WSO-GIVEUP). +// +// # An engine that stops reading, which the drain would hide +// +// The drain has a cost. celeris#607 was a recv the engine chained behind a +// send (IOSQE_IO_LINK) on a connection whose client had stopped reading: the +// send waited on the client's closed window, so the connection could not +// read at all. The drain is exactly what reopens that window, so a client +// that drains gets through such a connection and never gives up. +// +// So the rig also watches the server side of every connection, every +// wsoWatchEvery, for its whole life (watch): a connection whose server +// socket holds unread bytes while nothing asked the engine to stop reading +// it (the chanReader is not paused and holds nothing, the handler has not +// exited) is one the engine should be reading. If the engine reads nothing +// from it for wsoRecvStall, and the process was not starved of CPU meanwhile, +// the test fails with that stretch's state at both ends (WSO-STALL), whatever +// the client does afterwards. That is the celeris#607 class judged by what it +// does, not by the mechanism: a chained recv and a recv arm passed over +// behind an outstanding send look the same here. +// +// # The test binary's deadline +// +// A wedge costs wsoWriteIdle or more per subtest, and CI runs both oracles +// in one binary under -timeout. So every wait also ends at the binary's +// deadline less wsoReserve, which is kept for the verdicts, the engine +// shutdowns and the tests after it. A wait that runs into it gives up with +// reason "budget" and its WSO-GIVEUP record, and a subtest that starts with +// less than wsoMinBudget of it left fails at once, saying so, instead of +// starting a flood it cannot finish. So a run that would have timed out ends +// with each subtest's own verdict instead of a goroutine dump, and a slow +// run that would have finished in time still does: the budget never ends a +// wait before the binary itself would have. +// +// # What the timeline shows +// +// The client joins itself to the handler that serves it by sending its local +// port as ?cid= and the handler reading it back (Conn.Query). Every sample +// reads, for that connection: +// +// - the client's socket: send and receive queue, the windows it knows, +// RTO backoff; +// - the handler: where it is (ReadMessage, WriteMessage, exited), frames +// echoed, time since the last echo; +// - the chanReader: depth, spill, whether it holds the engine paused, and +// how many pauses and resumes it applied (the callbacks are wrapped); +// - the server's socket, read through the engine's own descriptor: bytes +// received and still unread (so, what the engine has read), bytes the +// kernel has accepted from the engine and how many the peer ACKed, the +// windows, retransmissions; +// - derived from those: the bytes the engine has read, and the bytes the +// handler wrote that the engine has not yet handed to the kernel. Both +// count from the connection's first byte, net of the upgrade request and +// the 101 response, whose sizes the client reads from its own socket +// once the handshake is done. +// +// That is the engine's recv and send progress, its resumes and its pending +// write bytes, observed from outside it, the same way on epoll and io_uring. + +import ( + "errors" + "fmt" + "io" + "net" + "os" + "sort" + "strconv" + "strings" + "sync" + "sync/atomic" + "syscall" + "testing" + "time" + + "golang.org/x/sys/unix" +) + +const ( + // wsoSlice is one write slice, and so how often a waiting client + // empties its receive queue and records a timeline sample. + wsoSlice = time.Second + // wsoWriteIdle: a write gives up after this long without progress. + // The longest healthy gap measured was 50.6 s (coverage, arm64 io_uring). + wsoWriteIdle = 90 * time.Second + // wsoCloseIdle: the close wait gives up after this long without a byte. + wsoCloseIdle = 30 * time.Second + // wsoWaitCap bounds any one wait, progress or not. + wsoWaitCap = 240 * time.Second + // wsoCloseSlow: a close wait longer than this is reported. It is the old + // absolute deadline, so these are the waits the old oracle failed. + wsoCloseSlow = 10 * time.Second + // wsoFirstSamples and wsoLastSamples: the timeline samples kept per + // connection, the first ones (the onset) and the most recent ones. + wsoFirstSamples = 8 + wsoLastSamples = 24 + // wsoWatchEvery: how often watch samples the server side of every + // connection. + wsoWatchEvery = 250 * time.Millisecond + // wsoWatchGap: a watcher run this long after its tick was due says the + // process was starved of CPU, and so perhaps the engine too: every open + // stretch starts over (see watch). + wsoWatchGap = time.Second + // wsoRecvStall: a connection whose engine read nothing for this long + // from a socket holding unread bytes, while nothing asked it to stop, + // fails the test. celeris#607 held such a connection for up to 14 s. + // A healthy engine's longest such stretch, measured on GitHub's runners + // in CI's shape, was 2.3 s on epoll and 1.7 s on io_uring. + wsoRecvStall = 5 * time.Second + // wsoQuiet: after its flood the #482 test's client neither reads nor + // writes for this long, holding the backpressure it built. A healthy + // engine reads what the socket holds, or pauses; one that stopped + // reading the connection behind its blocked echo (celeris#607) cannot, + // so that stretch outlasts wsoRecvStall. With #607 re-introduced such a + // stretch began up to 1.6 s after the flood ended; at 6 s of quiet the + // shortest reached 4.5 s. + wsoQuiet = 8 * time.Second + // wsoReserve: what a rig leaves of the test binary's -timeout for the + // verdicts, the engine shutdowns (up to 40 s) and the tests after it. + wsoReserve = 60 * time.Second + // wsoMinBudget: a subtest that starts with less of the budget left than + // this fails at once (its flood and quiet alone take about 20 s). + wsoMinBudget = 30 * time.Second +) + +// wsoRig joins each client connection to the handler that serves it and +// keeps what each side saw. +type wsoRig struct { + t0 time.Time + end time.Time // the wait budget from t.Deadline(); zero: none + srvPort atomic.Int32 // set once the server listens; read by handler goroutines + srvs sync.Map // client local port -> *wsoSrv + clisBy sync.Map // client local port -> *wsoCli + mu sync.Mutex + clis []*wsoCli + hs []*wsoSrv // every handler, joined or not + stallMax time.Duration // the longest stretch watch saw, set when it stops + longest wsoStall // that stretch (see watch) + // watchResets and watchMaxLag: how often, and for how long at most, the + // watcher was not run when its tick was due (see watch). + watchResets int + watchMaxLag time.Duration +} + +// newWSORig starts a subtest's rig. Its waits end, at the latest, at the +// test binary's deadline less wsoReserve (see the file comment); with less +// than wsoMinBudget of that left, the subtest fails here without starting. +func newWSORig(t *testing.T) *wsoRig { + t.Helper() + g := &wsoRig{t0: time.Now()} + if dl, ok := t.Deadline(); ok { + g.end = dl.Add(-wsoReserve) + if left := g.end.Sub(g.t0); left < wsoMinBudget { + t.Fatalf("only %v is left of the test binary's -timeout for this subtest's waits (its deadline less %v "+ + "kept for the verdicts and the shutdowns): not starting a flood that cannot finish. An earlier "+ + "subtest used the time; its own verdict says why", left.Round(time.Second), wsoReserve) + } + } + return g +} + +// overBudget reports whether now is past the rig's wait budget. +func (g *wsoRig) overBudget(now time.Time) bool { return !g.end.IsZero() && !now.Before(g.end) } + +// budgetNote describes the budget for a verdict message. +func (g *wsoRig) budgetNote() string { + if g.end.IsZero() { + return "no -timeout budget" + } + return "the -timeout budget ended at +" + wsoSec(int64(g.end.Sub(g.t0))) +} + +// since is the rig's clock: nanoseconds since the subtest started. +func (g *wsoRig) since() int64 { return int64(time.Since(g.t0)) } + +func wsoSec(ns int64) string { return strconv.FormatFloat(float64(ns)/1e9, 'f', 3, 64) + "s" } + +// setServer records the server's port from its ws:// address. +func (g *wsoRig) setServer(hostPort string) { + if i := strings.LastIndexByte(hostPort, ':'); i >= 0 { + p, _ := strconv.Atoi(hostPort[i+1:]) + g.srvPort.Store(int32(p)) + } +} + +// ---------------------------------------------------------------- server side + +type wsoSrvErr struct { + kind string // "write" or "read" + err error + ns int64 +} + +// wsoSrv is one handler's view of its connection. +type wsoSrv struct { + g *wsoRig + port int // the client's local port, or -1 when the client sent none + r *chanReader + frames atomic.Int64 + wroteB atomic.Int64 // bytes of every frame WriteMessage accepted + phase atomic.Int32 // 1 ReadMessage, 2 WriteMessage, 3 exited + echoNs atomic.Int64 // rig time of the last echo + pauses, resumes atomic.Int64 + pauseNs, resumeNs atomic.Int64 + + mu sync.Mutex + exit string + exitNs int64 + errs []wsoSrvErr + fd int + ino uint64 + gone bool // the socket was found once and its descriptor no longer is it +} + +// attach registers the handler's connection under the client's port and +// wraps the chanReader's pause and resume callbacks so each one applied to +// the engine is counted and timed. Call it first thing in the handler, on +// the handler goroutine: Read reads r.resume without the lock on that +// goroutine, and requestPause reads both under pausedMu, which is held here. +func (g *wsoRig) attach(c *Conn) *wsoSrv { + s := &wsoSrv{g: g, port: -1, r: c.engineReader, fd: -1} + s.phase.Store(1) + if p, err := strconv.Atoi(c.Query("cid")); err == nil && p > 0 { + s.port = p + g.srvs.Store(p, s) + } + g.mu.Lock() + g.hs = append(g.hs, s) + g.mu.Unlock() + if r := s.r; r != nil { + r.pausedMu.Lock() + if pause, resume := r.pause, r.resume; pause != nil && resume != nil { + r.pause = func() { + s.pauses.Add(1) + s.pauseNs.Store(g.since()) + pause() + } + r.resume = func() { + s.resumes.Add(1) + s.resumeNs.Store(g.since()) + resume() + } + } + r.pausedMu.Unlock() + } + g.srvSockFD(s) // find the engine's descriptor while the connection is young + return s +} + +// echoed records a frame the handler echoed. +func (s *wsoSrv) echoed(payload int) { + s.frames.Add(1) + s.wroteB.Add(int64(wsoFrameLen(payload))) + s.echoNs.Store(s.g.since()) +} + +// wsoFrameLen is the size on the wire of an unmasked server frame. +func wsoFrameLen(payload int) int { + switch { + case payload < 126: + return 2 + payload + case payload < 1<<16: + return 4 + payload + } + return 10 + payload +} + +// noteErr records a server-side error with its time, to be judged against +// the time its own client closed (see serverErrs). +func (s *wsoSrv) noteErr(kind string, err error) { + s.mu.Lock() + s.errs = append(s.errs, wsoSrvErr{kind: kind, err: err, ns: s.g.since()}) + s.mu.Unlock() +} + +// exited records how the handler's loop ended. +func (s *wsoSrv) exited(what string, err error) { + s.mu.Lock() + if s.exit == "" { + s.exit, s.exitNs = what+": "+wsoErrStr(err), s.g.since() + } + s.mu.Unlock() + s.phase.Store(3) +} + +func wsoErrStr(err error) string { + if err == nil { + return "nil" + } + return err.Error() +} + +// ------------------------------------------------------------ kernel sockets + +// wsoProcTCP finds the /proc/net/tcp row of the IPv4 socket with local port +// lport and remote port rport: a compact rendering and its inode. +func wsoProcTCP(lport, rport int) (string, uint64) { + b, err := os.ReadFile("/proc/net/tcp") + if err != nil { + return "proc-err", 0 + } + wl, wr := fmt.Sprintf(":%04X", lport), fmt.Sprintf(":%04X", rport) + lines := strings.Split(string(b), "\n") + for _, line := range lines[1:] { + f := strings.Fields(line) + if len(f) < 10 || !strings.HasSuffix(f[1], wl) || !strings.HasSuffix(f[2], wr) { + continue + } + q := strings.SplitN(f[4], ":", 2) + tx, rx := int64(-1), int64(-1) + if len(q) == 2 { + tx, _ = strconv.ParseInt(q[0], 16, 64) + rx, _ = strconv.ParseInt(q[1], 16, 64) + } + ino, _ := strconv.ParseUint(f[9], 10, 64) + return fmt.Sprintf("st=%s tx=%d rx=%d tm=%s retr=%s", f[3], tx, rx, f[5], f[6]), ino + } + return "absent", 0 +} + +// wsoFDByInode finds this process's descriptor for socket inode ino. +func wsoFDByInode(ino uint64) int { + if ino == 0 { + return -1 + } + want := "socket:[" + strconv.FormatUint(ino, 10) + "]" + ents, err := os.ReadDir("/proc/self/fd") + if err != nil { + return -1 + } + for _, e := range ents { + if l, err := os.Readlink("/proc/self/fd/" + e.Name()); err == nil && l == want { + fd, _ := strconv.Atoi(e.Name()) + return fd + } + } + return -1 +} + +// srvSockFD returns the engine's descriptor for this connection's server +// socket. It is found once through /proc and re-checked with fstat before +// every use, so a descriptor the engine closed and the kernel handed to +// another socket is not read (short of a reuse between the fstat and the +// read, which would only misreport a diagnostic). Once that check fails the +// socket is gone for good and is not searched for again (watch asks every +// wsoWatchEvery). Only read-only calls (fstat, getsockopt TCP_INFO, ioctl +// SIOCINQ/SIOCOUTQ) are made on it. +func (g *wsoRig) srvSockFD(s *wsoSrv) int { + srvPort := int(g.srvPort.Load()) + if s.port <= 0 || srvPort <= 0 { + return -1 + } + s.mu.Lock() + defer s.mu.Unlock() + if s.gone { + return -1 + } + if s.fd >= 0 { + var st unix.Stat_t + if unix.Fstat(s.fd, &st) == nil && st.Ino == s.ino { + return s.fd + } + s.fd, s.gone = -1, true + return -1 + } + _, ino := wsoProcTCP(srvPort, s.port) + if fd := wsoFDByInode(ino); fd >= 0 { + var st unix.Stat_t + if unix.Fstat(fd, &st) == nil && st.Ino == ino { + s.fd, s.ino = fd, ino + return fd + } + } + return -1 +} + +func wsoSockInfoFD(fd int) (ti *unix.TCPInfo, inq, outq int, ok bool) { + inq, outq = -1, -1 + v, err := unix.GetsockoptTCPInfo(fd, unix.SOL_TCP, unix.TCP_INFO) + if err != nil { + return nil, inq, outq, false + } + if n, err := unix.IoctlGetInt(fd, unix.SIOCINQ); err == nil { + inq = n + } + if n, err := unix.IoctlGetInt(fd, unix.SIOCOUTQ); err == nil { + outq = n + } + return v, inq, outq, true +} + +func (g *wsoRig) srvSockInfo(s *wsoSrv) (*unix.TCPInfo, int, int, bool) { + fd := g.srvSockFD(s) + if fd < 0 { + return nil, -1, -1, false + } + return wsoSockInfoFD(fd) +} + +// wsoCliSockInfo reads the client socket's TCP_INFO, receive queue (SIOCINQ) +// and send queue (SIOCOUTQ: bytes the server's kernel has not ACKed). +func wsoCliSockInfo(c net.Conn) (ti *unix.TCPInfo, inq, outq int) { + inq, outq = -1, -1 + tc, ok := c.(*net.TCPConn) + if !ok { + return nil, inq, outq + } + sc, err := tc.SyscallConn() + if err != nil { + return nil, inq, outq + } + _ = sc.Control(func(fd uintptr) { + ti, inq, outq, _ = wsoSockInfoFD(int(fd)) + }) + return ti, inq, outq +} + +// wsoDrainNow empties the client's receive queue without blocking. It +// returns the bytes read and how the read ended: nil when the queue was +// empty, io.EOF when the server had closed its side, or the socket's error. +// A raw read clears the socket's pending error (a reset reads as +// ECONNRESET once and then no more), so it is returned here, where the +// caller can still name it. +func wsoDrainNow(c net.Conn, buf []byte) (int64, error) { + tc, ok := c.(*net.TCPConn) + if !ok { + return 0, nil + } + sc, err := tc.SyscallConn() + if err != nil { + return 0, err + } + var total int64 + var end error + if err := sc.Read(func(fd uintptr) bool { + for { + n, e := unix.Read(int(fd), buf) + switch { + case n > 0: + total += int64(n) + continue + case e == unix.EINTR: + continue + case e == nil: + end = io.EOF + case e != unix.EAGAIN: + end = os.NewSyscallError("read", e) + } + return true // empty, EOF or an error: stop, never wait + } + }); err != nil && end == nil { + end = err + } + return total, end +} + +func wsoTCPShort(ti *unix.TCPInfo) string { + if ti == nil { + return "tcpinfo=nil" + } + return fmt.Sprintf("wnd=%d/%d unacked=%d notsent=%d backoff=%d probes=%d retr=%d ackAgo=%dms", + ti.Snd_wnd, ti.Rcv_wnd, ti.Unacked, ti.Notsent_bytes, ti.Backoff, ti.Probes, ti.Total_retrans, ti.Last_ack_recv) +} + +func wsoTCPFull(ti *unix.TCPInfo) string { + if ti == nil { + return "tcpinfo=nil" + } + return fmt.Sprintf("state=%d snd_wnd=%d rcv_wnd=%d rcv_space=%d rcv_ssthresh=%d snd_mss=%d rcv_mss=%d "+ + "unacked=%d notsent=%d probes=%d backoff=%d rto_ms=%d retrans=%d total_retrans=%d "+ + "bytes_sent=%d bytes_acked=%d bytes_received=%d data_segs_out=%d data_segs_in=%d "+ + "last_data_sent_ms=%d last_data_recv_ms=%d last_ack_recv_ms=%d rwnd_limited_ms=%d sndbuf_limited_ms=%d", + ti.State, ti.Snd_wnd, ti.Rcv_wnd, ti.Rcv_space, ti.Rcv_ssthresh, ti.Snd_mss, ti.Rcv_mss, + ti.Unacked, ti.Notsent_bytes, ti.Probes, ti.Backoff, ti.Rto/1000, ti.Retrans, ti.Total_retrans, + ti.Bytes_sent, ti.Bytes_acked, ti.Bytes_received, ti.Data_segs_out, ti.Data_segs_in, + ti.Last_data_sent, ti.Last_data_recv, ti.Last_ack_recv, ti.Rwnd_limited/1000, ti.Sndbuf_limited/1000) +} + +// ------------------------------------------------------------ server snapshot + +// wsoSrvView is one reading of a connection's server side. +type wsoSrvView struct { + found bool + phase int32 + frames int64 + echoAgo int64 + depth int + spill int64 + paused, closed bool + pauses, resumes int64 + pauseAgo, resumeAgo int64 + exit string + sockOK bool + ti *unix.TCPInfo + inq, outq int + rawRead int64 // bytes the engine has read from the socket, the upgrade request included + engRead, pendingWrite int64 +} + +func (g *wsoRig) srvView(port int) wsoSrvView { + x, ok := g.srvs.Load(port) + if !ok { + return wsoSrvView{} + } + return g.srvViewOf(x.(*wsoSrv)) +} + +func (g *wsoRig) srvViewOf(s *wsoSrv) wsoSrvView { + var v wsoSrvView + now := g.since() + v.found = true + v.phase = s.phase.Load() + v.frames = s.frames.Load() + v.echoAgo = -1 + if e := s.echoNs.Load(); e > 0 { + v.echoAgo = now - e + } + v.pauses, v.resumes = s.pauses.Load(), s.resumes.Load() + v.pauseAgo, v.resumeAgo = -1, -1 + if p := s.pauseNs.Load(); p > 0 { + v.pauseAgo = now - p + } + if p := s.resumeNs.Load(); p > 0 { + v.resumeAgo = now - p + } + v.depth, v.spill = -1, -1 + if r := s.r; r != nil { + v.depth = len(r.ch) + v.spill = r.spillLen.Load() + r.pausedMu.Lock() + v.paused = r.pausedState + r.pausedMu.Unlock() + v.closed = r.closed.Load() + } + wrote := s.wroteB.Load() + v.ti, v.inq, v.outq, v.sockOK = g.srvSockInfo(s) + s.mu.Lock() + v.exit = s.exit + s.mu.Unlock() + v.rawRead, v.engRead, v.pendingWrite = -1, -1, -1 + if !v.sockOK { + return v + } + // The server socket counts from its first payload byte (a passive open + // starts snd_una and rcv_nxt past the SYN), so bytes received minus + // unread is everything the engine has read, and bytes ACKed plus the + // send queue is everything it has handed to the kernel. Net of the + // upgrade request and the 101, both are the handler's frames. + v.rawRead = int64(v.ti.Bytes_received) - int64(v.inq) + if x, ok := g.clisBy.Load(s.port); ok { + cl := x.(*wsoCli) + if reqB, respB := cl.hsOut.Load(), cl.hsIn.Load(); reqB >= 0 && respB >= 0 { + v.engRead = v.rawRead - reqB + v.pendingWrite = wrote + respB - (int64(v.ti.Bytes_acked) + int64(v.outq)) + } + } + return v +} + +func wsoAgo(ns int64) string { + if ns < 0 { + return "never" + } + return wsoSec(ns) +} + +func (v wsoSrvView) String() string { + if !v.found { + return "srv{unjoined}" + } + ph := map[int32]string{1: "read", 2: "write", 3: "exited"}[v.phase] + s := fmt.Sprintf("srv{%s frames=%d echoAgo=%s depth=%d spill=%d paused=%t closed=%t pauses=%d resumes=%d pauseAgo=%s resumeAgo=%s", + ph, v.frames, wsoAgo(v.echoAgo), v.depth, v.spill, v.paused, v.closed, v.pauses, v.resumes, wsoAgo(v.pauseAgo), wsoAgo(v.resumeAgo)) + if v.exit != "" { + s += " exit=" + strconv.Quote(v.exit) + } + s += "}" + if !v.sockOK { + return s + " sock{gone}" + } + return s + fmt.Sprintf(" sock{inq=%d outq=%d %s} eng{read=%d pendingWrite=%d}", + v.inq, v.outq, wsoTCPShort(v.ti), v.engRead, v.pendingWrite) +} + +// wsoShape names the state a give-up was in, from both ends. It is a hint for +// the reader; the timeline and the raw state are printed with it. +func wsoShape(cli *unix.TCPInfo, cliInq int, v wsoSrvView) string { + switch { + case !v.found: + return "unjoined" + case v.phase == 3: + return "handler-exited" + case v.paused && v.depth == 0 && v.spill > 0: + return "WEDGE: engine paused, channel empty, chunks in the spill (celeris#705 shape)" + case v.paused && v.depth == 0: + return "WEDGE: engine paused with nothing buffered (celeris#672 shape)" + case v.sockOK && v.inq > 0 && !v.paused && v.phase == 1 && v.depth == 0: + return "WEDGE: socket readable, engine not paused, nothing delivered (recv not armed or not completing)" + case v.sockOK && v.pendingWrite > 0 && v.outq == 0: + return "WEDGE: handler bytes the engine never handed to the kernel (send not submitted or not completing)" + case cli != nil && cli.Snd_wnd == 0 && cliInq > 0 && v.sockOK && v.inq == 0 && v.depth == 0 && !v.paused: + return "lost window update: the client's full receive queue discards the server's ACKs (Linux >= 6.17 tcp_sequence)" + case v.paused || v.depth > 0: + return "server behind: handler has buffered input to echo" + } + return "unclassified" +} + +// ------------------------------------------------------------- client side + +// wsoCli is one client connection's record: milestones, the first timeline +// samples and a ring of the most recent ones. +type wsoCli struct { + g *wsoRig + port int + c net.Conn + // hsOut and hsIn: the upgrade request's and the 101 response's sizes, + // read from the client's socket once the handshake is done; -1 before. + hsOut, hsIn atomic.Int64 + mu sync.Mutex + marks []string + first []string + last [wsoLastSamples]string + nSamp int + closeNs int64 // rig time the client closed its socket; 0 while open + gaveUp bool + close wsoClose + fc, cw wsoWrite +} + +// client registers a dialed connection. Close it with closeClient, which +// records when the client let go (see serverErrs). +func (g *wsoRig) client(c net.Conn) *wsoCli { + cl := &wsoCli{g: g, c: c} + cl.hsOut.Store(-1) + cl.hsIn.Store(-1) + if a, ok := c.LocalAddr().(*net.TCPAddr); ok { + cl.port = a.Port + g.clisBy.Store(a.Port, cl) + } + g.mu.Lock() + g.clis = append(g.clis, cl) + g.mu.Unlock() + return cl +} + +// handshake upgrades the connection with the join key in its path, and then +// records the upgrade's size in each direction from the client's socket: the +// server sends nothing past its 101 until frames arrive, and it ACKed the +// whole request with that 101, so here bytes received are the 101 and bytes +// ACKed are the request. +func (cl *wsoCli) handshake(c net.Conn, hostPort string) error { + if err := wsHandshakePath(c, hostPort, cl.path()); err != nil { + return err + } + if ti, _, _ := wsoCliSockInfo(c); ti != nil { + cl.hsOut.Store(int64(ti.Bytes_acked)) + cl.hsIn.Store(int64(ti.Bytes_received)) + } + return nil +} + +func (g *wsoRig) closeClient(cl *wsoCli, c net.Conn) { + cl.mu.Lock() + cl.closeNs = g.since() + cl.mu.Unlock() + _ = c.Close() +} + +// path is the upgrade path that carries the join key. +func (cl *wsoCli) path() string { return "/ws?cid=" + strconv.Itoa(cl.port) } + +func (cl *wsoCli) mark(format string, args ...any) { + m := "+" + wsoSec(cl.g.since()) + " " + fmt.Sprintf(format, args...) + cl.mu.Lock() + cl.marks = append(cl.marks, m) + cl.mu.Unlock() +} + +// sample appends one timeline sample: both ends of the connection, now. +func (cl *wsoCli) sample(c net.Conn, label string) { + ti, inq, outq := wsoCliSockInfo(c) + line := fmt.Sprintf("+%s %s cli{inq=%d outq=%d %s} %s", wsoSec(cl.g.since()), label, inq, outq, wsoTCPShort(ti), cl.g.srvView(cl.port)) + cl.mu.Lock() + if len(cl.first) < wsoFirstSamples { + cl.first = append(cl.first, line) + } else { + cl.last[(cl.nSamp-wsoFirstSamples)%wsoLastSamples] = line + } + cl.nSamp++ + cl.mu.Unlock() +} + +// wsoWrite is how one bounded client write went. +type wsoWrite struct { + site string + ok bool + need int + left int + err error + reason string // "" (completed), "noprog", "cap", "budget", "rst", "eof", "err" + dur time.Duration + slices int + progress int + maxNoProg time.Duration + drained int64 + outq0 int + outq1 int +} + +func (w wsoWrite) String() string { + return fmt.Sprintf("%s: ok=%t left=%d/%d reason=%q err=%q dur=%s slices=%d progress=%d maxNoProg=%s drained=%d outq %d->%d", + w.site, w.ok, w.left, w.need, w.reason, wsoErrStr(w.err), w.dur.Round(time.Millisecond), w.slices, w.progress, + w.maxNoProg.Round(time.Millisecond), w.drained, w.outq0, w.outq1) +} + +// write writes every byte of buf, waiting on progress: 1 s slices; after +// each slice that times out, a timeline sample and a non-blocking drain of +// the client's receive queue. Progress is a byte written or a shrink of the +// client's send queue (the server's kernel ACKed: the engine is reading). +// It gives up after wsoWriteIdle without progress, at wsoWaitCap in total, +// at the rig's budget, on any error other than the slice's own deadline, or +// when a drain finds the connection reset ("rst") or closed ("eof"). +func (cl *wsoCli) write(c net.Conn, site string, buf, dbuf []byte) wsoWrite { + w := wsoWrite{site: site, need: len(buf)} + // The drain below goes through the conn's poller, which refuses at once + // under an expired read deadline, and the caller's read loop leaves one + // behind: clear it, or every drain is a no-op. + _ = c.SetReadDeadline(time.Time{}) + start := time.Now() + lastProg := start + _, _, w.outq0 = wsoCliSockInfo(c) + lastQ := w.outq0 + for len(buf) > 0 { + _ = c.SetWriteDeadline(time.Now().Add(wsoSlice)) + n, err := c.Write(buf) + buf = buf[n:] + w.slices++ + if len(buf) == 0 { + break + } + if err == nil { + continue + } + if !errors.Is(err, os.ErrDeadlineExceeded) { + w.err, w.reason = err, "err" + if errors.Is(err, syscall.ECONNRESET) { + w.reason = "rst" + } + break + } + now := time.Now() + _, _, q := wsoCliSockInfo(c) + if n > 0 || (q >= 0 && lastQ >= 0 && q < lastQ) { + if gap := now.Sub(lastProg); gap > w.maxNoProg { + w.maxNoProg = gap + } + lastProg = now + w.progress++ + } + lastQ = q + cl.sample(c, site+" slice "+strconv.Itoa(w.slices)) + if now.Sub(lastProg) >= wsoWriteIdle { + w.err, w.reason = err, "noprog" + break + } + if now.Sub(start) >= wsoWaitCap { + w.err, w.reason = err, "cap" + break + } + if cl.g.overBudget(now) { + w.err, w.reason = err, "budget" + break + } + n2, derr := wsoDrainNow(c, dbuf) + w.drained += n2 + if derr != nil { + w.err, w.reason = derr, "err" + switch { + case errors.Is(derr, syscall.ECONNRESET): + w.reason = "rst" + case errors.Is(derr, io.EOF): + w.reason = "eof" + } + break + } + } + _ = c.SetWriteDeadline(time.Time{}) + w.left = len(buf) + w.ok = w.left == 0 + w.dur = time.Since(start) + if gap := time.Since(lastProg); !w.ok && gap > w.maxNoProg { + w.maxNoProg = gap + } + _, _, w.outq1 = wsoCliSockInfo(c) + cl.mu.Lock() + if site == "fc" { + cl.fc = w + } else { + cl.cw = w + } + cl.mu.Unlock() + return w +} + +// wsoClose is how the wait for the server's close went. +type wsoClose struct { + outcome string // "eof", "rst", "idle", "cap", "budget", "err" + err error + dur time.Duration // Close sent -> outcome + bytes int64 // received during the wait + lastByte time.Duration // Close sent -> last byte received + maxGap time.Duration // the longest stretch with no byte, up to the outcome +} + +func (r wsoClose) String() string { + return fmt.Sprintf("close wait: outcome=%s err=%q dur=%s bytes=%d lastByte=%s maxGap=%s", + r.outcome, wsoErrStr(r.err), r.dur.Round(time.Millisecond), r.bytes, r.lastByte.Round(time.Millisecond), r.maxGap.Round(time.Millisecond)) +} + +// closeWait reads until the server closes. The give-up is re-armed by every +// byte received: it comes after wsoCloseIdle without one, at wsoWaitCap, or +// at the rig's budget. While it waits in silence it records a timeline +// sample every wsoSlice. +func (cl *wsoCli) closeWait(c net.Conn, buf []byte) wsoClose { + var r wsoClose + start := time.Now() + last := start + for { + dl := time.Now().Add(wsoSlice) + if lim := last.Add(wsoCloseIdle); dl.After(lim) { + dl = lim + } + if lim := start.Add(wsoWaitCap); dl.After(lim) { + dl = lim + } + if lim := cl.g.end; !lim.IsZero() && dl.After(lim) { + dl = lim + } + _ = c.SetReadDeadline(dl) + n, err := c.Read(buf) + now := time.Now() + if n > 0 { + if gap := now.Sub(last); gap > r.maxGap { + r.maxGap = gap + } + last = now + r.bytes += int64(n) + } + if err == nil { + continue + } + switch { + case errors.Is(err, io.EOF): + r.outcome = "eof" + case errors.Is(err, syscall.ECONNRESET): + r.outcome = "rst" + case errors.Is(err, os.ErrDeadlineExceeded): + switch { + case now.Sub(start) >= wsoWaitCap: + r.outcome = "cap" + case now.Sub(last) >= wsoCloseIdle: + r.outcome = "idle" + case cl.g.overBudget(now): + r.outcome = "budget" + default: + cl.sample(c, "close wait") + continue + } + default: + r.outcome = "err" + } + r.err = err + if gap := now.Sub(last); gap > r.maxGap { + r.maxGap = gap + } + r.dur = now.Sub(start) + r.lastByte = last.Sub(start) + cl.mu.Lock() + cl.close = r + cl.mu.Unlock() + return r + } +} + +// giveUp prints one WSO-GIVEUP record: the connection's milestones, its +// timeline, and the full state of both ends now, with the shape it is in. +// One t.Logf call, so records from concurrent clients never interleave, and +// printed at once, so it survives a test binary that later times out. +func (cl *wsoCli) giveUp(t *testing.T, c net.Conn, kind, what string) { + cl.mu.Lock() + cl.gaveUp = true + cl.mu.Unlock() + ti, inq, outq := wsoCliSockInfo(c) + v := cl.g.srvView(cl.port) + srvPort := int(cl.g.srvPort.Load()) + srvRow, _ := wsoProcTCP(srvPort, cl.port) + cliRow, _ := wsoProcTCP(cl.port, srvPort) + var b strings.Builder + fmt.Fprintf(&b, "WSO-GIVEUP %s conn=127.0.0.1:%d %s\n", kind, cl.port, what) + fmt.Fprintf(&b, " shape: %s\n", wsoShape(ti, inq, v)) + cl.mu.Lock() + fmt.Fprintf(&b, " milestones: %s\n", strings.Join(cl.marks, "; ")) + n := cl.nSamp + fmt.Fprintf(&b, " timeline (%d samples, one per second spent waiting; the first %d and the last %d are kept):\n", + n, wsoFirstSamples, wsoLastSamples) + for _, l := range cl.first { + fmt.Fprintf(&b, " %s\n", l) + } + from := wsoFirstSamples + if n-wsoLastSamples > from { + from = n - wsoLastSamples + fmt.Fprintf(&b, " ... %d samples not kept ...\n", from-wsoFirstSamples) + } + for i := from; i < n; i++ { + fmt.Fprintf(&b, " %s\n", cl.last[(i-wsoFirstSamples)%wsoLastSamples]) + } + cl.mu.Unlock() + fmt.Fprintf(&b, " now: cli{inq=%d outq=%d %s} /proc{%s}\n", inq, outq, wsoTCPFull(ti), cliRow) + fmt.Fprintf(&b, " now: %s /proc{%s}", v, srvRow) + if v.ti != nil { + fmt.Fprintf(&b, "\n now: srv tcpinfo{%s}", wsoTCPFull(v.ti)) + } + t.Log(b.String()) +} + +// --------------------------------------------------------------- the verdicts + +// wsoJudged is one server-side error, rendered with its timing against its +// own client's close. +type wsoJudged struct { + kind string // "write" or "read" + line string +} + +// serverErrs splits every server-side error into the judged and the +// excused. An error is excused only when its own client gave up on the +// connection or was reset (giveUp), which fails the test on its own, and the +// error came after that client closed its socket: then it is the client's +// teardown echoing back (a socket closed with unread data sends an RST; +// celeris#633: none of 1,165 such errors came before their client's close). +// Every other error is judged, whenever it came: on a connection whose +// client read the server's close (EOF), the client's receive queue was +// empty and its close sent a FIN, so nothing it did can explain a server +// error; and a handler notes an engine error only when its next read or +// write returns, so the time it was noted says nothing about its cause. An +// error whose connection the rig could not join is judged. +func (g *wsoRig) serverErrs() (judged, excused []wsoJudged) { + type end struct { + closeNs int64 + gaveUp bool + } + ends := map[int]end{} + g.mu.Lock() + for _, cl := range g.clis { + cl.mu.Lock() + ends[cl.port] = end{cl.closeNs, cl.gaveUp} + cl.mu.Unlock() + } + hs := append([]*wsoSrv(nil), g.hs...) + g.mu.Unlock() + for _, s := range hs { + s.mu.Lock() + errs := append([]wsoSrvErr(nil), s.errs...) + s.mu.Unlock() + for _, e := range errs { + en, ok := ends[s.port] + line := fmt.Sprintf("conn 127.0.0.1:%d %s error %q at %s", s.port, e.kind, wsoErrStr(e.err), wsoSec(e.ns)) + switch { + case ok && en.gaveUp && en.closeNs > 0 && e.ns >= en.closeNs: + excused = append(excused, wsoJudged{e.kind, line + fmt.Sprintf(", %s after its client gave up and closed", wsoSec(e.ns-en.closeNs))}) + continue + case !ok: + line += ", its connection was never joined" + case en.closeNs == 0: + line += ", its client had not closed" + case e.ns < en.closeNs: + line += fmt.Sprintf(", %s BEFORE its client closed", wsoSec(en.closeNs-e.ns)) + default: + line += fmt.Sprintf(", %s after its client closed on the server's close (EOF)", wsoSec(e.ns-en.closeNs)) + } + judged = append(judged, wsoJudged{e.kind, line}) + } + } + sort.Slice(judged, func(i, j int) bool { return judged[i].line < judged[j].line }) + sort.Slice(excused, func(i, j int) bool { return excused[i].line < excused[j].line }) + return judged, excused +} + +// ------------------------------------------------------------------ the watch + +// wsoStall is one stretch in which a connection's engine read nothing from +// a server socket holding unread bytes while nothing asked it to stop. +type wsoStall struct { + port int + fromNs, toNs int64 // rig time of the stretch's first and last sample + samples int + onset, last string // the server side at those samples + onsetC, lastC string // the client side at those samples + marks string // the client's milestones at the last sample + open bool // still going when the watch stopped +} + +func (st *wsoStall) String() string { + still := "" + if st.open { + still = ", still going when the watch stopped" + } + return fmt.Sprintf("WSO-STALL conn=127.0.0.1:%d the engine read nothing for %s (%d samples, +%s to +%s%s) "+ + "from a server socket holding unread bytes, while the chanReader was neither paused nor holding "+ + "anything and the handler had not exited\n onset: %s\n %s\n last: %s\n %s\n client milestones: %s", + st.port, wsoSec(st.toNs-st.fromNs), st.samples, wsoSec(st.fromNs), wsoSec(st.toNs), still, + st.onset, st.onsetC, st.last, st.lastC, st.marks) +} + +// longestStretch renders the longest stretch watch saw, for the verdict +// line's neighbour: how close a run came to wsoRecvStall. +func (g *wsoRig) longestStretch() string { + g.mu.Lock() + defer g.mu.Unlock() + if g.longest.samples == 0 { + return "longest unread stretch: none" + } + st := g.longest + return fmt.Sprintf("longest unread stretch (limit %v): %s", wsoRecvStall, strings.TrimPrefix(st.String(), "WSO-STALL ")) +} + +// wsoStallable reports whether the engine should be reading a connection +// now: its server socket holds unread bytes, and nothing asked the engine to +// stop (the chanReader is not paused, is not closed and holds nothing), and +// the handler has not exited. +func wsoStallable(v wsoSrvView) bool { + return v.found && v.sockOK && v.inq > 0 && !v.paused && !v.closed && v.depth == 0 && v.spill == 0 && v.phase != 3 +} + +// cliNow renders a joined connection's client side now, for a WSO-STALL. +func (g *wsoRig) cliNow(port int) string { + x, ok := g.clisBy.Load(port) + if !ok { + return "cli{unjoined}" + } + cl := x.(*wsoCli) + ti, inq, outq := wsoCliSockInfo(cl.c) + cl.mu.Lock() + last := "none" + if n := len(cl.marks); n > 0 { + last = strconv.Quote(cl.marks[n-1]) + } + cl.mu.Unlock() + return fmt.Sprintf("cli{inq=%d outq=%d %s lastMilestone=%s}", inq, outq, wsoTCPShort(ti), last) +} + +func (g *wsoRig) cliMarks(port int) string { + x, ok := g.clisBy.Load(port) + if !ok { + return "" + } + cl := x.(*wsoCli) + cl.mu.Lock() + defer cl.mu.Unlock() + return strings.Join(cl.marks, "; ") +} + +// watch samples the server side of every joined connection every +// wsoWatchEvery until stop is called. A stretch is a run of samples in which +// the connection is wsoStallable and the engine read nothing (bytes received +// minus unread did not move) and the chanReader applied no pause or resume. +// Every stretch that reaches wsoRecvStall is recorded, from its first sample +// to its last. stop returns those. The longest stretch seen on any +// connection is kept too, whether or not it reached the limit +// (longestStretch), so a healthy run shows how close it came. +// +// The samples between a stretch's first and last are not what proves it: +// bytes read, pauses and resumes only ever grow, a chanReader holding +// nothing that is handed nothing stays empty, and an unread socket that is +// not read stays unread, so equal readings at both ends mean the engine read +// nothing in between while nothing asked it to stop. What the samples guard +// against is the process itself not running: if the watcher waits more than +// wsoWatchGap past a due tick, the process was starved of CPU, the engine +// may have been too, and every open stretch starts over (counted as +// watchResets on the waits line). +// +// Start it once the server listens and stop it once the clients are done: +// the verdict is the test's, on the test goroutine. +func (g *wsoRig) watch() (stop func() []*wsoStall) { + type episode struct { + fromNs, raw, pauses, resumes int64 + samples int + onset, onsetC string + rec *wsoStall + } + done, fin := make(chan struct{}), make(chan struct{}) + var stalls []*wsoStall + var longest int64 + go func() { + defer close(fin) + eps := map[*wsoSrv]*episode{} + tick := time.NewTicker(wsoWatchEvery) + defer tick.Stop() + idleFrom := g.since() // when the watcher last started waiting for a tick + for { + select { + case <-done: + for _, ep := range eps { + if ep.rec != nil { + ep.rec.open = true + } + } + return + case <-tick.C: + } + // A tick is due at most wsoWatchEvery after the watcher starts + // waiting (sooner when a pass overran). Waiting longer than that + // plus wsoWatchGap means this goroutine was runnable and not run. + if lag := time.Duration(g.since()-idleFrom) - wsoWatchEvery; lag > 0 { + g.mu.Lock() + g.watchMaxLag = max(g.watchMaxLag, lag) + if lag > wsoWatchGap { + g.watchResets++ + clear(eps) + } + g.mu.Unlock() + } + g.mu.Lock() + hs := append([]*wsoSrv(nil), g.hs...) + g.mu.Unlock() + for _, s := range hs { + if s.port <= 0 || s.phase.Load() == 3 { + delete(eps, s) + continue + } + v := g.srvViewOf(s) + now := g.since() + if !wsoStallable(v) { + delete(eps, s) + continue + } + ep := eps[s] + if ep == nil || v.rawRead != ep.raw || v.pauses != ep.pauses || v.resumes != ep.resumes { + eps[s] = &episode{fromNs: now, raw: v.rawRead, pauses: v.pauses, resumes: v.resumes, + samples: 1, onset: v.String(), onsetC: g.cliNow(s.port)} + continue + } + ep.samples++ + d := now - ep.fromNs + if d > longest { + longest = d + st := wsoStall{port: s.port, fromNs: ep.fromNs, toNs: now, samples: ep.samples, + onset: ep.onset, onsetC: ep.onsetC, last: v.String(), lastC: g.cliNow(s.port), marks: g.cliMarks(s.port)} + g.mu.Lock() + g.longest = st + g.mu.Unlock() + } + if d < int64(wsoRecvStall) { + continue + } + if ep.rec == nil { + ep.rec = &wsoStall{port: s.port, fromNs: ep.fromNs, onset: ep.onset, onsetC: ep.onsetC} + stalls = append(stalls, ep.rec) + } + ep.rec.toNs, ep.rec.samples = now, ep.samples + ep.rec.last, ep.rec.lastC, ep.rec.marks = v.String(), g.cliNow(s.port), g.cliMarks(s.port) + } + idleFrom = g.since() + } + }() + return func() []*wsoStall { + close(done) + <-fin + g.mu.Lock() + g.stallMax = time.Duration(longest) + g.mu.Unlock() + return stalls + } +} + +// wsoAssertNoStalls prints every WSO-STALL record and fails the test if +// there is one (see watch). +func wsoAssertNoStalls(t *testing.T, name string, stalls []*wsoStall) { + t.Helper() + for _, st := range stalls { + t.Log(st.String()) + } + if len(stalls) != 0 { + t.Errorf("%s: %d connection(s) left unread for %v or longer: the server socket held bytes, nothing asked the "+ + "engine to stop reading (the chanReader was neither paused nor holding anything) and the engine read "+ + "nothing. That is the celeris#607 class, a recv held back behind a send the client's closed window "+ + "blocks; the client's own reads get it through, so no give-up shows it. Each WSO-STALL record above has "+ + "the stretch's state at both ends", name, len(stalls), wsoRecvStall) + } +} + +// wsoAssertNoLinkedRecv judges io_uring's linked-recv witness. These servers +// serve nothing but WebSocket connections, and each is detached by its +// upgrade (Context.Detach publishes it before the upgrade handler returns), +// before the engine flushes anything on it. So no recv may ever be chained +// behind a send (IOSQE_IO_LINK): on a detached connection that chain is +// celeris#607 itself (TestFlushSendLinkNeverChainsOnDetachedConn guards the +// function; this guards every path that reaches it). +func wsoAssertNoLinkedRecv(t *testing.T, name string, arms, blockedMaxNs uint64) { + t.Helper() + if arms != 0 { + t.Errorf("%s: the engine chained a recv behind a send %d time(s) (the longest waited %d ms) on a server "+ + "whose every connection is a detached WebSocket: celeris#607's mechanism", name, arms, blockedMaxNs/1e6) + } +} + +// wsoWaitStats summarises every connection's waits for the verdict line. +type wsoWaitStats struct { + conns int + fcMax, cwMax, closeMax time.Duration + closeP50, closeP99 time.Duration + closeSlow int // close waits over wsoCloseSlow: the old oracle failed these + closeSilentSlow int // close waits with a silent stretch over wsoCloseSlow + maxCloseGap, maxNoProg time.Duration + fcSlices, drainedBytes int64 + recvStallMax time.Duration // the longest stretch watch saw (see watch) + watchResets int + watchMaxLag time.Duration +} + +func (g *wsoRig) waitStats() wsoWaitStats { + var st wsoWaitStats + var closes []time.Duration + g.mu.Lock() + defer g.mu.Unlock() + st.recvStallMax, st.watchResets, st.watchMaxLag = g.stallMax, g.watchResets, g.watchMaxLag + for _, cl := range g.clis { + cl.mu.Lock() + st.conns++ + st.fcMax = max(st.fcMax, cl.fc.dur) + st.cwMax = max(st.cwMax, cl.cw.dur) + st.maxNoProg = max(st.maxNoProg, cl.fc.maxNoProg, cl.cw.maxNoProg) + st.fcSlices += int64(cl.fc.slices) + st.drainedBytes += cl.fc.drained + cl.cw.drained + if cl.close.outcome != "" { + closes = append(closes, cl.close.dur) + st.closeMax = max(st.closeMax, cl.close.dur) + st.maxCloseGap = max(st.maxCloseGap, cl.close.maxGap) + if cl.close.dur > wsoCloseSlow { + st.closeSlow++ + } + if cl.close.maxGap > wsoCloseSlow { + st.closeSilentSlow++ + } + } + cl.mu.Unlock() + } + if len(closes) > 0 { + sort.Slice(closes, func(i, j int) bool { return closes[i] < closes[j] }) + st.closeP50 = closes[len(closes)/2] + st.closeP99 = closes[(len(closes)*99)/100] + } + return st +} + +func (st wsoWaitStats) String() string { + r := func(d time.Duration) string { return d.Round(time.Millisecond).String() } + return fmt.Sprintf("waits: conns=%d fcMax=%s cwMax=%s maxNoProg=%s fcSlices=%d drainedWhileWaiting=%d closeP50=%s closeP99=%s closeMax=%s closeOver%s=%d closeSilentOver%s=%d maxCloseGap=%s recvStallMax=%s watchResets=%d watchMaxLag=%s", + st.conns, r(st.fcMax), r(st.cwMax), r(st.maxNoProg), st.fcSlices, st.drainedBytes, r(st.closeP50), r(st.closeP99), r(st.closeMax), + wsoCloseSlow, st.closeSlow, wsoCloseSlow, st.closeSilentSlow, r(st.maxCloseGap), r(st.recvStallMax), st.watchResets, r(st.watchMaxLag)) +} diff --git a/middleware/websocket/inbound_sequence_linux_test.go b/middleware/websocket/inbound_sequence_linux_test.go index 9f21f427..4b6142dc 100644 --- a/middleware/websocket/inbound_sequence_linux_test.go +++ b/middleware/websocket/inbound_sequence_linux_test.go @@ -9,6 +9,7 @@ import ( "fmt" "io" "net" + "os" "strings" "sync" "sync/atomic" @@ -21,21 +22,33 @@ import ( "github.com/goceleris/celeris/probe" ) +// closeHandshakeSlow is the point past which the client's wait for the +// server's close is worth reporting. See the client loop for the measured +// distribution behind it (celeris#566). +const closeHandshakeSlow = 10 * time.Second + // TestBackpressureInboundSequenceIntegrity (celeris#484 oracle): every inbound // frame carries a strictly increasing 64-bit sequence number and connection index, // so a lost, corrupted, or reordered frame is detected immediately. The client // floods without reading (forcing repeated pause/resume) and the server handler // validates sequence continuity per connection, payload content integrity, and // verifies the tail sequence number transmitted in the Close frame. -// closeHandshakeBudget is how long a client waits for the server to close -// after sending its Close frame, and closeHandshakeSlow is the point past -// which that wait is worth reporting. See the client loop for the measured -// distribution behind these numbers (celeris#566). -const ( - closeHandshakeBudget = 30 * time.Second - closeHandshakeSlow = 10 * time.Second -) - +// +// The frame count measures DELIVERY: a handler whose echo write fails stops +// echoing but keeps reading until the stream ends, so frames the engine did +// deliver are never reported lost because the handler walked away from them +// (celeris#611). The echo failure itself is judged, with its error and its +// connection, when it came before its own client closed. +// +// The client finishes its writes and waits for the close on progress, reading +// while it waits, like TestBackpressurePauseDoesNotCancelInflightSend's +// (backpressure_oracle_linux_test.go, celeris#633), and the rig's watch fails +// a connection the engine stops reading (WSO-STALL, the celeris#607 class). A +// write timeout during the flood is the backpressure the test creates, not a +// failure: it ends that burst, the next one resumes at the same wire +// position, and the connection stays in every verdict (it used to be taken +// out of the frame count). floodDeadlines on the verdict line counts those +// bursts, so a run shows whether it exercised that path at all. func TestBackpressureInboundSequenceIntegrity(t *testing.T) { conns := envInt("WS484_CONNS", 96) // bpBuf is the chanReader backpressure buffer capacity (default 256, matching the @@ -102,7 +115,11 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { detailMu.Unlock() } - var gaps, parseErr, overflowErr, protoErr, framesIn, framesSent, echoErr, slowClose atomic.Int64 + rig := newWSORig(t) + var gaps, parseErr, overflowErr, protoErr, framesIn, framesSent, echoErr, slowClose, readAfterEchoErr atomic.Int64 + // floodDeadlines: bursts a write deadline ended (the burst's + // backpressure); floodDeadlineConns: connections with one. + var floodDeadlines, floodDeadlineConns atomic.Int64 var closedOK, closeTimeout, dialFail, hsFail, clientCloseFail atomic.Int64 // An RST is NOT a clean close (celeris#530). Linux emits one // when a socket is closed with unread data still in its @@ -124,7 +141,11 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { connEchoed := make([]atomic.Int64, conns) connEchoErr := make([]atomic.Int64, conns) connProtoErr := make([]atomic.Int64, conns) + // Frames the handler read after its own echo write failed + // (celeris#611): delivered, and counted in connFramesIn. + connReadAfterEchoErr := make([]atomic.Int64, conns) connSent := make([]atomic.Int64, conns) + connPort := make([]atomic.Int64, conns) // the client's local port: the connection's address clientFailed := make([]atomic.Bool, conns) // Handlers must finish before the engine goes away — see settle(). @@ -136,10 +157,21 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { Handler: func(c *Conn) { handlerWG.Add(1) defer handlerWG.Done() + s := rig.attach(c) + defer s.exited("returned", nil) var myConnIdx = -1 + // echoing goes false when an echo write fails. The + // handler then stops writing and keeps READING to the + // end of the stream, so the frame count measures what + // the engine delivered, not where the handler gave up + // (celeris#611: 99.8-99.9% of the frames it reported + // missing were still queued, unread, at handler exit). + echoing := true for { + s.phase.Store(1) mt, msg, err := c.ReadMessage() if err != nil { + s.exited("read", err) if ce, ok := err.(*CloseError); ok { if len(ce.Text) == 8 && myConnIdx >= 0 { finalCount := int64(binary.BigEndian.Uint64([]byte(ce.Text))) @@ -160,6 +192,7 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { serverRST.Add(1) } else if !isCloseErr(err) { protoErr.Add(1) + s.noteErr("read", err) if myConnIdx >= 0 { connProtoErr[myConnIdx].Add(1) } @@ -207,18 +240,29 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { connFramesIn[cIdx].Add(1) framesIn.Add(1) + if !echoing { + connReadAfterEchoErr[cIdx].Add(1) + readAfterEchoErr.Add(1) + continue + } + s.phase.Store(2) if err := c.WriteMessage(mt, msg); err != nil { // The handler's own echo failure used to be // swallowed here. A client reporting a truncated // stream could not be told apart from a server // that stopped writing, which is the ambiguity - // celeris#562 spent several rounds inside. + // celeris#562 spent several rounds inside. It is + // recorded with its time and judged against its + // client's close; the handler keeps reading. echoErr.Add(1) connEchoErr[cIdx].Store(1) + s.noteErr("write", err) note("conn %d: echo write failed after %d frames: %v", cIdx, connEchoed[cIdx].Load(), err) - return + echoing = false + continue } + s.echoed(len(msg)) connEchoed[cIdx].Add(1) } }, @@ -264,6 +308,8 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { defer settle() hostPort := strings.TrimSuffix(strings.TrimPrefix(addr, "ws://"), "/ws") + rig.setServer(hostPort) + stopWatch := rig.watch() dialer := net.Dialer{Timeout: 3 * time.Second, Control: func(_, _ string, rc syscall.RawConn) error { var serr error _ = rc.Control(func(fd uintptr) { serr = syscall.SetsockoptInt(int(fd), syscall.SOL_SOCKET, syscall.SO_RCVBUF, 32<<10) }) @@ -282,8 +328,10 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { clientFailed[connID].Store(true) return } - defer func() { _ = c.Close() }() - if err := wsHandshake(c, hostPort); err != nil { + cl := rig.client(c) + connPort[connID].Store(int64(cl.port)) + defer rig.closeClient(cl, c) + if err := cl.handshake(c, hostPort); err != nil { hsFail.Add(1) clientFailed[connID].Store(true) return @@ -308,6 +356,7 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { var seq uint64 cur, off := encode(0), 0 + deadlines := 0 for b := 0; b < bursts; b++ { budget := perBurst * (6 + plen) for budget > 0 { @@ -320,17 +369,37 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { cur, off = encode(seq), 0 } if err != nil { + if errors.Is(err, os.ErrDeadlineExceeded) { + // Backpressure, which is what the + // test creates: this burst ends here + // and the next resumes at the same + // wire position. It used to fail the + // connection, taking exactly the + // backpressured connections out of + // the frame-count verdict. + deadlines++ + floodDeadlines.Add(1) + break + } clientFailed[connID].Store(true) clientCloseFail.Add(1) - break + cl.giveUp(t, c, testName, fmt.Sprintf("flood write failed: %v", err)) + return } } time.Sleep(200 * time.Millisecond) } + if deadlines > 0 { + floodDeadlineConns.Add(1) + } + cl.mark("flood end seq=%d off=%d burstsEndedByDeadline=%d", seq, off, deadlines) + cl.sample(c, "flood end") + buf := make([]byte, 64<<10) if off > 0 { - if !writeAll(c, cur[off:], 15*time.Second) { + if w := cl.write(c, "fc", cur[off:], buf); !w.ok { clientFailed[connID].Store(true) clientCloseFail.Add(1) + cl.giveUp(t, c, testName, "gave up finishing its last frame: "+w.String()) return } seq++ @@ -338,7 +407,6 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { connSent[connID].Store(int64(seq)) framesSent.Add(int64(seq)) - buf := make([]byte, 64<<10) for { _ = c.SetReadDeadline(time.Now().Add(1 * time.Second)) if _, err := c.Read(buf); err != nil { @@ -346,11 +414,14 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { } } - if !writeAll(c, maskedCloseFrameWithCount(seq), 10*time.Second) { + cl.mark("drained; writing Close") + if w := cl.write(c, "cw", maskedCloseFrameWithCount(seq), buf); !w.ok { clientFailed[connID].Store(true) clientCloseFail.Add(1) + cl.giveUp(t, c, testName, "gave up writing its Close frame: "+w.String()) return } + cl.mark("Close sent") // The server routinely takes ~12s to close after the // client's Close frame on this workload: 96 connections @@ -363,32 +434,39 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { // its side" for a server that closes two seconds later // (celeris#566). // - // The budget is generous enough that a genuine hang is - // still what fails, and the latency is recorded so a - // regression shows up as a number rather than as a - // boolean that silently starts tripping. - closeSentAt := time.Now() - _ = c.SetReadDeadline(closeSentAt.Add(closeHandshakeBudget)) - for { - if _, err := c.Read(buf); err != nil { - if d := time.Since(closeSentAt); d > closeHandshakeSlow { - slowClose.Add(1) - note("conn %d: server closed %v after the client's Close", - connID, d.Round(time.Millisecond)) - } - if errors.Is(err, syscall.ECONNRESET) { - clientRST.Add(1) - } else if errors.Is(err, io.EOF) { - closedOK.Add(1) - } else { - closeTimeout.Add(1) - } - return - } + // The wait gives up only after wsoCloseIdle with no + // byte from the server (re-armed by every byte), at + // wsoWaitCap, or at the rig's -timeout budget, so a + // genuine hang is still what fails, + // and the latency is recorded so a regression shows up + // as a number rather than as a boolean that silently + // starts tripping. + r := cl.closeWait(c, buf) + if r.dur > closeHandshakeSlow { + slowClose.Add(1) + note("conn %d: server closed %v after the client's Close (%s)", + connID, r.dur.Round(time.Millisecond), r) + } + // A connection that did not close cleanly fails below with + // its record, and leaves the frame count: frames still in + // the client's own send queue when it gave up, or thrown + // away with a reset, were never the engine's to deliver. + switch r.outcome { + case "eof": + closedOK.Add(1) + case "rst": + clientFailed[connID].Store(true) + clientRST.Add(1) + cl.giveUp(t, c, testName, "was RESET instead of closed: "+r.String()) + default: + clientFailed[connID].Store(true) + closeTimeout.Add(1) + cl.giveUp(t, c, testName, "gave up waiting for the server's close: "+r.String()) } }() } wg.Wait() + stalls := stopWatch() if dialFail.Load()+hsFail.Load() > 0 { t.Fatalf("environment: %d dial/handshake failures", dialFail.Load()+hsFail.Load()) @@ -404,8 +482,21 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { // whose summary read parseErr=0 went on to fail the parse-error // assertion. Anyone comparing two builds by grepping this line // is reading stale numbers. - t.Logf("%s: conns=%d framesSent=%d framesIn=%d seqGaps=%d parseErr=%d overflowErr=%d protocolErrors=%d clientCloseFail=%d closedOK=%d clientRST=%d serverRST=%d closeTimeout=%d dialFail=%d hsFail=%d", - testName, conns, framesSent.Load(), framesIn.Load(), gaps.Load(), parseErr.Load(), overflowErr.Load(), protoErr.Load(), clientCloseFail.Load(), closedOK.Load(), clientRST.Load(), serverRST.Load(), closeTimeout.Load(), dialFail.Load(), hsFail.Load()) + judged, excused := rig.serverErrs() + var echoJudged int64 + for _, e := range judged { + if e.kind == "write" { + echoJudged++ + } + } + t.Logf("%s: conns=%d framesSent=%d framesIn=%d seqGaps=%d parseErr=%d overflowErr=%d protocolErrors=%d clientCloseFail=%d closedOK=%d clientRST=%d serverRST=%d closeTimeout=%d dialFail=%d hsFail=%d echoErrJudged=%d readAfterEchoErr=%d serverErrsExcused=%d floodDeadlines=%d floodDeadlineConns=%d recvStalls=%d", + testName, conns, framesSent.Load(), framesIn.Load(), gaps.Load(), parseErr.Load(), overflowErr.Load(), protoErr.Load(), clientCloseFail.Load(), closedOK.Load(), clientRST.Load(), serverRST.Load(), closeTimeout.Load(), dialFail.Load(), hsFail.Load(), + echoJudged, readAfterEchoErr.Load(), len(excused), floodDeadlines.Load(), floodDeadlineConns.Load(), len(stalls)) + t.Logf("%s: %s", testName, rig.waitStats()) + t.Logf("%s: %s", testName, rig.longestStretch()) + for _, e := range excused { + t.Logf("%s: not judged, its client gave up: %s", testName, e.line) + } detailMu.Lock() for _, d := range details { @@ -413,8 +504,8 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { } detailMu.Unlock() - t.Logf("%s: echoWriteErrors=%d slowCloses=%d (>%v, budget %v)", - testName, echoErr.Load(), slowClose.Load(), closeHandshakeSlow, closeHandshakeBudget) + t.Logf("%s: echoWriteErrors=%d slowCloses=%d (>%v; the close wait gives up after %v without a byte, or at %v)", + testName, echoErr.Load(), slowClose.Load(), closeHandshakeSlow, wsoCloseIdle, wsoWaitCap) // Recv-arming witnesses (celeris#586), read AFTER settle() so // every worker has left its loop: the counters are direct @@ -433,6 +524,9 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { t.Logf("%s: RECVARM bp=%d resumeWhileCancelPending=%d resumeWhileRecvInFlight=%d armDeclined=%d doubleArmed=%d cqeUnaccounted=%d parseErr=%d", testName, bpBuf, m.RecvResumeWhileCancelPending, m.RecvResumeWhileRecvInFlight, m.RecvArmDeclined, m.RecvDoubleArmed, m.RecvCQEUnaccounted, parseErr.Load()) + t.Logf("%s: LINKBLOCK arms=%d blockedTotalMs=%d blockedMaxMs=%d", + testName, m.RecvLinkedArms, m.RecvLinkedBlockedNanos/1e6, m.RecvLinkedBlockedMaxNanos/1e6) + wsoAssertNoLinkedRecv(t, testName, m.RecvLinkedArms, m.RecvLinkedBlockedMaxNanos) if m.RecvDoubleArmed != 0 { t.Errorf("%s: RecvDoubleArmed=%d — a second recv SQE was placed on a connection that already had one (celeris#484)", testName, m.RecvDoubleArmed) @@ -444,9 +538,38 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { } for i := range conns { if connEchoErr[i].Load() != 0 { - t.Logf("%s: conn %d echo state: echoed=%d in=%d sent=%d", - testName, i, connEchoed[i].Load(), connFramesIn[i].Load(), connSent[i].Load()) + t.Logf("%s: conn %d (127.0.0.1:%d) echo state: echoed=%d in=%d (read after the echo failed %d) sent=%d", + testName, i, connPort[i].Load(), connEchoed[i].Load(), connFramesIn[i].Load(), connReadAfterEchoErr[i].Load(), connSent[i].Load()) + } + } + // celeris#611: every mismatched connection in 34 failing runs + // was also an echo-write failure, logged and never judged, so + // three investigations went looking at the read path. An echo + // failure is the engine closing or failing a connection it + // should have kept: judged here, with its error and its + // connection's address, unless its client gave up on the + // connection (which fails below) and it came after that + // client closed: then it is that client's teardown (printed + // above). + if len(judged) != 0 { + for _, e := range judged { + t.Logf("%s: %s", testName, e.line) } + t.Errorf("%s: %d handler error(s) (%d echo write, %d read) on connections whose client did not give up, or "+ + "before it did: the engine failed or tore down a connection it should have kept alive; each is listed "+ + "above with its error and address", testName, len(judged), echoJudged, int64(len(judged))-echoJudged) + } + wsoAssertNoStalls(t, testName, stalls) + // A connection that never finished its writes was counted + // and not judged, like the sibling oracle's (celeris#623). + if n := clientCloseFail.Load(); n != 0 { + t.Errorf("%s: %d client(s) never finished writing: a flood write failed, or nothing moved for %v "+ + "(or the wait ran past %v, or %s) finishing the last frame or the Close frame; each WSO-GIVEUP "+ + "record above has the connection's timeline from both ends", testName, n, wsoWriteIdle, wsoWaitCap, rig.budgetNote()) + } + if n := clientRST.Load(); n != 0 { + t.Errorf("%s: %d connection(s) were RESET rather than closed cleanly after the client's Close "+ + "(celeris#530); each WSO-GIVEUP record above has its state", testName, n) } if parseErr.Load() != 0 { t.Errorf("%s: %d frame parse error(s) observed — frames were corrupted", testName, parseErr.Load()) @@ -455,7 +578,9 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { t.Errorf("%s: %d sequence gap(s) observed — frames were dropped or reordered", testName, gaps.Load()) } if closeTimeout.Load() != 0 { - t.Errorf("%s: %d connection(s) timed out waiting for Close handshake", testName, closeTimeout.Load()) + t.Errorf("%s: %d connection(s) timed out waiting for Close handshake: no byte and no close for %v "+ + "(or the wait ran past %v, or %s); each WSO-GIVEUP record above has the state of both ends", + testName, closeTimeout.Load(), wsoCloseIdle, wsoWaitCap, rig.budgetNote()) } if bpBuf >= 256 && overflowErr.Load() != 0 { t.Errorf("%s: %d channel overflow error(s) observed at buffer capacity %d", testName, overflowErr.Load(), bpBuf) @@ -466,7 +591,13 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { in := connFramesIn[i].Load() sent := connSent[i].Load() if in != sent { - t.Errorf("%s conn %d: frame count mismatch: in=%d, sent=%d", testName, i, in, sent) + // The handler reads to the end of the stream even + // after its echo fails, so these frames were never + // delivered: the engine lost them (celeris#611 + // separates the two populations). + t.Errorf("%s conn %d (127.0.0.1:%d): frame count mismatch: sent=%d delivered=%d, never delivered %d "+ + "(of the delivered, %d were read after the handler's echo failed)", + testName, i, connPort[i].Load(), sent, in, sent-in, connReadAfterEchoErr[i].Load()) } if pe := connProtoErr[i].Load(); pe != 0 { t.Errorf("%s conn %d: protocol error on unfailed client: %d", testName, i, pe) diff --git a/middleware/websocket/pause_cancel_linux_test.go b/middleware/websocket/pause_cancel_linux_test.go index 578d1d68..d38a126b 100644 --- a/middleware/websocket/pause_cancel_linux_test.go +++ b/middleware/websocket/pause_cancel_linux_test.go @@ -8,7 +8,6 @@ import ( "encoding/base64" "errors" "fmt" - "io" "net" "os" "strconv" @@ -40,19 +39,37 @@ import ( // crosses high-water and requests the pause). A small backpressure buffer // makes the pause fire often per burst. // -// Three oracles, all must hold on every engine: +// Oracles, all must hold on every engine: // 1. the handler never observes ECANCELED from WriteMessage; -// 2. every conn still completes a Close handshake afterwards (a paused -// conn whose recv was never re-armed cannot read the Close frame and -// the client times out instead of seeing EOF); -// 3. engine shutdown completes within a bound (paused zombies block it). +// 2. every conn finishes its last frame and its Close frame, and then +// completes the Close handshake (a paused conn whose recv was never +// re-armed cannot read them, and the client gives up instead); +// 3. no handler sees a write or read error, unless its own client gave up +// on the connection (which fails 2) and the error came after that +// client closed; +// 4. engine shutdown completes within a bound (paused zombies block it); +// 5. the engine never leaves a connection unread for wsoRecvStall while its +// socket holds bytes and nothing asked it to stop (the celeris#607 +// class, WSO-STALL); +// 6. io_uring: the engine never chains a recv behind a send on these +// connections (RecvLinkedArms == 0; celeris#607's mechanism). +// +// The client waits on progress, not on a clock, and reads while it waits: +// see backpressure_oracle_linux_test.go for why (celeris#633). A connection +// that never finishes its writes used to be counted and not judged +// (clientCloseFail, celeris#623); it now fails the test, and every give-up +// prints both ends' timeline (WSO-GIVEUP). Because the client's reads would +// also get it through a connection the engine had stopped reading (that is +// how celeris#607 stalled), 5 and 6 judge that directly; the client holds +// its backpressure idle for wsoQuiet after the flood, so a connection the +// engine stopped reading stays unread long enough for 5 to see it. // // NOTE: inbound stream integrity is verified by the sequence-oracle test // (TestBackpressureInboundSequenceIntegrity, #484); this test asserts the #482 // fix: no ECANCELED on in-flight sends, clean close handshakes, and clean shutdown. func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { if testing.Short() { - t.Skip("needs ~20s of loopback flood") + t.Skip("needs ~45s of loopback flood") } conns := envInt("WS482_CONNS", 96) bpBuf := envInt("WS482_BP", 256) @@ -62,44 +79,45 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { for _, kind := range engineKinds(t) { kind := kind t.Run(kind.String(), func(t *testing.T) { + rig := newWSORig(t) var ecanceled, otherWriteErr, protoErr atomic.Int64 // protoErr counts every read error that is not a close, so on its // own it cannot distinguish a frame the engine mis-delivered from // a chanReader that hit ErrReadLimit because the asynchronous - // pause overshot its buffer. Those want opposite fixes. Record - // the errors and print them with the verdict; handler goroutines - // outlive the subtest, so they must never touch t directly. - var readErrMu sync.Mutex - readErrs := map[string]int{} - noteReadErr := func(err error) { - readErrMu.Lock() - if len(readErrs) < 16 { - readErrs[err.Error()]++ - } - readErrMu.Unlock() - } + // pause overshot its buffer. Those want opposite fixes. The rig + // records each error with its time and its connection, and they + // are printed with the verdict; handler goroutines outlive the + // subtest, so they must never touch t directly. addr, shutdownEngine, srv := startNativeServerWithHandle(t, kind, Config{ CheckOrigin: func(*celeris.Context) bool { return true }, ReadLimit: 256 * 1024, MaxBackpressureBuffer: bpBuf, // realistic buffer; headroom (cap-highWater) must exceed async pause-apply latency, else Append drops a chunk (ErrReadLimit) -- a config artifact, not engine reordering Handler: func(c *Conn) { + s := rig.attach(c) + defer s.exited("returned", nil) for { + s.phase.Store(1) mt, msg, err := c.ReadMessage() if err != nil { if !isCloseErr(err) { protoErr.Add(1) - noteReadErr(err) + s.noteErr("read", err) } + s.exited("read", err) return } + s.phase.Store(2) if err := c.WriteMessage(mt, msg); err != nil { if errors.Is(err, syscall.ECANCELED) { ecanceled.Add(1) } else if !errors.Is(err, ErrWriteClosed) { otherWriteErr.Add(1) + s.noteErr("write", err) } + s.exited("write", err) return } + s.echoed(len(msg)) } }, }) @@ -119,12 +137,12 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { // Recv-arming witnesses, read after shutdown so every worker // has left its loop and no stall episode is still open. They // are direct atomics, so nothing is stranded in a - // per-iteration batch. RECVSTALL is the celeris#607 witness: + // per-iteration batch. RECVSTALL is a celeris#607 witness: // a connection owed a recv arm that the dirty-list retry // passed over because a SEND was outstanding. Logged, not - // asserted -- an episode is normal SQ-ring pressure; it is - // the DURATION, joined against closeTimeout above, that - // carries the verdict. + // asserted -- an episode is normal SQ-ring pressure; what + // one does to a connection is judged by the watch + // (WSO-STALL), whatever the mechanism. if kind != celeris.IOUring { return } @@ -142,10 +160,13 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { t.Logf("%s: LINKBLOCK arms=%d blockedTotalMs=%d blockedMaxMs=%d", kind, m.RecvLinkedArms, m.RecvLinkedBlockedNanos/1e6, m.RecvLinkedBlockedMaxNanos/1e6) + wsoAssertNoLinkedRecv(t, kind.String(), m.RecvLinkedArms, m.RecvLinkedBlockedMaxNanos) }() hostPort := strings.TrimPrefix(addr, "ws://") hostPort = strings.TrimSuffix(hostPort, "/ws") + rig.setServer(hostPort) + stopWatch := rig.watch() var closedOK, closeTimeout, dialFail, hsFail, clientCloseFail, clientMisaligned, framesSent atomic.Int64 // An RST is NOT a clean close (celeris#530). Linux emits one when @@ -174,8 +195,9 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { dialFail.Add(1) return } - defer func() { _ = c.Close() }() - if err := wsHandshake(c, hostPort); err != nil { + cl := rig.client(c) + defer rig.closeClient(cl, c) + if err := cl.handshake(c, hostPort); err != nil { hsFail.Add(1) return } @@ -197,20 +219,30 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { } time.Sleep(200 * time.Millisecond) } - // Complete the current frame so the wire is a whole number of frames. writeAll - // retries to completion: once the flood stops the server drains and the send - // buffer empties. A conn the server wrongly killed surfaces as an error here. + cl.mark("flood end wrote=%d", wrote) + cl.sample(c, "flood end") + // Hold the backpressure with the client idle (wsoQuiet): a + // healthy engine reads what this connection's socket + // holds or pauses it; one that stopped reading it behind + // the blocked echo SEND (celeris#607) cannot, and the + // watch sees that stretch pass wsoRecvStall. + time.Sleep(wsoQuiet) + cl.mark("quiet end") + buf := make([]byte, 64<<10) + // Complete the current frame so the wire is a whole number of frames, + // waiting on progress and reading while it waits. A conn the server + // wrongly killed or wedged surfaces here, with its timeline. if rem := wrote % 126; rem != 0 { need := 126 - rem start := wrote % len(batch) - if !writeAll(c, batch[start:start+need], 15*time.Second) { + if w := cl.write(c, "fc", batch[start:start+need], buf); !w.ok { clientCloseFail.Add(1) + cl.giveUp(t, c, kind.String(), "gave up finishing its last frame: "+w.String()) return } wrote += need } // Drain the echo backlog (a real WS client reads); relieves backpressure. - buf := make([]byte, 64<<10) for { _ = c.SetReadDeadline(time.Now().Add(1 * time.Second)) if _, err := c.Read(buf); err != nil { @@ -221,30 +253,30 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { clientMisaligned.Add(1) } framesSent.Add(int64(wrote / 126)) - if !writeAll(c, maskedCloseFrame(), 10*time.Second) { + cl.mark("drained; writing Close") + if w := cl.write(c, "cw", maskedCloseFrame(), buf); !w.ok { clientCloseFail.Add(1) + cl.giveUp(t, c, kind.String(), "gave up writing its Close frame: "+w.String()) return } - _ = c.SetReadDeadline(time.Now().Add(10 * time.Second)) - tClose := time.Now() - for { - _, err := c.Read(buf) - if err == nil { - continue - } - if errors.Is(err, syscall.ECONNRESET) { - clientRST.Add(1) - } else if errors.Is(err, io.EOF) { - closedOK.Add(1) - } else { - closeTimeout.Add(1) - t.Logf("close-timeout %s: %v after Close sent", c.LocalAddr(), time.Since(tClose).Round(time.Millisecond)) + cl.mark("Close sent") + switch r := cl.closeWait(c, buf); r.outcome { + case "eof": + closedOK.Add(1) + if r.maxGap > wsoCloseSlow { + t.Logf("%s: close after a silent stretch longer than %v, conn %s: %s", kind, wsoCloseSlow, c.LocalAddr(), r) } - return + case "rst": + clientRST.Add(1) + cl.giveUp(t, c, kind.String(), "was RESET instead of closed: "+r.String()) + default: + closeTimeout.Add(1) + cl.giveUp(t, c, kind.String(), "gave up waiting for the server's close: "+r.String()) } }() } wg.Wait() + stalls := stopWatch() t.Logf("%s: clientMisaligned=%d framesSent=%d (if misaligned>0 the client truncated; if 0 while protocol errors>0 the engine mis-delivered)", kind, clientMisaligned.Load(), framesSent.Load()) @@ -252,8 +284,30 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { t.Errorf("%d client conn(s) ended mid-frame -- test client bug, not a server verdict", clientMisaligned.Load()) } - t.Logf("%s: conns=%d protoErr=%d clientCloseFail=%d ecanceled=%d otherWriteErr=%d closedOK=%d clientRST=%d closeTimeout=%d dialFail=%d hsFail=%d", - kind, conns, protoErr.Load(), clientCloseFail.Load(), ecanceled.Load(), otherWriteErr.Load(), closedOK.Load(), clientRST.Load(), closeTimeout.Load(), dialFail.Load(), hsFail.Load()) + // A server error is excused only on a connection whose client + // gave up or was reset, which fails the test below, and only + // after that client closed: that is the client's teardown + // echoing back (celeris#633: 1,165 of 1,165 such errors came + // after). Every other one is judged; the excused are printed. + // protoErr and otherWriteErr count every error, as they always + // have; the ...Judged fields are the ones judged. + judged, excused := rig.serverErrs() + var writeJudged, readJudged int64 + for _, e := range judged { + if e.kind == "write" { + writeJudged++ + } else { + readJudged++ + } + } + t.Logf("%s: conns=%d protoErr=%d clientCloseFail=%d ecanceled=%d otherWriteErr=%d closedOK=%d clientRST=%d closeTimeout=%d dialFail=%d hsFail=%d protoErrJudged=%d otherWriteErrJudged=%d serverErrsExcused=%d recvStalls=%d", + kind, conns, protoErr.Load(), clientCloseFail.Load(), ecanceled.Load(), otherWriteErr.Load(), closedOK.Load(), clientRST.Load(), closeTimeout.Load(), dialFail.Load(), hsFail.Load(), + readJudged, writeJudged, len(excused), len(stalls)) + t.Logf("%s: %s", kind, rig.waitStats()) + t.Logf("%s: %s", kind, rig.longestStretch()) + for _, e := range excused { + t.Logf("%s: not judged, its client gave up: %s", kind, e.line) + } if dialFail.Load()+hsFail.Load() > 0 { t.Fatalf("%d conns failed to dial/handshake -- environment problem, not a verdict", dialFail.Load()+hsFail.Load()) } @@ -263,17 +317,17 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { } // protoErr and otherWriteErr were logged but never asserted, so a // server that killed a healthy connection mid-write passed on - // these counters alone (celeris#530). Measured zero on both - // engines across repeated runs before this assertion was added. - if n := protoErr.Load(); n != 0 { - readErrMu.Lock() - for msg, count := range readErrs { - t.Logf("%s: read error x%d: %s", kind, count, msg) + // these counters alone (celeris#530). Judged unless excused + // (see above). + if len(judged) != 0 { + for _, e := range judged { + t.Logf("%s: %s", kind, e.line) } - readErrMu.Unlock() - t.Errorf("%d handler(s) saw a protocol error on the read side: the engine "+ - "mis-delivered or tore down a healthy connection", n) + t.Errorf("%d handler error(s) (%d read, %d write) on connections whose client did not give up, or "+ + "before it did: the engine mis-delivered, or tore down or failed a connection it should have kept alive", + len(judged), readJudged, writeJudged) } + wsoAssertNoStalls(t, kind.String(), stalls) // Measured zero on both engines across repeated runs before this // was asserted: nothing was actually being hidden inside closedOK // here, unlike the sibling inbound oracle where the same folding @@ -284,25 +338,28 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { "socket is closed with unread data still queued, which is a connection torn "+ "down mid-stream, not a clean close (celeris#530)", n) } - if n := otherWriteErr.Load(); n != 0 { - t.Errorf("%d handler(s) saw a non-ECANCELED write error: the engine failed a "+ - "WriteMessage on a connection it should have kept alive", n) + // celeris#623: counted and printed for months and never judged, + // while it was the larger of the two stall populations. The + // client now waits on progress and reads while it waits, so a + // give-up here is a connection on which nothing moved for + // wsoWriteIdle (or that ran past wsoWaitCap), not a slow server. + if n := clientCloseFail.Load(); n != 0 { + t.Errorf("%d conn(s) never finished writing their last frame or their Close frame: nothing "+ + "moved for %v (or the wait ran past %v, or %s). Each WSO-GIVEUP record above has the "+ + "connection's timeline from both ends and the shape it stalled in (celeris#623)", + n, wsoWriteIdle, wsoWaitCap, rig.budgetNote()) } if n := closeTimeout.Load(); n != 0 { - // Deliberately does NOT name a cause. This counter only says - // the server never closed its side within the client's - // window; it cannot distinguish a conn whose recv stayed - // paused (celeris#482) from one whose send accounting - // desynchronised so the dirty-list flush skipped it forever - // (celeris#519), and asserting the former sent three separate - // investigations down the wrong path -- in the #519 failures - // nothing was paused at all. - t.Errorf("%d conn(s) never completed the Close handshake within the client's read window "+ - "after sending Close: "+ - "the server never closed its side. Cause is NOT implied by this counter -- dump the "+ - "engine's per-conn state (recvPaused/recvArmed/sending/dirty/closing) to tell a "+ - "stuck pause (celeris#482) from stranded send accounting (celeris#519)", - n) + // Deliberately does NOT name a cause in the message: a + // connection whose recv stayed paused (celeris#482) and one + // whose send accounting desynchronised so the dirty-list flush + // skipped it forever (celeris#519) look the same here, and + // asserting the former sent three separate investigations down + // the wrong path. The WSO-GIVEUP record carries the state. + t.Errorf("%d conn(s) never completed the Close handshake: no byte and no close from the "+ + "server for %v after the client's Close (or the wait ran past %v, or %s). Each WSO-GIVEUP "+ + "record above has the connection's timeline and the state of both ends", + n, wsoCloseIdle, wsoWaitCap, rig.budgetNote()) } }) } @@ -318,9 +375,15 @@ func envInt(k string, def int) int { } func wsHandshake(c net.Conn, hostPort string) error { + return wsHandshakePath(c, hostPort, "/ws") +} + +// wsHandshakePath is wsHandshake for a request target other than "/ws" (the +// backpressure oracles put their join key in the query). +func wsHandshakePath(c net.Conn, hostPort, path string) error { key := make([]byte, 16) _, _ = rand.Read(key) - req := "GET /ws HTTP/1.1\r\nHost: " + hostPort + "\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n" + + req := "GET " + path + " HTTP/1.1\r\nHost: " + hostPort + "\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n" + "Sec-WebSocket-Key: " + base64.StdEncoding.EncodeToString(key) + "\r\nSec-WebSocket-Version: 13\r\n\r\n" _ = c.SetDeadline(time.Now().Add(5 * time.Second)) defer func() { _ = c.SetDeadline(time.Time{}) }() @@ -359,23 +422,6 @@ func maskedTextFrames(n, plen int) []byte { return out } -// writeAll writes every byte of buf, retrying partial writes until done or the -// overall deadline. A frame written this way is never truncated, so the stream -// stays well-formed regardless of backpressure timing. Returns false on error -// or deadline (e.g. the server killed the conn -- the celeris#482 symptom). -func writeAll(c net.Conn, buf []byte, within time.Duration) bool { - end := time.Now().Add(within) - for len(buf) > 0 { - _ = c.SetWriteDeadline(end) - n, err := c.Write(buf) - buf = buf[n:] - if err != nil { - return len(buf) == 0 - } - } - return true -} - // maskedCloseFrame is a client->server Close (opcode 8) with status 1000. func maskedCloseFrame() []byte { m := [4]byte{0x11, 0x22, 0x33, 0x44} diff --git a/middleware/websocket/writeall_closeprobe_linux_test.go b/middleware/websocket/writeall_closeprobe_linux_test.go new file mode 100644 index 00000000..e5f99d60 --- /dev/null +++ b/middleware/websocket/writeall_closeprobe_linux_test.go @@ -0,0 +1,25 @@ +//go:build linux && celeris_closeprobe + +package websocket + +import ( + "net" + "time" +) + +// writeAll writes every byte of buf, retrying partial writes until done or the +// overall deadline. A frame written this way is never truncated, so the stream +// stays well-formed regardless of backpressure timing. Returns false on error +// or deadline (e.g. the server killed the conn -- the celeris#482 symptom). +func writeAll(c net.Conn, buf []byte, within time.Duration) bool { + end := time.Now().Add(within) + for len(buf) > 0 { + _ = c.SetWriteDeadline(end) + n, err := c.Write(buf) + buf = buf[n:] + if err != nil { + return len(buf) == 0 + } + } + return true +}