mirror of
https://github.com/soypat/lneto.git
synced 2026-08-20 14:39:02 +00:00
begin adding TxQueue failing multipacket test
This commit is contained in:
+1
-2
@@ -158,7 +158,6 @@ func (rtx *ringTx) RecvACK(ack Value) error {
|
|||||||
pkt0 := rtx.pkt(first)
|
pkt0 := rtx.pkt(first)
|
||||||
if ack.LessThanEq(pkt0.seq) {
|
if ack.LessThanEq(pkt0.seq) {
|
||||||
return fmt.Errorf("incoming ack %d older than first packet seq %d", ack, pkt0.seq)
|
return fmt.Errorf("incoming ack %d older than first packet seq %d", ack, pkt0.seq)
|
||||||
// return errors.New("old packet")
|
|
||||||
}
|
}
|
||||||
// lastAckedPkt stores last fully acked packet.
|
// lastAckedPkt stores last fully acked packet.
|
||||||
var lastAckedPkt *ringidx
|
var lastAckedPkt *ringidx
|
||||||
@@ -195,7 +194,7 @@ func (rtx *ringTx) RecvACK(ack Value) error {
|
|||||||
acked := int(ack - pkt.seq)
|
acked := int(ack - pkt.seq)
|
||||||
pring := rtx.ring(pkt.off, pkt.end)
|
pring := rtx.ring(pkt.off, pkt.end)
|
||||||
buffered := pring.Buffered()
|
buffered := pring.Buffered()
|
||||||
if acked > buffered || acked < minBufferSize {
|
if acked > buffered {
|
||||||
panic("unreachable")
|
panic("unreachable")
|
||||||
}
|
}
|
||||||
off := rtx.addOff(pkt.off, acked)
|
off := rtx.addOff(pkt.off, acked)
|
||||||
|
|||||||
@@ -8,6 +8,76 @@ import (
|
|||||||
"testing"
|
"testing"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
func TestTxQueue_multipacket(t *testing.T) {
|
||||||
|
const mtu = 256
|
||||||
|
const iss = 1
|
||||||
|
const maxPkts = 3
|
||||||
|
const maxWrites = 20
|
||||||
|
const maxWriteSize = mtu / maxWrites
|
||||||
|
var rtx ringTx
|
||||||
|
internalbuff := make([]byte, mtu)
|
||||||
|
rng := rand.New(rand.NewSource(1))
|
||||||
|
var wbuf, rbuf [mtu]byte
|
||||||
|
for itest := 0; itest < 32; itest++ {
|
||||||
|
err := rtx.Reset(internalbuff, maxPkts, iss)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
numWrites := rng.Intn(maxWrites) + 1
|
||||||
|
total := 0
|
||||||
|
woff := 0
|
||||||
|
for iw := 0; iw < numWrites; iw++ {
|
||||||
|
wlen := rng.Intn(maxWriteSize) + 1
|
||||||
|
towrite := wbuf[woff : woff+wlen]
|
||||||
|
rng.Read(towrite)
|
||||||
|
n, err := rtx.Write(towrite)
|
||||||
|
testQueueSanity(t, &rtx)
|
||||||
|
woff += n
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
} else if n != wlen {
|
||||||
|
t.Fatal("expected wlen==n", wlen, n)
|
||||||
|
}
|
||||||
|
total += wlen
|
||||||
|
}
|
||||||
|
npkt := rng.Intn(maxPkts) + 1
|
||||||
|
roff := 0
|
||||||
|
seq := Value(iss)
|
||||||
|
for ipkt := 0; ipkt < npkt; ipkt++ {
|
||||||
|
maxToPacket := min(total-roff, maxWriteSize)
|
||||||
|
pktlen := rng.Intn(maxToPacket) + 1
|
||||||
|
pkt := rbuf[roff : roff+pktlen]
|
||||||
|
expectPkt := wbuf[roff : roff+pktlen]
|
||||||
|
ngot, err := rtx.MakePacket(pkt, seq)
|
||||||
|
testQueueSanity(t, &rtx)
|
||||||
|
roff += ngot
|
||||||
|
seq += Value(ngot)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
} else if pktlen != ngot && roff != total {
|
||||||
|
t.Fatal(err)
|
||||||
|
} else if !bytes.Equal(expectPkt, pkt) {
|
||||||
|
t.Fatal("mismatched data written", expectPkt, pkt)
|
||||||
|
}
|
||||||
|
if roff == total {
|
||||||
|
break // made packet from all data.
|
||||||
|
}
|
||||||
|
}
|
||||||
|
acked := 0
|
||||||
|
for acked < roff {
|
||||||
|
maxToack := min(roff-acked, maxWriteSize)
|
||||||
|
toack := rng.Intn(maxToack) + 1
|
||||||
|
t.Log("\n", rtx.string())
|
||||||
|
err = rtx.RecvACK(iss + Value(acked+toack))
|
||||||
|
testQueueSanity(t, &rtx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
acked += toack
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestTxQueue(t *testing.T) {
|
func TestTxQueue(t *testing.T) {
|
||||||
const bufsize = 1024
|
const bufsize = 1024
|
||||||
var msgBuf, ringBuf, readBuf, aux [bufsize]byte
|
var msgBuf, ringBuf, readBuf, aux [bufsize]byte
|
||||||
|
|||||||
+11
-6
@@ -34,16 +34,21 @@ func TestStackAsyncTCP_multipacket(t *testing.T) {
|
|||||||
tst.TestTCPSetupAndEstablish(sv, client, svconn, clconn, svPort, 1337)
|
tst.TestTCPSetupAndEstablish(sv, client, svconn, clconn, svPort, 1337)
|
||||||
tst.TestTCPClose(client, sv, clconn, svconn)
|
tst.TestTCPClose(client, sv, clconn, svconn)
|
||||||
var buf [MTU]byte
|
var buf [MTU]byte
|
||||||
for i := 0; i < 30; i++ {
|
for i := 0; i < 1; i++ {
|
||||||
payloadSize := rng.Intn(maxPktLen) + 1
|
payloadSize := rng.Intn(maxPktLen) + 1
|
||||||
tst.TestTCPSetupAndEstablish(sv, client, svconn, clconn, svPort, 1337)
|
tst.TestTCPSetupAndEstablish(sv, client, svconn, clconn, svPort, 1337)
|
||||||
npkt := rng.Intn(maxNPkt-1) + 2
|
// npkt := rng.Intn(maxNPkt-1) + 2
|
||||||
for ipkt := 0; ipkt < npkt; ipkt++ {
|
a, _ := rng.Read(buf[:payloadSize])
|
||||||
a, _ := rng.Read(buf[:payloadSize])
|
tst.TestTCPEstablishedSingleData(sv, client, svconn, clconn, buf[:a])
|
||||||
tst.TestTCPEstablishedSingleData(sv, client, svconn, clconn, buf[:a])
|
a, _ = rng.Read(buf[:payloadSize])
|
||||||
}
|
tst.TestTCPEstablishedSingleData(sv, client, svconn, clconn, buf[:a])
|
||||||
|
// for ipkt := 0; ipkt < npkt; ipkt++ {
|
||||||
|
// a, _ := rng.Read(buf[:payloadSize])
|
||||||
|
// tst.TestTCPEstablishedSingleData(sv, client, svconn, clconn, buf[:a])
|
||||||
|
// }
|
||||||
tst.TestTCPClose(client, sv, clconn, svconn)
|
tst.TestTCPClose(client, sv, clconn, svconn)
|
||||||
if t.Failed() {
|
if t.Failed() {
|
||||||
|
t.Error("multi failed")
|
||||||
t.FailNow()
|
t.FailNow()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user