diff --git a/.gitignore b/.gitignore index 66391f2..e34276b 100644 --- a/.gitignore +++ b/.gitignore @@ -1,2 +1,4 @@ *.bin *.work* +*.pprof +*.log diff --git a/cmd/beamng.go b/cmd/beamng.go index 8307e7a..488029b 100644 --- a/cmd/beamng.go +++ b/cmd/beamng.go @@ -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 } } diff --git a/constants/constants.go b/constants/constants.go new file mode 100644 index 0000000..d4e7a87 --- /dev/null +++ b/constants/constants.go @@ -0,0 +1,3 @@ +package constants + +const ProgramName = "TelemetryMockserver" diff --git a/go.mod b/go.mod index bb45902..0c58ce3 100644 --- a/go.mod +++ b/go.mod @@ -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 ) diff --git a/internal/mockservers/beamng/reader.go b/internal/mockservers/beamng/reader.go deleted file mode 100644 index 36bd87c..0000000 --- a/internal/mockservers/beamng/reader.go +++ /dev/null @@ -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 -} diff --git a/internal/mockservers/beamng/record.go b/internal/mockservers/beamng/record.go index 6a11622..dcf8b03 100644 --- a/internal/mockservers/beamng/record.go +++ b/internal/mockservers/beamng/record.go @@ -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! - } } } } diff --git a/internal/mockservers/beamng/replay.go b/internal/mockservers/beamng/replay.go index 4430b4d..f47264c 100644 --- a/internal/mockservers/beamng/replay.go +++ b/internal/mockservers/beamng/replay.go @@ -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 in a UDP server : 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! - } - } } } diff --git a/internal/mockservers/beamng/server.go b/internal/mockservers/beamng/server.go deleted file mode 100644 index 771a14f..0000000 --- a/internal/mockservers/beamng/server.go +++ /dev/null @@ -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() -} diff --git a/internal/mockservers/beamng/views.go b/internal/mockservers/beamng/views.go index d508b10..7b5cf9e 100644 --- a/internal/mockservers/beamng/views.go +++ b/internal/mockservers/beamng/views.go @@ -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, "}") }