mirror of
https://github.com/soypat/lneto.git
synced 2026-08-19 22:24:03 +00:00
@@ -194,6 +194,10 @@ func (h *Handler) Recv(incomingPacket []byte) error {
|
|||||||
return err
|
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 {
|
if segIncoming.Flags.HasAny(FlagSYN) && h.remotePort == 0 {
|
||||||
// Remote reached out and has given us their port, set it on our side.
|
// 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)))
|
h.debug("tcp.Handler:rx-remoteport-set", slog.Uint64("port", uint64(h.localPort)), slog.Uint64("remoteport", uint64(remotePort)))
|
||||||
|
|||||||
@@ -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
|
// TestBufferNotClearedOnPassiveClose tests that data remains readable after
|
||||||
// the TCP connection is closed by the remote peer. This is a regression test
|
// 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
|
// for a bug where the receive buffer was cleared when the connection transitioned
|
||||||
|
|||||||
+3
-3
@@ -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).
|
// Start of buffer will be SENT, end of buffer will be UNSENT(or empty).
|
||||||
// Packet generated has offset at old unsentOff.
|
// Packet generated has offset at old unsentOff.
|
||||||
size := rtx.Size()
|
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) {
|
if pkt.off != oldUnsentOff || pkt.end != addEnd(pkt.off, n, size) {
|
||||||
panic("invalid generated packet")
|
panic("invalid generated packet")
|
||||||
}
|
}
|
||||||
@@ -295,7 +295,7 @@ func (sl *sentlist) Free() int {
|
|||||||
return cap(sl.pkts) - len(sl.pkts)
|
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()
|
free := sl.Free()
|
||||||
if free == 0 {
|
if free == 0 {
|
||||||
panic("pkt buffer full")
|
panic("pkt buffer full")
|
||||||
@@ -307,7 +307,7 @@ func (sl *sentlist) AddPacket(datalen, off, bufsize int) *ringidx {
|
|||||||
sl.pkts = append(sl.pkts, ringidx{
|
sl.pkts = append(sl.pkts, ringidx{
|
||||||
off: off,
|
off: off,
|
||||||
end: addEnd(off, datalen, bufsize),
|
end: addEnd(off, datalen, bufsize),
|
||||||
seq: sl.EndSeq(),
|
seq: seq,
|
||||||
size: Size(datalen),
|
size: Size(datalen),
|
||||||
})
|
})
|
||||||
return &sl.pkts[len(sl.pkts)-1]
|
return &sl.pkts[len(sl.pkts)-1]
|
||||||
|
|||||||
+7
-7
@@ -138,17 +138,17 @@ func TestSentlist_multi(t *testing.T) {
|
|||||||
sl.Reset(3, 0)
|
sl.Reset(3, 0)
|
||||||
|
|
||||||
// Test multi packet x2.
|
// Test multi packet x2.
|
||||||
p1 := sl.AddPacket(5, 0, bufsize)
|
p1 := sl.AddPacket(5, 0, bufsize, 0)
|
||||||
p2 := sl.AddPacket(5, p1.end, bufsize)
|
p2 := sl.AddPacket(5, p1.end, bufsize, p1.endSeq())
|
||||||
sl.RecvAck(Value(p2.size+p1.size), bufsize)
|
sl.RecvAck(Value(p2.size+p1.size), bufsize)
|
||||||
if sl.Oldest() != nil {
|
if sl.Oldest() != nil {
|
||||||
t.Fatal("expected full ack")
|
t.Fatal("expected full ack")
|
||||||
}
|
}
|
||||||
// multi packet x3.
|
// multi packet x3.
|
||||||
sl.Reset(3, 0)
|
sl.Reset(3, 0)
|
||||||
p1 = sl.AddPacket(3, 0, bufsize)
|
p1 = sl.AddPacket(3, 0, bufsize, 0)
|
||||||
p2 = sl.AddPacket(3, p1.end, bufsize)
|
p2 = sl.AddPacket(3, p1.end, bufsize, p1.endSeq())
|
||||||
p3 := sl.AddPacket(4, p2.end, bufsize)
|
p3 := sl.AddPacket(4, p2.end, bufsize, p2.endSeq())
|
||||||
sl.RecvAck(2, bufsize)
|
sl.RecvAck(2, bufsize)
|
||||||
oldest := sl.Oldest()
|
oldest := sl.Oldest()
|
||||||
if oldest != p1 {
|
if oldest != p1 {
|
||||||
@@ -167,7 +167,7 @@ func TestSentlist_simple(t *testing.T) {
|
|||||||
// Test full ack.
|
// Test full ack.
|
||||||
const bufsize = 16
|
const bufsize = 16
|
||||||
const pkt = 10
|
const pkt = 10
|
||||||
sl.AddPacket(pkt, 0, bufsize)
|
sl.AddPacket(pkt, 0, bufsize, 0)
|
||||||
if sl.Oldest() == nil || sl.Newest() != sl.Oldest() {
|
if sl.Oldest() == nil || sl.Newest() != sl.Oldest() {
|
||||||
t.Error("expected same oldest/newest non-nil packet")
|
t.Error("expected same oldest/newest non-nil packet")
|
||||||
}
|
}
|
||||||
@@ -179,7 +179,7 @@ func TestSentlist_simple(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Test partial ack.
|
// Test partial ack.
|
||||||
sl.AddPacket(pkt, 0, bufsize)
|
sl.AddPacket(pkt, 0, bufsize, sl.ssn)
|
||||||
for i := Value(0); i < pkt-1; i++ {
|
for i := Value(0); i < pkt-1; i++ {
|
||||||
ack++
|
ack++
|
||||||
sl.RecvAck(ack, bufsize)
|
sl.RecvAck(ack, bufsize)
|
||||||
|
|||||||
Reference in New Issue
Block a user