From c98e7ba21ea719e0e83c31624f851490ea023130 Mon Sep 17 00:00:00 2001 From: "Tomi P. Hakala" <7030001+tphakala@users.noreply.github.com> Date: Sat, 26 Sep 2026 09:39:06 +0300 Subject: [PATCH 1/2] fix(etcp): count progress ahead of a probe as liveness and keep the catchup until the server speaks Liveness: the watcher took only inbound frames as proof of life, and the probe queues behind any unsent backlog. On a slow uplink with the server silent, a healthy link was declared dead; each reconnect replayed the backlog and the probe answer queued behind it again, so the upload could loop without finishing. At 1 KiB/s, 10 of 96 KiB arrived in an hour of fake time. The writer now sends each batch in 4 KiB pieces and signals liveness for a piece written while a probe still waits behind it: the echo cannot come before the probe is sent, and once the send buffer is full a write completes only as the peer acknowledges data. Writes with no probe behind them, the probe itself included, never count, so neither keystrokes nor large packets written into a dead link keep it alive. Two limits remain and are documented on Dialer.KeepAlive: an uplink under about 400 B/s, and data already in the kernel send buffer, which drains out of sight. Replay trim: recover trimmed the ring to ReplayLimit right after writing our catchup, but the server counts a catchup only once it has decoded the whole message. A link lost in between made the next recover fail with ErrReplayExceeded. The ring now holds everything from the sequence the peer acknowledged until the server's first frame on the new link, which upstream sends only after decoding our catchup (BackedWriter::write waits on the recover mutex, src/base/BackedWriter.cpp:17-18 and src/base/Connection.cpp:109,134-142 at et-v7.0.0). The hold is capped at twice ReplayLimit of written bytes, beyond which the next recover fails with ErrReplayExceeded as before. BenchmarkDrainBacklog shows no significant change against main (benchstat p=0.57). Verified: go test -race ./... (etcp also -count=20), golangci-lint (linux and windows), go fix -diff, wasm vet, and the e2e suite against etserver on localhost. Each new test was checked by removing the production line it pins and watching it fail. --- internal/etcp/dialer.go | 31 ++++--- internal/etcp/drain_bench_test.go | 2 +- internal/etcp/helpers_test.go | 53 ++++++++++++ internal/etcp/link.go | 98 ++++++++++++++++++---- internal/etcp/link_internal_test.go | 4 +- internal/etcp/liveness_test.go | 122 +++++++++++++++++++++++++++- internal/etcp/recover.go | 7 +- internal/etcp/recover_test.go | 108 ++++++++++++++++++++++-- internal/etcp/ring.go | 10 ++- internal/etcp/ring_test.go | 17 ++++ internal/etcp/throttle_test.go | 20 +++-- 11 files changed, 424 insertions(+), 48 deletions(-) diff --git a/internal/etcp/dialer.go b/internal/etcp/dialer.go index af79354..a5eb0f6 100644 --- a/internal/etcp/dialer.go +++ b/internal/etcp/dialer.go @@ -69,10 +69,15 @@ type Dialer struct { DialContext(ctx context.Context, network, address string) (net.Conn, error) } // KeepAlive is the quiet period after which a probe is sent; after two - // quiet periods the link is declared dead. Only inbound frames count as - // proof of life, and the probe queues behind any unsent backlog, so an - // upload that takes longer than two periods to drain while the server - // sends nothing costs a reconnect (the data survives it). Zero means 5 s + // quiet periods the link is declared dead. Inbound frames count as proof + // of life, and so do our writes while a probe waits behind them: the + // probe queues behind the unsent backlog, so a slow upload while the + // server sends nothing would otherwise be taken for a dead link. Writes + // are observed in 4 KiB pieces, so an uplink too slow to send one piece + // in two periods (under about 400 B/s at the default) still costs a + // reconnect, and so does data the kernel has already accepted, which + // drains out of sight: a send buffer holding more than two periods' + // worth ahead of the probe delays its echo past the deadline. Zero means 5 s // (upstream's maximum, src/base/Headers.hpp:180 at et-v7.0.0); Dial // refuses a negative value or one below 100 ms. Probing starts only // after the first WritePacket, because etserver aborts when a session's @@ -90,13 +95,17 @@ type Dialer struct { // bound applied separately to the two kinds of retained packets: packets // already written to a socket are trimmed down to it, and WritePacket // blocks while the not-yet-written backlog exceeds it (a single packet - // may take the backlog over, and probes are never blocked). While - // disconnected the ring can therefore hold about twice ReplayLimit plus - // one packet. Packets count as sent once written to the socket, and - // written packets are trimmed first, so a window smaller than what the - // kernel and the network can hold in flight turns a reconnect into - // ErrReplayExceeded; so does a catchup too large for one handshake - // message. Small values are for tests only. + // may take the backlog over, and probes are never blocked). After a + // reconnect, written packets the server has not acknowledged are kept + // until it first sends on the new link, up to twice ReplayLimit: the + // server counts our catchup only once it has decoded all of it, so a + // link lost before then must be able to replay it again. The ring can + // therefore hold about three times ReplayLimit plus one packet. Packets + // count as sent once written to the socket, and written packets are + // trimmed first, so a window smaller than what the kernel and the network + // can hold in flight turns a reconnect into ErrReplayExceeded; so does a + // catchup too large for one handshake message. Small values are for + // tests only. ReplayLimit int // Logger receives connection events. Nil discards them. Logger *slog.Logger diff --git a/internal/etcp/drain_bench_test.go b/internal/etcp/drain_bench_test.go index d4e65c5..84f31a6 100644 --- a/internal/etcp/drain_bench_test.go +++ b/internal/etcp/drain_bench_test.go @@ -36,7 +36,7 @@ func BenchmarkDrainBacklog(b *testing.B) { } ctx, cancel := context.WithCancel(b.Context()) w := &countingWriter{want: entries * (4 + 18), done: cancel} - if err := c.writeLoop(ctx, w); err == nil || ctx.Err() == nil { + if err := c.writeLoop(ctx, newLink(), w); err == nil || ctx.Err() == nil { b.Fatalf("writeLoop = %v before draining %d bytes (took %d)", err, w.want, w.n) } c.cancel(nil) diff --git a/internal/etcp/helpers_test.go b/internal/etcp/helpers_test.go index 016d1e2..ffaa1fe 100644 --- a/internal/etcp/helpers_test.go +++ b/internal/etcp/helpers_test.go @@ -297,3 +297,56 @@ func (s *rawServer) open(b []byte) (protocol.Packet, error) { } return protocol.Packet{Header: h, Payload: plain}, nil } + +// lostCatchup completes the exchange from the client's point of view, then +// drops the link as if it died before the server counted the client's +// catchup. +func (s *rawServer) lostCatchup() { + if !s.startRecover() { + return + } + var theirs protocol.CatchupBuffer + _ = wire.ReadMessage(s.br, &theirs) +} + +// recoverAll completes the exchange and returns the numbers of the packets +// in the client's catchup, or nil when the exchange fails. +func (s *rawServer) recoverAll() []int { + if !s.startRecover() { + return nil + } + var theirs protocol.CatchupBuffer + if wire.ReadMessage(s.br, &theirs) != nil { + return nil + } + nums := []int{} + for _, b := range theirs.GetBuffer() { + p, err := s.open(b) + if err != nil { + return nil + } + n, err := number(p) + if err != nil { + return nil + } + nums = append(nums, n) + } + return nums +} + +// startRecover answers RETURNING_CLIENT and runs the server's half of the +// recover exchange up to reading the client's catchup, in upstream's order, +// claiming to have received nothing and sending an empty catchup. +func (s *rawServer) startRecover() bool { + if s.respond(protocol.ConnectStatus_RETURNING_CLIENT) != nil { + return false + } + if wire.WriteMessage(s.conn, &protocol.SequenceHeader{}) != nil { + return false + } + var mine protocol.SequenceHeader + if wire.ReadMessage(s.br, &mine) != nil { + return false + } + return wire.WriteMessage(s.conn, &protocol.CatchupBuffer{}) == nil +} diff --git a/internal/etcp/link.go b/internal/etcp/link.go index 8fb212f..7ddb58c 100644 --- a/internal/etcp/link.go +++ b/internal/etcp/link.go @@ -25,6 +25,10 @@ const ( // frame more than this many into writeBufSize. Taking the whole backlog // instead would copy it under c.mu on every pass, quadratic in its length. maxBatchEntries = writeBufSize/(4+2+secretbox.Overhead) + 1 + + // progressChunk is the piece size the writer sends a batch in, so that + // a slow uplink shows progress well within one keepAlive period. + progressChunk = 4 << 10 ) var ( @@ -34,10 +38,12 @@ var ( // link is the per-TCP-connection state shared by its goroutines. type link struct { - alive chan struct{} // cap 1: a frame arrived + alive chan struct{} // cap 1: a frame arrived, or a write progressed ahead of a waiting probe delivering atomic.Bool // the reader is blocked handing a packet to ReadPacket } +func newLink() *link { return &link{alive: make(chan struct{}, 1)} } + // runLink runs a reader, a writer and a liveness watcher on nc until one of // them fails, then joins all three and returns the cause. func (c *Conn) runLink(nc net.Conn, catchup [][]byte) error { @@ -46,10 +52,10 @@ func (c *Conn) runLink(nc net.Conn, catchup [][]byte) error { stop := context.AfterFunc(ctx, func() { _ = nc.Close() }) defer stop() - l := &link{alive: make(chan struct{}, 1)} + l := newLink() var wg sync.WaitGroup wg.Go(func() { cancel(c.readLoop(ctx, l, nc, catchup)) }) - wg.Go(func() { cancel(c.writeLoop(ctx, nc)) }) + wg.Go(func() { cancel(c.writeLoop(ctx, l, nc)) }) wg.Go(func() { cancel(c.watch(ctx, l)) }) wg.Wait() _ = nc.Close() @@ -57,7 +63,8 @@ func (c *Conn) runLink(nc net.Conn, catchup [][]byte) error { } // readLoop hands over packets an earlier link left pending, then delivers the -// peer's catchup, then frames from the link. +// peer's catchup, then frames from the link. The first frame ends the hold +// recover placed on our catchup (see releaseHold). func (c *Conn) readLoop(ctx context.Context, l *link, r io.Reader, catchup [][]byte) error { for len(c.pending) > 0 { if err := c.handOver(ctx, l, c.pending[0]); err != nil { @@ -73,7 +80,10 @@ func (c *Conn) readLoop(ctx context.Context, l *link, r io.Reader, catchup [][]b } } br := bufio.NewReaderSize(r, readBufSize) - var buf []byte + var ( + buf []byte + released bool + ) for { frame, err := wire.ReadFrame(br, buf) if err != nil { @@ -84,6 +94,10 @@ func (c *Conn) readLoop(ctx context.Context, l *link, r io.Reader, catchup [][]b } buf = frame signal(l.alive) + if !released { + c.releaseHold() + released = true + } if err := c.deliver(ctx, l, frame); err != nil { return err } @@ -138,20 +152,44 @@ func (c *Conn) handOver(ctx context.Context, l *link, p protocol.Packet) error { } } +// releaseHold ends the hold recover placed on our catchup and trims the ring +// back to ReplayLimit. It is called on a link's first frame: upstream writes +// on a recovered socket only after it has decoded our whole catchup, since +// BackedWriter::write waits on the recover mutex that Connection::recover +// holds until then (src/base/BackedWriter.cpp:17-18 and +// src/base/Connection.cpp:109,134-142 at et-v7.0.0). +func (c *Conn) releaseHold() { + c.mu.Lock() + defer c.mu.Unlock() + if c.ring.held { + c.ring.held = false + c.ring.trim(c.flushed, c.unsent) + } +} + // writeLoop sends ring entries from flushed onwards, trimming the ring and // releasing blocked writers as the backlog drains. It frames a batch of // entries (at least one, then up to writeBufSize bytes) into one reused -// buffer and sends it with a single Write, so the steady state allocates +// buffer and sends it in progressChunk pieces, so the steady state allocates // nothing per packet. Entries are immutable once sealed, so they are framed // outside the lock. -func (c *Conn) writeLoop(ctx context.Context, w io.Writer) error { +// +// A piece written while a probe still waits behind it signals l.alive. The +// probe's echo cannot come before the probe is sent, and once the socket's +// send buffer is full a write completes only as the peer acknowledges data, +// so a backlog draining ahead of the probe proves the path is alive while +// the server sends nothing. Other writes prove nothing (a dead link's send +// buffer still takes them), so data written with no probe behind it, the +// probe itself included, never counts. +func (c *Conn) writeLoop(ctx context.Context, l *link, w io.Writer) error { var ( batch [][]byte buf []byte ) for { c.mu.Lock() - batch = c.ring.appendRange(batch[:0], c.flushed, min(c.ring.next(), c.flushed+maxBatchEntries)) + base, probe := c.flushed, c.lastProbe + batch = c.ring.appendRange(batch[:0], base, min(c.ring.next(), base+maxBatchEntries)) c.mu.Unlock() if len(batch) == 0 { select { @@ -163,10 +201,14 @@ func (c *Conn) writeLoop(ctx context.Context, w io.Writer) error { } buf = buf[:0] sent, n := 0, 0 + probeAt := -1 // offset in buf of the probe's frame, if it is in this batch for _, f := range batch { if sent > 0 && len(buf)+4+len(f) > writeBufSize { break } + if base+int64(sent) == probe { + probeAt = len(buf) + } var err error // WritePacket refuses packets above wire.MaxFrameSize, so an // error here means the ring is corrupt and no link can help. @@ -177,13 +219,20 @@ func (c *Conn) writeLoop(ctx context.Context, w io.Writer) error { n += len(f) } clear(batch) // drop references so trimmed entries can be collected - // A short count without an error breaks the io.Writer contract; - // counting the whole batch as sent would drop its tail from replay. - if m, err := w.Write(buf); err != nil || m != len(buf) { - if err == nil { - err = io.ErrShortWrite + for written := 0; written < len(buf); { + piece := buf[written:min(len(buf), written+progressChunk)] + // A short count without an error breaks the io.Writer contract; + // counting the whole batch as sent would drop its tail from replay. + if m, err := w.Write(piece); err != nil || m != len(piece) { + if err == nil { + err = io.ErrShortWrite + } + return fmt.Errorf("etcp: write: %w", err) + } + written += len(piece) + if c.probeBehind(base+int64(sent), probeAt, written) { + signal(l.alive) } - return fmt.Errorf("etcp: write: %w", err) } c.mu.Lock() c.flushed += int64(sent) @@ -197,10 +246,25 @@ func (c *Conn) writeLoop(ctx context.Context, w io.Writer) error { } } +// probeBehind reports whether a probe still waits behind the first written +// bytes of the current batch: in the batch past them (probeAt is its frame's +// offset, or -1), or queued after the batch, whose entries end before +// sequence end. Only the watcher queues probes, so one queued after the batch +// was framed is always behind it. +func (c *Conn) probeBehind(end int64, probeAt, written int) bool { + if probeAt >= 0 { + return written <= probeAt + } + c.mu.Lock() + defer c.mu.Unlock() + return c.lastProbe >= end +} + // watch declares the link dead after two quiet keepAlive periods, sending the -// probe after the first. Time the reader spends blocked on a slow ReadPacket -// caller does not count as silence, and neither does the time before the -// caller's first packet, when no probe may be sent. +// probe after the first. A frame from the server, or writer progress ahead +// of a waiting probe (see writeLoop), breaks the quiet. Time the reader spends blocked on +// a slow ReadPacket caller does not count as silence, and neither does the +// time before the caller's first packet, when no probe may be sent. func (c *Conn) watch(ctx context.Context, l *link) error { t := time.NewTimer(c.keepAlive) defer t.Stop() diff --git a/internal/etcp/link_internal_test.go b/internal/etcp/link_internal_test.go index d0c7e49..a210c45 100644 --- a/internal/etcp/link_internal_test.go +++ b/internal/etcp/link_internal_test.go @@ -35,7 +35,7 @@ func TestWriteLoopShortWrite(t *testing.T) { // more data forever. ctx, cancel := context.WithTimeout(t.Context(), time.Minute) defer cancel() - err := c.writeLoop(ctx, shortWriter{}) + err := c.writeLoop(ctx, newLink(), shortWriter{}) if !errors.Is(err, io.ErrShortWrite) { t.Fatalf("writeLoop = %v, want io.ErrShortWrite", err) } @@ -104,7 +104,7 @@ func TestWriteLoopKeepsWrittenWhileBacklogFull(t *testing.T) { ctx, cancel := context.WithTimeout(t.Context(), time.Minute) defer cancel() done := make(chan error, 1) - go func() { done <- c.writeLoop(ctx, &backlogWriter{ctx: ctx, c: c}) }() + go func() { done <- c.writeLoop(ctx, newLink(), &backlogWriter{ctx: ctx, c: c}) }() synctest.Wait() // the first batch is written, the second is stuck c.mu.Lock() diff --git a/internal/etcp/liveness_test.go b/internal/etcp/liveness_test.go index 2f4a6fd..c5f875d 100644 --- a/internal/etcp/liveness_test.go +++ b/internal/etcp/liveness_test.go @@ -141,6 +141,84 @@ func TestLivenessDetectsDeadLink(t *testing.T) { }) } +// Writes that complete with nothing behind them prove nothing: a dead link's +// send buffer still takes them. Typing into a link whose server never answers +// must not keep it alive. +func TestLivenessTypingDoesNotHideDeadLink(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + h := newHarness(t, etcp.Dialer{}) + defer h.close() + h.srv.EchoKeepAlive(false) + + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan struct{}) + defer func() { + cancel() + <-done + }() + go func() { // a keystroke every 2 s + defer close(done) + for i := 0; ctx.Err() == nil; i++ { + if h.conn.WritePacket(ctx, numbered(i, 10)) != nil { + return + } + select { + case <-time.After(2 * time.Second): + case <-ctx.Done(): + } + } + }() + + // Probe at 5 s, dead at 10 s, immediate redial. + synctest.Sleep(9 * time.Second) + if got := h.net.Dials(); got != 1 { + t.Fatalf("Dials() = %d at 9 s, want 1", got) + } + synctest.Sleep(2 * time.Second) + if got := h.net.Dials(); got != 2 { + t.Fatalf("Dials() = %d at 11 s, want 2: keystrokes kept a dead link alive", got) + } + }) +} + +// Packets larger than one write piece prove nothing either while no probe +// waits behind them: the send buffer takes every piece at once. +func TestLivenessLargeWritesDoNotHideDeadLink(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + h := newHarness(t, etcp.Dialer{}) + defer h.close() + h.srv.EchoKeepAlive(false) + + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan struct{}) + defer func() { + cancel() + <-done + }() + go func() { // a 5 KiB packet every 3 s + defer close(done) + for i := 0; ctx.Err() == nil; i++ { + if h.conn.WritePacket(ctx, numbered(i, 5<<10)) != nil { + return + } + select { + case <-time.After(3 * time.Second): + case <-ctx.Done(): + } + } + }() + + synctest.Sleep(9 * time.Second) + if got := h.net.Dials(); got != 1 { + t.Fatalf("Dials() = %d at 9 s, want 1", got) + } + synctest.Sleep(2 * time.Second) + if got := h.net.Dials(); got != 2 { + t.Fatalf("Dials() = %d at 11 s, want 2: large writes kept a dead link alive", got) + } + }) +} + // While the reader is blocked on a caller that is not reading, the link is // not declared dead: that silence is ours, not the network's. func TestLivenessIgnoresSlowReader(t *testing.T) { @@ -189,10 +267,9 @@ func TestDialRejectsOversizedProbe(t *testing.T) { } // A long upload over a slow uplink queues the probe behind the backlog, so -// its echo comes late while the server itself sends nothing, and the -// watcher may drop the link (the KeepAlive doc states this limit). Whatever -// reconnects that costs, every packet must still arrive exactly once and in -// order. +// its echo comes late while the server itself sends nothing. Whatever +// reconnects that may cost, every packet must still arrive exactly once and +// in order. func TestLivenessSlowUploadDeliversEverything(t *testing.T) { synctest.Test(t, func(t *testing.T) { srv := etservertest.NewServer(testID, testKey) @@ -218,3 +295,40 @@ func TestLivenessSlowUploadDeliversEverything(t *testing.T) { } }) } + +// A slow uplink that is still moving is a live link, even when one socket +// write takes longer than two keepalive periods and the server says nothing +// meanwhile. Without counting write progress each reconnect replays the +// backlog, the probe's echo queues behind it again, and the upload may never +// finish. +func TestLivenessSlowUploadKeepsLink(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + srv := etservertest.NewServer(testID, testKey) + nw := etservertest.NewNetwork(srv) + defer nw.Close() + d := etcp.Dialer{NetDialer: throttledDialer{inner: nw, perKiB: time.Second}} + conn, err := d.Dial(t.Context(), testAddr, testID, testKey) + if err != nil { + t.Fatalf("Dial: %v", err) + } + defer func() { _ = conn.Close() }() + + // 160 KiB at 1 KiB/s: one 64 KiB batch alone takes over a minute, and + // the probe from the first quiet period, queued after all of it, is + // framed behind tens of KiB of its own batch. + const n = 160 + for i := range n { + if err := conn.WritePacket(t.Context(), numbered(i, 1024)); err != nil { + t.Fatalf("WritePacket %d: %v", i, err) + } + } + ctx, cancel := context.WithTimeout(t.Context(), time.Hour) + defer cancel() + if err := expectNumbered(ctx, n, srv.Recv); err != nil { + t.Fatalf("server side: %v", err) + } + if got := nw.Dials(); got != 1 { + t.Fatalf("Dials() = %d, want 1: a link that keeps accepting data is alive", got) + } + }) +} diff --git a/internal/etcp/recover.go b/internal/etcp/recover.go index 195c53f..bae0478 100644 --- a/internal/etcp/recover.go +++ b/internal/etcp/recover.go @@ -70,7 +70,8 @@ func (c *Conn) reconnect(b *backoff) (net.Conn, [][]byte, error) { // the connection cannot buffer what both sides write (at once over net.Pipe, // and over TCP once both catchups exceed the socket buffers). Whichever side // fails first closes conn to unblock the other; readerErrWins decides which -// error is reported. +// error is reported. On success the ring holds our catchup until the server +// first speaks on the new link (see releaseHold). func (c *Conn) recover(conn net.Conn) ([][]byte, error) { var ( peer protocol.SequenceHeader @@ -109,6 +110,10 @@ func (c *Conn) recover(conn net.Conn) ([][]byte, error) { c.mu.Lock() c.unsent -= c.ring.bytesBetween(c.flushed, snap) c.flushed = snap + // The server counts our catchup only once it has decoded the whole + // message, so until it speaks on this link the next SequenceHeader may + // ask for it again: hold everything from the position it acknowledged. + c.ring.held, c.ring.hold = true, int64(peer.GetSequenceNumber()) // Trim as writeLoop does after a Write: a link that dies before its // first Write would otherwise leave the catchup in the ring while // WritePacket admits another ReplayLimit. diff --git a/internal/etcp/recover_test.go b/internal/etcp/recover_test.go index 6630c48..5d20a26 100644 --- a/internal/etcp/recover_test.go +++ b/internal/etcp/recover_test.go @@ -299,9 +299,9 @@ func TestWritePacketRacingRecovery(t *testing.T) { } // A link that recovers and then dies before writing anything must not let the -// replay ring grow past ReplayLimit: recover counts our catchup as sent, which +// replay ring grow without bound: recover counts our catchup as sent, which // frees WritePacket to admit another ReplayLimit of packets, so recover must -// also trim what it has now written. The server here reports its +// also trim what the server has acknowledged. The server here reports its // true received count in each SequenceHeader and drops every link right after // the exchange, while the caller keeps writing. func TestFlappingLinkKeepsRingBounded(t *testing.T) { @@ -378,10 +378,108 @@ func TestFlappingLinkKeepsRingBounded(t *testing.T) { if lastCatchup.Load() == 0 { t.Fatal("writer stalled: the last recover carried no catchup") } - // At most ReplayLimit of written entries survive a trim, plus the - // unsent backlog WritePacket admits: ReplayLimit and one packet. - if got, bound := etcp.RingBytes(conn), 2*limit+sealed; got > bound { + // This server decodes every catchup and acknowledges it on the next + // link, so recover's hold keeps only the latest catchup, which is at + // most the backlog WritePacket admits (ReplayLimit and one packet). + // Add the new backlog admitted since. + if got, bound := etcp.RingBytes(conn), 2*(limit+sealed); got > bound { t.Fatalf("ring holds %d bytes after %d links, want at most %d", got, s.dials.Load(), bound) } }) } + +// A recover succeeds on our side once the server has read our catchup, but +// the server counts it only after decoding the whole message, and it starts +// writing on the new link only after that (src/base/Connection.cpp:134-142 +// and src/base/BackedWriter.cpp:17-18 at et-v7.0.0). If the link dies in +// between, the next SequenceHeader asks for the same catchup again, so it +// must still be in the ring even though it exceeds ReplayLimit. +func TestCatchupKeptUntilServerSpeaks(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + const ( + limit = 1 << 10 + n = 5 // 5 packets of about 218 sealed bytes: over the limit + ) + got := make(chan []int, 1) + s := &scripted{handle: func(i int, c *rawServer) { + switch i { + case 0: + // Accept the link and read nothing, so every packet below + // is still unsent when it drops. + if c.respond(protocol.ConnectStatus_NEW_CLIENT) == nil { + time.Sleep(time.Second) + } + case 1: + c.lostCatchup() + case 2: + got <- c.recoverAll() + c.drain() + } + }} + d := etcp.Dialer{NetDialer: s, ReplayLimit: limit} + conn, err := d.Dial(t.Context(), testAddr, testID, testKey) + if err != nil { + t.Fatalf("Dial: %v", err) + } + defer func() { + _ = conn.Close() + s.wg.Wait() + }() + for i := range n { + if err := conn.WritePacket(t.Context(), numbered(i, 200)); err != nil { + t.Fatalf("WritePacket %d: %v", i, err) + } + } + nums := within(t, got) + if nums == nil { + // The exchange failed on our side; the Conn says why. + _, err := readPacket(t, conn) + t.Fatalf("second recover failed; ReadPacket = %v", err) + } + if len(nums) != n || nums[0] != 0 || nums[n-1] != n-1 { + t.Fatalf("second recover carried packets %v, want 0..%d", nums, n-1) + } + }) +} + +// The server writes on a recovered link only after it has decoded our whole +// catchup (src/base/BackedWriter.cpp:17-18 at et-v7.0.0), so its first +// packet there ends the hold, and the ring goes back to ReplayLimit. +func TestServerPacketReleasesCatchup(t *testing.T) { + synctest.Test(t, func(t *testing.T) { + const limit = 1 << 10 + h := newHarness(t, etcp.Dialer{ReplayLimit: limit}) + defer h.close() + + synctest.Wait() + h.net.SetRefuse(true) + h.net.CutAll() + synctest.Wait() + const n = 5 // about 1090 sealed bytes: over the limit + for i := range n { + if err := h.conn.WritePacket(t.Context(), numbered(i, 200)); err != nil { + t.Fatalf("WritePacket %d: %v", i, err) + } + } + h.net.SetRefuse(false) + ctx, cancel := context.WithTimeout(t.Context(), time.Second) + defer cancel() + if err := expectNumbered(ctx, n, h.srv.Recv); err != nil { + t.Fatalf("server side: %v", err) + } + synctest.Wait() + if got := etcp.RingBytes(h.conn); got <= limit { + t.Fatalf("ring holds %d bytes before the server spoke, want the whole catchup (over %d)", got, limit) + } + if err := h.srv.Send(ctx, numbered(0, 10)); err != nil { + t.Fatalf("Send: %v", err) + } + if err := expectNumbered(ctx, 1, h.conn.ReadPacket); err != nil { + t.Fatalf("client side: %v", err) + } + synctest.Wait() + if got := etcp.RingBytes(h.conn); got > limit { + t.Fatalf("ring holds %d bytes after the server spoke, want at most %d", got, limit) + } + }) +} diff --git a/internal/etcp/ring.go b/internal/etcp/ring.go index 061a359..251756b 100644 --- a/internal/etcp/ring.go +++ b/internal/etcp/ring.go @@ -10,6 +10,10 @@ type ring struct { entries [][]byte bytes int limit int + // While held, trim keeps entries at or after hold, the sequence the + // peer last acknowledged, until written entries exceed twice limit. + held bool + hold int64 } // next is the sequence number the next pushed packet gets, which is also the @@ -25,10 +29,14 @@ func (r *ring) push(data []byte) { // bytes, never dropping an entry at or after keep (not yet written to any // link). unsent is the size of the entries from keep on; the limit applies // to written entries alone, so a full backlog cannot squeeze out the replay -// copies of packets still in flight. +// copies of packets still in flight. While held, entries at or after hold +// are dropped only once written entries exceed twice limit. func (r *ring) trim(keep int64, unsent int) { n := 0 for n < len(r.entries) && r.bytes-unsent > r.limit && r.first+int64(n) < keep { + if r.held && r.first+int64(n) >= r.hold && r.bytes-unsent <= 2*r.limit { + break + } r.bytes -= len(r.entries[n]) r.entries[n] = nil n++ diff --git a/internal/etcp/ring_test.go b/internal/etcp/ring_test.go index a3f5d28..6b1b90b 100644 --- a/internal/etcp/ring_test.go +++ b/internal/etcp/ring_test.go @@ -36,6 +36,23 @@ func TestRingTrimKeepsUnsent(t *testing.T) { } } +// While held, trim keeps the entries the peer has not acknowledged, but only +// up to twice limit of written bytes: a server that never catches up must +// not grow the ring without bound. +func TestRingTrimHoldCeiling(t *testing.T) { + r := filledRing(25, 10, 10) // 100 bytes held, limit 25 + r.held, r.hold = true, 3 + r.trim(10, 0) + if r.first != 5 || r.bytes != 50 { + t.Fatalf("first, bytes = %d, %d; want 5, 50 (trim unacknowledged entries only down to twice the limit)", r.first, r.bytes) + } + r.hold = 7 + r.trim(10, 0) + if r.first != 7 { + t.Fatalf("first = %d, want 7 (acknowledged entries trim down to the limit, not below the hold)", r.first) + } +} + func TestRingSince(t *testing.T) { r := filledRing(25, 10, 10) r.trim(10, 0) // retains 8 and 9 diff --git a/internal/etcp/throttle_test.go b/internal/etcp/throttle_test.go index 24ea488..ad0673f 100644 --- a/internal/etcp/throttle_test.go +++ b/internal/etcp/throttle_test.go @@ -11,14 +11,15 @@ import ( "github.com/tphakala/et-go/internal/etservertest" ) -// throttledDialer slows every client-side write to 1 KiB per 125 ms of fake -// time (8 KiB/s), like a congested uplink, or with downlink set every -// client-side read instead, like a congested downlink. +// throttledDialer slows every client-side write to 1 KiB per perKiB of fake +// time (zero means 125 ms, 8 KiB/s), like a congested uplink, or with +// downlink set every client-side read instead, like a congested downlink. type throttledDialer struct { inner interface { DialContext(ctx context.Context, network, address string) (net.Conn, error) } downlink bool + perKiB time.Duration } func (d throttledDialer) DialContext(ctx context.Context, network, address string) (net.Conn, error) { @@ -29,7 +30,11 @@ func (d throttledDialer) DialContext(ctx context.Context, network, address strin if d.downlink { return slowReadConn{c}, nil } - return throttledConn{c}, nil + perKiB := d.perKiB + if perKiB == 0 { + perKiB = 125 * time.Millisecond + } + return throttledConn{Conn: c, perKiB: perKiB}, nil } // slowReadConn reads at most 1 KiB per 125 ms of fake time. @@ -43,7 +48,10 @@ func (c slowReadConn) Read(p []byte) (int, error) { return n, err } -type throttledConn struct{ net.Conn } +type throttledConn struct { + net.Conn + perKiB time.Duration +} func (c throttledConn) Write(p []byte) (int, error) { var n int @@ -54,7 +62,7 @@ func (c throttledConn) Write(p []byte) (int, error) { return n, err } p = p[m:] - time.Sleep(125 * time.Millisecond) + time.Sleep(c.perKiB) } return n, nil } From 753ce5255d700c5e636f5526e606832f1fb8a852 Mon Sep 17 00:00:00 2001 From: "Tomi P. Hakala" <7030001+tphakala@users.noreply.github.com> Date: Sat, 26 Sep 2026 10:00:04 +0300 Subject: [PATCH 2/2] docs(etcp): clarify the throttle rate and document newLink --- internal/etcp/link.go | 1 + internal/etcp/throttle_test.go | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/internal/etcp/link.go b/internal/etcp/link.go index 7ddb58c..39ef6cb 100644 --- a/internal/etcp/link.go +++ b/internal/etcp/link.go @@ -42,6 +42,7 @@ type link struct { delivering atomic.Bool // the reader is blocked handing a packet to ReadPacket } +// newLink returns the shared state for one TCP connection's goroutines. func newLink() *link { return &link{alive: make(chan struct{}, 1)} } // runLink runs a reader, a writer and a liveness watcher on nc until one of diff --git a/internal/etcp/throttle_test.go b/internal/etcp/throttle_test.go index ad0673f..e0d74d1 100644 --- a/internal/etcp/throttle_test.go +++ b/internal/etcp/throttle_test.go @@ -11,7 +11,7 @@ import ( "github.com/tphakala/et-go/internal/etservertest" ) -// throttledDialer slows every client-side write to 1 KiB per perKiB of fake +// throttledDialer slows every client-side write to 1 KiB every perKiB of fake // time (zero means 125 ms, 8 KiB/s), like a congested uplink, or with // downlink set every client-side read instead, like a congested downlink. type throttledDialer struct {