mirror of
https://github.com/soypat/lneto.git
synced 2026-08-20 06:29:03 +00:00
lneto rework: Require explicit BackoffStrategy in all APIs (#116)
* make backoff explicit API * fix tests to use explicit backoff * fix examples with explicit tcp
This commit is contained in:
@@ -8,7 +8,6 @@ import (
|
||||
|
||||
"github.com/soypat/lneto"
|
||||
"github.com/soypat/lneto/dhcp/dhcpv4"
|
||||
"github.com/soypat/lneto/internal"
|
||||
"github.com/soypat/lneto/tcp"
|
||||
)
|
||||
|
||||
@@ -21,6 +20,9 @@ var (
|
||||
)
|
||||
|
||||
func (s *StackAsync) StackBlocking(stackProtoBackoff lneto.BackoffStrategy) StackBlocking {
|
||||
if stackProtoBackoff == nil {
|
||||
panic("nil backoff to StackBlocking")
|
||||
}
|
||||
return StackBlocking{
|
||||
async: s,
|
||||
_backoff: stackProtoBackoff,
|
||||
@@ -206,9 +208,5 @@ func (s StackBlocking) backoff(consecutiveBackoffs uint) {
|
||||
}
|
||||
|
||||
func backoff(bo lneto.BackoffStrategy, consecutiveBackoffs uint) {
|
||||
if bo != nil {
|
||||
bo.Do(consecutiveBackoffs)
|
||||
} else {
|
||||
internal.BackoffStackProto(consecutiveBackoffs)
|
||||
}
|
||||
bo.Do(consecutiveBackoffs)
|
||||
}
|
||||
|
||||
@@ -25,6 +25,9 @@ type StackGoConfig struct {
|
||||
}
|
||||
|
||||
func (s *StackAsync) StackGo(stackProtoBackoff lneto.BackoffStrategy, cfg StackGoConfig) StackGo {
|
||||
if stackProtoBackoff == nil || cfg.ListenerPoolConfig.NewBackoff == nil {
|
||||
panic("nil backoff to StackGo")
|
||||
}
|
||||
return s.StackBlocking(stackProtoBackoff).StackGo(cfg)
|
||||
}
|
||||
|
||||
@@ -142,9 +145,11 @@ func (s StackGo) SocketNetip(ctx context.Context, network string, family, sotype
|
||||
var conn tcp.Conn
|
||||
// DIAL TCP: active connection a.k.a TCP Client branch.
|
||||
err = conn.Configure(tcp.ConnConfig{
|
||||
// TODO(pato): Eventually add UDP configuration. we use TCP for now for simplicity's sake.
|
||||
TxBuf: make([]byte, s.plcfg.TxBufSize),
|
||||
RxBuf: make([]byte, s.plcfg.RxBufSize),
|
||||
TxPacketQueueSize: s.plcfg.QueueSize,
|
||||
RWBackoff: s.plcfg.NewBackoff(),
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
||||
@@ -10,6 +10,9 @@ import (
|
||||
)
|
||||
|
||||
func (s *StackAsync) StackRetrying(stackProtoBackoff lneto.BackoffStrategy) StackRetrying {
|
||||
if stackProtoBackoff == nil {
|
||||
panic("nil backoff to StackRetrying")
|
||||
}
|
||||
return StackRetrying{
|
||||
block: s.StackBlocking(stackProtoBackoff),
|
||||
}
|
||||
|
||||
@@ -70,6 +70,7 @@ func newUDPConn6(t testing.TB) *udp.Conn {
|
||||
TxBuf: make([]byte, bufSize),
|
||||
RxQueueSize: 4,
|
||||
TxQueueSize: 4,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal("UDP Configure:", err)
|
||||
}
|
||||
@@ -85,6 +86,7 @@ func newTCPConn6(t testing.TB) *tcp.Conn {
|
||||
RxBuf: make([]byte, bufSize),
|
||||
TxBuf: make([]byte, bufSize),
|
||||
TxPacketQueueSize: 4,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal("TCP Configure:", err)
|
||||
}
|
||||
|
||||
+4
-4
@@ -53,7 +53,7 @@ type TCPPoolConfig struct {
|
||||
ClosingTimeout time.Duration
|
||||
// NewUserData is used to create user data used for each individual TCP connection and returned on GetTCP.
|
||||
NewUserData func() any
|
||||
// NewBackoff returns the backoff to use for every newly configured TCP connection.
|
||||
// NewBackoff returns the backoff to use for every newly configured TCP connection. Must be non-nil.
|
||||
// This should always return a static(non-method) function unless you know what you are doing.
|
||||
NewBackoff func() lneto.BackoffStrategy
|
||||
}
|
||||
@@ -61,6 +61,8 @@ type TCPPoolConfig struct {
|
||||
func NewTCPPool(cfg TCPPoolConfig) (*TCPPool, error) {
|
||||
if cfg.EstablishedTimeout <= 0 || cfg.ClosingTimeout <= 0 {
|
||||
return nil, lneto.ErrInvalidConfig
|
||||
} else if cfg.NewBackoff == nil {
|
||||
return nil, lneto.ErrMissingHALConfig
|
||||
}
|
||||
n := int(cfg.PoolSize)
|
||||
pool := &TCPPool{
|
||||
@@ -84,9 +86,7 @@ func NewTCPPool(cfg TCPPoolConfig) (*TCPPool, error) {
|
||||
TxBuf: bufSpace[txOff : txOff+cfg.TxBufSize],
|
||||
TxPacketQueueSize: cfg.QueueSize,
|
||||
Logger: cfg.ConnLogger,
|
||||
}
|
||||
if cfg.NewBackoff != nil {
|
||||
conncfg.RWBackoff = cfg.NewBackoff()
|
||||
RWBackoff: cfg.NewBackoff(),
|
||||
}
|
||||
err := pool.conns[i].Configure(conncfg)
|
||||
if err != nil {
|
||||
|
||||
@@ -53,6 +53,7 @@ func TestTCPListener_ConcurrentEcho(t *testing.T) {
|
||||
RxBufSize: 512,
|
||||
EstablishedTimeout: 5 * time.Second,
|
||||
ClosingTimeout: 5 * time.Second,
|
||||
NewBackoff: func() lneto.BackoffStrategy { return backoffYield },
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -96,6 +97,7 @@ func TestTCPListener_ConcurrentEcho(t *testing.T) {
|
||||
RxBuf: connBufs[bufOff : bufOff+tcpBufSize],
|
||||
TxBuf: connBufs[bufOff+tcpBufSize : bufOff+2*tcpBufSize],
|
||||
TxPacketQueueSize: 4,
|
||||
RWBackoff: backoffYield,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("client %d conn configure: %v", i, err)
|
||||
@@ -327,7 +329,7 @@ func testCloseTransmitsPending(tst *tester, s1, s2 *StackAsync, c1, c2 *tcp.Conn
|
||||
RxBuf: nil,
|
||||
TxBuf: make([]byte, tx1Buf),
|
||||
TxPacketQueueSize: queueSize,
|
||||
RWBackoff: backoffGosched,
|
||||
RWBackoff: backoffYield,
|
||||
Logger: logger,
|
||||
})
|
||||
if err != nil {
|
||||
@@ -431,7 +433,3 @@ func testCloseTransmitsPending(tst *tester, s1, s2 *StackAsync, c1, c2 *tcp.Conn
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func backoffGosched(consecutiveBackoffs uint) (sleep time.Duration) {
|
||||
return lneto.BackoffFlagGosched
|
||||
}
|
||||
|
||||
@@ -318,6 +318,7 @@ func testStackSeeded(t *testing.T, seed1, seed2 int64) {
|
||||
RxBuf: make([]byte, bufsize),
|
||||
TxBuf: make([]byte, bufsize),
|
||||
TxPacketQueueSize: 1 + int(uint16(seed1)%10),
|
||||
RWBackoff: backoffYield,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -326,6 +327,7 @@ func testStackSeeded(t *testing.T, seed1, seed2 int64) {
|
||||
RxBuf: make([]byte, bufsize),
|
||||
TxBuf: make([]byte, bufsize),
|
||||
TxPacketQueueSize: 1 + int(uint16(seed2)%10),
|
||||
RWBackoff: backoffYield,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -335,6 +337,7 @@ func testStackSeeded(t *testing.T, seed1, seed2 int64) {
|
||||
TxBuf: make([]byte, bufsize),
|
||||
RxQueueSize: int(1 + uint16(seed1>>32)%10),
|
||||
TxQueueSize: int(1 + uint16(seed1>>32)%10),
|
||||
RWBackoff: backoffYield,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -344,6 +347,7 @@ func testStackSeeded(t *testing.T, seed1, seed2 int64) {
|
||||
TxBuf: make([]byte, bufsize),
|
||||
RxQueueSize: int(1 + uint16(seed2>>32)%10),
|
||||
TxQueueSize: int(1 + uint16(seed2>>32)%10),
|
||||
RWBackoff: backoffYield,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/soypat/lneto"
|
||||
"github.com/soypat/lneto/ethernet"
|
||||
"github.com/soypat/lneto/tcp"
|
||||
)
|
||||
@@ -49,6 +50,7 @@ func TestStackAsyncListener_SingleConnection(t *testing.T) {
|
||||
RxBuf: make([]byte, MTU),
|
||||
TxBuf: make([]byte, MTU),
|
||||
TxPacketQueueSize: 4,
|
||||
RWBackoff: backoffYield,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -62,6 +64,7 @@ func TestStackAsyncListener_SingleConnection(t *testing.T) {
|
||||
RxBufSize: MTU,
|
||||
EstablishedTimeout: 10e9,
|
||||
ClosingTimeout: 10e9,
|
||||
NewBackoff: func() lneto.BackoffStrategy { return backoffYield },
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -150,6 +153,7 @@ func TestStackAsyncListener_MultiSequentialConn(t *testing.T) {
|
||||
RxBufSize: bufsize,
|
||||
EstablishedTimeout: 10e9,
|
||||
ClosingTimeout: 10e9,
|
||||
NewBackoff: func() lneto.BackoffStrategy { return backoffYield },
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -188,6 +192,7 @@ func TestStackAsyncListener_MultiSequentialConn(t *testing.T) {
|
||||
RxBuf: make([]byte, bufsize),
|
||||
TxBuf: make([]byte, bufsize),
|
||||
TxPacketQueueSize: 4,
|
||||
RWBackoff: backoffYield,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -239,6 +244,7 @@ func TestListener_Close(t *testing.T) {
|
||||
RxBufSize: 512,
|
||||
EstablishedTimeout: 10e9,
|
||||
ClosingTimeout: 10e9,
|
||||
NewBackoff: func() lneto.BackoffStrategy { return backoffYield },
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -278,6 +284,7 @@ func TestListener_ResetAfterClose(t *testing.T) {
|
||||
RxBufSize: 512,
|
||||
EstablishedTimeout: 10e9,
|
||||
ClosingTimeout: 10e9,
|
||||
NewBackoff: func() lneto.BackoffStrategy { return backoffYield },
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"net/netip"
|
||||
"testing"
|
||||
|
||||
"github.com/soypat/lneto"
|
||||
"github.com/soypat/lneto/ethernet"
|
||||
"github.com/soypat/lneto/ipv4"
|
||||
"github.com/soypat/lneto/tcp"
|
||||
@@ -73,6 +74,7 @@ func TestStackAsync_ListenerSynAckAddressedToClient(t *testing.T) {
|
||||
RxBufSize: mtu,
|
||||
EstablishedTimeout: 10e9,
|
||||
ClosingTimeout: 10e9,
|
||||
NewBackoff: func() lneto.BackoffStrategy { return backoffYield },
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -104,6 +106,7 @@ func TestStackAsync_ListenerSynAckAddressedToClient(t *testing.T) {
|
||||
if err = clConn.Configure(tcp.ConnConfig{
|
||||
RxBuf: make([]byte, mtu), TxBuf: make([]byte, mtu),
|
||||
TxPacketQueueSize: 4,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
@@ -208,6 +208,7 @@ func newTCPStacks(t testing.TB, randSeed int64, mtu int) (s1, s2 *StackAsync, c1
|
||||
RxBuf: buf[:mtu],
|
||||
TxBuf: buf[mtu : mtu*2],
|
||||
TxPacketQueueSize: 4,
|
||||
RWBackoff: backoffYield,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -216,6 +217,7 @@ func newTCPStacks(t testing.TB, randSeed int64, mtu int) (s1, s2 *StackAsync, c1
|
||||
RxBuf: buf[2*mtu : 3*mtu],
|
||||
TxBuf: buf[3*mtu : 4*mtu],
|
||||
TxPacketQueueSize: 4,
|
||||
RWBackoff: backoffYield,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -1045,3 +1047,7 @@ func getTCPFrame(etherFrame []byte) (tcp.Frame, bool) {
|
||||
}
|
||||
return tfrm, true
|
||||
}
|
||||
|
||||
func backoffYield(consecutiveBackoffs uint) time.Duration {
|
||||
return lneto.BackoffFlagGosched
|
||||
}
|
||||
|
||||
@@ -60,6 +60,7 @@ func TestStackAsyncRegisterListenerUDP_ReceiveData(t *testing.T) {
|
||||
if err := pc.Configure(udp.PacketConnConfig{
|
||||
RxBuf: make([]byte, testUDPBufSize), TxBuf: make([]byte, testUDPBufSize),
|
||||
RxQueueSize: testUDPQueueSize, TxQueueSize: testUDPQueueSize,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal("pc Configure:", err)
|
||||
}
|
||||
@@ -74,6 +75,7 @@ func TestStackAsyncRegisterListenerUDP_ReceiveData(t *testing.T) {
|
||||
if err := conn.Configure(udp.ConnConfig{
|
||||
RxBuf: make([]byte, testUDPBufSize), TxBuf: make([]byte, testUDPBufSize),
|
||||
RxQueueSize: testUDPQueueSize, TxQueueSize: testUDPQueueSize,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal("conn Configure:", err)
|
||||
}
|
||||
@@ -119,6 +121,7 @@ func TestStackAsyncRegisterListenerUDP_ReplyToClient(t *testing.T) {
|
||||
if err := pc.Configure(udp.PacketConnConfig{
|
||||
RxBuf: make([]byte, testUDPBufSize), TxBuf: make([]byte, testUDPBufSize),
|
||||
RxQueueSize: testUDPQueueSize, TxQueueSize: testUDPQueueSize,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal("pc Configure:", err)
|
||||
}
|
||||
@@ -133,6 +136,7 @@ func TestStackAsyncRegisterListenerUDP_ReplyToClient(t *testing.T) {
|
||||
if err := conn.Configure(udp.ConnConfig{
|
||||
RxBuf: make([]byte, testUDPBufSize), TxBuf: make([]byte, testUDPBufSize),
|
||||
RxQueueSize: testUDPQueueSize, TxQueueSize: testUDPQueueSize,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal("conn Configure:", err)
|
||||
}
|
||||
@@ -228,6 +232,7 @@ func TestStackAsyncRegisterListenerUDP_MultiSource(t *testing.T) {
|
||||
if err := pc.Configure(udp.PacketConnConfig{
|
||||
RxBuf: make([]byte, testUDPBufSize), TxBuf: make([]byte, testUDPBufSize),
|
||||
RxQueueSize: testUDPQueueSize, TxQueueSize: testUDPQueueSize,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal("pc Configure:", err)
|
||||
}
|
||||
@@ -243,6 +248,7 @@ func TestStackAsyncRegisterListenerUDP_MultiSource(t *testing.T) {
|
||||
if err := conn1.Configure(udp.ConnConfig{
|
||||
RxBuf: make([]byte, testUDPBufSize), TxBuf: make([]byte, testUDPBufSize),
|
||||
RxQueueSize: testUDPQueueSize, TxQueueSize: testUDPQueueSize,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal("conn1 Configure:", err)
|
||||
}
|
||||
@@ -260,6 +266,7 @@ func TestStackAsyncRegisterListenerUDP_MultiSource(t *testing.T) {
|
||||
if err := conn2.Configure(udp.ConnConfig{
|
||||
RxBuf: make([]byte, testUDPBufSize), TxBuf: make([]byte, testUDPBufSize),
|
||||
RxQueueSize: testUDPQueueSize, TxQueueSize: testUDPQueueSize,
|
||||
RWBackoff: backoffYield,
|
||||
}); err != nil {
|
||||
t.Fatal("conn2 Configure:", err)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user