Compare commits
4
Commits
v1.1.3
...
better-sdk
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ba6b2c700b | ||
|
|
ec2ed17f43 | ||
|
|
67b9cf617c | ||
|
|
a7fbf74f0a |
@@ -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
|
||||
```
|
||||
|
||||
@@ -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
|
||||
|
||||
+29
-13
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
+135
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
Reference in New Issue
Block a user