A better SDK with more features and more consistency
This commit is contained in:
@@ -4,8 +4,11 @@ package bngsdk
|
|||||||
import (
|
import (
|
||||||
"encoding/binary"
|
"encoding/binary"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"log/slog"
|
||||||
"math"
|
"math"
|
||||||
"net"
|
"time"
|
||||||
|
"unsafe"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -15,72 +18,207 @@ const (
|
|||||||
DefaultUDPPort = 4444
|
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 {
|
type BeamNGSDK struct {
|
||||||
Addr *net.UDPAddr
|
Opts Options
|
||||||
Conn *net.UDPConn
|
reader BngImporter
|
||||||
Buffer []byte
|
writer BngExporter
|
||||||
Data Outgauge
|
Data Outgauge
|
||||||
|
buffer []byte
|
||||||
|
receiver *Receiver
|
||||||
}
|
}
|
||||||
|
|
||||||
func createUDPConnection(ip string, port int) (*net.UDPConn, *net.UDPAddr, error) {
|
func NewBngSDK(opts Options) (*BeamNGSDK, error) {
|
||||||
// Define the IP address and port to listen on
|
sdk := BeamNGSDK{
|
||||||
addr := &net.UDPAddr{
|
Opts: opts,
|
||||||
IP: net.ParseIP(ip),
|
buffer: make([]byte, unsafe.Sizeof(Outgauge{})),
|
||||||
Port: port,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Create a UDP socket
|
// Set the passed logger as the default logger
|
||||||
conn, err := net.ListenUDP("udp", addr)
|
slog.SetDefault(sdk.Opts.Logger)
|
||||||
|
|
||||||
|
var err error
|
||||||
|
|
||||||
|
err = sdk.openReader()
|
||||||
if err != nil {
|
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) openReader() error {
|
||||||
func (sdk *BeamNGSDK) parseData() error {
|
switch sdk.Opts.SourceType {
|
||||||
sdk.Data.Time = binary.LittleEndian.Uint32(sdk.Buffer[0:4])
|
case BinaryFile:
|
||||||
copy(sdk.Data.Car[:], sdk.Buffer[4:8])
|
reader, err := NewBinaryImporter(sdk.Opts.BinSourcePath)
|
||||||
sdk.Data.Flags = binary.LittleEndian.Uint16(sdk.Buffer[8:10])
|
if err != nil {
|
||||||
sdk.Data.Gear = int8(sdk.Buffer[10])
|
slog.Error("failed to create BinaryImporter",
|
||||||
sdk.Data.Plid = int8(sdk.Buffer[11])
|
"path", sdk.Opts.BinSourcePath, "err", err)
|
||||||
sdk.Data.Speed = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[12:16]))
|
return err
|
||||||
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]))
|
slog.Info(
|
||||||
sdk.Data.Fuel = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[28:32]))
|
"created BinaryImporter",
|
||||||
sdk.Data.OilPressure = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[32:36]))
|
"path", sdk.Opts.BinSourcePath,
|
||||||
sdk.Data.OilTemp = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[36:40]))
|
)
|
||||||
sdk.Data.DashLights = binary.LittleEndian.Uint32(sdk.Buffer[40:44])
|
sdk.reader = reader
|
||||||
sdk.Data.ShowLights = binary.LittleEndian.Uint32(sdk.Buffer[44:48])
|
case UDPData:
|
||||||
sdk.Data.Throttle = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[48:52]))
|
reader, err := NewSocketImporter(sdk.Opts.ImportUDPAddress, sdk.Opts.ImportUDPPort)
|
||||||
sdk.Data.Brake = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[52:56]))
|
if err != nil {
|
||||||
sdk.Data.Throttle = math.Float32frombits(binary.LittleEndian.Uint32(sdk.Buffer[56:60]))
|
slog.Error(
|
||||||
copy(sdk.Data.Display1[:], sdk.Buffer[60:76])
|
"failed to create SocketImporter",
|
||||||
copy(sdk.Data.Display2[:], sdk.Buffer[76:92])
|
"address", sdk.Opts.ImportUDPAddress, "port", sdk.Opts.ImportUDPPort, "err", err,
|
||||||
sdk.Data.ID = int32(binary.LittleEndian.Uint32(sdk.Buffer[92:96]))
|
)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
slog.Info(
|
||||||
|
"created SocketImporter",
|
||||||
|
"address", sdk.Opts.ImportUDPAddress, "port", sdk.Opts.ImportUDPPort,
|
||||||
|
)
|
||||||
|
sdk.reader = reader
|
||||||
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// ReadData will read new data from the UDP server
|
func (sdk *BeamNGSDK) openWriter() error {
|
||||||
func (sdk *BeamNGSDK) ReadData() error {
|
// if the user didn't request data export we don't need a writer
|
||||||
// Receive data from the socket
|
if !sdk.Opts.ExportData {
|
||||||
n, _, err := sdk.Conn.ReadFromUDP(sdk.Buffer)
|
return nil
|
||||||
if err != nil {
|
|
||||||
fmt.Println("Error reading from UDP:", err)
|
|
||||||
return err
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Check if enough data was received to fill our struct
|
switch sdk.Opts.ExportDataType {
|
||||||
if n < outgaugeSize {
|
case BinaryFile:
|
||||||
fmt.Println("Received packet too small for Outgauge struct")
|
writer, err := NewBinaryExporter(sdk.Opts.ExportDataPath)
|
||||||
return err
|
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
|
// Read the binary data into the struct
|
||||||
err = sdk.parseData()
|
err := sdk.parseData(buf)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
fmt.Println("Error decoding UDP packet:", err)
|
fmt.Println("Error decoding UDP packet:", err)
|
||||||
return err
|
return err
|
||||||
@@ -90,27 +228,39 @@ func (sdk *BeamNGSDK) ReadData() error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sdk *BeamNGSDK) Close() error {
|
func (sdk *BeamNGSDK) GetBuffer() []byte {
|
||||||
return sdk.Conn.Close()
|
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
|
// Init initializes a BeamNG SDK struct
|
||||||
// NOTE: Change this to output a *BeamNGSDK
|
// NOTE: Change this to output a *BeamNGSDK
|
||||||
func Init(ip string, port int) (BeamNGSDK, error) {
|
// func Init(ip string, port int) (BeamNGSDK, error) {
|
||||||
var err error
|
// var err error
|
||||||
sdk := BeamNGSDK{}
|
// sdk := BeamNGSDK{}
|
||||||
|
//
|
||||||
// Create the connection to the OutGauge server
|
// // Create the connection to the OutGauge server
|
||||||
sdk.Conn, sdk.Addr, err = createUDPConnection(ip, port)
|
// sdk.Conn, sdk.Addr, err = createUDPConnection(ip, port, int(unsafe.Sizeof(sdk.Buffer)))
|
||||||
if err != nil {
|
// if err != nil {
|
||||||
return BeamNGSDK{}, err
|
// return BeamNGSDK{}, err
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
// Initiate the data variables
|
// // Initiate the data variables
|
||||||
sdk.Buffer = make([]byte, 1024)
|
// sdk.Buffer = make([]byte, 1024)
|
||||||
|
//
|
||||||
return sdk, nil
|
// return sdk, nil
|
||||||
}
|
// }
|
||||||
|
|
||||||
// SDK utilities
|
// SDK utilities
|
||||||
|
|
||||||
|
|||||||
@@ -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 / 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,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
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user