diff --git a/CHANGELOG.md b/CHANGELOG.md index ea3083b8..c1055de8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -277,6 +277,85 @@ Detailed per-release notes are on the `~/.pilot/update-state.json`. ### Fixed +- **Stream segments fit one packet.** A full stream segment was 4096 bytes, + about 4.2 KB on the wire and three IP fragments on a 1500-byte path. NATs, + firewalls and some virtual networks drop fragments, so on those paths + handshakes, pings and messages under ~1.4 KB worked while every full + segment was lost: larger replies never arrived, `send-file` failed, and + `bench` reported a few kbit/s until its retransmissions gave up. Segments + are now at most 1152 bytes (`SendSegmentSize`), a 1231-byte datagram on the + relay path — within what fits IPv6's minimum MTU, so within any path that + carries IP. This is a sender-side change: receivers accept segments up to + the old size, so old and new nodes interoperate in both directions. With + non-first fragments dropped at the receiver, a 3 KB message, a 60 KB + message and 64 KB and 1 MB file transfers all failed before and all + succeed now. + - The peer's receive window is advertised in segments, not bytes, and is + now enforced as such. A sender of short segments used to put more of + them in flight than the peer had slots for. + - A short write no longer stalls its sender. Nagle's algorithm held a + short segment until everything before it was acknowledged, and the + writer was blocked until it had gone out. The ACK it waited for is + delayed by the peer when the full segments before it are an odd number + (5ms; 40ms up to v1.14.1), nothing could join a held write because its + writer could not write again, so a stream of short writes moved at one + write per round trip, and since a client's sends are handled in order on + its IPC connection, everything else that client asked for waited too. + With 1152-byte segments every 4 KB write ends in a short tail, so this + had to go. Now: + - the tail of a write of a segment or more leaves with it; + - a write shorter than a segment waits only for an earlier short + segment, which is what lets small writes coalesce; + - a held write does not block: it stays buffered, later writes join it, + and it is sent when that earlier segment is acknowledged, after 40ms + at the latest, or ahead of the FIN when the connection is closed. + + This also removes the stall per file chunk that held transfers back. A + 60 MB file between two nodes on one host took 4.8s and takes 1.2s; from + an upgraded node into v1.13.5 it takes 1.8s, where v1.13.5 to itself + takes 22.7s. + - A receiver acknowledges a lone segment at once for the first 32 + segments of a connection, and again after each time its delayed-ACK + timer has had to fire. A peer that waits for an ACK before sending its + next small write no longer stalls for the timer each time; the first + exchange on a connection goes from 5.4ms to 0.3ms. + - The congestion window grows and backs off over the same number of bytes + as before. + + Mixed versions were checked in both directions against v1.13.5, v1.14.1 + and the previous main: messages up to 100 KB, files up to 60 MB, four + concurrent transfers, echo and pub/sub, on a clean network, under a stock + kernel's socket buffer limit, at 1% loss over a 30ms path, and through the + relay. Two things to know while a network is part upgraded: + - Only an upgraded sender stops fragmenting. A reply of more than a + packet from a node that is not upgraded is still lost on a path that + drops fragments. + - A node that is not upgraded and writes back what it reads piece by + piece (the echo service, so `pilotctl bench` against it) now reads + 1152-byte pieces, each below its own 4096-byte segment size, and its + own Nagle rule sends one per round trip: 1 MB echoed over a 30ms path + takes about 35s. Its replies to requests, files and pub/sub are not + affected, and upgrading that node removes it (1.5–2s). +- **A transfer no longer collapses after a burst of losses.** Three faults + in loss recovery, found at 2% packet loss on a 30ms path, where 2 of 5 + 5 MB transfers failed after several minutes: + - In fast recovery, segments the peer had already reported (SACK) were + left out of the amount in flight while the window was also grown by one + segment per duplicate ACK, counting each twice: the amount outstanding + doubled every round trip for as long as the first loss stayed + unrepaired. + - The peer holds at most 128 out-of-order segments and silently drops the + rest; the sender did not know and ran past it. It now keeps at most 128 + segments unacknowledged. + - After a retransmission timeout, every further segment missing from that + window waited for a timeout of its own, and the timeout doubles each + time (1s, 2s, 4s … 10s). The next missing segment is now retransmitted + as soon as an ACK shows the previous one arrived. + + Measured, 5 MB over a 30ms path, five runs each: at 2% loss, 9.7–12.4s + with 2 of 5 failing before, 10.3–11.8s with none failing after; at 0.5% + loss 9.4–10.5s before, 7.3–7.6s after; without loss 8.2–8.4s before, + 7.6–7.7s after. - **`-advertise-endpoint` survives a re-registration.** When the daemon re-registered (registry reconnect, transport watchdog recovery) it sent the tunnel socket's local address instead of the advertised endpoint, so diff --git a/pkg/daemon/daemon.go b/pkg/daemon/daemon.go index 6b0d6a86..a21c617c 100644 --- a/pkg/daemon/daemon.go +++ b/pkg/daemon/daemon.go @@ -3392,7 +3392,7 @@ func (d *Daemon) handleStreamPacket(pkt *protocol.Packet) { // Process peer's receive window from SYN (H9 fix: always update, including Window==0) conn.RetxMu.Lock() prevWin := conn.PeerRecvWin - conn.PeerRecvWin = int(pkt.Window) * MaxSegmentSize + conn.PeerRecvWin = int(pkt.Window) * SendSegmentSize // prevWin==-1 is the sentinel (no advertisement yet); don't signal // window-opened on the first transition (unknown→zero would fire). winOpened := prevWin != -1 && conn.PeerRecvWin > prevWin && conn.WindowAvailable() @@ -3481,7 +3481,7 @@ func (d *Daemon) handleStreamPacket(pkt *protocol.Packet) { // Process peer's receive window from SYN-ACK (H9 fix: always update) conn.RetxMu.Lock() prevWin := conn.PeerRecvWin - conn.PeerRecvWin = int(pkt.Window) * MaxSegmentSize + conn.PeerRecvWin = int(pkt.Window) * SendSegmentSize winOpened := prevWin != -1 && conn.PeerRecvWin > prevWin && conn.WindowAvailable() conn.RetxMu.Unlock() if winOpened && conn.WindowCh != nil { @@ -3595,7 +3595,7 @@ func (d *Daemon) handleStreamPacket(pkt *protocol.Packet) { // Update peer's receive window (H9 fix: always update, honor Window==0) conn.RetxMu.Lock() prevPeerWin := conn.PeerRecvWin - conn.PeerRecvWin = int(pkt.Window) * MaxSegmentSize + conn.PeerRecvWin = int(pkt.Window) * SendSegmentSize peerWinOpened := prevPeerWin != -1 && conn.PeerRecvWin > prevPeerWin && conn.WindowAvailable() conn.RetxMu.Unlock() if peerWinOpened && conn.WindowCh != nil { @@ -3661,13 +3661,24 @@ func (d *Daemon) handleStreamPacket(pkt *protocol.Packet) { conn.AckMu.Unlock() d.sendDelayedACK(conn) } else { - // Delayed ACK: batch up to 2 segments or 40ms + // Delayed ACK: every second segment at once, a lone one at + // once while the quick-ACK budget lasts, otherwise after + // DelayedACKTimeout. conn.PendingACKs++ if conn.PendingACKs >= DelayedACKThreshold { conn.AckMu.Unlock() d.sendDelayedACK(conn) + } else if conn.QuickACKs > 0 { + conn.QuickACKs-- + conn.AckMu.Unlock() + d.sendDelayedACK(conn) } else if conn.ACKTimer == nil { conn.ACKTimer = time.AfterFunc(DelayedACKTimeout, func() { + // Nothing followed the lone segment: its sender is + // probably waiting for this ACK before sending more. + conn.AckMu.Lock() + conn.QuickACKs = QuickACKBudget + conn.AckMu.Unlock() d.sendDelayedACK(conn) }) conn.AckMu.Unlock() @@ -4201,22 +4212,45 @@ const NagleTimeout = 40 * time.Millisecond // DelayedACKTimeout is the max time to delay an ACK (RFC 1122 suggests 500ms max). // -// It is also how long a write can stall with both ends idle. Nagle holds a -// short write — a frame's body after its header, the tail of a large write — -// until the data before it is ACKed, and the ACK of a lone or odd segment -// waits for this timer. At 40ms that was 40ms on the first exchange of every +// It is also how long a sender that waits for an ACK before it sends more +// can stall with both ends idle: a short write behind another short write +// (see nagleHoldsTail), once the receiver's quick-ACK budget is spent (see +// QuickACKBudget). At 40ms, and with every short tail held until all the +// data before it was ACKed, that was 40ms on the first exchange of every // connection and one stall per 48KB file chunk (about 1.5 MB/s on any link). // -// The timer is shortened rather than removed from the path. ACKing short -// segments at once, or sending tails without waiting, takes the pause between -// writes away entirely; several streams then burst into the peer's socket -// buffer faster than a stock kernel's can hold, and four concurrent 20MB -// transfers measured 34-40s instead of 1.4s. +// An earlier attempt to take the pause between writes away entirely had +// several streams burst into the peer's socket buffer faster than a stock +// kernel's can hold: four concurrent 20MB transfers measured 34-40s instead +// of 1.4s. Two things changed since. A sender keeps at most +// MaxSegmentsOutstanding segments unacknowledged, which bounds the burst, +// and a burst of losses is repaired on the ACK clock instead of one +// retransmission timeout per segment. With no pause between writes the same +// four transfers take 1.4s under a stock kernel's limit (rmem_max 212992). const DelayedACKTimeout = 5 * time.Millisecond // DelayedACKThreshold is the number of segments to receive before sending an ACK immediately. const DelayedACKThreshold = 2 +// QuickACKBudget is how many lone segments are acknowledged at once, without +// the delayed-ACK timer, at the start of a connection and again after each +// time the timer has had to fire. +// +// A sender that holds its next small write until the previous one is +// acknowledged (Nagle) sends one segment and waits. Each of those is a lone +// segment here, and delaying its ACK stalls that sender for the whole timer. +// The usual case is the first exchange on a connection — a frame header, then +// its body. The costly one is a peer up to v1.14.1 relaying what it reads: +// it reads this node's 1152-byte segments, writes each back as a write below +// its own 4096-byte segment size, and waits for the ACK of every one. bench +// against such a peer took 3.2s for 1 MB at 5ms a segment. +// +// While the budget lasts every arriving segment is acknowledged on its own. +// A bulk transfer uses it up within its first 32 segments and from then on +// is acknowledged every second segment as before: the timer, which restores +// the budget, only fires when a segment arrives and nothing follows it. +const QuickACKBudget = 32 + // SendData sends data over an established connection. // Implements Nagle's algorithm: small writes are coalesced into MSS-sized // segments unless NoDelay is set. Large writes (>= MSS) are sent immediately. @@ -4256,91 +4290,189 @@ func (d *Daemon) SendData(conn *Connection, data []byte) error { conn.NagleBuf = append(conn.NagleBuf, data...) conn.NagleMu.Unlock() - return d.nagleFlush(conn) + // A write of a segment or more takes its tail with it: it is the end of + // a message, not one of a run of small writes that would coalesce. + _, err := d.flushNagle(conn, len(data) >= SendSegmentSize, false) + return err } -// nagleFlush sends buffered data according to Nagle's algorithm: -// - Full MSS segments are always sent -// - Sub-MSS data is sent only if no unacknowledged data exists or timeout +// nagleFlush sends what Nagle's algorithm allows from the connection's +// buffer: every full segment, and the short remainder unless it is held. +// +// A remainder is held when the write that left it was shorter than a +// segment and an earlier short segment is still unacknowledged (see +// nagleHoldsTail). It stays in the buffer, where the next write joins it, and +// flushHeldTail sends it once that segment is acknowledged or NagleTimeout +// has passed. The caller does not wait for that. +// +// It used to: this function returned only once the buffer was empty. A +// writer whose short write was held could not write again until it went out, +// so nothing ever joined it, and a stream of short writes moved at one write +// per round trip — a node echoing back 4 KB writes it had received as three +// full segments and a tail managed 136 KB/s at 30ms. IPC sends are handled +// inline in the client's read loop, so everything else that client asked +// for waited as well. +// +// SendData does not hold the remainder of a write of a segment or more. func (d *Daemon) nagleFlush(conn *Connection) error { + _, err := d.flushNagle(conn, false, false) + return err +} + +// flushNagle is nagleFlush with two switches: force sends the remainder even +// if Nagle would hold it, and asFlusher marks the call as coming from +// flushHeldTail. +// +// flushHeldTail sends only a held remainder, and sends it without waiting for +// the window: it is less than one segment, and a goroutine that waited here +// would hold SendMu for as long as the window stayed shut, with the +// connection unable to close in order behind it. It does not start another +// copy of itself, and its bookkeeping (Connection.tailFlusher) is cleared +// here, under NagleMu, at the moment nothing is left for it to send — so a +// writer that holds a new remainder right afterwards finds no flusher and +// starts one. +// +// SendMu is held from taking bytes out of the buffer until they are handed +// on to be sent, so segments leave in the order the bytes were written +// whichever goroutine sends them. +func (d *Daemon) flushNagle(conn *Connection, force, asFlusher bool) (held bool, err error) { + conn.SendMu.Lock() + defer conn.SendMu.Unlock() + for { conn.NagleMu.Lock() if len(conn.NagleBuf) == 0 { + if asFlusher { + conn.tailFlusher = false + } conn.NagleMu.Unlock() - return nil + return false, nil } // If we have at least MSS bytes, send a full segment - if len(conn.NagleBuf) >= MaxSegmentSize { - segment := make([]byte, MaxSegmentSize) - copy(segment, conn.NagleBuf[:MaxSegmentSize]) - conn.NagleBuf = conn.NagleBuf[MaxSegmentSize:] + if len(conn.NagleBuf) >= SendSegmentSize { + if asFlusher { + // A writer has added to the buffer and is on its way here + // to send it; full segments wait for the window, which is + // the writer's place to wait, not this goroutine's. + conn.NagleMu.Unlock() + return true, nil + } + segment := make([]byte, SendSegmentSize) + copy(segment, conn.NagleBuf[:SendSegmentSize]) + conn.NagleBuf = conn.NagleBuf[SendSegmentSize:] conn.NagleMu.Unlock() if err := d.sendSegment(conn, segment); err != nil { - return err + return false, err } continue } // Sub-MSS data: check if we can send now (check under NagleMu). - // Use BytesInFlight() rather than len(Unacked): SACKed entries - // stay in Unacked until cumulative ACK removes them, but they - // are already at the peer and should not delay a flush. conn.RetxMu.Lock() - hasUnacked := conn.BytesInFlight() > 0 + hold := !force && conn.nagleHoldsTail() conn.RetxMu.Unlock() - if !hasUnacked { - // No data in flight — send immediately (Nagle allows this) - segment := make([]byte, len(conn.NagleBuf)) - copy(segment, conn.NagleBuf) - conn.NagleBuf = conn.NagleBuf[:0] + if hold { + if !asFlusher && !conn.tailFlusher { + conn.tailFlusher = true + go d.flushHeldTail(conn) + } conn.NagleMu.Unlock() + return true, nil + } - return d.sendSegment(conn, segment) + segment := make([]byte, len(conn.NagleBuf)) + copy(segment, conn.NagleBuf) + conn.NagleBuf = conn.NagleBuf[:0] + if asFlusher { + conn.tailFlusher = false } conn.NagleMu.Unlock() - // Data in flight — wait for ACK or timeout - nagleTimer := time.NewTimer(NagleTimeout) + if asFlusher { + return false, d.transmitSegment(conn, segment) + } + return false, d.sendSegment(conn, segment) + } +} + +// flushHeldTail sends a remainder that nagleFlush left held: as soon as the +// short segment ahead of it is acknowledged (NagleCh), and regardless once +// NagleTimeout has passed since it was first held. One runs per connection +// at most, and only while something is held. +func (d *Daemon) flushHeldTail(conn *Connection) { + timer := time.NewTimer(NagleTimeout) + defer timer.Stop() + force := false + for { select { case <-conn.NagleCh: - nagleTimer.Stop() - // All data ACKed — flush now - case <-nagleTimer.C: - // Timeout — flush regardless + // May be a signal left over from an earlier ACK; flushNagle + // looks again and keeps holding if it must. + case <-timer.C: + // Held long enough. If a writer turns out to be sending the + // buffer itself, look again after another interval. + force = true + timer.Reset(NagleTimeout) case <-conn.RetxStop: - nagleTimer.Stop() - return protocol.ErrConnClosed - } - - // Re-check under lock after waking - conn.NagleMu.Lock() - if len(conn.NagleBuf) == 0 { + conn.NagleMu.Lock() + conn.tailFlusher = false conn.NagleMu.Unlock() - return nil + return } - - // Send whatever we have (might have reached MSS now) - if len(conn.NagleBuf) >= MaxSegmentSize { - conn.NagleMu.Unlock() - continue // loop back to send full segments + if held, _ := d.flushNagle(conn, force, true); !held { + return } + } +} - segment := make([]byte, len(conn.NagleBuf)) - copy(segment, conn.NagleBuf) - conn.NagleBuf = conn.NagleBuf[:0] - conn.NagleMu.Unlock() - - return d.sendSegment(conn, segment) +// sendHeldTailBeforeClose sends a remainder Nagle is still holding, so that it +// is on the wire with a sequence number below the FIN's. It does not wait +// for the window — a close must not block on a peer that has stopped +// reading, and the remainder is less than one segment. +// +// SendMu is what orders this against flushHeldTail, which holds it only for +// the instant it takes to hand a remainder on. A writer can hold it for as +// long as the window stays shut; a close does not wait for that, and goes +// ahead as it always has. +func (d *Daemon) sendHeldTailBeforeClose(conn *Connection) { + conn.Mu.Lock() + established := conn.State == StateEstablished + conn.Mu.Unlock() + if !established { + return + } + locked := conn.SendMu.TryLock() + for i := 0; !locked && i < 20; i++ { + time.Sleep(time.Millisecond) + locked = conn.SendMu.TryLock() + } + if !locked { + return + } + defer conn.SendMu.Unlock() + conn.NagleMu.Lock() + tail := conn.NagleBuf + conn.NagleBuf = nil + conn.NagleMu.Unlock() + for len(tail) > 0 { + n := len(tail) + if n > SendSegmentSize { + n = SendSegmentSize + } + if err := d.transmitSegment(conn, tail[:n]); err != nil { + return + } + tail = tail[n:] } } // sendDataImmediate sends data in MSS-sized segments without Nagle coalescing. func (d *Daemon) sendDataImmediate(conn *Connection, data []byte) error { for offset := 0; offset < len(data); { - end := offset + MaxSegmentSize + end := offset + SendSegmentSize if end > len(data) { end = len(data) } @@ -4437,6 +4569,12 @@ func (d *Daemon) sendSegment(conn *Connection, data []byte) error { } } + return d.transmitSegment(conn, data) +} + +// transmitSegment puts one segment on the wire and tracks it for +// retransmission, without consulting the window. +func (d *Daemon) transmitSegment(conn *Connection, data []byte) error { // v1.9.1: reserve seq atomically with the read — pre-incrementing SendSeq // inside the same Mu critical section prevents two concurrent sendSegment // callers from reading the same SendSeq and emitting packets with identical @@ -4727,6 +4865,7 @@ func (d *Daemon) CloseConnection(conn *Connection) { // concurrent sendSegment cannot read the same value and produce a data // segment with the same seq as the FIN sentinel (iter 24 fix, same // pattern as the sendSegment pre-increment fix in iter 23). + d.sendHeldTailBeforeClose(conn) conn.Mu.Lock() st := conn.State sendSeq := conn.SendSeq diff --git a/pkg/daemon/ports.go b/pkg/daemon/ports.go index 2d669d86..d682d35a 100644 --- a/pkg/daemon/ports.go +++ b/pkg/daemon/ports.go @@ -206,16 +206,56 @@ type recvSegment struct { data []byte } +// SendSegmentSize is the most stream payload this node puts in one segment. +// +// A segment travels as one UDP datagram: 34 bytes of stream header, 36 bytes of +// tunnel framing (magic, sender, nonce, AEAD tag) and, when relayed, 9 more +// for the beacon header. At the old size, 4096, every full segment was a +// datagram of about 4.2 KB — three IP fragments on a 1500-byte path. NATs, +// firewalls and some virtual networks drop fragments, so on those paths the +// handshake, pings and small messages worked while every full segment was +// lost: a reply of more than ~1.4 KB never arrived, send-file failed, and +// bench "ran" at 0.017 Mbps until its retransmits gave up. +// +// 1152 keeps the datagram at 1231 bytes on the relay path, within the 1232 +// bytes that fit IPv6's minimum MTU (1280) and so within any path that can +// carry IP at all. It is a sender-side choice: a receiver takes segments of +// any size up to MaxSegmentSize, so peers running the old size interoperate +// in both directions. +// +// Which constant goes where: +// +// - SendSegmentSize wherever a segment is counted: cutting data into +// segments, converting the peer's window (advertised in segments) to +// bytes, and the congestion-window adjustments that stand for one segment +// leaving the network (one per duplicate ACK, three on entering fast +// recovery, the add-back on a partial ACK). +// - MaxSegmentSize for how fast the congestion window moves and how low it +// may go: the initial window, the growth per ACK in slow start and +// congestion avoidance, the ssthresh floor, the window after a timeout. +// These are unchanged, so a connection ramps up and backs off over the +// same number of bytes as before. Scaling them down with the segment made +// a lossy path markedly slower (5 MB at 30 ms RTT and 0.5% loss: 17s +// against 10s) for no gain: the old values were already what this +// transport ran with, in 4 KB datagrams. +const SendSegmentSize = 1152 + // Default window parameters const ( InitialCongWin = 10 * MaxSegmentSize // 40 KB initial congestion window (IW10, RFC 6928) MaxCongWin = 1024 * 1024 // 1 MB max congestion window - MaxSegmentSize = 4096 // MTU for virtual segments + MaxSegmentSize = 4096 // largest stream segment accepted from a peer; unit of the byte limits here RecvBufSize = 512 // receive buffer channel capacity (segments) MaxRecvWin = RecvBufSize * MaxSegmentSize // 2 MB max receive window MaxOOOBuf = 128 // max out-of-order segments buffered per connection - AcceptQueueLen = 64 // listener accept channel capacity - SendBufLen = 256 // send buffer channel capacity (segments) + // MaxSegmentsOutstanding is the most segments a sender keeps unacknowledged. + // When the oldest of them is lost, every later one waits in the peer's + // reorder buffer, which holds MaxOOOBuf segments and silently drops the + // rest. A sender past that limit loses the excess even on a clean path, + // and finds out one retransmission timeout at a time. + MaxSegmentsOutstanding = MaxOOOBuf + AcceptQueueLen = 64 // listener accept channel capacity + SendBufLen = 256 // send buffer channel capacity (segments) // MaxNagleBuf caps the per-connection NagleBuf at 64 segments // (256 KB). v1.9.1 fix: SendData previously appended without bound, @@ -275,19 +315,27 @@ type Connection struct { WindowCh chan struct{} // signaled when window opens up PeerRecvWin int // peer's advertised receive window (-1 = not yet received, 0 = explicit zero-window) // Nagle algorithm (write coalescing) - NagleBuf []byte // pending small write data - NagleMu sync.Mutex // protects NagleBuf - NagleCh chan struct{} // signaled when Nagle should flush - DialCh chan struct{} // signaled when an outbound dial leaves SYN_SENT - NoDelay bool // if true, disable Nagle (send immediately) + NagleBuf []byte // pending small write data + NagleMu sync.Mutex // protects NagleBuf and tailFlusher + // SendMu is held from taking bytes out of NagleBuf until they are handed + // to sendSegment, so segments leave in the order the bytes were written + // whichever goroutine sends them. Taken before NagleMu. + SendMu sync.Mutex + // tailFlusher is true while a flushHeldTail goroutine is waiting to send + // a remainder that Nagle is holding. + tailFlusher bool + NagleCh chan struct{} // signaled when Nagle should flush + DialCh chan struct{} // signaled when an outbound dial leaves SYN_SENT + NoDelay bool // if true, disable Nagle (send immediately) // Receive window (reassembly) RecvMu sync.Mutex ExpectedSeq uint32 // next in-order seq expected OOOBuf []*recvSegment // out-of-order buffer // Delayed ACK - AckMu sync.Mutex // protects PendingACKs and ACKTimer + AckMu sync.Mutex // protects PendingACKs, ACKTimer and QuickACKs PendingACKs int // count of unacked received segments ACKTimer *time.Timer // delayed ACK timer + QuickACKs int // lone segments still to be ACKed at once (see QuickACKBudget) // Keepalive dead-peer detection KeepaliveUnacked int // consecutive unanswered keepalive probes // Close @@ -500,6 +548,7 @@ func (pm *PortManager) NewConnection(localPort uint16, remoteAddr protocol.Addr, SSThresh: MaxCongWin / 2, WindowCh: make(chan struct{}, 1), NagleCh: make(chan struct{}, 1), + QuickACKs: QuickACKBudget, DialCh: make(chan struct{}, 1), PeerRecvWin: -1, // sentinel: no window advertisement received yet // Initialize RetxStop here (instead of in startRetxLoop) so it is @@ -748,6 +797,40 @@ func (c *Connection) BytesInFlight() int { return total } +// nagleHoldsTail reports whether a write shorter than a segment must wait +// before it is sent. Must be called with RetxMu held. (The tail of a write +// of a segment or more is never held; see SendData.) +// +// The rule is Nagle's with Minshall's refinement: a short segment waits only +// while an earlier short segment is unacknowledged. Full segments in flight +// do not hold it. +// +// Nagle's original rule — wait while anything is unacknowledged — meets the +// peer's delayed ACK badly. A peer acknowledges every second segment at once +// and a lone or odd one on a timer (5ms; 40ms up to v1.14.1), so the tail +// behind an odd number of full segments waited for that timer. With 4096-byte +// segments only writes that were not a multiple of 4096 paid it. With +// SendSegmentSize every 4096-byte write is three full segments and a tail: +// a program streaming 4 KB writes stalled 5ms per write against a current +// peer and 40ms against an older one (1 MB took 3.6s and 10s). +// +// A stream of tiny writes still coalesces: the second waits for the first to +// be acknowledged. +// +// SACKed entries stay in Unacked until a cumulative ACK removes them, but +// they are already at the peer and do not count. +func (c *Connection) nagleHoldsTail() bool { + for _, e := range c.Unacked { + if e.sacked || e.isFIN || len(e.data) == 0 { + continue + } + if len(e.data) < SendSegmentSize { + return true + } + } + return false +} + // EffectiveWindow returns the effective send window (minimum of congestion // window and peer's advertised receive window). // Must be called with RetxMu held. @@ -764,8 +847,41 @@ func (c *Connection) EffectiveWindow() int { // WindowAvailable returns true if the effective window allows more data. // Must be called with RetxMu held. +// +// Three limits apply. +// +// Bytes: what is in flight against the smaller of the congestion window and +// the peer's window. A SACKed segment has left the network and normally does +// not count. In fast recovery it does: there the congestion window is +// inflated by one segment per duplicate ACK, which already accounts for the +// segments that have left. Leaving them out as well counted each one twice, +// so every duplicate ACK released two new segments and the amount in flight +// doubled each round trip for as long as the hole stayed open. +// +// Segments against the peer's window: the peer advertises free slots in its +// receive buffer, one per segment whatever the segment's size. Short segments +// — the tail of a write, a small message — fill slots without filling bytes, +// so a sender that only counted bytes sent more segments than the peer had +// room for; the peer's packet loop then blocks for up to a second per segment +// on its full buffer, stalling every connection on that node. SACKed segments +// count here: they wait in the peer's reorder buffer and take a slot each the +// moment the hole before them is filled. +// +// Segments against the peer's reorder buffer: see MaxSegmentsOutstanding. func (c *Connection) WindowAvailable() bool { - return c.BytesInFlight() < c.EffectiveWindow() + bytes := 0 + for _, e := range c.Unacked { + if !e.sacked || c.FastRecovery { + bytes += len(e.data) + } + } + if bytes >= c.EffectiveWindow() { + return false + } + if c.PeerRecvWin >= 0 && len(c.Unacked) >= c.PeerRecvWin/SendSegmentSize { + return false + } + return len(c.Unacked) < MaxSegmentsOutstanding } // TrackSend adds a sent data segment to the retransmission buffer. @@ -846,7 +962,7 @@ func (c *Connection) ProcessAck(ack uint32, pureACK bool) { // (InRecovery=true, FastRecovery=true) must inflate cwnd by SMSS regardless // of the current count. if c.DupAckCount < 3 && c.InRecovery && c.FastRecovery && len(c.Unacked) > 0 { - c.CongWin += MaxSegmentSize + c.CongWin += SendSegmentSize if c.CongWin > MaxCongWin { c.CongWin = MaxCongWin } @@ -897,7 +1013,7 @@ func (c *Connection) ProcessAck(ack uint32, pureACK bool) { // §3 step 6 retransmit + deflation even when DupAckCount was reset // to 0 by a prior partial ACK. c.FastRecovery = true - c.CongWin = c.SSThresh + 3*MaxSegmentSize + c.CongWin = c.SSThresh + 3*SendSegmentSize // For small CongWin (< 6*SMSS), SSThresh+3*MSS > old CongWin, so // the window may have opened. Signal the sender so it doesn't stall. if c.WindowCh != nil && c.WindowAvailable() { @@ -911,7 +1027,7 @@ func (c *Connection) ProcessAck(ack uint32, pureACK bool) { // ACK reset DupAckCount to 0. fastRetransmit (above) is the RFC // 6582 §3 step 6a retransmit; the step-5 per-dup cwnd inflation // must also fire — it is not gated on newEpisode. - c.CongWin += MaxSegmentSize + c.CongWin += SendSegmentSize if c.CongWin > MaxCongWin { c.CongWin = MaxCongWin } @@ -934,7 +1050,7 @@ func (c *Connection) ProcessAck(ack uint32, pureACK bool) { return } // Inflate window for each additional dup ACK (only in recovery) - c.CongWin += MaxSegmentSize + c.CongWin += SendSegmentSize if c.CongWin > MaxCongWin { c.CongWin = MaxCongWin } @@ -1062,8 +1178,8 @@ func (c *Connection) ProcessAck(ack uint32, pureACK bool) { // so that a 1-SMSS partial ACK is neutral, and cwnd never falls // below ssthresh. c.CongWin -= bytesAcked - if bytesAcked >= MaxSegmentSize { - c.CongWin += MaxSegmentSize + if bytesAcked >= SendSegmentSize { + c.CongWin += SendSegmentSize } if c.CongWin < c.SSThresh { c.CongWin = c.SSThresh @@ -1073,6 +1189,16 @@ func (c *Connection) ProcessAck(ack uint32, pureACK bool) { // the 100ms RTO tick — up to one full RTO of unnecessary delay. c.fastRetransmit(recvAck) } + } else if wasInRecovery && c.InRecovery { + // Partial ACK in timeout recovery. The cumulative ACK stopped at the + // segment now at the head, and that segment was sent before the + // timeout, so the peer does not have it: resend it now. Waiting for + // its own timer instead costs a full RTO per missing segment, and + // the RTO doubles each time (Karn: no RTT sample comes from a + // retransmitted segment) — a window with thirty segments missing + // took minutes and the connection was given up on. One segment per + // ACK keeps the retransmissions clocked by what the peer receives. + c.fastRetransmit(recvAck) } // Congestion window growth (Appropriate Byte Counting, RFC 3465). @@ -1125,11 +1251,10 @@ func (c *Connection) ProcessAck(ack uint32, pureACK bool) { } } - // Signal Nagle flush when no data is genuinely in flight. - // Use BytesInFlight()==0 rather than len(c.Unacked)==0: after iter-49, - // sacked entries remain in Unacked across partial ACKs, so len>0 even - // when the peer already has every outstanding byte. - if c.BytesInFlight() == 0 && c.NagleCh != nil { + // Signal Nagle flush when no short segment is left in flight (see + // nagleHoldsTail). SACKed entries remain in Unacked across partial ACKs + // but are at the peer already, and do not count. + if c.NagleCh != nil && !c.nagleHoldsTail() { select { case c.NagleCh <- struct{}{}: default: @@ -1479,11 +1604,10 @@ func (c *Connection) ProcessSACK(blocks []SACKBlock) { } } - // When all outstanding bytes are now sacked, BytesInFlight()==0 means the - // peer has received everything in flight. Signal NagleCh so a nagleFlush - // goroutine waiting for "all data ACKed" wakes immediately instead of + // A SACK can be what tells us the short segment in flight has arrived. + // Signal NagleCh so a nagleFlush waiting on it wakes now instead of // stalling for NagleTimeout (40 ms) until the cumulative ACK arrives. - if c.NagleCh != nil && c.BytesInFlight() == 0 { + if c.NagleCh != nil && !c.nagleHoldsTail() { select { case c.NagleCh <- struct{}{}: default: diff --git a/pkg/daemon/zz_daemon_senddata_test.go b/pkg/daemon/zz_daemon_senddata_test.go index 41aae1a8..8481b742 100644 --- a/pkg/daemon/zz_daemon_senddata_test.go +++ b/pkg/daemon/zz_daemon_senddata_test.go @@ -125,7 +125,7 @@ func TestSendDataNagleFullMSSSegmentEmitsSingleFrame(t *testing.T) { t.Parallel() d, peer, conn := setupSendDataConn(t) - payload := bytes.Repeat([]byte{0x61}, MaxSegmentSize) // exactly one MSS + payload := bytes.Repeat([]byte{0x61}, SendSegmentSize) // exactly one MSS if err := d.SendData(conn, payload); err != nil { t.Fatalf("SendData: %v", err) } @@ -137,8 +137,8 @@ func TestSendDataNagleFullMSSSegmentEmitsSingleFrame(t *testing.T) { if err != nil { t.Fatal(err) } - if len(pkt.Payload) != MaxSegmentSize { - t.Errorf("Payload size=%d, want %d", len(pkt.Payload), MaxSegmentSize) + if len(pkt.Payload) != SendSegmentSize { + t.Errorf("Payload size=%d, want %d", len(pkt.Payload), SendSegmentSize) } } @@ -227,7 +227,7 @@ func TestNagleFlushWaitsForNagleChWhenUnackedPresent(t *testing.T) { } } -func TestNagleFlushSubMSSWithUnackedRetxStopAborts(t *testing.T) { +func TestHeldTailIsDroppedWhenConnectionStops(t *testing.T) { t.Parallel() d, _, conn := setupSendDataConn(t) conn.NagleMu.Lock() @@ -238,20 +238,34 @@ func TestNagleFlushSubMSSWithUnackedRetxStopAborts(t *testing.T) { conn.RetxMu.Unlock() conn.RetxStop = make(chan struct{}) - errCh := make(chan error, 1) - go func() { errCh <- d.nagleFlush(conn) }() + // The short write is held behind the unacknowledged short segment; the + // caller is not made to wait for it. + if err := d.nagleFlush(conn); err != nil { + t.Fatalf("nagleFlush: %v", err) + } + flusher := func() bool { + conn.NagleMu.Lock() + defer conn.NagleMu.Unlock() + return conn.tailFlusher + } + if !flusher() { + t.Fatal("a held tail should leave a flusher waiting to send it") + } - // Close RetxStop to trigger the ErrConnClosed branch. - time.Sleep(20 * time.Millisecond) + // The connection stopping ends the wait without sending. close(conn.RetxStop) - - select { - case err := <-errCh: - if err != protocol.ErrConnClosed { - t.Errorf("err=%v, want ErrConnClosed", err) + deadline := time.Now().Add(2 * time.Second) + for flusher() { + if time.Now().After(deadline) { + t.Fatal("the flusher did not exit after RetxStop closed") } - case <-time.After(2 * time.Second): - t.Fatal("nagleFlush did not return after RetxStop close") + time.Sleep(time.Millisecond) + } + conn.RetxMu.Lock() + sent := len(conn.Unacked) + conn.RetxMu.Unlock() + if sent != 1 { + t.Errorf("%d segments tracked, want only the one that was already in flight", sent) } } @@ -265,27 +279,24 @@ func TestNagleFlushNagleTimeoutFlushesEventually(t *testing.T) { conn.Unacked = []*retxEntry{{seq: 1, data: []byte("x"), sentAt: time.Now(), attempts: 1}} conn.RetxMu.Unlock() + // nagleFlush leaves the held write in the buffer and returns; it goes + // out when NagleTimeout has passed with the short segment ahead of it + // still unacknowledged. start := time.Now() - errCh := make(chan error, 1) - go func() { errCh <- d.nagleFlush(conn) }() - - select { - case err := <-errCh: - if err != nil { - t.Fatalf("nagleFlush: %v", err) - } - case <-time.After(1 * time.Second): - t.Fatal("nagleFlush did not fire after NagleTimeout") + if err := d.nagleFlush(conn); err != nil { + t.Fatalf("nagleFlush: %v", err) } - elapsed := time.Since(start) - if elapsed < NagleTimeout-5*time.Millisecond { - t.Errorf("nagleFlush returned too early: %v (want >= %v)", elapsed, NagleTimeout) + if returned := time.Since(start); returned >= NagleTimeout { + t.Errorf("nagleFlush took %v to return: it waited for the held write", returned) } frame := readOneFrame(t, peer) if frame == nil { t.Fatal("no frame received after Nagle timeout") } + if elapsed := time.Since(start); elapsed < NagleTimeout-5*time.Millisecond { + t.Errorf("held write was sent too early: %v (want >= %v)", elapsed, NagleTimeout) + } pkt, err := protocol.Unmarshal(frame[4:]) if err != nil { t.Fatal(err) diff --git a/pkg/daemon/zz_dup_ack_fresh_recovery_per_ack_inflation_bug_test.go b/pkg/daemon/zz_dup_ack_fresh_recovery_per_ack_inflation_bug_test.go index 9dfbd57f..63be8f8e 100644 --- a/pkg/daemon/zz_dup_ack_fresh_recovery_per_ack_inflation_bug_test.go +++ b/pkg/daemon/zz_dup_ack_fresh_recovery_per_ack_inflation_bug_test.go @@ -94,7 +94,9 @@ func TestFreshFastRecoveryPerDupAckInflation(t *testing.T) { congWin := c.CongWin c.RetxMu.Unlock() - expected := congWinBefore + MaxSegmentSize // 5*MSS + MSS = 6*MSS = 24576 + // The inflation is one segment as sent: SendSegmentSize, which is what a + // duplicate ACK says has left the network. + expected := congWinBefore + SendSegmentSize if congWin < expected { t.Errorf("4th dup ACK in fresh fast recovery: CongWin=%d, want >=%d (%d+MSS=%d); "+ "RFC 5681 §3.2 step 5 requires cwnd += SMSS per additional dup ACK while in "+ diff --git a/pkg/daemon/zz_fast_recovery_post_partial_ack_dup_inflation_bug_test.go b/pkg/daemon/zz_fast_recovery_post_partial_ack_dup_inflation_bug_test.go index 2a076339..5c62c1e1 100644 --- a/pkg/daemon/zz_fast_recovery_post_partial_ack_dup_inflation_bug_test.go +++ b/pkg/daemon/zz_fast_recovery_post_partial_ack_dup_inflation_bug_test.go @@ -95,7 +95,9 @@ func TestFastRecoveryPostPartialAckDupAckInflation(t *testing.T) { // RFC 5681 §3.2 step 5 ("while in fast retransmit"): cwnd += SMSS. c.ProcessAck(seqB, true) - if c.CongWin < congWinBefore+MaxSegmentSize { + // The inflation is one segment as sent: SendSegmentSize, which is what a + // duplicate ACK says has left the network. + if c.CongWin < congWinBefore+SendSegmentSize { t.Errorf("1st dup ACK after partial-ACK DupAckCount reset: CongWin=%d, want >=%d "+ "(%d+MSS=%d); RFC 5681 §3.2 step 5 requires cwnd += SMSS for every dup ACK "+ "'while in fast retransmit' (InRecovery=true, FastRecovery=true); "+ @@ -105,7 +107,7 @@ func TestFastRecoveryPostPartialAckDupAckInflation(t *testing.T) { "the 3-dup threshold is the ENTRY condition only, not a per-inflation gate; "+ "fix: add 'if DupAckCount<3 && InRecovery && FastRecovery' branch before the "+ "DupAckCount==3 check to inflate cwnd += MSS for counts 1-2 inside fast recovery", - c.CongWin, congWinBefore+MaxSegmentSize, - congWinBefore, congWinBefore+MaxSegmentSize) + c.CongWin, congWinBefore+SendSegmentSize, + congWinBefore, congWinBefore+SendSegmentSize) } } diff --git a/pkg/daemon/zz_fast_recovery_third_dup_ack_same_episode_inflation_bug_test.go b/pkg/daemon/zz_fast_recovery_third_dup_ack_same_episode_inflation_bug_test.go index 2d092fdb..92c379a0 100644 --- a/pkg/daemon/zz_fast_recovery_third_dup_ack_same_episode_inflation_bug_test.go +++ b/pkg/daemon/zz_fast_recovery_third_dup_ack_same_episode_inflation_bug_test.go @@ -103,7 +103,9 @@ func TestFastRecoveryThirdDupAckSameEpisodeInflation(t *testing.T) { // RFC 5681 §3.2 step 5: also inflate cwnd += SMSS. c.ProcessAck(seqB, true) - if c.CongWin < congWinBefore+MaxSegmentSize { + // The inflation is one segment as sent: SendSegmentSize, which is what a + // duplicate ACK says has left the network. + if c.CongWin < congWinBefore+SendSegmentSize { t.Errorf("3rd dup ACK (DupAckCount 2→3) same-episode fast recovery: "+ "CongWin=%d, want >=%d (%d+MSS=%d); "+ "RFC 5681 §3.2 step 5 requires cwnd += SMSS for every dup ACK while in "+ @@ -114,8 +116,8 @@ func TestFastRecoveryThirdDupAckSameEpisodeInflation(t *testing.T) { "the +MSS inflation must fire regardless of newEpisode when already in fast "+ "recovery; fix: add 'else if InRecovery && FastRecovery { CongWin += MSS }' "+ "after the newEpisode block inside the DupAckCount==3 branch", - c.CongWin, congWinBefore+MaxSegmentSize, - congWinBefore, congWinBefore+MaxSegmentSize, + c.CongWin, congWinBefore+SendSegmentSize, + congWinBefore, congWinBefore+SendSegmentSize, recoveryPoint) } } diff --git a/pkg/daemon/zz_loss_recovery_segments_test.go b/pkg/daemon/zz_loss_recovery_segments_test.go new file mode 100644 index 00000000..c60d4424 --- /dev/null +++ b/pkg/daemon/zz_loss_recovery_segments_test.go @@ -0,0 +1,142 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package daemon + +import ( + "testing" + "time" +) + +// In fast recovery the congestion window grows by one segment per duplicate +// ACK, which is how the segments that have left the network are accounted +// for. Leaving the SACKed ones out of the amount in flight as well counted +// each of them twice: every duplicate ACK released two new segments, the +// amount outstanding doubled each round trip while the hole stayed open, and +// the peer's reorder buffer overflowed. +func TestWindowCountsSACKedBytesInFastRecovery(t *testing.T) { + t.Parallel() + const segs = 10 + unacked := func() []*retxEntry { + out := make([]*retxEntry, segs) + for i := range out { + out[i] = &retxEntry{seq: uint32(1000 + i*SendSegmentSize), data: make([]byte, SendSegmentSize), attempts: 1} + out[i].sacked = i > 0 // the first is the hole, the rest are at the peer + } + return out + } + // The window as fast recovery leaves it after segs-1 duplicate ACKs: + // exactly what is outstanding. + c := &Connection{PeerRecvWin: 1 << 20, CongWin: segs * SendSegmentSize, Unacked: unacked()} + + c.InRecovery, c.FastRecovery = true, true + if c.WindowAvailable() { + t.Fatalf("fast recovery, %d segments outstanding against a window of %d: window should be closed", segs, segs) + } + c.CongWin += SendSegmentSize // one more duplicate ACK + if !c.WindowAvailable() { + t.Fatal("fast recovery: one more duplicate ACK should release one more segment") + } + + // Outside fast recovery nothing inflates the window, and a SACKed + // segment is simply no longer in flight. + c.InRecovery, c.FastRecovery = false, false + c.CongWin = 2 * SendSegmentSize + if !c.WindowAvailable() { + t.Fatal("outside fast recovery, one unSACKed segment against a window of two: window should be open") + } +} + +// When the oldest outstanding segment is lost, every later one waits in the +// peer's reorder buffer. That buffer holds MaxOOOBuf segments and drops the +// rest without a word, so a sender past it loses segments on a clean path and +// then recovers them one timeout at a time. +func TestWindowClosesAtPeerReorderCapacity(t *testing.T) { + t.Parallel() + outstanding := func(n int) []*retxEntry { + out := make([]*retxEntry, n) + for i := range out { + out[i] = &retxEntry{seq: uint32(1000 + i*10), data: make([]byte, 10), attempts: 1} + } + return out + } + c := &Connection{PeerRecvWin: MaxRecvWin, CongWin: MaxCongWin} + + c.Unacked = outstanding(MaxSegmentsOutstanding - 1) + if !c.WindowAvailable() { + t.Fatalf("%d segments outstanding: window should be open", MaxSegmentsOutstanding-1) + } + c.Unacked = outstanding(MaxSegmentsOutstanding) + if c.WindowAvailable() { + t.Fatalf("%d segments outstanding, as many as the peer can hold out of order: window should be closed", MaxSegmentsOutstanding) + } + if MaxSegmentsOutstanding > MaxOOOBuf { + t.Fatalf("MaxSegmentsOutstanding %d exceeds the reorder buffer, %d", MaxSegmentsOutstanding, MaxOOOBuf) + } +} + +// After a retransmission timeout, each further segment missing from that +// window used to wait for a timeout of its own, and the timeout doubles each +// time. Thirty missing segments took minutes; the connection was abandoned +// first. +func TestPartialAckInTimeoutRecoveryRetransmitsNextSegment(t *testing.T) { + t.Parallel() + const ( + seqA = uint32(1000) + seqB = seqA + SendSegmentSize + seqC = seqB + SendSegmentSize + end = seqC + SendSegmentSize + ) + setup := func() (*Connection, *capturedSender) { + c := newAckTestConn(t) + cs := &capturedSender{} + c.RetxSend = cs.send + c.LastAck = seqA + // Timeout recovery: entered by the retransmission timer, not by + // duplicate ACKs. + c.InRecovery, c.FastRecovery = true, false + c.RecoveryPoint = end + c.SSThresh = 2 * SendSegmentSize + c.CongWin = SendSegmentSize + sent := time.Now().Add(-2 * time.Second) + c.Unacked = []*retxEntry{ + {seq: seqA, data: make([]byte, SendSegmentSize), attempts: 2, sentAt: time.Now()}, // the timer's retransmission + {seq: seqB, data: make([]byte, SendSegmentSize), attempts: 1, sentAt: sent}, + {seq: seqC, data: make([]byte, SendSegmentSize), attempts: 1, sentAt: sent}, + } + return c, cs + } + + // The ACK for the retransmission stops at seqB: the peer does not have it. + c, cs := setup() + c.ProcessAck(seqB, true) + pkts := cs.all() + if len(pkts) != 1 || pkts[0].Seq != seqB { + t.Fatalf("partial ACK in timeout recovery: retransmitted %d segments, want exactly the next missing one (seq %d)", len(pkts), seqB) + } + if !c.InRecovery { + t.Fatal("partial ACK should leave the connection in recovery") + } + // The next partial ACK moves on to seqC. + c.ProcessAck(seqC, true) + if pkts = cs.all(); len(pkts) != 2 || pkts[1].Seq != seqC { + t.Fatalf("second partial ACK: %d retransmissions in total, want 2 with the last for seq %d", len(pkts), seqC) + } + + // An ACK for the whole window ends recovery and retransmits nothing. + c, cs = setup() + c.ProcessAck(end, true) + if n := len(cs.all()); n != 0 { + t.Fatalf("full ACK in timeout recovery retransmitted %d segments, want none", n) + } + if c.InRecovery { + t.Fatal("full ACK should end recovery") + } + + // Outside recovery an ACK retransmits nothing. + c, cs = setup() + c.InRecovery = false + c.ProcessAck(seqB, true) + if n := len(cs.all()); n != 0 { + t.Fatalf("ACK outside recovery retransmitted %d segments, want none", n) + } +} diff --git a/pkg/daemon/zz_ports_logic_test.go b/pkg/daemon/zz_ports_logic_test.go index 0f53d5a0..9354d2f8 100644 --- a/pkg/daemon/zz_ports_logic_test.go +++ b/pkg/daemon/zz_ports_logic_test.go @@ -474,13 +474,14 @@ func TestProcessAckThirdDupACKTriggersFastRetransmit(t *testing.T) { } // Multiplicative decrease: SSThresh = max(FlightSize/2, 2*SMSS). // FlightSize = len("one") = 3; FlightSize/2 = 1 < 2*MSS = 8192 → floor applies. - // CongWin = SSThresh + 3*MSS (fast-recovery inflation). + // CongWin = SSThresh + three segments as sent (fast-recovery inflation: + // each duplicate ACK is one SendSegmentSize segment leaving the network). wantSSThresh := 2 * MaxSegmentSize // max(3/2=1, 2*MSS=8192) if c.SSThresh != wantSSThresh { t.Fatalf("SSThresh = %d, want %d (max(FlightSize/2, 2*SMSS))", c.SSThresh, wantSSThresh) } - if c.CongWin != wantSSThresh+3*MaxSegmentSize { - t.Fatalf("CongWin = %d, want %d", c.CongWin, wantSSThresh+3*MaxSegmentSize) + if c.CongWin != wantSSThresh+3*SendSegmentSize { + t.Fatalf("CongWin = %d, want %d", c.CongWin, wantSSThresh+3*SendSegmentSize) } if c.Stats.FastRetx != 1 { t.Fatalf("Stats.FastRetx = %d, want 1", c.Stats.FastRetx) diff --git a/pkg/daemon/zz_segment_fits_packet_test.go b/pkg/daemon/zz_segment_fits_packet_test.go new file mode 100644 index 00000000..0fde1ab1 --- /dev/null +++ b/pkg/daemon/zz_segment_fits_packet_test.go @@ -0,0 +1,293 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package daemon + +import ( + "net" + "testing" + "time" + + "github.com/pilot-protocol/common/protocol" +) + +// What a stream segment gains on its way to the wire, and the room it has. +const ( + // [PILS][sender node ID][nonce] before the ciphertext, GCM tag after it. + testTunnelFraming = 4 + 4 + 12 + 16 + // The plaintext frame the fixture below sends: [PILP] and nothing else. + testPlainFraming = 4 + // [MsgRelay][sender node ID][destination node ID] when the beacon relays. + testRelayHeader = 1 + 4 + 4 + // UDP payload that fits IPv6's minimum MTU (1280) after the IPv6 and UDP + // headers. Any path that carries IP carries a datagram this size whole. + testSafeDatagram = 1280 - 40 - 8 +) + +// A full stream segment used to be 4096 bytes, about 4.2 KB on the wire and +// three IP fragments on a 1500-byte path. Where fragments are dropped — NATs, +// firewalls, some virtual networks — pings and small messages worked and +// every full segment was lost. +func TestSendSegmentSizeFitsOneUnfragmentedDatagram(t *testing.T) { + t.Parallel() + onWire := protocol.PacketHeaderSize() + SendSegmentSize + testTunnelFraming + testRelayHeader + if onWire > testSafeDatagram { + t.Fatalf("a full segment is a %d-byte datagram on the relay path, want at most %d", onWire, testSafeDatagram) + } + if SendSegmentSize > MaxSegmentSize { + t.Fatalf("SendSegmentSize %d exceeds MaxSegmentSize %d, the most a peer accepts", SendSegmentSize, MaxSegmentSize) + } +} + +// readDatagramSizes collects the sizes of the datagrams arriving on pc until +// want bytes of stream payload have been seen. +func readDatagramSizes(t *testing.T, pc *net.UDPConn, want int) (sizes []int, payload []byte) { + t.Helper() + buf := make([]byte, 65535) + for len(payload) < want { + pc.SetReadDeadline(time.Now().Add(2 * time.Second)) + n, _, err := pc.ReadFromUDP(buf) + if err != nil { + t.Fatalf("read after %d of %d payload bytes: %v", len(payload), want, err) + } + pkt, err := protocol.Unmarshal(buf[testPlainFraming:n]) + if err != nil { + t.Fatalf("unmarshal: %v", err) + } + sizes = append(sizes, n) + payload = append(payload, pkt.Payload...) + } + return sizes, payload +} + +func TestStreamWritesLeaveAsDatagramsThatFitOnePacket(t *testing.T) { + t.Parallel() + for _, tc := range []struct { + name string + noDelay bool + node uint32 + }{ + {"nagle", false, 0xD7D70001}, + {"nodelay", true, 0xD7D70002}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + d, pc, conn := newSendDataFixture(t, tc.node) + conn.NoDelay = tc.noDelay + + data := make([]byte, 20000) + for i := range data { + data[i] = byte(i * 7) + } + if err := d.SendData(conn, data); err != nil { + t.Fatalf("SendData: %v", err) + } + sizes, payload := readDatagramSizes(t, pc, len(data)) + if string(payload) != string(data) { + t.Fatalf("payload reassembled from %d datagrams differs from what was written", len(sizes)) + } + for i, n := range sizes { + // The fixture's tunnel is plaintext; an encrypted, relayed + // one adds the difference. + onWire := n - testPlainFraming + testTunnelFraming + testRelayHeader + if onWire > testSafeDatagram { + t.Fatalf("datagram %d of %d would be %d bytes encrypted and relayed, want at most %d", i+1, len(sizes), onWire, testSafeDatagram) + } + } + }) + } +} + +// The peer advertises its window as free slots in its receive buffer, one per +// segment whatever the segment's size. Counting only bytes, a sender of short +// segments put more of them in flight than the peer had slots; the peer's +// packet loop then blocks on its full buffer, up to a second per segment, +// with every other connection on that node waiting behind it. +func TestWindowAvailableCountsSegmentsAgainstPeerSlots(t *testing.T) { + t.Parallel() + short := func(n int) []*retxEntry { + out := make([]*retxEntry, n) + for i := range out { + out[i] = &retxEntry{seq: uint32(1000 + 10*i), data: make([]byte, 10), attempts: 1} + } + return out + } + const slots = 4 + c := &Connection{CongWin: MaxCongWin, PeerRecvWin: slots * SendSegmentSize} + + c.Unacked = short(slots - 1) + if !c.WindowAvailable() { + t.Fatalf("%d segments in flight against %d free slots: window should be open", slots-1, slots) + } + c.Unacked = short(slots) + if c.WindowAvailable() { + t.Fatalf("%d segments in flight against %d free slots: window should be closed (only %d bytes are in flight, but every slot is taken)", slots, slots, c.BytesInFlight()) + } + // A SACKed segment still counts: it waits in the peer's reorder buffer + // and takes a slot the moment the hole before it is filled. + c.Unacked[0].sacked = true + if c.WindowAvailable() { + t.Fatalf("%d segments outstanding, one of them SACKed, against %d free slots: window should be closed", slots, slots) + } + // No advertisement yet (-1): only the congestion window applies. + c.PeerRecvWin = -1 + c.Unacked = short(100) + if !c.WindowAvailable() { + t.Fatal("no peer window advertised yet: window should be open") + } +} + +func TestNagleHoldsTail(t *testing.T) { + t.Parallel() + full := func() *retxEntry { return &retxEntry{data: make([]byte, SendSegmentSize)} } + short := func() *retxEntry { return &retxEntry{data: make([]byte, 100)} } + sacked := func(e *retxEntry) *retxEntry { e.sacked = true; return e } + fin := &retxEntry{data: []byte{0}, isFIN: true} + for _, tc := range []struct { + name string + inFlight []*retxEntry + hold bool + }{ + {"nothing in flight", nil, false}, + {"a short segment in flight", []*retxEntry{short()}, true}, + {"2000-byte write: one full segment ahead", []*retxEntry{full()}, false}, + {"4096-byte write: three full segments ahead", []*retxEntry{full(), full(), full()}, false}, + {"bulk: forty full segments ahead", func() []*retxEntry { + out := make([]*retxEntry, 40) + for i := range out { + out[i] = full() + } + return out + }(), false}, + {"a short segment behind full ones", []*retxEntry{full(), short(), full()}, true}, + {"the short segment in flight is SACKed", []*retxEntry{sacked(short()), full()}, false}, + {"only the FIN sentinel", []*retxEntry{fin}, false}, + } { + c := &Connection{Unacked: tc.inFlight} + if got := c.nagleHoldsTail(); got != tc.hold { + t.Errorf("%s: tail held = %v, want %v", tc.name, got, tc.hold) + } + } +} + +// The tail of a write of a segment or more leaves with it, even while an +// earlier short segment is unacknowledged: it is the end of a message, and +// holding it delays the message by a round trip. +func TestSendDataSendsTailOfMultiSegmentWriteAtOnce(t *testing.T) { + t.Parallel() + const write = MaxSegmentSize // three full segments and a 640-byte tail + d, pc, conn := newSendDataFixture(t, 0xD8D80001) + conn.NoDelay = false + conn.RetxMu.Lock() + conn.Unacked = append(conn.Unacked, &retxEntry{data: []byte("an earlier short segment"), seq: 500, sentAt: time.Now(), attempts: 1}) + conn.RetxMu.Unlock() + + if err := d.SendData(conn, make([]byte, write)); err != nil { + t.Fatalf("SendData: %v", err) + } + // Nothing ACKs here, so a tail left held would still be in the buffer. + conn.NagleMu.Lock() + left := len(conn.NagleBuf) + conn.NagleMu.Unlock() + if left != 0 { + t.Fatalf("%d bytes of a %d-byte write were held back", left, write) + } + if sizes, _ := readDatagramSizes(t, pc, write); len(sizes) != 4 { + t.Fatalf("a %d-byte write left as %d datagrams, want 4", write, len(sizes)) + } +} + +// readSegments reads stream packets off pc until n have arrived. +func readSegments(t *testing.T, pc *net.UDPConn, n int) []*protocol.Packet { + t.Helper() + var out []*protocol.Packet + for len(out) < n { + out = append(out, readOneSegment(t, pc, 2*time.Second)) + } + return out +} + +// A short write that Nagle holds does not hold its writer: the write +// returns, later writes join it in the buffer, full segments leave as they +// fill, and what is left goes out when the short segment ahead is ACKed. +// Holding the writer instead meant nothing could ever join a held write, and +// a stream of short writes moved at one write per round trip. +func TestShortWritesBehindAHeldTailCoalesce(t *testing.T) { + t.Parallel() + d, pc, conn := newSendDataFixture(t, 0xD9D90001) + conn.NoDelay = false + conn.RetxMu.Lock() + conn.Unacked = append(conn.Unacked, &retxEntry{data: []byte("an earlier short segment"), seq: 500, sentAt: time.Now(), attempts: 1}) + conn.RetxMu.Unlock() + + const piece, pieces = 300, 5 // 1500 bytes: one full segment and 348 left + var want []byte + start := time.Now() + for i := 0; i < pieces; i++ { + b := make([]byte, piece) + for j := range b { + b[j] = byte('a' + i) + } + want = append(want, b...) + if err := d.SendData(conn, b); err != nil { + t.Fatalf("SendData %d: %v", i, err) + } + } + if took := time.Since(start); took >= NagleTimeout { + t.Fatalf("%d short writes took %v: a held write made its writer wait", pieces, took) + } + + // The fourth write filled a segment; it left without waiting. + first := readOneSegment(t, pc, time.Second) + if len(first.Payload) != SendSegmentSize { + t.Fatalf("first segment carries %d bytes, want a full one (%d) coalesced from the short writes", len(first.Payload), SendSegmentSize) + } + conn.NagleMu.Lock() + left := len(conn.NagleBuf) + conn.NagleMu.Unlock() + if left != piece*pieces-SendSegmentSize { + t.Fatalf("%d bytes left in the buffer, want %d held", left, piece*pieces-SendSegmentSize) + } + + // The earlier short segment is acknowledged: the rest follows. + conn.RetxMu.Lock() + conn.Unacked = conn.Unacked[1:] + conn.RetxMu.Unlock() + select { + case conn.NagleCh <- struct{}{}: + default: + } + rest := readOneSegment(t, pc, time.Second) + got := append(append([]byte{}, first.Payload...), rest.Payload...) + if string(got) != string(want) { + t.Fatalf("received %d bytes that differ from the %d written", len(got), len(want)) + } + if rest.Seq != first.Seq+uint32(len(first.Payload)) { + t.Fatalf("second segment seq %d, want %d (directly after the first)", rest.Seq, first.Seq+uint32(len(first.Payload))) + } +} + +// Closing a connection sends a write Nagle is still holding, ahead of the FIN. +func TestCloseSendsHeldTailBeforeFIN(t *testing.T) { + t.Parallel() + d, pc, conn := newSendDataFixture(t, 0xD9D90002) + conn.NoDelay = false + conn.RetxMu.Lock() + conn.Unacked = append(conn.Unacked, &retxEntry{data: []byte("an earlier short segment"), seq: 500, sentAt: time.Now(), attempts: 1}) + conn.RetxMu.Unlock() + + if err := d.SendData(conn, []byte("last words")); err != nil { + t.Fatalf("SendData: %v", err) + } + d.CloseConnection(conn) + + pkts := readSegments(t, pc, 2) + if string(pkts[0].Payload) != "last words" { + t.Fatalf("first packet after close carries %q, want the held write", pkts[0].Payload) + } + if pkts[1].Flags&protocol.FlagFIN == 0 { + t.Fatalf("second packet after close has flags %#x, want FIN", pkts[1].Flags) + } + if want := pkts[0].Seq + uint32(len(pkts[0].Payload)); pkts[1].Seq != want { + t.Fatalf("FIN seq %d, want %d (directly after the held write)", pkts[1].Seq, want) + } +} diff --git a/pkg/daemon/zz_send_test.go b/pkg/daemon/zz_send_test.go index e40a96e2..9d82145d 100644 --- a/pkg/daemon/zz_send_test.go +++ b/pkg/daemon/zz_send_test.go @@ -325,7 +325,7 @@ func TestSendDataImmediateSplitsIntoMSSSegments(t *testing.T) { conn.RecvAck = 0 // Payload = 2 full MSS + a small tail. - payload := make([]byte, 2*MaxSegmentSize+10) + payload := make([]byte, 2*SendSegmentSize+10) for i := range payload { payload[i] = byte(i % 251) } @@ -356,7 +356,7 @@ func TestSendDataImmediateSplitsIntoMSSSegments(t *testing.T) { if len(sizes) != 3 { t.Fatalf("got %d segments, want 3; sizes=%v", len(sizes), sizes) } - if sizes[0] != MaxSegmentSize || sizes[1] != MaxSegmentSize || sizes[2] != 10 { + if sizes[0] != SendSegmentSize || sizes[1] != SendSegmentSize || sizes[2] != 10 { t.Fatalf("segment sizes = %v, want [MSS, MSS, 10]", sizes) } } diff --git a/pkg/daemon/zz_senddata_test.go b/pkg/daemon/zz_senddata_test.go index 9d53c86e..4f2c4225 100644 --- a/pkg/daemon/zz_senddata_test.go +++ b/pkg/daemon/zz_senddata_test.go @@ -120,7 +120,7 @@ func TestNagleFlushFullMSSSendsSegment(t *testing.T) { const peerNode uint32 = 0xD2D2D2D2 d, pc, conn := newSendDataFixture(t, peerNode) conn.NoDelay = false - conn.NagleBuf = make([]byte, MaxSegmentSize) + conn.NagleBuf = make([]byte, SendSegmentSize) for i := range conn.NagleBuf { conn.NagleBuf[i] = byte('A' + (i % 26)) } @@ -129,8 +129,8 @@ func TestNagleFlushFullMSSSendsSegment(t *testing.T) { t.Fatalf("nagleFlush: %v", err) } pkt := readOneSegment(t, pc, 500*time.Millisecond) - if len(pkt.Payload) != MaxSegmentSize { - t.Fatalf("payload len = %d, want %d", len(pkt.Payload), MaxSegmentSize) + if len(pkt.Payload) != SendSegmentSize { + t.Fatalf("payload len = %d, want %d", len(pkt.Payload), SendSegmentSize) } conn.NagleMu.Lock() remaining := len(conn.NagleBuf) @@ -185,14 +185,14 @@ func TestNagleFlushSubMSSWithUnackedBlocksThenNagleChFlushes(t *testing.T) { } } -func TestNagleFlushRetxStopReturnsErrConnClosed(t *testing.T) { +func TestNagleFlushDoesNotWaitForAHeldWrite(t *testing.T) { t.Parallel() const peerNode uint32 = 0xD4D4D4D4 d, _, conn := newSendDataFixture(t, peerNode) conn.NoDelay = false - conn.NagleBuf = []byte("blocked") + conn.NagleBuf = []byte("held") conn.RetxStop = make(chan struct{}) - // Create in-flight so nagleFlush waits. + // A short segment in flight, so the write is held. conn.RetxMu.Lock() conn.Unacked = append(conn.Unacked, &retxEntry{ data: []byte("x"), seq: 1, sentAt: time.Now(), attempts: 1, @@ -201,18 +201,21 @@ func TestNagleFlushRetxStopReturnsErrConnClosed(t *testing.T) { errCh := make(chan error, 1) go func() { errCh <- d.nagleFlush(conn) }() - - time.Sleep(10 * time.Millisecond) - close(conn.RetxStop) - select { case err := <-errCh: - if err != protocol.ErrConnClosed { - t.Fatalf("err = %v, want protocol.ErrConnClosed", err) + if err != nil { + t.Fatalf("nagleFlush: %v", err) } case <-time.After(1 * time.Second): - t.Fatal("nagleFlush did not return after RetxStop close") + t.Fatal("nagleFlush waited for the held write instead of returning") + } + conn.NagleMu.Lock() + buffered := string(conn.NagleBuf) + conn.NagleMu.Unlock() + if buffered != "held" { + t.Fatalf("buffer = %q, want the held write still in it", buffered) } + close(conn.RetxStop) // lets the flusher go } // --- Stop / doStop --- diff --git a/pkg/daemon/zz_streampacket_test.go b/pkg/daemon/zz_streampacket_test.go index a4708d17..c791caf7 100644 --- a/pkg/daemon/zz_streampacket_test.go +++ b/pkg/daemon/zz_streampacket_test.go @@ -410,6 +410,11 @@ func TestHandleStreamACKWithDataDeliversAndSchedulesDelayedACK(t *testing.T) { conn.State = StateEstablished conn.ExpectedSeq = 1000 conn.Mu.Unlock() + // A connection starts with a budget of lone segments it ACKs at once + // (QuickACKBudget); this test is about what happens once it is spent. + conn.AckMu.Lock() + conn.QuickACKs = 0 + conn.AckMu.Unlock() data := streamPacket(protocol.FlagACK, peerNode, d.NodeID(), 443, 55555, 1000, 1) data.Payload = []byte("hello") @@ -474,3 +479,76 @@ func TestHandleStreamACKWithDataNonEstablishedIsIgnoredPayloadNotDelivered(t *te t.Fatalf("non-Established conn received data: bytes=%d segs=%d", bytesRecv, segsRecv) } } + +// A sender that holds its next small write until the previous one is ACKed +// sends one segment and waits. Delaying that lone segment's ACK stalled it +// for the whole delayed-ACK timer: once on the first exchange of every +// connection, and once per segment for a peer up to v1.14.1 relaying what it +// read from this node (bench against one took 3.2s for 1 MB). +func TestLoneSegmentIsAckedAtOnceWithinQuickACKBudget(t *testing.T) { + t.Parallel() + d, peerNode, _ := setupDaemonWithPeer(t, Config{Public: true}) + d.setNodeID_testhelper(0xABCD0010) + + conn := d.ports.NewConnection(55556, protocol.Addr{Network: 0, Node: peerNode}, 443) + conn.LocalAddr = protocol.Addr{Network: 0, Node: d.NodeID()} + conn.Mu.Lock() + conn.State = StateEstablished + conn.ExpectedSeq = 1000 + conn.Mu.Unlock() + + ackState := func() (pending, quick int, timer bool) { + conn.AckMu.Lock() + defer conn.AckMu.Unlock() + return conn.PendingACKs, conn.QuickACKs, conn.ACKTimer != nil + } + if _, quick, _ := ackState(); quick != QuickACKBudget { + t.Fatalf("a new connection has a quick-ACK budget of %d, want %d", quick, QuickACKBudget) + } + seq := uint32(1000) + lone := func() { + pkt := streamPacket(protocol.FlagACK, peerNode, d.NodeID(), 443, 55556, seq, 1) + pkt.Payload = []byte("hello") + seq += 5 + d.handleStreamPacket(pkt) + } + + // Within the budget: ACKed at once, nothing left pending, no timer. + lone() + if pending, quick, timer := ackState(); pending != 0 || timer || quick != QuickACKBudget-1 { + t.Fatalf("first lone segment: pending=%d timer=%v budget=%d, want 0, false, %d", pending, timer, quick, QuickACKBudget-1) + } + + // Budget spent: the lone segment's ACK is delayed. + conn.AckMu.Lock() + conn.QuickACKs = 0 + conn.AckMu.Unlock() + lone() + if pending, _, timer := ackState(); pending != 1 || !timer { + t.Fatalf("lone segment with the budget spent: pending=%d timer=%v, want 1, true", pending, timer) + } + + // The timer firing means nothing followed that segment, so its sender is + // likely waiting on the ACK: the budget is restored. + deadline := time.Now().Add(2 * time.Second) + for { + pending, quick, timer := ackState() + if pending == 0 && !timer && quick == QuickACKBudget { + break + } + if time.Now().After(deadline) { + t.Fatalf("after the delayed-ACK timer: pending=%d timer=%v budget=%d, want 0, false, %d", pending, timer, quick, QuickACKBudget) + } + time.Sleep(time.Millisecond) + } + + // A pair of segments is ACKed on the second and spends none of it. + conn.AckMu.Lock() + conn.QuickACKs = 0 + conn.AckMu.Unlock() + lone() + lone() + if pending, quick, timer := ackState(); pending != 0 || timer || quick != 0 { + t.Fatalf("two segments with the budget spent: pending=%d timer=%v budget=%d, want 0, false, 0", pending, timer, quick) + } +} diff --git a/pkg/daemon/zz_window_update_wakeup_bug_test.go b/pkg/daemon/zz_window_update_wakeup_bug_test.go index 00b96613..e053560a 100644 --- a/pkg/daemon/zz_window_update_wakeup_bug_test.go +++ b/pkg/daemon/zz_window_update_wakeup_bug_test.go @@ -104,7 +104,9 @@ func TestWindowUpdateDoesNotWakeSender(t *testing.T) { } // Peer sends a window-update ACK: same cumulative ACK (1000 == LastAck), - // but Window=1 (receiver window now open). + // but Window=2 (receiver window now open). The window is a count of free + // receive slots and the unacked segment above will take one of them, so + // two is the smallest update that leaves room for a new segment. // This is a dup-ACK from ProcessAck's point of view; ProcessAck will NOT // signal WindowCh on the dup-ACK path. windowUpdatePkt := &protocol.Packet{ @@ -117,7 +119,7 @@ func TestWindowUpdateDoesNotWakeSender(t *testing.T) { DstPort: localPort, Seq: 500, Ack: 1000, // == LastAck — dup-ACK path in ProcessAck - Window: 1, // non-zero: peer's window just opened + Window: 2, // room for the segment in flight plus one more } d.handleStreamPacket(windowUpdatePkt) @@ -127,7 +129,7 @@ func TestWindowUpdateDoesNotWakeSender(t *testing.T) { case <-conn.WindowCh: // good — sender wakes up promptly case <-time.After(100 * time.Millisecond): - t.Errorf("window-update ACK (Ack=LastAck, Window=1 with PeerRecvWin=0) did not " + + t.Errorf("window-update ACK (Ack=LastAck, Window=2 with PeerRecvWin=0) did not " + "signal conn.WindowCh within 100ms; " + "handleStreamPacket updates PeerRecvWin but never signals WindowCh; " + "ProcessAck is called with ack=LastAck (dup-ACK path) which returns " + diff --git a/tests/zz_stream_latency_test.go b/tests/zz_stream_latency_test.go index a4ec4794..28c4f21b 100644 --- a/tests/zz_stream_latency_test.go +++ b/tests/zz_stream_latency_test.go @@ -3,6 +3,7 @@ package tests import ( + "encoding/binary" "io" "testing" "time" @@ -148,3 +149,147 @@ func TestLargeWritesDoNotStallOnDelayedACK(t *testing.T) { t.Errorf("%d chunks took %v at best, want under 200ms: each large write is stalling on a delayed ACK", chunks, fastest) } } + +// A message of a few kilobytes crosses as two to four segments, the last one +// short. Held until the ones before it are ACKed, that tail waits for the +// receiver's delayed-ACK timer whenever the full segments are an odd number — +// one segment for a 2 KB message — so the message arrives a timer late: 5ms +// from a current peer, 40ms from an older one. +// +// Measured on a connection's first write, the one-message-per-connection +// case. Later writes on the same connection can be let through by a wake-up +// left over from the previous exchange's last ACK. +func TestFirstWriteOfAFewSegmentsIsNotDelayed(t *testing.T) { + requireRealNetwork(t) + env := NewTestEnv(t) + a := env.AddDaemon() + b := env.AddDaemon() + + const port = 5103 + ln, err := a.Driver.Listen(port) + if err != nil { + t.Fatalf("listen: %v", err) + } + defer ln.Close() + + // The first four bytes of a request say how long it is. + go func() { + for { + conn, err := ln.Accept() + if err != nil { + return + } + go func() { + defer conn.Close() + head := make([]byte, 4) + if _, err := io.ReadFull(conn, head); err != nil { + return + } + size := int(binary.BigEndian.Uint32(head)) + if _, err := io.ReadFull(conn, make([]byte, size-4)); err != nil { + return + } + _, _ = conn.Write([]byte{1}) + }() + } + }() + + // Judged on the fastest of six connections per size: a timer in the path + // puts a floor under every one of them, load only adds to some. + reply := make([]byte, 1) + for _, size := range []int{2000, 3000, 4096} { + request := make([]byte, size) + binary.BigEndian.PutUint32(request, uint32(size)) + fastest := time.Hour + for round := 0; round < 6; round++ { + conn, err := b.Driver.DialAddr(a.Daemon.Addr(), port) + if err != nil { + t.Fatalf("dial: %v", err) + } + start := time.Now() + if _, err := conn.Write(request); err != nil { + t.Fatalf("write %d bytes: %v", size, err) + } + if _, err := io.ReadFull(conn, reply); err != nil { + t.Fatalf("read reply to %d bytes: %v", size, err) + } + if d := time.Since(start); d < fastest { + fastest = d + } + conn.Close() + } + t.Logf("fastest first %d-byte write and 1-byte reply: %v", size, fastest) + // The delayed-ACK timer is 5ms; an exchange that waits for it + // cannot finish sooner. + if fastest > 4*time.Millisecond { + t.Errorf("fastest first %d-byte exchange took %v, want under 4ms: its last segment is waiting for a delayed ACK", size, fastest) + } + } +} + +// A program streaming 4 KB writes — a default bufio.Writer does — sends three +// full segments and a short tail per write. Held until everything before it +// was ACKed, each tail waited for the receiver's delayed-ACK timer on the +// third segment: 5ms a write against a current peer, 40ms against an older +// one, so 1 MB took 1.3s or 10s. The receiver here only reads, so no data of +// its own carries the ACK back sooner. +func TestStreamOf4KWritesDoesNotStallOnDelayedACK(t *testing.T) { + requireRealNetwork(t) + env := NewTestEnv(t) + a := env.AddDaemon() + b := env.AddDaemon() + + const port = 5104 + const writeSize, writes = 4096, 100 + ln, err := a.Driver.Listen(port) + if err != nil { + t.Fatalf("listen: %v", err) + } + defer ln.Close() + + received := make(chan error, 1) + go func() { + for { + conn, err := ln.Accept() + if err != nil { + return + } + _, err = io.CopyN(io.Discard, conn, writeSize*writes) + conn.Close() + received <- err + } + }() + + // Judged on the fastest of three runs: a timer per write puts a floor of + // 500ms under every run, load only adds. + fastest := time.Hour + buf := make([]byte, writeSize) + for run := 0; run < 3; run++ { + conn, err := b.Driver.DialAddr(a.Daemon.Addr(), port) + if err != nil { + t.Fatalf("dial: %v", err) + } + start := time.Now() + for i := 0; i < writes; i++ { + if _, err := conn.Write(buf); err != nil { + t.Fatalf("write %d: %v", i, err) + } + } + select { + case err := <-received: + if err != nil { + t.Fatalf("receive: %v", err) + } + case <-time.After(60 * time.Second): + t.Fatal("timeout receiving") + } + if d := time.Since(start); d < fastest { + fastest = d + } + conn.Close() + } + t.Logf("fastest run: %d writes of %d bytes in %v", writes, writeSize, fastest) + if fastest > 150*time.Millisecond { + t.Errorf("%d writes of %d bytes took %v at best, want under 150ms: each write is stalling on a delayed ACK", writes, writeSize, fastest) + } +}