From e13bb4585442d7bc524f092e62bf07a34bdfefa5 Mon Sep 17 00:00:00 2001 From: Pat Whittingslow Date: Mon, 19 Jan 2026 23:22:54 -0300 Subject: [PATCH] attempt to fix #22 (#23) --- tcp/handler.go | 4 ++ tcp/handler_test.go | 119 ++++++++++++++++++++++++++++++++++++++++++++ tcp/txqueue.go | 6 +-- tcp/txqueue_test.go | 14 +++--- 4 files changed, 133 insertions(+), 10 deletions(-) diff --git a/tcp/handler.go b/tcp/handler.go index ce553d0..33ffe13 100644 --- a/tcp/handler.go +++ b/tcp/handler.go @@ -194,6 +194,10 @@ func (h *Handler) Recv(incomingPacket []byte) error { return err } } + if segIncoming.Flags.HasAny(FlagACK) { + // Update TX ring buffer to free up acked data. + h.bufTx.RecvACK(segIncoming.ACK) + } if segIncoming.Flags.HasAny(FlagSYN) && h.remotePort == 0 { // Remote reached out and has given us their port, set it on our side. h.debug("tcp.Handler:rx-remoteport-set", slog.Uint64("port", uint64(h.localPort)), slog.Uint64("remoteport", uint64(remotePort))) diff --git a/tcp/handler_test.go b/tcp/handler_test.go index f799eda..c5f12c3 100644 --- a/tcp/handler_test.go +++ b/tcp/handler_test.go @@ -148,6 +148,125 @@ func clear[E any, T []E](s T) { } } +// TestTxBufferFreedOnACK tests that the TX buffer is freed when ACKs are received. +// This is a regression test for https://github.com/soypat/lneto/issues/22 +// where ringTx.sentoff and ringTx.sentend were not being updated when ACKs +// were received, causing AvailableOutput() to return 0 indefinitely after +// the initial buffer was consumed. +func TestTxBufferFreedOnACK(t *testing.T) { + const mtu = 256 + const maxpackets = 4 + const txBufSize = 128 // Small TX buffer to easily fill it + rng := rand.New(rand.NewSource(42)) + + // Create handlers with small TX buffers to easily trigger the issue. + client := new(Handler) + server := new(Handler) + err := client.SetBuffers(make([]byte, txBufSize), make([]byte, mtu), maxpackets) + if err != nil { + t.Fatal(err) + } + err = server.SetBuffers(make([]byte, txBufSize), make([]byte, mtu), maxpackets) + if err != nil { + t.Fatal(err) + } + + // Setup and establish connection. + err = server.OpenListen(uint16(rng.Uint32()), 0) + if err != nil { + t.Fatal(err) + } + err = client.OpenActive(uint16(rng.Uint32()), server.LocalPort(), 0) + if err != nil { + t.Fatal(err) + } + + var rawbuf [mtu]byte + establish(t, client, server, rawbuf[:]) + + // Record initial available space. + initialAvailable := client.AvailableOutput() + if initialAvailable == 0 { + t.Fatal("expected non-zero initial available output") + } + + // Write data to fill a significant portion of the TX buffer. + data := make([]byte, txBufSize/2) + for i := range data { + data[i] = byte(i) + } + n, err := client.Write(data) + if err != nil { + t.Fatal("client write:", err) + } else if n != len(data) { + t.Fatalf("expected to write %d bytes, wrote %d", len(data), n) + } + + // Available space should have decreased. + afterWriteAvailable := client.AvailableOutput() + if afterWriteAvailable >= initialAvailable { + t.Fatalf("expected available to decrease after write: before=%d, after=%d", + initialAvailable, afterWriteAvailable) + } + + // Client sends DATA packet. + clear(rawbuf[:]) + n, err = client.Send(rawbuf[:]) + if err != nil { + t.Fatal("client sending data:", err) + } + if n < len(data)+sizeHeaderTCP { + t.Fatal("expected client to send full data packet") + } + dataPacket := append([]byte(nil), rawbuf[:n]...) + + // After sending, data moves from "unsent" to "sent" - available should still be reduced + // until we receive an ACK. + afterSendAvailable := client.AvailableOutput() + + // Server receives DATA. + err = server.Recv(dataPacket) + if err != nil { + t.Fatal("server receiving data:", err) + } + + // Server sends ACK. + clear(rawbuf[:]) + n, err = server.Send(rawbuf[:]) + if err != nil { + t.Fatal("server sending ACK:", err) + } + ackPacket := append([]byte(nil), rawbuf[:n]...) + + // Client receives ACK - this is where the bug manifests. + // Without the fix, the TX buffer's sentoff/sentend are not updated, + // so AvailableOutput() remains low. + err = client.Recv(ackPacket) + if err != nil { + t.Fatal("client receiving ACK:", err) + } + + // THE BUG: After receiving ACK, the TX buffer should be freed. + // Without the fix, AvailableOutput() stays at the post-send value. + afterAckAvailable := client.AvailableOutput() + + if afterAckAvailable <= afterSendAvailable { + t.Fatalf("BUG (issue #22): TX buffer not freed after receiving ACK\n"+ + "AvailableOutput() after send: %d\n"+ + "AvailableOutput() after ACK: %d\n"+ + "Expected available space to increase after ACK is received.\n"+ + "The ringTx.sentoff and ringTx.sentend fields are not being updated\n"+ + "because ringTx.RecvACK() is not called when ACKs are received.", + afterSendAvailable, afterAckAvailable) + } + + // Should be back to (approximately) initial available space. + if afterAckAvailable < initialAvailable-10 { // Allow small margin for overhead + t.Fatalf("expected available to return close to initial: initial=%d, afterAck=%d", + initialAvailable, afterAckAvailable) + } +} + // TestBufferNotClearedOnPassiveClose tests that data remains readable after // the TCP connection is closed by the remote peer. This is a regression test // for a bug where the receive buffer was cleared when the connection transitioned diff --git a/tcp/txqueue.go b/tcp/txqueue.go index c1ec4c0..f4ef732 100644 --- a/tcp/txqueue.go +++ b/tcp/txqueue.go @@ -139,7 +139,7 @@ func (rtx *ringTx) MakePacket(b []byte, currentSeq Value) (int, error) { // Start of buffer will be SENT, end of buffer will be UNSENT(or empty). // Packet generated has offset at old unsentOff. size := rtx.Size() - pkt := rtx.slist.AddPacket(n, oldUnsentOff, size) + pkt := rtx.slist.AddPacket(n, oldUnsentOff, size, currentSeq) if pkt.off != oldUnsentOff || pkt.end != addEnd(pkt.off, n, size) { panic("invalid generated packet") } @@ -295,7 +295,7 @@ func (sl *sentlist) Free() int { return cap(sl.pkts) - len(sl.pkts) } -func (sl *sentlist) AddPacket(datalen, off, bufsize int) *ringidx { +func (sl *sentlist) AddPacket(datalen, off, bufsize int, seq Value) *ringidx { free := sl.Free() if free == 0 { panic("pkt buffer full") @@ -307,7 +307,7 @@ func (sl *sentlist) AddPacket(datalen, off, bufsize int) *ringidx { sl.pkts = append(sl.pkts, ringidx{ off: off, end: addEnd(off, datalen, bufsize), - seq: sl.EndSeq(), + seq: seq, size: Size(datalen), }) return &sl.pkts[len(sl.pkts)-1] diff --git a/tcp/txqueue_test.go b/tcp/txqueue_test.go index 78613fe..f10d02c 100644 --- a/tcp/txqueue_test.go +++ b/tcp/txqueue_test.go @@ -138,17 +138,17 @@ func TestSentlist_multi(t *testing.T) { sl.Reset(3, 0) // Test multi packet x2. - p1 := sl.AddPacket(5, 0, bufsize) - p2 := sl.AddPacket(5, p1.end, bufsize) + p1 := sl.AddPacket(5, 0, bufsize, 0) + p2 := sl.AddPacket(5, p1.end, bufsize, p1.endSeq()) sl.RecvAck(Value(p2.size+p1.size), bufsize) if sl.Oldest() != nil { t.Fatal("expected full ack") } // multi packet x3. sl.Reset(3, 0) - p1 = sl.AddPacket(3, 0, bufsize) - p2 = sl.AddPacket(3, p1.end, bufsize) - p3 := sl.AddPacket(4, p2.end, bufsize) + p1 = sl.AddPacket(3, 0, bufsize, 0) + p2 = sl.AddPacket(3, p1.end, bufsize, p1.endSeq()) + p3 := sl.AddPacket(4, p2.end, bufsize, p2.endSeq()) sl.RecvAck(2, bufsize) oldest := sl.Oldest() if oldest != p1 { @@ -167,7 +167,7 @@ func TestSentlist_simple(t *testing.T) { // Test full ack. const bufsize = 16 const pkt = 10 - sl.AddPacket(pkt, 0, bufsize) + sl.AddPacket(pkt, 0, bufsize, 0) if sl.Oldest() == nil || sl.Newest() != sl.Oldest() { t.Error("expected same oldest/newest non-nil packet") } @@ -179,7 +179,7 @@ func TestSentlist_simple(t *testing.T) { } // Test partial ack. - sl.AddPacket(pkt, 0, bufsize) + sl.AddPacket(pkt, 0, bufsize, sl.ssn) for i := Value(0); i < pkt-1; i++ { ack++ sl.RecvAck(ack, bufsize)