From ec2ed17f43b23edda8cca842566f3cb65f0b86c4 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Wed, 16 Sep 2026 14:36:04 +0100 Subject: [PATCH] made the mechanism resilient enough to detect stalled data --- bngsdk.go | 2 +- example/read_socket_data/main.go | 2 +- reader_socket.go | 13 ++++---- transport.go | 52 ++++++++++++++++++++++++++++---- 4 files changed, 54 insertions(+), 15 deletions(-) diff --git a/bngsdk.go b/bngsdk.go index 10b9e77..1299864 100644 --- a/bngsdk.go +++ b/bngsdk.go @@ -181,7 +181,7 @@ func (sdk *BeamNGSDK) Update() (int, error) { nBytes, err = sdk.reader.Next(sdk.buffer) } if err != nil { - return nBytes, err + return 0, err } if sdk.Opts.ExportData { diff --git a/example/read_socket_data/main.go b/example/read_socket_data/main.go index a3e04b8..1b99a24 100644 --- a/example/read_socket_data/main.go +++ b/example/read_socket_data/main.go @@ -35,7 +35,7 @@ func main() { } defer sdk.Close() - ticker := time.NewTicker(time.Second / 60) + ticker := time.NewTicker(time.Second) defer ticker.Stop() for { diff --git a/reader_socket.go b/reader_socket.go index b1ccf8f..5b99bb0 100644 --- a/reader_socket.go +++ b/reader_socket.go @@ -1,14 +1,15 @@ package bngsdk import ( - "fmt" - "log/slog" + "errors" ) type OgUDPReader struct { udpConnection *UDPTransport } +var ErrInvalidOutgaugeData = errors.New("data is of different size than outgauge") + func NewOgUDPReader(ip string, port int) (*OgUDPReader, error) { conn, err := NewUDPReader(ip, port) if err != nil { @@ -34,17 +35,15 @@ func (ogr *OgUDPReader) Reset() error { } func (ogr *OgUDPReader) Next(buffer []byte) (int, error) { - slog.Debug("Reading") - nBytes, err := ogr.udpConnection.Read(buffer) if err != nil { return 0, err } // Check if enough data was received to fill our struct - if nBytes < outgaugeSize { - return 0, fmt.Errorf("received data is smaller than outgauge size") + if nBytes != outgaugeSize { + return 0, ErrInvalidOutgaugeData } - return 0, nil + return nBytes, nil } diff --git a/transport.go b/transport.go index f1f8690..cc91469 100644 --- a/transport.go +++ b/transport.go @@ -1,12 +1,23 @@ package bngsdk import ( + "errors" + "log/slog" "net" + "time" + "unsafe" +) + +// ErrReaderDisabled = errors.New("reader is currently turned off") +var ( + ErrNoData = errors.New("no new data available") + readTimeout = (time.Second / 60) * 5 // N missed frames at 60fps ) type UDPTransport struct { address *net.UDPAddr connection *net.UDPConn + dataChan chan []byte } func NewUDPReader(ip string, port int) (*UDPTransport, error) { @@ -21,10 +32,15 @@ func NewUDPReader(ip string, port int) (*UDPTransport, error) { return nil, err } - return &UDPTransport{ + udpT := UDPTransport{ address: addr, connection: conn, - }, nil + dataChan: make(chan []byte, 1), + } + + go udpT.udpSink() + + return &udpT, nil } func NewUDPWriter(ip string, port int) (*UDPTransport, error) { @@ -45,17 +61,41 @@ func NewUDPWriter(ip string, port int) (*UDPTransport, error) { }, nil } +// udpSink is a loop that will consume the most recent packets on the port +// instead of letting them pile up +func (ut *UDPTransport) udpSink() { + defer close(ut.dataChan) + + for { + ut.connection.SetReadDeadline(time.Now().Add(readTimeout)) + + buf := make([]byte, unsafe.Sizeof(Outgauge{})) + nBytes, _, err := ut.connection.ReadFromUDP(buf) + if err != nil { + return + } + + select { + case ut.dataChan <- buf[:nBytes]: + default: + <-ut.dataChan + ut.dataChan <- buf[:nBytes] + } + } +} + func (ut *UDPTransport) Write(data []byte) (int, error) { return ut.connection.Write(data) } func (ut *UDPTransport) Read(buffer []byte) (int, error) { - n, _, err := ut.connection.ReadFromUDP(buffer) - if err != nil { - return 0, err + data, ok := <-ut.dataChan + if !ok { + slog.Error(ErrNoData.Error()) + return 0, ErrNoData } - return n, nil + return copy(buffer, data), nil } func (ut *UDPTransport) Close() error {