Update beamng to 2.0.0 #3

Merged
esilva merged 2 commits from update-beamng-to-2.0.0 into master 2026-09-17 00:10:11 +01:00
9 changed files with 112 additions and 286 deletions
+2
View File
@@ -1,2 +1,4 @@
*.bin
*.work*
*.pprof
*.log
+6 -3
View File
@@ -3,6 +3,7 @@ package cmd
import (
"context"
"fmt"
"log/slog"
beamng "github.com/ESilva15/TelemetryMockserver/internal/mockservers/beamng"
@@ -59,14 +60,16 @@ func replayAction(cmd *cobra.Command, args []string) {
replayer, err := beamng.NewReplayer(address, port, inputFile)
if err != nil {
fmt.Printf("Something went wrong setting up the player: %+v", err)
slog.Error("Something went wrong setting up the player", "err", err)
return
}
// NOTE: is this doing anything at all??
ctx := context.Background()
if err := replayer.Replay(ctx, loop); err != nil {
fmt.Printf("Something went wrong while playing the file: %v", err)
err = replayer.Replay(ctx, loop)
if err != nil {
slog.Error("Something went wrong while playing the file", "err", err)
return
}
}
+3
View File
@@ -0,0 +1,3 @@
package constants
const ProgramName = "TelemetryMockserver"
+1 -1
View File
@@ -3,7 +3,7 @@ module github.com/ESilva15/TelemetryMockserver
go 1.23.2
require (
github.com/ESilva15/gobngsdk v1.1.3
github.com/ESilva15/gobngsdk v2.2.0
github.com/ESilva15/goirsdk v0.3.0
github.com/spf13/cobra v1.10.1
)
-45
View File
@@ -1,45 +0,0 @@
package mockserver
import (
"io"
"unsafe"
sdk "github.com/ESilva15/gobngsdk"
)
type GobReader struct {
TotalRead int64
File io.ReadSeeker
Buf []byte
}
func NewGobReader(r io.ReadSeeker) *GobReader {
return &GobReader{
TotalRead: 0,
File: r,
Buf: make([]byte, unsafe.Sizeof(sdk.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) error {
_, err := io.ReadFull(g.File, buffer)
if err != nil {
return err
}
pos, _ := g.File.Seek(0, io.SeekCurrent)
g.TotalRead = pos
return nil
}
+22 -54
View File
@@ -3,82 +3,59 @@ package mockserver
import (
"bytes"
"context"
"encoding/binary"
"fmt"
"log/slog"
"os"
"sync"
"time"
"github.com/ESilva15/TelemetryMockserver/constants"
bngsdk "github.com/ESilva15/gobngsdk"
)
type recorderViewData struct {
TotalBytes int
SDK *bngsdk.BeamNGSDK
TotalBytes int64
Og bngsdk.Outgauge
}
type Recorder struct {
SDK bngsdk.BeamNGSDK
OutputFile *os.File
TotalBytes int
SDK *bngsdk.BeamNGSDK
// Views
mut sync.RWMutex
viewDataMut sync.RWMutex
viewData recorderViewData
viewCh chan *recorderViewData
recorderCh chan []byte
}
func NewRecorder(fp string, address string, port int) (*Recorder, error) {
var recorder Recorder
var err error
recorder.SDK, err = bngsdk.Init(address, port)
recorder.SDK, err = bngsdk.NewBngSDK(bngsdk.Options{
Logger: slog.Default().With("SDK", "BeamNG"),
SourceType: bngsdk.UDPData,
ImportUDPAddress: address,
ImportUDPPort: port,
ExportData: true,
ExportDataType: bngsdk.BinaryFile,
ExportDataPath: fp,
})
if err != nil {
return &Recorder{}, err
}
recorder.OutputFile, err = os.Create(fp)
if err != nil {
return &Recorder{}, err
return nil, err
}
recorder.viewData = recorderViewData{}
recorder.viewCh = make(chan *recorderViewData, 1)
recorder.recorderCh = make(chan []byte, 1)
return &recorder, nil
}
func (r *Recorder) Close() {
r.SDK.Close()
if r.OutputFile != nil {
r.OutputFile.Close()
}
}
func (r *Recorder) record(ctx context.Context) {
for {
select {
case <-ctx.Done():
return
case data := <-r.recorderCh:
r.mut.Lock()
err := binary.Write(r.OutputFile, binary.LittleEndian, r.SDK.Data)
r.TotalBytes += len(data)
r.mut.Unlock()
if err != nil {
// NOTE: find a way of logging this somehow
}
}
}
}
func (r *Recorder) view(ctx context.Context) {
var buf bytes.Buffer
var nBytes int
buf.Grow(2048)
@@ -89,16 +66,15 @@ func (r *Recorder) view(ctx context.Context) {
case viewData := <-r.viewCh:
buf.Reset()
buf.WriteString("\x1b[2J\x1b[H")
fmt.Fprintf(&buf, "\x1b]0;%s - Recording ", ProgramName)
stringifyRecordingProgress(&buf, nBytes)
fmt.Fprintf(&buf, "\x1b]0;%s - Recording ", constants.ProgramName)
stringifyRecordingProgress(&buf, viewData.TotalBytes)
fmt.Fprintf(&buf, "\x07")
stringifyRecordingProgress(&buf, nBytes)
stringifyRecordingProgress(&buf, viewData.TotalBytes)
fmt.Fprintf(&buf, "\n\n")
r.viewDataMut.RLock()
nBytes = viewData.TotalBytes
stringifyOutgaugeData(&buf, &r.SDK)
stringifyOutgaugeData(&buf, &viewData.Og)
r.viewDataMut.RUnlock()
_, _ = buf.WriteTo(os.Stdout)
@@ -111,7 +87,6 @@ func (r *Recorder) Record(ctx context.Context) error {
ticker := time.NewTicker(time.Second / 60)
defer ticker.Stop()
go r.record(ctx)
go r.view(ctx)
for {
@@ -119,12 +94,13 @@ func (r *Recorder) Record(ctx context.Context) error {
case <-ctx.Done():
return nil
case <-ticker.C:
err := r.SDK.ReadData()
og, err := r.SDK.Update()
if err != nil {
return err
}
r.viewData.TotalBytes = r.TotalBytes
r.viewData.TotalBytes = r.SDK.GetTotalWritten()
r.viewData.Og = *og
// Send the data to the view
select {
@@ -133,14 +109,6 @@ func (r *Recorder) Record(ctx context.Context) error {
default:
// Dropped the frame!
}
// Write the data to the file
select {
case r.recorderCh <- r.SDK.Buffer:
// Sent the data
default:
// Dropped the frame!
}
}
}
}
+29 -93
View File
@@ -5,32 +5,30 @@ package mockserver
import (
"bytes"
"context"
"encoding/binary"
"fmt"
"io"
"log/slog"
"math"
"os"
"sync"
"time"
"unsafe"
"github.com/ESilva15/TelemetryMockserver/constants"
bngsdk "github.com/ESilva15/gobngsdk"
)
type ViewData struct {
SDK *bngsdk.BeamNGSDK
Outgauge bngsdk.Outgauge
SizeRead int64
FileSize int64
}
// Replayer does the replaying
// Should we make a "player" struct that can record and replay?
type Replayer struct {
SDK bngsdk.BeamNGSDK
DataSourcePath string
Socket *UDPTransport
SDK *bngsdk.BeamNGSDK
// Streams
dataViewCh chan ViewData
socketCh chan []byte
// Mut
mut sync.RWMutex
@@ -40,21 +38,23 @@ type Replayer struct {
}
func NewReplayer(address string, port int, fp string) (*Replayer, error) {
udp, err := NewUDPTransport(address, port)
sdk, err := bngsdk.NewBngSDK(bngsdk.Options{
Logger: slog.Default().With("SDK", "BeamNG"),
SourceType: bngsdk.BinaryFile,
BinSourcePath: fp,
ExportData: true,
ExportDataType: bngsdk.UDPData,
ExportUDPAddress: address,
ExportUDPPort: port,
Loop: true,
})
if err != nil {
return nil, err
}
replayer := &Replayer{
DataSourcePath: fp,
SDK: bngsdk.BeamNGSDK{
Data: bngsdk.Outgauge{},
Buffer: make([]byte, unsafe.Sizeof(bngsdk.Outgauge{})),
},
Socket: udp,
viewData: ViewData{},
SDK: sdk,
dataViewCh: make(chan ViewData, 1),
socketCh: make(chan []byte, 1),
}
return replayer, nil
@@ -62,15 +62,7 @@ func NewReplayer(address string, port int, fp string) (*Replayer, error) {
// renderToTerminal will render the data for the users viewing pleasure
func (r *Replayer) renderToTerminal(ctx context.Context) {
fileInfo, err := os.Stat(r.DataSourcePath)
if err != nil {
// NOTE: learn how to handle this error
// return fmt.Errorf("error stating file: %v", err)
}
var buf bytes.Buffer
var bytesReader bytes.Reader
buf.Grow(2048)
for {
@@ -79,60 +71,25 @@ func (r *Replayer) renderToTerminal(ctx context.Context) {
return
case data := <-r.dataViewCh:
// Reset to the start of the terminal
percent := int(float64(data.SizeRead) / float64(fileInfo.Size()) * 100)
// percent := int(float64(data.SizeRead) / float64(data.FileSize) * 100)
percent := int(math.Round((float64(data.SizeRead) / float64(data.FileSize)) * 100))
buf.Reset()
buf.WriteString("\x1b[2J\x1b[H")
fmt.Fprintf(&buf, "\x1b]0;%s - Replaying %d%%\x07", ProgramName, percent)
bytesReader.Reset(data.SDK.Buffer)
err := binary.Read(&bytesReader, binary.LittleEndian, &r.SDK.Data)
if err != nil {
fmt.Fprintf(&buf, "FAILED TO PARSE DATA\nError: %+v", err)
_, _ = buf.WriteTo(os.Stdout)
continue
}
fmt.Fprintf(&buf, "\x1b]0;%s - Replaying %d%%\x07", constants.ProgramName, percent)
fmt.Fprintf(&buf, "Replayed: %d%%\n", percent)
stringifyOutgaugeData(&buf, data.SDK)
stringifyOutgaugeData(&buf, &data.Outgauge)
buf.WriteTo(os.Stdout)
}
}
}
// writeToUDPSocket will write the telemetry data to the UDP socket
func (r *Replayer) writeToUDPSocket(ctx context.Context) {
for {
select {
case <-ctx.Done():
return
case data := <-r.socketCh:
r.mut.RLock()
_, err := r.Socket.Send(data)
r.mut.RUnlock()
if err != nil {
panic(fmt.Sprintf("error writing buffer to socket: %+v", err))
continue
// NOTE: log the error somewhere maybe
// return err
}
}
}
}
// Replay replays a given file <fp> in a UDP server <addr>:<port>
func (r *Replayer) Replay(ctx context.Context, loop bool) error {
bin, err := os.Open(r.DataSourcePath)
if err != nil {
return fmt.Errorf("error opening file: %v", err)
}
reader := NewGobReader(bin)
go r.renderToTerminal(ctx)
go r.writeToUDPSocket(ctx)
ticker := time.NewTicker(time.Second / 60)
defer ticker.Stop()
@@ -142,30 +99,18 @@ func (r *Replayer) Replay(ctx context.Context, loop bool) error {
return ctx.Err()
case <-ticker.C:
r.mut.Lock()
err := reader.Next(r.SDK.Buffer)
r.mut.Unlock()
if err == io.EOF {
if !loop {
return r.Socket.Close()
}
err = reader.Reset()
if err != nil {
return err
}
continue
}
og, err := r.SDK.Update()
if err != nil {
// TODO: log here
slog.Error("an error occurred when updating", "err", err)
return err
}
// NOTE: Really like this???
r.viewData.SDK = &r.SDK
r.viewData.SizeRead = reader.TotalRead
r.mut.Lock()
r.viewData.SizeRead = r.SDK.GetTotalRead()
r.viewData.FileSize = r.SDK.GetSourceSize()
r.viewData.Outgauge = *og
r.mut.Unlock()
// Send the data to the view
select {
@@ -174,15 +119,6 @@ func (r *Replayer) Replay(ctx context.Context, loop bool) error {
default:
// Dropped the frame!
}
// Send the data to the UDP socket
select {
case r.socketCh <- r.SDK.Buffer:
// Sent the data
default:
// Dropped the frame!
}
}
}
}
-41
View File
@@ -1,41 +0,0 @@
package mockserver
import (
"fmt"
"net"
)
const ProgramName = "TelemetryMockerserver"
type Transport interface {
Send(data []byte) error
Close() error
}
type UDPTransport struct {
Conn *net.UDPConn
Addr *net.UDPAddr
}
func NewUDPTransport(address string, port int) (*UDPTransport, error) {
addr, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", address, port))
if err != nil {
return nil, err
}
conn, err := net.ListenUDP("udp", nil)
if err != nil {
return nil, err
}
return &UDPTransport{Conn: conn, Addr: addr}, nil
}
// Send will send a byte array of data trough the UDP server
func (u *UDPTransport) Send(data []byte) (int, error) {
return u.Conn.WriteToUDP(data, u.Addr)
}
func (u *UDPTransport) Close() error {
return u.Conn.Close()
}
+49 -49
View File
@@ -14,7 +14,7 @@ const (
GiB
)
func stringifyRecordingProgress(s *bytes.Buffer, nBytes int) {
func stringifyRecordingProgress(s *bytes.Buffer, nBytes int64) {
if nBytes < KiB {
fmt.Fprintf(s, "%d B", nBytes)
} else if nBytes < MiB {
@@ -26,64 +26,64 @@ func stringifyRecordingProgress(s *bytes.Buffer, nBytes int) {
}
}
func stringifyOutgaugeData(s *bytes.Buffer, sdk *bngsdk.BeamNGSDK) {
func stringifyOutgaugeData(s *bytes.Buffer, og *bngsdk.Outgauge) {
// NOTE: write a string serialization function on the SDK itself
fmt.Fprint(s, "Outgauge {\n")
fmt.Fprintf(s, " Time: %d ms\n", sdk.Data.Time)
fmt.Fprintf(s, " Car: %s\n", sdk.Data.Car)
fmt.Fprintf(s, " Flags: %b\n", sdk.Data.Flags)
fmt.Fprintf(s, " Gear: %d\n", sdk.Data.Gear)
fmt.Fprintf(s, " Plid: %d\n", sdk.Data.Plid)
fmt.Fprintf(s, " Speed: %f m/s\n", sdk.Data.Speed)
fmt.Fprintf(s, " RPM: %f RPM\n", sdk.Data.RPM)
fmt.Fprintf(s, " Turbo: %f Bar\n", sdk.Data.Turbo)
fmt.Fprintf(s, " EngTemp: %f °C\n", sdk.Data.EngTemp)
fmt.Fprintf(s, " Fuel: %f\n", sdk.Data.Fuel)
fmt.Fprintf(s, " OilPressure: %f Bar\n", sdk.Data.OilPressure)
fmt.Fprintf(s, " OilTemp: %f °C\n", sdk.Data.OilTemp)
fmt.Fprintf(s, " DashLights: %b\n", sdk.Data.DashLights)
fmt.Fprintf(s, " ShowLights: %b\n", sdk.Data.ShowLights)
fmt.Fprintf(s, " Throttle: %f\n", sdk.Data.Throttle)
fmt.Fprintf(s, " Brakes: %f\n", sdk.Data.Brake)
fmt.Fprintf(s, " Clutch: %f\n", sdk.Data.Clutch)
fmt.Fprintf(s, " Display1: %s\n", sdk.Data.Display1)
fmt.Fprintf(s, " Display2: %s\n", sdk.Data.Display2)
fmt.Fprintf(s, " ID: %d\n", sdk.Data.ID)
fmt.Fprintf(s, " Time: %d ms\n", og.Time)
fmt.Fprintf(s, " Car: %s\n", og.Car)
fmt.Fprintf(s, " Flags: %b\n", og.Flags)
fmt.Fprintf(s, " Gear: %d\n", og.Gear)
fmt.Fprintf(s, " Plid: %d\n", og.Plid)
fmt.Fprintf(s, " Speed: %f m/s\n", og.Speed)
fmt.Fprintf(s, " RPM: %f RPM\n", og.RPM)
fmt.Fprintf(s, " Turbo: %f Bar\n", og.Turbo)
fmt.Fprintf(s, " EngTemp: %f °C\n", og.EngTemp)
fmt.Fprintf(s, " Fuel: %f\n", og.Fuel)
fmt.Fprintf(s, " OilPressure: %f Bar\n", og.OilPressure)
fmt.Fprintf(s, " OilTemp: %f °C\n", og.OilTemp)
fmt.Fprintf(s, " DashLights: %b\n", og.DashLights)
fmt.Fprintf(s, " ShowLights: %b\n", og.ShowLights)
fmt.Fprintf(s, " Throttle: %f\n", og.Throttle)
fmt.Fprintf(s, " Brakes: %f\n", og.Brake)
fmt.Fprintf(s, " Clutch: %f\n", og.Clutch)
fmt.Fprintf(s, " Display1: %s\n", og.Display1)
fmt.Fprintf(s, " Display2: %s\n", og.Display2)
fmt.Fprintf(s, " ID: %d\n", og.ID)
fmt.Fprint(s, "}\n\n")
fmt.Fprint(s, "DashLights {\n")
fmt.Fprintf(s, " DL_SHIFT: %t\n", sdk.HasShiftLight())
fmt.Fprintf(s, " DL_FULLBEAM: %t\n", sdk.HasHighBeamLight())
fmt.Fprintf(s, " DL_HANDBRAKE: %t\n", sdk.HasHandbrakeLight())
fmt.Fprintf(s, " DL_PITSPEED: %t\n", sdk.HasPitspeed())
fmt.Fprintf(s, " DL_TC: %t\n", sdk.HasTractionControlLight())
fmt.Fprintf(s, " DL_SIGNAL_L: %t\n", sdk.HasLeftIndicatorLight())
fmt.Fprintf(s, " DL_SIGNAL_R: %t\n", sdk.HasRightIndicatorLight())
fmt.Fprintf(s, " DL_SIGNAL_ANY: %t\n", sdk.HasAnyIndicatorLight())
fmt.Fprintf(s, " DL_OILWARN: %t\n", sdk.HasOilLight())
fmt.Fprintf(s, " DL_BATTERY: %t\n", sdk.HasBatteryLight())
fmt.Fprintf(s, " DL_ABS: %t\n", sdk.HasABSLight())
fmt.Fprintf(s, " DL_SPARE: %t\n", sdk.Data.DashLights&bngsdk.DL_SPARE != 0)
fmt.Fprintf(s, " DL_SHIFT: %t\n", og.HasShiftLight())
fmt.Fprintf(s, " DL_FULLBEAM: %t\n", og.HasHighBeamLight())
fmt.Fprintf(s, " DL_HANDBRAKE: %t\n", og.HasHandbrakeLight())
fmt.Fprintf(s, " DL_PITSPEED: %t\n", og.HasPitspeed())
fmt.Fprintf(s, " DL_TC: %t\n", og.HasTractionControlLight())
fmt.Fprintf(s, " DL_SIGNAL_L: %t\n", og.HasLeftIndicatorLight())
fmt.Fprintf(s, " DL_SIGNAL_R: %t\n", og.HasRightIndicatorLight())
fmt.Fprintf(s, " DL_SIGNAL_ANY: %t\n", og.HasAnyIndicatorLight())
fmt.Fprintf(s, " DL_OILWARN: %t\n", og.HasOilLight())
fmt.Fprintf(s, " DL_BATTERY: %t\n", og.HasBatteryLight())
fmt.Fprintf(s, " DL_ABS: %t\n", og.HasABSLight())
fmt.Fprintf(s, " DL_SPARE: %t\n", og.HasSpare())
fmt.Fprint(s, "}\n\n")
fmt.Fprint(s, "ShowLights {\n") // Fixed typo "ShowLigths"
fmt.Fprintf(s, " DL_SHIFT: %t\n", sdk.ShiftLight())
fmt.Fprintf(s, " DL_FULLBEAM: %t\n", sdk.HighBeam())
fmt.Fprintf(s, " DL_HANDBRAKE: %t\n", sdk.Handbrake())
fmt.Fprintf(s, " DL_PITSPEED: %t\n", sdk.Pitspeed())
fmt.Fprintf(s, " DL_TC: %t\n", sdk.TractionControl())
fmt.Fprintf(s, " DL_SIGNAL_L: %t\n", sdk.LeftIndicator())
fmt.Fprintf(s, " DL_SIGNAL_R: %t\n", sdk.RightIndicator())
fmt.Fprintf(s, " DL_SIGNAL_ANY: %t\n", sdk.AnyIndicator())
fmt.Fprintf(s, " DL_OILWARN: %t\n", sdk.OilLight())
fmt.Fprintf(s, " DL_BATTERY: %t\n", sdk.BatteryLight())
fmt.Fprintf(s, " DL_ABS: %t\n", sdk.ABS())
fmt.Fprintf(s, " DL_SPARE: %t\n", sdk.Data.ShowLights&bngsdk.DL_SPARE != 0)
fmt.Fprintf(s, " DL_SHIFT: %t\n", og.ShiftLight())
fmt.Fprintf(s, " DL_FULLBEAM: %t\n", og.HighBeam())
fmt.Fprintf(s, " DL_HANDBRAKE: %t\n", og.Handbrake())
fmt.Fprintf(s, " DL_PITSPEED: %t\n", og.Pitspeed())
fmt.Fprintf(s, " DL_TC: %t\n", og.TractionControl())
fmt.Fprintf(s, " DL_SIGNAL_L: %t\n", og.LeftIndicator())
fmt.Fprintf(s, " DL_SIGNAL_R: %t\n", og.RightIndicator())
fmt.Fprintf(s, " DL_SIGNAL_ANY: %t\n", og.AnyIndicator())
fmt.Fprintf(s, " DL_OILWARN: %t\n", og.OilLight())
fmt.Fprintf(s, " DL_BATTERY: %t\n", og.BatteryLight())
fmt.Fprintf(s, " DL_ABS: %t\n", og.ABS())
fmt.Fprintf(s, " DL_SPARE: %t\n", og.Spare())
fmt.Fprint(s, "}\n\n")
fmt.Fprint(s, "Flags {\n")
fmt.Fprintf(s, " OG_TURBO (Has Turbo): %t\n", sdk.HasTurbo())
fmt.Fprintf(s, " OG_KM (Is Metric): %t\n", sdk.PrefersKm())
fmt.Fprintf(s, " OG_BAR (Pressure): %t\n", sdk.PrefersBAR())
fmt.Fprintf(s, " OG_TURBO (Has Turbo): %t\n", og.HasTurbo())
fmt.Fprintf(s, " OG_KM (Is Metric): %t\n", og.PrefersKm())
fmt.Fprintf(s, " OG_BAR (Pressure): %t\n", og.PrefersBAR())
fmt.Fprint(s, "}")
}