From 13ca047e885503b8fa996dd65e11b2a17012fbe7 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 20:08:25 +0200 Subject: [PATCH 1/3] test(websocket): judge every connection the #482 oracle cannot finish, with waits that give up on no progress and both ends' timeline (celeris#633, celeris#623) TestBackpressurePauseDoesNotCancelInflightSend finished its last frame and its Close frame with fixed 15 s and 10 s deadlines, waited a fixed 10 s for the server's close, never read while it waited, and counted a connection that never finished its writes (clientCloseFail) without judging it. celeris#633 measured what those deadlines judged on GitHub's runners: a server that was behind but moving (457 of 459 close-timeouts were still receiving bytes), a client that stalled itself because its full receive queue makes Linux 6.17 discard the server's ACKs, and server errors that were the echo of the client's own give-up (1,165 of 1,165 after that client's close). The client now writes in 1 s slices, empties its receive queue without blocking after every slice that times out, and gives up only after 90 s with no progress (no byte written, no shrink of its send queue) or at 240 s; the close wait is re-armed by every byte and gives up after 30 s of silence or at 240 s. SO_RCVBUF stays 32 KiB and the flood never reads, so the echo SEND still blocks, which is what the test is for. In the same change, clientCloseFail is asserted (celeris#623), and a server error is judged only when it came before its own client closed. Every give-up prints a WSO-GIVEUP record: the connection's milestones and a 1 s timeline of both ends (client socket; handler phase and echoes; chanReader depth, spill, pause state and the pauses and resumes it applied; the server socket read through the engine's descriptor: bytes read by the engine, bytes accepted by the kernel and ACKed, the handler bytes the engine has not handed to the kernel), and the shape it stalled in. The client joins its handler through ?cid= and Conn.Query. --- .../backpressure_oracle_linux_test.go | 898 ++++++++++++++++++ .../websocket/pause_cancel_linux_test.go | 179 ++-- 2 files changed, 1006 insertions(+), 71 deletions(-) create mode 100644 middleware/websocket/backpressure_oracle_linux_test.go diff --git a/middleware/websocket/backpressure_oracle_linux_test.go b/middleware/websocket/backpressure_oracle_linux_test.go new file mode 100644 index 00000000..f44d9816 --- /dev/null +++ b/middleware/websocket/backpressure_oracle_linux_test.go @@ -0,0 +1,898 @@ +//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). +// +// # 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. +// +// 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 +) + +// wsoRig joins each client connection to the handler that serves it and +// keeps what each side saw. +type wsoRig struct { + t0 time.Time + srvPort atomic.Int32 // set once the server listens; read by handler goroutines + srvs sync.Map // client local port -> *wsoSrv + mu sync.Mutex + clis []*wsoCli + hs []*wsoSrv // every handler, joined or not +} + +func newWSORig() *wsoRig { return &wsoRig{t0: time.Now()} } + +// 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 + baseOut int64 // server socket: bytes ACKed + queued when the handler started + baseIn int64 // server socket: bytes received - unread when the handler started + baseOK bool + exit string + exitNs int64 + errs []wsoSrvErr + fd int + ino uint64 +} + +// 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() + } + if ti, inq, outq, ok := g.srvSockInfo(s); ok { + s.mu.Lock() + s.baseOut = int64(ti.Bytes_acked) + int64(outq) + s.baseIn = int64(ti.Bytes_received) - int64(inq) + s.baseOK = true + s.mu.Unlock() + } + 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). 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.fd >= 0 { + var st unix.Stat_t + if unix.Fstat(s.fd, &st) == nil && st.Ino == s.ino { + return s.fd + } + s.fd = -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. +func wsoDrainNow(c net.Conn, buf []byte) int64 { + tc, ok := c.(*net.TCPConn) + if !ok { + return 0 + } + sc, err := tc.SyscallConn() + if err != nil { + return 0 + } + var total int64 + _ = sc.Read(func(fd uintptr) bool { + for { + n, e := unix.Read(int(fd), buf) + if n > 0 { + total += int64(n) + continue + } + if e == unix.EINTR { + continue + } + return true // EAGAIN (empty), EOF or an error: stop, never wait + } + }) + return total +} + +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 + engRead, pendingWrite int64 +} + +func (g *wsoRig) srvView(port int) wsoSrvView { + var v wsoSrvView + x, ok := g.srvs.Load(port) + if !ok { + return v + } + s := x.(*wsoSrv) + 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 + baseOK, baseOut, baseIn := s.baseOK, s.baseOut, s.baseIn + s.mu.Unlock() + v.engRead, v.pendingWrite = -1, -1 + if v.sockOK && baseOK { + v.engRead = int64(v.ti.Bytes_received) - int64(v.inq) - baseIn + v.pendingWrite = wrote - (int64(v.ti.Bytes_acked) + int64(v.outq) - baseOut) + } + 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 + 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} + if a, ok := c.LocalAddr().(*net.TCPAddr); ok { + cl.port = a.Port + } + g.mu.Lock() + g.clis = append(g.clis, cl) + g.mu.Unlock() + return cl +} + +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", "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, or +// on any error other than the slice's own deadline. +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" + 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 + } + w.drained += wsoDrainNow(c, dbuf) + } + _ = 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", "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, or at wsoWaitCap. +// 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 + } + _ = 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" + 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 those that came BEFORE +// their own client closed its socket, which the engine did, and those after, +// which are the client's own teardown echoing back (a socket closed with +// unread data sends an RST). An error whose connection the rig could not +// join is counted as before: it cannot be excused. +func (g *wsoRig) serverErrs() (before, after []wsoJudged) { + closeAt := map[int]int64{} + g.mu.Lock() + for _, cl := range g.clis { + cl.mu.Lock() + closeAt[cl.port] = cl.closeNs + 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 { + ca, ok := closeAt[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)) + if ok && ca > 0 && e.ns >= ca { + after = append(after, wsoJudged{e.kind, line + fmt.Sprintf(", %s after its client closed", wsoSec(e.ns-ca))}) + continue + } + if ok && ca > 0 { + line += fmt.Sprintf(", %s BEFORE its client closed", wsoSec(ca-e.ns)) + } else { + line += ", its client had not closed" + } + before = append(before, wsoJudged{e.kind, line}) + } + } + sort.Slice(before, func(i, j int) bool { return before[i].line < before[j].line }) + sort.Slice(after, func(i, j int) bool { return after[i].line < after[j].line }) + return before, after +} + +// 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 +} + +func (g *wsoRig) waitStats() wsoWaitStats { + var st wsoWaitStats + var closes []time.Duration + g.mu.Lock() + defer g.mu.Unlock() + 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", + 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)) +} diff --git a/middleware/websocket/pause_cancel_linux_test.go b/middleware/websocket/pause_cancel_linux_test.go index 578d1d68..b22817e6 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,12 +39,19 @@ 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 before its own client closed; +// 4. engine shutdown completes within a bound (paused zombies block it). +// +// 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). // // NOTE: inbound stream integrity is verified by the sequence-oracle test // (TestBackpressureInboundSequenceIntegrity, #484); this test asserts the #482 @@ -62,44 +68,45 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { for _, kind := range engineKinds(t) { kind := kind t.Run(kind.String(), func(t *testing.T) { + rig := newWSORig() 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)) } }, }) @@ -146,6 +153,7 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { hostPort := strings.TrimPrefix(addr, "ws://") hostPort = strings.TrimSuffix(hostPort, "/ws") + rig.setServer(hostPort) 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 +182,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 := wsHandshakePath(c, hostPort, cl.path()); err != nil { hsFail.Add(1) return } @@ -197,20 +206,23 @@ 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") + 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,26 +233,25 @@ 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()) } }() } @@ -252,8 +263,26 @@ 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 that came after its own client closed is that + // client's teardown echoing back, not the engine's doing + // (celeris#633: 1,165 of 1,165 such errors came after). Only the + // ones before are judged; the rest are printed. + errsBefore, errsAfter := rig.serverErrs() + var writeBefore, readBefore int64 + for _, e := range errsBefore { + if e.kind == "write" { + writeBefore++ + } else { + readBefore++ + } + } + t.Logf("%s: conns=%d protoErr=%d clientCloseFail=%d ecanceled=%d otherWriteErr=%d closedOK=%d clientRST=%d closeTimeout=%d dialFail=%d hsFail=%d serverErrsAfterClientClose=%d protoErrAll=%d otherWriteErrAll=%d", + kind, conns, readBefore, clientCloseFail.Load(), ecanceled.Load(), writeBefore, closedOK.Load(), clientRST.Load(), closeTimeout.Load(), dialFail.Load(), hsFail.Load(), + len(errsAfter), protoErr.Load(), otherWriteErr.Load()) + t.Logf("%s: %s", kind, rig.waitStats()) + for _, e := range errsAfter { + t.Logf("%s: not judged, after its client's close: %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,16 +292,15 @@ 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 only when the error + // came before its own client closed (see above). + if len(errsBefore) != 0 { + for _, e := range errsBefore { + 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) before their own client closed: the engine "+ + "mis-delivered, or tore down or failed a connection it should have kept alive", + len(errsBefore), readBefore, writeBefore) } // Measured zero on both engines across repeated runs before this // was asserted: nothing was actually being hidden inside closedOK @@ -284,25 +312,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). Each WSO-GIVEUP record above has the "+ + "connection's timeline from both ends and the shape it stalled in (celeris#623)", + n, wsoWriteIdle, wsoWaitCap) } 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). Each WSO-GIVEUP "+ + "record above has the connection's timeline and the state of both ends", + n, wsoCloseIdle, wsoWaitCap) } }) } @@ -318,9 +349,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{}) }() From 3421840468c6c28d216aaead92be71283d2d35fa Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Sun, 27 Sep 2026 20:08:26 +0200 Subject: [PATCH 2/3] test(websocket): count delivered frames, not where the handler gave up, and judge the echo failure with its connection (celeris#611) TestBackpressureInboundSequenceIntegrity's handler returned when its echo write failed, leaving frames the engine had delivered unread, and the test then reported them lost (99.8-99.9% of the "missing" frames were still queued at handler exit); the echo failure itself was logged and never judged. The handler now stops echoing and keeps reading to the end of the stream, so the frame count measures delivery; frames read after the echo failed are counted separately, and a mismatch names the connection and both populations. An echo write or read error before its own client closed fails the test with the error and the connection's address; after the client's close it is that client's teardown and is printed, not judged. Its client gets the same wait/read pattern as the #482 oracle's (celeris#633): a flood write that times out is the backpressure the test creates and ends the burst, instead of taking the connection out of the frame-count verdict; the last frame and the Close frame are written on progress with drains between slices; the close wait is re-armed per byte. A client that cannot finish, and a close that arrives as an RST, now fail the test with the connection's timeline, as clientCloseFail and clientRST were counted and not judged. writeAll, which only the celeris_closeprobe-tagged close-drain test still uses, moves to a file under that tag, so the default build does not carry an unused helper. --- .../websocket/inbound_sequence_linux_test.go | 198 +++++++++++++----- .../websocket/pause_cancel_linux_test.go | 17 -- .../writeall_closeprobe_linux_test.go | 25 +++ 3 files changed, 175 insertions(+), 65 deletions(-) create mode 100644 middleware/websocket/writeall_closeprobe_linux_test.go diff --git a/middleware/websocket/inbound_sequence_linux_test.go b/middleware/websocket/inbound_sequence_linux_test.go index 9f21f427..856d65f4 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,29 @@ 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). A write timeout during the +// flood is the backpressure the test creates, not a failure: it used to take +// the connection out of the frame-count verdict. func TestBackpressureInboundSequenceIntegrity(t *testing.T) { conns := envInt("WS484_CONNS", 96) // bpBuf is the chanReader backpressure buffer capacity (default 256, matching the @@ -102,7 +111,8 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { detailMu.Unlock() } - var gaps, parseErr, overflowErr, protoErr, framesIn, framesSent, echoErr, slowClose atomic.Int64 + rig := newWSORig() + var gaps, parseErr, overflowErr, protoErr, framesIn, framesSent, echoErr, slowClose, readAfterEchoErr 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 +134,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 +150,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 +185,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 +233,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 +301,7 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { defer settle() hostPort := strings.TrimSuffix(strings.TrimPrefix(addr, "ws://"), "/ws") + rig.setServer(hostPort) 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 +320,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 := wsHandshakePath(c, hostPort, cl.path()); err != nil { hsFail.Add(1) clientFailed[connID].Store(true) return @@ -320,17 +360,32 @@ 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. + 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) } + cl.mark("flood end seq=%d off=%d", seq, off) + 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 +393,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 +400,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,28 +420,27 @@ 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), or at + // wsoWaitCap, 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) + } + switch r.outcome { + case "eof": + closedOK.Add(1) + case "rst": + clientRST.Add(1) + cl.giveUp(t, c, testName, "was RESET instead of closed: "+r.String()) + default: + closeTimeout.Add(1) + cl.giveUp(t, c, testName, "gave up waiting for the server's close: "+r.String()) } }() } @@ -404,8 +460,20 @@ 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()) + errsBefore, errsAfter := rig.serverErrs() + var echoBefore int64 + for _, e := range errsBefore { + if e.kind == "write" { + echoBefore++ + } + } + 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 echoErrBeforeClientClose=%d readAfterEchoErr=%d serverErrsAfterClientClose=%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(), + echoBefore, readAfterEchoErr.Load(), len(errsAfter)) + t.Logf("%s: %s", testName, rig.waitStats()) + for _, e := range errsAfter { + t.Logf("%s: not judged, after its client's close: %s", testName, e.line) + } detailMu.Lock() for _, d := range details { @@ -413,8 +481,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 @@ -444,9 +512,35 @@ 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 before its own client closed is the engine closing + // or failing a connection it should have kept: judged here, + // with its error and its connection's address. After the + // client's close it is that client's teardown (printed above). + if len(errsBefore) != 0 { + for _, e := range errsBefore { + t.Logf("%s: %s", testName, e.line) } + t.Errorf("%s: %d handler error(s) before their own client closed (%d echo write, %d read): the engine "+ + "failed or tore down a connection it should have kept alive; each is listed above with its error and address", + testName, len(errsBefore), echoBefore, int64(len(errsBefore))-echoBefore) + } + // 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) 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) + } + 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 +549,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); each WSO-GIVEUP record above has the state of both ends", + testName, closeTimeout.Load(), wsoCloseIdle, wsoWaitCap) } 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 +562,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 b22817e6..da8c9687 100644 --- a/middleware/websocket/pause_cancel_linux_test.go +++ b/middleware/websocket/pause_cancel_linux_test.go @@ -396,23 +396,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 +} From 1ee81fc0882c94d9d1a6bfd9f9a1027d056f7cc0 Mon Sep 17 00:00:00 2001 From: Albert Bausili Date: Mon, 28 Sep 2026 00:09:57 +0200 Subject: [PATCH 3/3] test(websocket): fail a connection the engine stops reading, and judge every server error on a connection that closed cleanly (celeris#607 class, celeris#623) The progress-based client of the previous commits empties its receive queue between write slices. That is what reopens a client window a blocked server SEND waits on, so it also gets through a connection whose engine chained its RECV behind that SEND (celeris#607). Measured in CI's shape on GitHub's runners with #607 re-introduced (flushSendLink's Detached early-return deleted), 5 processes per arch: the #482 oracle before this PR counted such connections as clientCloseFail (19 to 55 of 96, in 6 of 10 processes) and never judged them; this PR's previous head finished every one of them (0 of 10 processes failed) while the engine held a chained recv for up to 6.5 s. So the rig now watches the server side of every connection every 250 ms for its whole life. A connection whose server socket holds unread bytes, whose chanReader is neither paused nor holding anything, and whose handler has not exited is one the engine should be reading. If the engine reads nothing from it for 5 s the test fails with that stretch's state at both ends (WSO-STALL), whatever the client does afterwards. Bytes read, pauses and resumes only grow, so equal readings at a stretch's two ends prove it; a watcher run more than 1 s after its tick was due means the process was starved, and every open stretch starts over. The healthy maximum measured in CI's shape was 2.3 s (epoll) and 1.7 s (io_uring). After its flood the #482 client now holds its backpressure idle for 8 s, so a stretch that begins late in the flood still reaches the limit. On io_uring both oracles also assert RecvLinkedArms == 0: every connection on these servers is a detached WebSocket, and a chain on one is #607's own mechanism. Server errors: an error is excused only on a connection whose client gave up or was reset (which already fails the test) and only after that client closed. An error on a connection whose client read the server's close is judged whenever the handler noted it: that client's queue was empty and its close sent a FIN, and an engine error is noted only when the handler next returns from a read or write. Also: - every wait also ends at the test binary's deadline less 60 s, and a subtest that starts with less than 30 s of that left fails at once, so a wedge under CI's -timeout=300s ends as each subtest's verdict and WSO-GIVEUP record, not a goroutine dump, while a slow run that would have finished in time still does; - the inbound verdict line counts bursts a flood write deadline ended (floodDeadlines, floodDeadlineConns), and a connection whose close wait gave up or was reset leaves the frame count, as a failed write already did: the frames still in its client's send queue were never the engine's to deliver, and the give-up fails the test with its record; - the timeline's engine-read and pending-write bytes count from the connection's first byte net of the upgrade, whose sizes the client reads from its own socket, not from a baseline sampled after the 101 was queued; a drain that finds the connection reset or closed ends the wait as "rst" or "eof"; - PauseCancel's protoErr and otherWriteErr keep main's meaning (every error); the judged ones are protoErrJudged and otherWriteErrJudged; the waits line adds recvStallMax, watchResets and watchMaxLag, and a "longest unread stretch" line shows how close a run came to the limit. --- .../backpressure_oracle_linux_test.go | 578 +++++++++++++++--- .../websocket/inbound_sequence_linux_test.go | 87 ++- .../websocket/pause_cancel_linux_test.go | 96 +-- 3 files changed, 611 insertions(+), 150 deletions(-) diff --git a/middleware/websocket/backpressure_oracle_linux_test.go b/middleware/websocket/backpressure_oracle_linux_test.go index f44d9816..7043723c 100644 --- a/middleware/websocket/backpressure_oracle_linux_test.go +++ b/middleware/websocket/backpressure_oracle_linux_test.go @@ -43,6 +43,38 @@ package websocket // 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 @@ -60,7 +92,10 @@ package websocket // 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. +// 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. @@ -101,20 +136,81 @@ const ( // 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 - srvPort atomic.Int32 // set once the server listens; read by handler goroutines - srvs sync.Map // client local port -> *wsoSrv - mu sync.Mutex - clis []*wsoCli - hs []*wsoSrv // every handler, joined or not + 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 } -func newWSORig() *wsoRig { return &wsoRig{t0: time.Now()} } +// 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)) } @@ -149,15 +245,13 @@ type wsoSrv struct { pauses, resumes atomic.Int64 pauseNs, resumeNs atomic.Int64 - mu sync.Mutex - baseOut int64 // server socket: bytes ACKed + queued when the handler started - baseIn int64 // server socket: bytes received - unread when the handler started - baseOK bool - exit string - exitNs int64 - errs []wsoSrvErr - fd int - ino uint64 + 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 @@ -191,13 +285,7 @@ func (g *wsoRig) attach(c *Conn) *wsoSrv { } r.pausedMu.Unlock() } - if ti, inq, outq, ok := g.srvSockInfo(s); ok { - s.mu.Lock() - s.baseOut = int64(ti.Bytes_acked) + int64(outq) - s.baseIn = int64(ti.Bytes_received) - int64(inq) - s.baseOK = true - s.mu.Unlock() - } + g.srvSockFD(s) // find the engine's descriptor while the connection is young return s } @@ -295,8 +383,10 @@ func wsoFDByInode(ino uint64) int { // 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). Only read-only calls -// (fstat, getsockopt TCP_INFO, ioctl SIOCINQ/SIOCOUTQ) are made on it. +// 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 { @@ -304,12 +394,16 @@ func (g *wsoRig) srvSockFD(s *wsoSrv) int { } 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 = -1 + s.fd, s.gone = -1, true + return -1 } _, ino := wsoProcTCP(srvPort, s.port) if fd := wsoFDByInode(ino); fd >= 0 { @@ -363,31 +457,43 @@ func wsoCliSockInfo(c net.Conn) (ti *unix.TCPInfo, inq, outq int) { return ti, inq, outq } -// wsoDrainNow empties the client's receive queue without blocking. -func wsoDrainNow(c net.Conn, buf []byte) int64 { +// 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 + return 0, nil } sc, err := tc.SyscallConn() if err != nil { - return 0 + return 0, err } var total int64 - _ = sc.Read(func(fd uintptr) bool { + var end error + if err := sc.Read(func(fd uintptr) bool { for { n, e := unix.Read(int(fd), buf) - if n > 0 { + switch { + case n > 0: total += int64(n) continue - } - if e == unix.EINTR { + case e == unix.EINTR: continue + case e == nil: + end = io.EOF + case e != unix.EAGAIN: + end = os.NewSyscallError("read", e) } - return true // EAGAIN (empty), EOF or an error: stop, never wait + return true // empty, EOF or an error: stop, never wait } - }) - return total + }); err != nil && end == nil { + end = err + } + return total, end } func wsoTCPShort(ti *unix.TCPInfo) string { @@ -429,16 +535,20 @@ type wsoSrvView struct { 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 { - var v wsoSrvView x, ok := g.srvs.Load(port) if !ok { - return v + return wsoSrvView{} } - s := x.(*wsoSrv) + 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() @@ -468,12 +578,23 @@ func (g *wsoRig) srvView(port int) wsoSrvView { v.ti, v.inq, v.outq, v.sockOK = g.srvSockInfo(s) s.mu.Lock() v.exit = s.exit - baseOK, baseOut, baseIn := s.baseOK, s.baseOut, s.baseIn s.mu.Unlock() - v.engRead, v.pendingWrite = -1, -1 - if v.sockOK && baseOK { - v.engRead = int64(v.ti.Bytes_received) - int64(v.inq) - baseIn - v.pendingWrite = wrote - (int64(v.ti.Bytes_acked) + int64(v.outq) - baseOut) + 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 } @@ -532,25 +653,32 @@ func wsoShape(cli *unix.TCPInfo, cliInq int, v wsoSrvView) string { // 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 - 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 + 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} + 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) @@ -558,6 +686,22 @@ func (g *wsoRig) client(c net.Conn) *wsoCli { 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() @@ -596,7 +740,7 @@ type wsoWrite struct { need int left int err error - reason string // "" (completed), "noprog", "cap", "err" + reason string // "" (completed), "noprog", "cap", "budget", "rst", "eof", "err" dur time.Duration slices int progress int @@ -616,8 +760,9 @@ func (w wsoWrite) String() string { // 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, or -// on any error other than the slice's own deadline. +// 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 @@ -641,6 +786,9 @@ func (cl *wsoCli) write(c net.Conn, site string, buf, dbuf []byte) wsoWrite { } 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() @@ -662,7 +810,22 @@ func (cl *wsoCli) write(c net.Conn, site string, buf, dbuf []byte) wsoWrite { w.err, w.reason = err, "cap" break } - w.drained += wsoDrainNow(c, dbuf) + 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) @@ -684,7 +847,7 @@ func (cl *wsoCli) write(c net.Conn, site string, buf, dbuf []byte) wsoWrite { // wsoClose is how the wait for the server's close went. type wsoClose struct { - outcome string // "eof", "rst", "idle", "cap", "err" + outcome string // "eof", "rst", "idle", "cap", "budget", "err" err error dur time.Duration // Close sent -> outcome bytes int64 // received during the wait @@ -698,8 +861,9 @@ func (r wsoClose) String() string { } // closeWait reads until the server closes. The give-up is re-armed by every -// byte received: it comes after wsoCloseIdle without one, or at wsoWaitCap. -// While it waits in silence it records a timeline sample every wsoSlice. +// 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() @@ -712,6 +876,9 @@ func (cl *wsoCli) closeWait(c net.Conn, buf []byte) wsoClose { 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() @@ -736,6 +903,8 @@ func (cl *wsoCli) closeWait(c net.Conn, buf []byte) wsoClose { 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 @@ -806,17 +975,28 @@ type wsoJudged struct { line string } -// serverErrs splits every server-side error into those that came BEFORE -// their own client closed its socket, which the engine did, and those after, -// which are the client's own teardown echoing back (a socket closed with -// unread data sends an RST). An error whose connection the rig could not -// join is counted as before: it cannot be excused. -func (g *wsoRig) serverErrs() (before, after []wsoJudged) { - closeAt := map[int]int64{} +// 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() - closeAt[cl.port] = cl.closeNs + ends[cl.port] = end{cl.closeNs, cl.gaveUp} cl.mu.Unlock() } hs := append([]*wsoSrv(nil), g.hs...) @@ -826,23 +1006,245 @@ func (g *wsoRig) serverErrs() (before, after []wsoJudged) { errs := append([]wsoSrvErr(nil), s.errs...) s.mu.Unlock() for _, e := range errs { - ca, ok := closeAt[s.port] + 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)) - if ok && ca > 0 && e.ns >= ca { - after = append(after, wsoJudged{e.kind, line + fmt.Sprintf(", %s after its client closed", wsoSec(e.ns-ca))}) + 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 - } - if ok && ca > 0 { - line += fmt.Sprintf(", %s BEFORE its client closed", wsoSec(ca-e.ns)) - } else { + 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)) } - before = append(before, wsoJudged{e.kind, line}) + judged = append(judged, wsoJudged{e.kind, line}) } } - sort.Slice(before, func(i, j int) bool { return before[i].line < before[j].line }) - sort.Slice(after, func(i, j int) bool { return after[i].line < after[j].line }) - return before, after + 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. @@ -854,6 +1256,9 @@ type wsoWaitStats struct { 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 { @@ -861,6 +1266,7 @@ func (g *wsoRig) waitStats() 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++ @@ -892,7 +1298,7 @@ func (g *wsoRig) waitStats() wsoWaitStats { 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", + 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)) + 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 856d65f4..4b6142dc 100644 --- a/middleware/websocket/inbound_sequence_linux_test.go +++ b/middleware/websocket/inbound_sequence_linux_test.go @@ -42,9 +42,13 @@ const closeHandshakeSlow = 10 * time.Second // // 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). A write timeout during the -// flood is the backpressure the test creates, not a failure: it used to take -// the connection out of the frame-count verdict. +// (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 @@ -111,8 +115,11 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { detailMu.Unlock() } - rig := newWSORig() + 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 @@ -302,6 +309,7 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { 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) }) @@ -323,7 +331,7 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { cl := rig.client(c) connPort[connID].Store(int64(cl.port)) defer rig.closeClient(cl, c) - if err := wsHandshakePath(c, hostPort, cl.path()); err != nil { + if err := cl.handshake(c, hostPort); err != nil { hsFail.Add(1) clientFailed[connID].Store(true) return @@ -348,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 { @@ -368,6 +377,8 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { // connection, taking exactly the // backpressured connections out of // the frame-count verdict. + deadlines++ + floodDeadlines.Add(1) break } clientFailed[connID].Store(true) @@ -378,7 +389,10 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { } time.Sleep(200 * time.Millisecond) } - cl.mark("flood end seq=%d off=%d", seq, off) + 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 { @@ -421,8 +435,9 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { // (celeris#566). // // The wait gives up only after wsoCloseIdle with no - // byte from the server (re-armed by every byte), or at - // wsoWaitCap, so a genuine hang is still what fails, + // 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. @@ -432,19 +447,26 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { 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()) @@ -460,19 +482,20 @@ 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. - errsBefore, errsAfter := rig.serverErrs() - var echoBefore int64 - for _, e := range errsBefore { + judged, excused := rig.serverErrs() + var echoJudged int64 + for _, e := range judged { if e.kind == "write" { - echoBefore++ + 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 echoErrBeforeClientClose=%d readAfterEchoErr=%d serverErrsAfterClientClose=%d", + 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(), - echoBefore, readAfterEchoErr.Load(), len(errsAfter)) + echoJudged, readAfterEchoErr.Load(), len(excused), floodDeadlines.Load(), floodDeadlineConns.Load(), len(stalls)) t.Logf("%s: %s", testName, rig.waitStats()) - for _, e := range errsAfter { - t.Logf("%s: not judged, after its client's close: %s", testName, e.line) + 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() @@ -501,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) @@ -519,24 +545,27 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { // 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 before its own client closed is the engine closing - // or failing a connection it should have kept: judged here, - // with its error and its connection's address. After the - // client's close it is that client's teardown (printed above). - if len(errsBefore) != 0 { - for _, e := range errsBefore { + // 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) before their own client closed (%d echo write, %d read): the engine "+ - "failed or tore down a connection it should have kept alive; each is listed above with its error and address", - testName, len(errsBefore), echoBefore, int64(len(errsBefore))-echoBefore) + 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) 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) + "(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 "+ @@ -550,8 +579,8 @@ func TestBackpressureInboundSequenceIntegrity(t *testing.T) { } if closeTimeout.Load() != 0 { 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); each WSO-GIVEUP record above has the state of both ends", - testName, closeTimeout.Load(), wsoCloseIdle, wsoWaitCap) + "(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) diff --git a/middleware/websocket/pause_cancel_linux_test.go b/middleware/websocket/pause_cancel_linux_test.go index da8c9687..d38a126b 100644 --- a/middleware/websocket/pause_cancel_linux_test.go +++ b/middleware/websocket/pause_cancel_linux_test.go @@ -44,21 +44,32 @@ import ( // 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 before its own client closed; -// 4. engine shutdown completes within a bound (paused zombies block it). +// 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). +// 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) @@ -68,7 +79,7 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { for _, kind := range engineKinds(t) { kind := kind t.Run(kind.String(), func(t *testing.T) { - rig := newWSORig() + 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 @@ -126,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 } @@ -149,11 +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 @@ -184,7 +197,7 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { } cl := rig.client(c) defer rig.closeClient(cl, c) - if err := wsHandshakePath(c, hostPort, cl.path()); err != nil { + if err := cl.handshake(c, hostPort); err != nil { hsFail.Add(1) return } @@ -208,6 +221,13 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { } 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 @@ -256,6 +276,7 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { }() } 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()) @@ -263,25 +284,29 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { t.Errorf("%d client conn(s) ended mid-frame -- test client bug, not a server verdict", clientMisaligned.Load()) } - // A server error that came after its own client closed is that - // client's teardown echoing back, not the engine's doing - // (celeris#633: 1,165 of 1,165 such errors came after). Only the - // ones before are judged; the rest are printed. - errsBefore, errsAfter := rig.serverErrs() - var writeBefore, readBefore int64 - for _, e := range errsBefore { + // 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" { - writeBefore++ + writeJudged++ } else { - readBefore++ + readJudged++ } } - t.Logf("%s: conns=%d protoErr=%d clientCloseFail=%d ecanceled=%d otherWriteErr=%d closedOK=%d clientRST=%d closeTimeout=%d dialFail=%d hsFail=%d serverErrsAfterClientClose=%d protoErrAll=%d otherWriteErrAll=%d", - kind, conns, readBefore, clientCloseFail.Load(), ecanceled.Load(), writeBefore, closedOK.Load(), clientRST.Load(), closeTimeout.Load(), dialFail.Load(), hsFail.Load(), - len(errsAfter), protoErr.Load(), otherWriteErr.Load()) + 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()) - for _, e := range errsAfter { - t.Logf("%s: not judged, after its client's close: %s", kind, e.line) + 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()) @@ -292,16 +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). Judged only when the error - // came before its own client closed (see above). - if len(errsBefore) != 0 { - for _, e := range errsBefore { + // 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) } - t.Errorf("%d handler error(s) (%d read, %d write) before their own client closed: the engine "+ - "mis-delivered, or tore down or failed a connection it should have kept alive", - len(errsBefore), readBefore, writeBefore) + 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 @@ -319,9 +345,9 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { // 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). Each WSO-GIVEUP record above has the "+ + "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) + n, wsoWriteIdle, wsoWaitCap, rig.budgetNote()) } if n := closeTimeout.Load(); n != 0 { // Deliberately does NOT name a cause in the message: a @@ -331,9 +357,9 @@ func TestBackpressurePauseDoesNotCancelInflightSend(t *testing.T) { // 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). Each WSO-GIVEUP "+ + "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) + n, wsoCloseIdle, wsoWaitCap, rig.budgetNote()) } }) }