Implemented hearthbit client side + tests
This commit is contained in:
@@ -99,6 +99,20 @@ func BuildDisconnectPacket() []byte {
|
||||
})
|
||||
}
|
||||
|
||||
// BuildAckPacket returns a 17-byte ACK datagram. Unicast clients send it
|
||||
// periodically as a keepalive: UDPSServer refreshes the client's last-seen
|
||||
// without re-sending CONFIG (which a repeated CONNECT would trigger).
|
||||
func BuildAckPacket() []byte {
|
||||
return buildHeader(PacketHeader{
|
||||
Magic: MagicUDPS,
|
||||
Type: PktACK,
|
||||
Counter: 0,
|
||||
FragmentIdx: 0,
|
||||
TotalFragments: 1,
|
||||
PayloadBytes: 0,
|
||||
})
|
||||
}
|
||||
|
||||
// ─── Signal descriptor (136 bytes) ───────────────────────────────────────────
|
||||
|
||||
// SignalInfo holds the parsed metadata for one signal.
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
package wshub
|
||||
|
||||
import (
|
||||
"net"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"marte2/common/udpsprotocol"
|
||||
)
|
||||
|
||||
// TestUDPClientSendsPeriodicKeepAliveAcks verifies that a unicast UDPClient
|
||||
// re-sends ACK datagrams from the SAME socket at keepAliveInterval. The
|
||||
// UDPSServer refreshes a unicast client's last-seen only on client->server
|
||||
// traffic; without this keepalive it evicts the client after ClientTimeout
|
||||
// (default 30 s) and the stream dies.
|
||||
func TestUDPClientSendsPeriodicKeepAliveAcks(t *testing.T) {
|
||||
srv, err := net.ListenUDP("udp4", &net.UDPAddr{IP: net.ParseIP("127.0.0.1")})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer srv.Close()
|
||||
|
||||
c := NewUDPClient(srv.LocalAddr().String(), "ka1", NewHub(), "", 0)
|
||||
c.keepAliveInterval = 150 * time.Millisecond
|
||||
go c.Run()
|
||||
defer c.Stop()
|
||||
|
||||
buf := make([]byte, 512)
|
||||
|
||||
// 1) CONNECT from the client's ephemeral socket.
|
||||
srv.SetReadDeadline(time.Now().Add(3 * time.Second))
|
||||
n, clientAddr, err := srv.ReadFromUDP(buf)
|
||||
if err != nil {
|
||||
t.Fatalf("expected CONNECT: %v", err)
|
||||
}
|
||||
hdr, err := udpsprotocol.ParseHeader(buf[:n])
|
||||
if err != nil {
|
||||
t.Fatalf("parse CONNECT: %v", err)
|
||||
}
|
||||
if hdr.Type != udpsprotocol.PktConnect {
|
||||
t.Fatalf("first packet type = %d, want CONNECT (%d)", hdr.Type, udpsprotocol.PktConnect)
|
||||
}
|
||||
|
||||
// 2) Collect ACKs for ~1 s: must be periodic and from the SAME socket
|
||||
// (a new ephemeral socket would be registered as a new client).
|
||||
deadline := time.Now().Add(time.Second)
|
||||
acks := 0
|
||||
for time.Now().Before(deadline) {
|
||||
srv.SetReadDeadline(deadline)
|
||||
n, addr, err := srv.ReadFromUDP(buf)
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
hdr, err := udpsprotocol.ParseHeader(buf[:n])
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
if hdr.Type != udpsprotocol.PktACK {
|
||||
t.Fatalf("unexpected packet type %d from %s", hdr.Type, addr)
|
||||
}
|
||||
if addr.String() != clientAddr.String() {
|
||||
t.Fatalf("ACK from %s, want same socket as CONNECT (%s)", addr, clientAddr)
|
||||
}
|
||||
acks++
|
||||
}
|
||||
if acks < 3 {
|
||||
t.Fatalf("expected >= 3 keepalive ACKs in 1 s, got %d", acks)
|
||||
}
|
||||
}
|
||||
@@ -172,27 +172,34 @@ const (
|
||||
reconnectDelay = 2 * time.Second
|
||||
readBufSize = 65536
|
||||
udpRcvBufSize = 8 * 1024 * 1024
|
||||
// keepAliveInterval is the unicast keepalive period. The UDPStreamer
|
||||
// server evicts silent unicast clients after its ClientTimeout (default
|
||||
// 30 s); an ACK from the same socket refreshes its last-seen without
|
||||
// triggering a CONFIG resend (a CONNECT would).
|
||||
keepAliveInterval = 15 * time.Second
|
||||
)
|
||||
|
||||
// UDPClient manages the connection to one MARTe2 streamer source.
|
||||
type UDPClient struct {
|
||||
serverAddr string
|
||||
sourceID string
|
||||
hub *Hub
|
||||
multicastGroup string
|
||||
dataPort int
|
||||
stopCh chan struct{}
|
||||
serverAddr string
|
||||
sourceID string
|
||||
hub *Hub
|
||||
multicastGroup string
|
||||
dataPort int
|
||||
keepAliveInterval time.Duration
|
||||
stopCh chan struct{}
|
||||
}
|
||||
|
||||
// NewUDPClient creates a UDPClient bound to a specific source ID.
|
||||
func NewUDPClient(serverAddr, sourceID string, hub *Hub, multicastGroup string, dataPort int) *UDPClient {
|
||||
return &UDPClient{
|
||||
serverAddr: serverAddr,
|
||||
sourceID: sourceID,
|
||||
hub: hub,
|
||||
multicastGroup: multicastGroup,
|
||||
dataPort: dataPort,
|
||||
stopCh: make(chan struct{}),
|
||||
serverAddr: serverAddr,
|
||||
sourceID: sourceID,
|
||||
hub: hub,
|
||||
multicastGroup: multicastGroup,
|
||||
dataPort: dataPort,
|
||||
keepAliveInterval: keepAliveInterval,
|
||||
stopCh: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -253,6 +260,20 @@ func (u *UDPClient) runSession() error {
|
||||
return err
|
||||
}
|
||||
log.Printf("[%s] udp: sent CONNECT", u.sourceID)
|
||||
lastData := time.Now()
|
||||
lastKeepAlive := time.Now()
|
||||
// sendKeepAliveIfDue sends an ACK if the keepalive interval has elapsed.
|
||||
// ACK refreshes the server's last-seen without re-sending CONFIG (which a
|
||||
// repeated CONNECT would trigger).
|
||||
sendKeepAliveIfDue := func() error {
|
||||
if u.keepAliveInterval > 0 && time.Since(lastKeepAlive) >= u.keepAliveInterval {
|
||||
if _, err := conn.WriteToUDP(udpsprotocol.BuildAckPacket(), serverAddr); err != nil {
|
||||
return err
|
||||
}
|
||||
lastKeepAlive = time.Now()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
reassembler := udpsprotocol.NewReassembler(2 * time.Second)
|
||||
buf := make([]byte, readBufSize)
|
||||
@@ -260,14 +281,34 @@ func (u *UDPClient) runSession() error {
|
||||
var currentPublishMode uint8
|
||||
|
||||
for {
|
||||
conn.SetReadDeadline(time.Now().Add(silenceTimeout))
|
||||
// Wake up at least every keepalive interval so ACKs are sent even
|
||||
// when the server is idle; the read deadline also doubles as the
|
||||
// silence detector (no data for silenceTimeout = server gone).
|
||||
wakeup := silenceTimeout
|
||||
if u.keepAliveInterval > 0 && u.keepAliveInterval < wakeup {
|
||||
wakeup = u.keepAliveInterval
|
||||
}
|
||||
conn.SetReadDeadline(time.Now().Add(wakeup))
|
||||
|
||||
n, _, err := conn.ReadFromUDP(buf)
|
||||
arrivalTime := time.Now()
|
||||
if err != nil {
|
||||
if ne, ok := err.(net.Error); ok && ne.Timeout() {
|
||||
if time.Since(lastData) >= silenceTimeout {
|
||||
// True silence: stream is dead — Run() reconnects.
|
||||
conn.WriteToUDP(udpsprotocol.BuildDisconnectPacket(), serverAddr)
|
||||
return err
|
||||
}
|
||||
// Short wakeup: keepalive if due, then keep waiting.
|
||||
if kaErr := sendKeepAliveIfDue(); kaErr != nil {
|
||||
return kaErr
|
||||
}
|
||||
continue
|
||||
}
|
||||
conn.WriteToUDP(udpsprotocol.BuildDisconnectPacket(), serverAddr)
|
||||
return err
|
||||
}
|
||||
lastData = arrivalTime
|
||||
|
||||
if n < udpsprotocol.HeaderSize {
|
||||
log.Printf("[%s] udp: short datagram (%d bytes), skipping", u.sourceID, n)
|
||||
@@ -334,6 +375,10 @@ func (u *UDPClient) runSession() error {
|
||||
return nil
|
||||
default:
|
||||
}
|
||||
|
||||
if kaErr := sendKeepAliveIfDue(); kaErr != nil {
|
||||
return kaErr
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user