From a7fbf74f0a40274c3c1a2363f788fa9aaeb74733 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Tue, 15 Sep 2026 13:19:48 +0100 Subject: [PATCH 1/4] A better SDK with more features and more consistency --- bngsdk.go | 280 ++++++++++++++++++++++++------- example/read_binary_file/main.go | 60 +++++++ example/read_socket_data/main.go | 56 +++++++ reader.go | 26 +++ reader_binary.go | 43 +++++ reader_socket.go | 42 +++++ transport.go | 67 ++++++++ writer.go | 23 +++ writer_binary.go | 30 ++++ writer_socket.go | 21 +++ 10 files changed, 583 insertions(+), 65 deletions(-) create mode 100644 example/read_binary_file/main.go create mode 100644 example/read_socket_data/main.go create mode 100644 reader.go create mode 100644 reader_binary.go create mode 100644 reader_socket.go create mode 100644 transport.go create mode 100644 writer.go create mode 100644 writer_binary.go create mode 100644 writer_socket.go diff --git a/bngsdk.go b/bngsdk.go index 9a46f8a..bb664ed 100644 --- a/bngsdk.go +++ b/bngsdk.go @@ -4,8 +4,11 @@ package bngsdk import ( "encoding/binary" "fmt" + "io" + "log/slog" "math" - "net" + "time" + "unsafe" ) const ( @@ -15,72 +18,207 @@ const ( DefaultUDPPort = 4444 ) +type TelemetryContainer int + +const ( + BinaryFile TelemetryContainer = iota + UDPData TelemetryContainer = iota +) + +type Options struct { + Logger *slog.Logger + // Import settings + SourceType TelemetryContainer // type of source data + BinSourcePath string // Path to source binary + ImportUDPAddress string + ImportUDPPort int + Loop bool // Wheter to loop when we reach the end of file + // Export settings + ExportUDPAddress string + ExportUDPPort int + ExportDataType TelemetryContainer // export type of telemetry: store .bin or replay in UDP + ExportDataPath string // path where to export the data + ExportData bool // whether to export the telemetry data +} + type BeamNGSDK struct { - Addr *net.UDPAddr - Conn *net.UDPConn - Buffer []byte - Data Outgauge + Opts Options + reader BngImporter + writer BngExporter + Data Outgauge + buffer []byte + receiver *Receiver } -func createUDPConnection(ip string, port int) (*net.UDPConn, *net.UDPAddr, error) { - // Define the IP address and port to listen on - addr := &net.UDPAddr{ - IP: net.ParseIP(ip), - Port: port, +func NewBngSDK(opts Options) (*BeamNGSDK, error) { + sdk := BeamNGSDK{ + Opts: opts, + buffer: make([]byte, unsafe.Sizeof(Outgauge{})), } - // Create a UDP socket - conn, err := net.ListenUDP("udp", addr) + // Set the passed logger as the default logger + slog.SetDefault(sdk.Opts.Logger) + + var err error + + err = sdk.openReader() if err != nil { - return nil, nil, err + slog.Error("failed to open reader", "err", err) + return nil, err } - return conn, addr, nil + err = sdk.openWriter() + if err != nil { + slog.Error("failed to open writer", "err", err) + return nil, err + } + + return &sdk, nil } -// parseData will parse the bytes read from the socket into the Outgauge struct -func (sdk *BeamNGSDK) parseData() error { - sdk.Data.Time = binary.LittleEndian.Uint32(sdk.Buffer[0:4]) - copy(sdk.Data.Car[:], sdk.Buffer[4:8]) - sdk.Data.Flags = binary.LittleEndian.Uint16(sdk.Buffer[8:10]) - sdk.Data.Gear = int8(sdk.Buffer[10]) - sdk.Data.Plid = int8(sdk.Buffer[11]) - sdk.Data.Speed = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[12:16])) - sdk.Data.RPM = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[16:20])) - sdk.Data.Turbo = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[20:24])) - sdk.Data.EngTemp = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[24:28])) - sdk.Data.Fuel = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[28:32])) - sdk.Data.OilPressure = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[32:36])) - sdk.Data.OilTemp = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[36:40])) - sdk.Data.DashLights = binary.LittleEndian.Uint32(sdk.Buffer[40:44]) - sdk.Data.ShowLights = binary.LittleEndian.Uint32(sdk.Buffer[44:48]) - sdk.Data.Throttle = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[48:52])) - sdk.Data.Brake = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[52:56])) - sdk.Data.Throttle = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[56:60])) - copy(sdk.Data.Display1[:], sdk.Buffer[60:76]) - copy(sdk.Data.Display2[:], sdk.Buffer[76:92]) - sdk.Data.ID = int32(binary.LittleEndian.Uint32(sdk.Buffer[92:96])) +func (sdk *BeamNGSDK) openReader() error { + switch sdk.Opts.SourceType { + case BinaryFile: + reader, err := NewBinaryImporter(sdk.Opts.BinSourcePath) + if err != nil { + slog.Error("failed to create BinaryImporter", + "path", sdk.Opts.BinSourcePath, "err", err) + return err + } + + slog.Info( + "created BinaryImporter", + "path", sdk.Opts.BinSourcePath, + ) + sdk.reader = reader + case UDPData: + reader, err := NewSocketImporter(sdk.Opts.ImportUDPAddress, sdk.Opts.ImportUDPPort) + if err != nil { + slog.Error( + "failed to create SocketImporter", + "address", sdk.Opts.ImportUDPAddress, "port", sdk.Opts.ImportUDPPort, "err", err, + ) + return err + } + + slog.Info( + "created SocketImporter", + "address", sdk.Opts.ImportUDPAddress, "port", sdk.Opts.ImportUDPPort, + ) + sdk.reader = reader + } return nil } -// ReadData will read new data from the UDP server -func (sdk *BeamNGSDK) ReadData() error { - // Receive data from the socket - n, _, err := sdk.Conn.ReadFromUDP(sdk.Buffer) - if err != nil { - fmt.Println("Error reading from UDP:", err) - return err +func (sdk *BeamNGSDK) openWriter() error { + // if the user didn't request data export we don't need a writer + if !sdk.Opts.ExportData { + return nil } - // Check if enough data was received to fill our struct - if n < outgaugeSize { - fmt.Println("Received packet too small for Outgauge struct") - return err + switch sdk.Opts.ExportDataType { + case BinaryFile: + writer, err := NewBinaryExporter(sdk.Opts.ExportDataPath) + if err != nil { + slog.Error("failed to create BinaryExporter", + "path", sdk.Opts.ExportDataPath, "err", err) + return err + } + + slog.Info( + "created BinaryExporter", + "path", sdk.Opts.ExportDataPath, + ) + sdk.writer = writer + case UDPData: + writer, err := NewSocketExporter(sdk.Opts.ExportUDPAddress, sdk.Opts.ExportUDPPort) + if err != nil { + slog.Error( + "failed to create SocketExporter", + "address", sdk.Opts.ExportUDPAddress, "port", sdk.Opts.ExportUDPPort, "err", err, + ) + return err + } + + slog.Info( + "created SocketExporter", + "address", sdk.Opts.ExportUDPAddress, "port", sdk.Opts.ExportUDPPort, + ) + sdk.writer = writer + } + + return nil +} + +func (sdk *BeamNGSDK) Update() (int, error) { + var nBytes int + var err error + + nBytes, err = sdk.reader.Next(sdk.buffer) + // We check this first because we want to know if we need to loop + if err == io.EOF { + if sdk.Opts.Loop { + err = sdk.reader.Reset() + if err != nil { + return 0, err + } + } + + // NOTE: could we make the reset return the next piece of data? + // We update to the start of the file since we had to reset + nBytes, err = sdk.reader.Next(sdk.buffer) + } + if err != nil { + return nBytes, err + } + + if sdk.Opts.ExportData { + sdk.writer.Write(sdk.buffer) + } + + return nBytes, sdk.parseData(sdk.buffer) +} + +func ParseData(ogData *Outgauge, buffer []byte) error { + ogData.Time = binary.LittleEndian.Uint32(buffer[0:4]) + copy(ogData.Car[:], buffer[4:8]) + ogData.Flags = binary.LittleEndian.Uint16(buffer[8:10]) + ogData.Gear = int8(buffer[10]) + ogData.Plid = int8(buffer[11]) + ogData.Speed = math.Float32frombits(binary.LittleEndian.Uint32(buffer[12:16])) + ogData.RPM = math.Float32frombits(binary.LittleEndian.Uint32(buffer[16:20])) + ogData.Turbo = math.Float32frombits(binary.LittleEndian.Uint32(buffer[20:24])) + ogData.EngTemp = math.Float32frombits(binary.LittleEndian.Uint32(buffer[24:28])) + ogData.Fuel = math.Float32frombits(binary.LittleEndian.Uint32(buffer[28:32])) + ogData.OilPressure = math.Float32frombits(binary.LittleEndian.Uint32(buffer[32:36])) + ogData.OilTemp = math.Float32frombits(binary.LittleEndian.Uint32(buffer[36:40])) + ogData.DashLights = binary.LittleEndian.Uint32(buffer[40:44]) + ogData.ShowLights = binary.LittleEndian.Uint32(buffer[44:48]) + ogData.Throttle = math.Float32frombits(binary.LittleEndian.Uint32(buffer[48:52])) + ogData.Brake = math.Float32frombits(binary.LittleEndian.Uint32(buffer[52:56])) + ogData.Clutch = math.Float32frombits(binary.LittleEndian.Uint32(buffer[56:60])) + copy(ogData.Display1[:], buffer[60:76]) + copy(ogData.Display2[:], buffer[76:92]) + ogData.ID = int32(binary.LittleEndian.Uint32(buffer[92:96])) + + return nil +} + +func (sdk *BeamNGSDK) parseData(buffer []byte) error { + return ParseData(&sdk.Data, buffer) +} + +// ReadData will read new data from the UDP server +func (sdk *BeamNGSDK) ReadData(timeout time.Duration) error { + buf := make([]byte, 2048) + nBytes := sdk.receiver.GetLatest(buf) + if nBytes < 0 { + return fmt.Errorf("no new data") } // Read the binary data into the struct - err = sdk.parseData() + err := sdk.parseData(buf) if err != nil { fmt.Println("Error decoding UDP packet:", err) return err @@ -90,27 +228,39 @@ func (sdk *BeamNGSDK) ReadData() error { return nil } -func (sdk *BeamNGSDK) Close() error { - return sdk.Conn.Close() +func (sdk *BeamNGSDK) GetBuffer() []byte { + buf := make([]byte, 2048) + sdk.receiver.GetLatest(buf) + return buf +} + +func (sdk *BeamNGSDK) GetBufferPtr() []byte { + return sdk.receiver.GetBufferPtr() +} + +func (sdk *BeamNGSDK) Close() { + if sdk.receiver != nil { + sdk.receiver.Close() + } } // Init initializes a BeamNG SDK struct // NOTE: Change this to output a *BeamNGSDK -func Init(ip string, port int) (BeamNGSDK, error) { - var err error - sdk := BeamNGSDK{} - - // Create the connection to the OutGauge server - sdk.Conn, sdk.Addr, err = createUDPConnection(ip, port) - if err != nil { - return BeamNGSDK{}, err - } - - // Initiate the data variables - sdk.Buffer = make([]byte, 1024) - - return sdk, nil -} +// func Init(ip string, port int) (BeamNGSDK, error) { +// var err error +// sdk := BeamNGSDK{} +// +// // Create the connection to the OutGauge server +// sdk.Conn, sdk.Addr, err = createUDPConnection(ip, port, int(unsafe.Sizeof(sdk.Buffer))) +// if err != nil { +// return BeamNGSDK{}, err +// } +// +// // Initiate the data variables +// sdk.Buffer = make([]byte, 1024) +// +// return sdk, nil +// } // SDK utilities diff --git a/example/read_binary_file/main.go b/example/read_binary_file/main.go new file mode 100644 index 0000000..5008ef3 --- /dev/null +++ b/example/read_binary_file/main.go @@ -0,0 +1,60 @@ +package main + +import ( + "fmt" + "log" + "log/slog" + "os" + "time" + + bngsdk "github.com/ESilva15/gobngsdk" +) + +func main() { + fmt.Println("Example of how to read a binary file with the SDK") + + output, err := os.OpenFile("./output.log", os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o755) + if err != nil { + log.Fatalf("Failed to open log file: %+v", err) + } + + logger := slog.New( + slog.NewTextHandler(output, &slog.HandlerOptions{ + Level: slog.LevelDebug, + }), + ) + + sdk, err := bngsdk.NewBngSDK(bngsdk.Options{ + Logger: logger.With("service", "bngsdk"), + SourceType: bngsdk.BinaryFile, + BinSourcePath: "../../sunburstManual.bin", + ExportData: true, + ExportDataType: bngsdk.UDPData, + ExportUDPAddress: "127.0.0.1", + ExportUDPPort: 4444, + Loop: true, + }) + if err != nil { + log.Fatalf("failed to open the sdk: %+v", err) + } + defer sdk.Close() + + ticker := time.NewTicker(time.Second / 60) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + _, err := sdk.Update() + if err != nil { + log.Fatalf("failed to update data: %+v", err) + } + + fmt.Printf("\033[?25l\033[2J\033[H") + fmt.Printf( + "Gear: %d, RPM: %f, Speed: %f\n", + sdk.Data.Gear, sdk.Data.RPM, sdk.Data.Speed, + ) + } + } +} diff --git a/example/read_socket_data/main.go b/example/read_socket_data/main.go new file mode 100644 index 0000000..a3e04b8 --- /dev/null +++ b/example/read_socket_data/main.go @@ -0,0 +1,56 @@ +package main + +import ( + "fmt" + "log" + "log/slog" + "os" + "time" + + bngsdk "github.com/ESilva15/gobngsdk" +) + +func main() { + fmt.Println("Example of how to read from the socket with the SDK") + + output, err := os.OpenFile("./output.log", os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0o755) + if err != nil { + log.Fatalf("Failed to open log file: %+v", err) + } + + logger := slog.New( + slog.NewTextHandler(output, &slog.HandlerOptions{ + Level: slog.LevelDebug, + }), + ) + + sdk, err := bngsdk.NewBngSDK(bngsdk.Options{ + Logger: logger.With("service", "bngsdk"), + SourceType: bngsdk.UDPData, + ImportUDPAddress: "127.0.0.1", + ImportUDPPort: 4444, + }) + if err != nil { + log.Fatalf("failed to open the sdk: %+v", err) + } + defer sdk.Close() + + ticker := time.NewTicker(time.Second / 60) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + _, err := sdk.Update() + if err != nil { + log.Fatalf("failed to update data: %+v", err) + } + + fmt.Printf("\033[?25l\033[2J\033[H") + fmt.Printf( + "Gear: %d, RPM: %f, Speed: %f\n", + sdk.Data.Gear, sdk.Data.RPM, sdk.Data.Speed, + ) + } + } +} diff --git a/reader.go b/reader.go new file mode 100644 index 0000000..e6f4867 --- /dev/null +++ b/reader.go @@ -0,0 +1,26 @@ +package bngsdk + +import "os" + +type BngImporter interface { + Reset() error + Next([]byte) (int, error) +} + +func NewBinaryImporter(file string) (BngImporter, error) { + bin, err := os.Open(file) + if err != nil { + return nil, err + } + + return OgBinReader(bin), nil +} + +func NewSocketImporter(address string, port int) (BngImporter, error) { + reader, err := NewOgUDPReader(address, port) + if err != nil { + return nil, err + } + + return reader, nil +} diff --git a/reader_binary.go b/reader_binary.go new file mode 100644 index 0000000..bf7d1b0 --- /dev/null +++ b/reader_binary.go @@ -0,0 +1,43 @@ +package bngsdk + +import ( + "io" + "unsafe" +) + +type GobReader struct { + TotalRead int64 + File io.ReadSeeker + Buf []byte +} + +func OgBinReader(r io.ReadSeeker) *GobReader { + return &GobReader{ + TotalRead: 0, + File: r, + Buf: make([]byte, unsafe.Sizeof(Outgauge{})), + } +} + +func (g *GobReader) Reset() error { + _, err := g.File.Seek(0, io.SeekStart) + if err != nil { + return err + } + + g.TotalRead = 0 + + return nil +} + +func (g *GobReader) Next(buffer []byte) (int, error) { + nBytes, err := io.ReadFull(g.File, buffer) + if err != nil { + return 0, err + } + + pos, _ := g.File.Seek(0, io.SeekCurrent) + g.TotalRead = pos + + return nBytes, nil +} diff --git a/reader_socket.go b/reader_socket.go new file mode 100644 index 0000000..cd4aea1 --- /dev/null +++ b/reader_socket.go @@ -0,0 +1,42 @@ +package bngsdk + +import ( + "fmt" + "log/slog" +) + +type OgUDPReader struct { + udpConnection *UDPTransport +} + +func NewOgUDPReader(ip string, port int) (*OgUDPReader, error) { + conn, err := NewUDPReader(ip, port) + if err != nil { + return nil, err + } + + return &OgUDPReader{ + udpConnection: conn, + }, nil +} + +func (ogr *OgUDPReader) Reset() error { + // NOTE: what to implement here? + return nil +} + +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") + } + + return 0, nil +} diff --git a/transport.go b/transport.go new file mode 100644 index 0000000..f1f8690 --- /dev/null +++ b/transport.go @@ -0,0 +1,67 @@ +package bngsdk + +import ( + "net" +) + +type UDPTransport struct { + address *net.UDPAddr + connection *net.UDPConn +} + +func NewUDPReader(ip string, port int) (*UDPTransport, error) { + // Define the IP address and port to listen on + addr := &net.UDPAddr{ + IP: net.ParseIP(ip), + Port: port, + } + + conn, err := net.ListenUDP("udp", addr) + if err != nil { + return nil, err + } + + return &UDPTransport{ + address: addr, + connection: conn, + }, nil +} + +func NewUDPWriter(ip string, port int) (*UDPTransport, error) { + // Define the IP address and port to listen on + addr := &net.UDPAddr{ + IP: net.ParseIP(ip), + Port: port, + } + + conn, err := net.DialUDP("udp", nil, addr) + if err != nil { + return nil, err + } + + return &UDPTransport{ + address: addr, + connection: conn, + }, nil +} + +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 + } + + return n, nil +} + +func (ut *UDPTransport) Close() error { + if ut.connection != nil { + return ut.connection.Close() + } + + return nil +} diff --git a/writer.go b/writer.go new file mode 100644 index 0000000..027a087 --- /dev/null +++ b/writer.go @@ -0,0 +1,23 @@ +package bngsdk + +type BngExporter interface { + Write([]byte) (int, error) +} + +func NewSocketExporter(address string, port int) (BngExporter, error) { + writer, err := NewSocketWriter(address, port) + if err != nil { + return nil, err + } + + return writer, nil +} + +func NewBinaryExporter(path string) (BngExporter, error) { + writer, err := NewOgBinWriter(path) + if err != nil { + return nil, err + } + + return writer, nil +} diff --git a/writer_binary.go b/writer_binary.go new file mode 100644 index 0000000..a4daa74 --- /dev/null +++ b/writer_binary.go @@ -0,0 +1,30 @@ +package bngsdk + +import ( + "encoding/binary" + "os" +) + +type OgBinWriter struct { + file *os.File +} + +func NewOgBinWriter(path string) (*OgBinWriter, error) { + outputFile, err := os.Create(path) + if err != nil { + return nil, err + } + + return &OgBinWriter{ + file: outputFile, + }, nil +} + +func (ogw *OgBinWriter) Write(data []byte) (int, error) { + err := binary.Write(ogw.file, binary.LittleEndian, data) + if err != nil { + return 0, err + } + + return len(data), nil +} diff --git a/writer_socket.go b/writer_socket.go new file mode 100644 index 0000000..c35ceac --- /dev/null +++ b/writer_socket.go @@ -0,0 +1,21 @@ +package bngsdk + +type SocketWriter struct { + udpConnection *UDPTransport +} + +func NewSocketWriter(ip string, port int) (*SocketWriter, error) { + conn, err := NewUDPWriter(ip, port) + if err != nil { + return nil, err + } + + return &SocketWriter{ + udpConnection: conn, + }, nil +} + +func (sw *SocketWriter) Write(data []byte) (int, error) { + // slog.Debug("Writing", "data", data) + return sw.udpConnection.Write(data) +} -- 2.54.0 From 67b9cf617c26b0b3da981032bba13ba5b43dfb34 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Tue, 15 Sep 2026 17:41:11 +0100 Subject: [PATCH 2/4] added Close() methods to the structs --- bngsdk.go | 80 ++++++++++++------------------------------------ reader.go | 10 ++---- reader_binary.go | 20 ++++++++++-- reader_socket.go | 8 +++++ writer.go | 1 + writer_binary.go | 8 +++++ writer_socket.go | 8 +++++ 7 files changed, 63 insertions(+), 72 deletions(-) diff --git a/bngsdk.go b/bngsdk.go index bb664ed..10b9e77 100644 --- a/bngsdk.go +++ b/bngsdk.go @@ -3,11 +3,9 @@ package bngsdk import ( "encoding/binary" - "fmt" "io" "log/slog" "math" - "time" "unsafe" ) @@ -42,12 +40,11 @@ type Options struct { } type BeamNGSDK struct { - Opts Options - reader BngImporter - writer BngExporter - Data Outgauge - buffer []byte - receiver *Receiver + Opts Options + reader BngImporter + writer BngExporter + Data Outgauge + buffer []byte } func NewBngSDK(opts Options) (*BeamNGSDK, error) { @@ -76,6 +73,20 @@ func NewBngSDK(opts Options) (*BeamNGSDK, error) { return &sdk, nil } +func (sdk *BeamNGSDK) Close() error { + sdk.buffer = nil + + if sdk.writer != nil { + sdk.writer.Close() + } + + if sdk.reader != nil { + sdk.reader.Close() + } + + return nil +} + func (sdk *BeamNGSDK) openReader() error { switch sdk.Opts.SourceType { case BinaryFile: @@ -209,59 +220,6 @@ func (sdk *BeamNGSDK) parseData(buffer []byte) error { return ParseData(&sdk.Data, buffer) } -// ReadData will read new data from the UDP server -func (sdk *BeamNGSDK) ReadData(timeout time.Duration) error { - buf := make([]byte, 2048) - nBytes := sdk.receiver.GetLatest(buf) - if nBytes < 0 { - return fmt.Errorf("no new data") - } - - // Read the binary data into the struct - err := sdk.parseData(buf) - if err != nil { - fmt.Println("Error decoding UDP packet:", err) - return err - } - - // this means there's new data - return nil -} - -func (sdk *BeamNGSDK) GetBuffer() []byte { - buf := make([]byte, 2048) - sdk.receiver.GetLatest(buf) - return buf -} - -func (sdk *BeamNGSDK) GetBufferPtr() []byte { - return sdk.receiver.GetBufferPtr() -} - -func (sdk *BeamNGSDK) Close() { - if sdk.receiver != nil { - sdk.receiver.Close() - } -} - -// Init initializes a BeamNG SDK struct -// NOTE: Change this to output a *BeamNGSDK -// func Init(ip string, port int) (BeamNGSDK, error) { -// var err error -// sdk := BeamNGSDK{} -// -// // Create the connection to the OutGauge server -// sdk.Conn, sdk.Addr, err = createUDPConnection(ip, port, int(unsafe.Sizeof(sdk.Buffer))) -// if err != nil { -// return BeamNGSDK{}, err -// } -// -// // Initiate the data variables -// sdk.Buffer = make([]byte, 1024) -// -// return sdk, nil -// } - // SDK utilities // ShowLights - functions to check if a given dash light is on [START] diff --git a/reader.go b/reader.go index e6f4867..0cb7963 100644 --- a/reader.go +++ b/reader.go @@ -1,19 +1,13 @@ package bngsdk -import "os" - type BngImporter interface { Reset() error Next([]byte) (int, error) + Close() error } func NewBinaryImporter(file string) (BngImporter, error) { - bin, err := os.Open(file) - if err != nil { - return nil, err - } - - return OgBinReader(bin), nil + return NewOgBinReader(file), nil } func NewSocketImporter(address string, port int) (BngImporter, error) { diff --git a/reader_binary.go b/reader_binary.go index bf7d1b0..7ff138e 100644 --- a/reader_binary.go +++ b/reader_binary.go @@ -2,23 +2,37 @@ package bngsdk import ( "io" + "os" "unsafe" ) type GobReader struct { TotalRead int64 - File io.ReadSeeker + File *os.File Buf []byte } -func OgBinReader(r io.ReadSeeker) *GobReader { +func NewOgBinReader(fp string) *GobReader { + bin, err := os.Open(fp) + if err != nil { + return nil + } + return &GobReader{ TotalRead: 0, - File: r, + File: bin, Buf: make([]byte, unsafe.Sizeof(Outgauge{})), } } +func (g *GobReader) Close() error { + if g.File != nil { + return g.File.Close() + } + + return nil +} + func (g *GobReader) Reset() error { _, err := g.File.Seek(0, io.SeekStart) if err != nil { diff --git a/reader_socket.go b/reader_socket.go index cd4aea1..b1ccf8f 100644 --- a/reader_socket.go +++ b/reader_socket.go @@ -20,6 +20,14 @@ func NewOgUDPReader(ip string, port int) (*OgUDPReader, error) { }, nil } +func (ogr *OgUDPReader) Close() error { + if ogr.udpConnection != nil { + return ogr.udpConnection.Close() + } + + return nil +} + func (ogr *OgUDPReader) Reset() error { // NOTE: what to implement here? return nil diff --git a/writer.go b/writer.go index 027a087..5d7eaad 100644 --- a/writer.go +++ b/writer.go @@ -2,6 +2,7 @@ package bngsdk type BngExporter interface { Write([]byte) (int, error) + Close() error } func NewSocketExporter(address string, port int) (BngExporter, error) { diff --git a/writer_binary.go b/writer_binary.go index a4daa74..0bc16d9 100644 --- a/writer_binary.go +++ b/writer_binary.go @@ -20,6 +20,14 @@ func NewOgBinWriter(path string) (*OgBinWriter, error) { }, nil } +func (ogw *OgBinWriter) Close() error { + if ogw.file != nil { + return ogw.file.Close() + } + + return nil +} + func (ogw *OgBinWriter) Write(data []byte) (int, error) { err := binary.Write(ogw.file, binary.LittleEndian, data) if err != nil { diff --git a/writer_socket.go b/writer_socket.go index c35ceac..9b347c3 100644 --- a/writer_socket.go +++ b/writer_socket.go @@ -15,6 +15,14 @@ func NewSocketWriter(ip string, port int) (*SocketWriter, error) { }, nil } +func (sw *SocketWriter) Close() error { + if sw.udpConnection != nil { + return sw.udpConnection.Close() + } + + return nil +} + func (sw *SocketWriter) Write(data []byte) (int, error) { // slog.Debug("Writing", "data", data) return sw.udpConnection.Write(data) -- 2.54.0 From ec2ed17f43b23edda8cca842566f3cb65f0b86c4 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Wed, 16 Sep 2026 14:36:04 +0100 Subject: [PATCH 3/4] 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 { -- 2.54.0 From ba6b2c700b0d0790bdd4e8ace76275b8560abf9a Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Wed, 16 Sep 2026 15:06:27 +0100 Subject: [PATCH 4/4] reduced allocations --- README.md | 11 +++++------ bngsdk_test.go | 42 +++++++++++++++++++++++++++++------------- transport.go | 46 +++++++++++++++++++++++++++++++++++++--------- 3 files changed, 71 insertions(+), 28 deletions(-) diff --git a/README.md b/README.md index 80cce45..eb35639 100644 --- a/README.md +++ b/README.md @@ -3,7 +3,7 @@ Simple SDK to interact with BeamNG.drive OutGauge data. ## Development ### Performance -`go test -bench=BenchmarkReadData -benchmem -memprofile=mem.pprof` +`go test -bench=BenchmarkUpdate -benchmem -memprofile=mem.pprof` replace the function to be tested Use `go tool pprof` to analyze the results @@ -15,18 +15,17 @@ goos: linux goarch: amd64 pkg: github.com/ESilva15/gobngsdk cpu: AMD Ryzen 7 5800X3D 8-Core Processor -BenchmarkReadData-16 320931 4058 ns/op 100 B/op 2 allocs/op +BenchmarkReadData-16 362514 3188 ns/op 4 B/op 1 allocs/op PASS -ok github.com/ESilva15/gobngsdk 1.345s +ok github.com/ESilva15/gobngsdk 1.194s # New footprint goos: linux goarch: amd64 pkg: github.com/ESilva15/gobngsdk cpu: AMD Ryzen 7 5800X3D 8-Core Processor -BenchmarkReadData-16 362514 3188 ns/op 4 B/op 1 allocs/op +BenchmarkUpdate-16 159170 7287 ns/op 4 B/op 1 allocs/op PASS -ok github.com/ESilva15/gobngsdk 1.194s - +ok github.com/ESilva15/gobngsdk 1.241s # Pretty good enough. I can finally go be productive instead of "procrastinating" here ``` diff --git a/bngsdk_test.go b/bngsdk_test.go index 4d17946..3ad85b5 100644 --- a/bngsdk_test.go +++ b/bngsdk_test.go @@ -3,35 +3,51 @@ package bngsdk import ( "bytes" "encoding/binary" + "io" + "log/slog" "net" "testing" ) -func BenchmarkReadData(b *testing.B) { - // Spin up an UDP server - sdk, err := Init("127.0.0.1", 0) +func BenchmarkUpdate(b *testing.B) { + // Silence logging output so slog calls don't pollute benchmark stats + slogger := slog.New(slog.NewTextHandler(io.Discard, nil)) + + // Initialize the SDK with port 0 to bind to an OS-assigned ephemeral port + sdk, err := NewBngSDK(Options{ + Logger: slogger, + SourceType: UDPData, + ImportUDPAddress: "127.0.0.1", + ImportUDPPort: 0, + }) if err != nil { b.Fatalf("Failed to initialize SDK: %v", err) } defer sdk.Close() - // Retrieve the actual assigned UDP address - localAddr := sdk.Conn.LocalAddr().(*net.UDPAddr) + // Access the underlying reader connection to determine the dynamically bound port + ogReader, ok := sdk.reader.(*OgUDPReader) + if !ok || ogReader.udpConnection == nil || ogReader.udpConnection.connection == nil { + b.Fatalf("Failed to retrieve underlying UDP connection") + } - // Start a client to stream data - clientConn, err := net.DialUDP("udp", nil, localAddr) + serverAddr := ogReader.udpConnection.connection.LocalAddr().(*net.UDPAddr) + + // Dial the UDP socket as a client to send test data + clientConn, err := net.DialUDP("udp", nil, serverAddr) if err != nil { - b.Fatalf("Failed to dial local UDP socket: %v", err) + b.Fatalf("Failed to dial UDP server: %v", err) } defer clientConn.Close() - // Pre serialize some data + // Pre-serialize a dummy Outgauge struct matching the required byte layout dummyOutgauge := Outgauge{ Time: 424242, Car: [4]byte{'P', 'E', 'R', 'F'}, Speed: 45.2, RPM: 3500.0, } + var buf bytes.Buffer if err := binary.Write(&buf, binary.LittleEndian, dummyOutgauge); err != nil { b.Fatalf("Failed to serialize dummy struct: %v", err) @@ -42,16 +58,16 @@ func BenchmarkReadData(b *testing.B) { b.ReportAllocs() for i := 0; i < b.N; i++ { - // Feed a packet into the network buffer right before reading it + // Feed a packet into the network transport socket _, err := clientConn.Write(packetBytes) if err != nil { b.Fatalf("Failed to write to UDP socket: %v", err) } - // Execute the target function - err = sdk.ReadData() + // Run the main API loop method + _, err = sdk.Update() if err != nil { - b.Fatalf("ReadData failed at iteration %d: %v", i, err) + b.Fatalf("Update failed at iteration %d: %v", i, err) } } } diff --git a/transport.go b/transport.go index cc91469..e3e9d1a 100644 --- a/transport.go +++ b/transport.go @@ -4,6 +4,7 @@ import ( "errors" "log/slog" "net" + "sync" "time" "unsafe" ) @@ -12,12 +13,25 @@ import ( var ( ErrNoData = errors.New("no new data available") readTimeout = (time.Second / 60) * 5 // N missed frames at 60fps + packetPool = sync.Pool{ + New: func() any { + var b packetBuffer + return &b + }, + } ) +type packetBuffer [unsafe.Sizeof(Outgauge{})]byte + +type frame struct { + Buf *packetBuffer + Len int +} + type UDPTransport struct { address *net.UDPAddr connection *net.UDPConn - dataChan chan []byte + dataChan chan frame } func NewUDPReader(ip string, port int) (*UDPTransport, error) { @@ -35,7 +49,7 @@ func NewUDPReader(ip string, port int) (*UDPTransport, error) { udpT := UDPTransport{ address: addr, connection: conn, - dataChan: make(chan []byte, 1), + dataChan: make(chan frame, 1), } go udpT.udpSink() @@ -69,17 +83,28 @@ func (ut *UDPTransport) udpSink() { for { ut.connection.SetReadDeadline(time.Now().Add(readTimeout)) - buf := make([]byte, unsafe.Sizeof(Outgauge{})) - nBytes, _, err := ut.connection.ReadFromUDP(buf) + bufPtr := packetPool.Get().(*packetBuffer) + + nBytes, _, err := ut.connection.ReadFromUDP(bufPtr[:]) if err != nil { return } + frame := frame{ + Buf: bufPtr, + Len: nBytes, + } + select { - case ut.dataChan <- buf[:nBytes]: + case ut.dataChan <- frame: + // Packet sent successfuly default: - <-ut.dataChan - ut.dataChan <- buf[:nBytes] + select { + case oldFrame := <-ut.dataChan: + packetPool.Put(oldFrame.Buf) + default: + } + ut.dataChan <- frame } } } @@ -89,13 +114,16 @@ func (ut *UDPTransport) Write(data []byte) (int, error) { } func (ut *UDPTransport) Read(buffer []byte) (int, error) { - data, ok := <-ut.dataChan + latestFrame, ok := <-ut.dataChan if !ok { slog.Error(ErrNoData.Error()) return 0, ErrNoData } - return copy(buffer, data), nil + nBytes := copy(buffer, latestFrame.Buf[:latestFrame.Len]) + packetPool.Put(latestFrame.Buf) + + return nBytes, nil } func (ut *UDPTransport) Close() error { -- 2.54.0