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.go b/bngsdk.go index 9a46f8a..1299864 100644 --- a/bngsdk.go +++ b/bngsdk.go @@ -3,9 +3,10 @@ package bngsdk import ( "encoding/binary" - "fmt" + "io" + "log/slog" "math" - "net" + "unsafe" ) const ( @@ -15,101 +16,208 @@ 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 + Opts Options + reader BngImporter + writer BngExporter Data Outgauge + buffer []byte } -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 -} - -// 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])) - - 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) + err = sdk.openWriter() if err != nil { - fmt.Println("Error reading from UDP:", err) - return err + slog.Error("failed to open writer", "err", err) + return nil, err } - // Check if enough data was received to fill our struct - if n < outgaugeSize { - fmt.Println("Received packet too small for Outgauge struct") - return err - } - - // Read the binary data into the struct - err = sdk.parseData() - if err != nil { - fmt.Println("Error decoding UDP packet:", err) - return err - } - - // this means there's new data - return nil + return &sdk, nil } func (sdk *BeamNGSDK) Close() error { - return sdk.Conn.Close() -} + sdk.buffer = nil -// 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 + if sdk.writer != nil { + sdk.writer.Close() } - // Initiate the data variables - sdk.Buffer = make([]byte, 1024) + if sdk.reader != nil { + sdk.reader.Close() + } - return sdk, nil + return nil +} + +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 +} + +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 + } + + 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 0, 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) } // SDK utilities 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/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..1b99a24 --- /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) + 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..0cb7963 --- /dev/null +++ b/reader.go @@ -0,0 +1,20 @@ +package bngsdk + +type BngImporter interface { + Reset() error + Next([]byte) (int, error) + Close() error +} + +func NewBinaryImporter(file string) (BngImporter, error) { + return NewOgBinReader(file), 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..7ff138e --- /dev/null +++ b/reader_binary.go @@ -0,0 +1,57 @@ +package bngsdk + +import ( + "io" + "os" + "unsafe" +) + +type GobReader struct { + TotalRead int64 + File *os.File + Buf []byte +} + +func NewOgBinReader(fp string) *GobReader { + bin, err := os.Open(fp) + if err != nil { + return nil + } + + return &GobReader{ + TotalRead: 0, + 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 { + 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..5b99bb0 --- /dev/null +++ b/reader_socket.go @@ -0,0 +1,49 @@ +package bngsdk + +import ( + "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 { + return nil, err + } + + return &OgUDPReader{ + udpConnection: conn, + }, 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 +} + +func (ogr *OgUDPReader) Next(buffer []byte) (int, error) { + 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, ErrInvalidOutgaugeData + } + + return nBytes, nil +} diff --git a/transport.go b/transport.go new file mode 100644 index 0000000..e3e9d1a --- /dev/null +++ b/transport.go @@ -0,0 +1,135 @@ +package bngsdk + +import ( + "errors" + "log/slog" + "net" + "sync" + "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 + 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 frame +} + +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 + } + + udpT := UDPTransport{ + address: addr, + connection: conn, + dataChan: make(chan frame, 1), + } + + go udpT.udpSink() + + return &udpT, 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 +} + +// 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)) + + bufPtr := packetPool.Get().(*packetBuffer) + + nBytes, _, err := ut.connection.ReadFromUDP(bufPtr[:]) + if err != nil { + return + } + + frame := frame{ + Buf: bufPtr, + Len: nBytes, + } + + select { + case ut.dataChan <- frame: + // Packet sent successfuly + default: + select { + case oldFrame := <-ut.dataChan: + packetPool.Put(oldFrame.Buf) + default: + } + ut.dataChan <- frame + } + } +} + +func (ut *UDPTransport) Write(data []byte) (int, error) { + return ut.connection.Write(data) +} + +func (ut *UDPTransport) Read(buffer []byte) (int, error) { + latestFrame, ok := <-ut.dataChan + if !ok { + slog.Error(ErrNoData.Error()) + return 0, ErrNoData + } + + nBytes := copy(buffer, latestFrame.Buf[:latestFrame.Len]) + packetPool.Put(latestFrame.Buf) + + return nBytes, 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..5d7eaad --- /dev/null +++ b/writer.go @@ -0,0 +1,24 @@ +package bngsdk + +type BngExporter interface { + Write([]byte) (int, error) + Close() 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..0bc16d9 --- /dev/null +++ b/writer_binary.go @@ -0,0 +1,38 @@ +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) 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 { + return 0, err + } + + return len(data), nil +} diff --git a/writer_socket.go b/writer_socket.go new file mode 100644 index 0000000..9b347c3 --- /dev/null +++ b/writer_socket.go @@ -0,0 +1,29 @@ +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) 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) +}