updated the recorder gobngsdk v2
This commit is contained in:
@@ -0,0 +1,3 @@
|
|||||||
|
package constants
|
||||||
|
|
||||||
|
const ProgramName = "TelemetryMockserver"
|
||||||
@@ -3,7 +3,7 @@ module github.com/ESilva15/TelemetryMockserver
|
|||||||
go 1.23.2
|
go 1.23.2
|
||||||
|
|
||||||
require (
|
require (
|
||||||
github.com/ESilva15/gobngsdk v1.1.3
|
github.com/ESilva15/gobngsdk v2.2.0
|
||||||
github.com/ESilva15/goirsdk v0.3.0
|
github.com/ESilva15/goirsdk v0.3.0
|
||||||
github.com/spf13/cobra v1.10.1
|
github.com/spf13/cobra v1.10.1
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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
|
|
||||||
}
|
|
||||||
@@ -3,83 +3,59 @@ package mockserver
|
|||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"log/slog"
|
||||||
"os"
|
"os"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/ESilva15/TelemetryMockserver/constants"
|
||||||
bngsdk "github.com/ESilva15/gobngsdk"
|
bngsdk "github.com/ESilva15/gobngsdk"
|
||||||
)
|
)
|
||||||
|
|
||||||
type recorderViewData struct {
|
type recorderViewData struct {
|
||||||
TotalBytes int
|
TotalBytes int64
|
||||||
SDK *bngsdk.BeamNGSDK
|
Og bngsdk.Outgauge
|
||||||
}
|
}
|
||||||
|
|
||||||
type Recorder struct {
|
type Recorder struct {
|
||||||
SDK bngsdk.BeamNGSDK
|
SDK *bngsdk.BeamNGSDK
|
||||||
OutputFile *os.File
|
|
||||||
TotalBytes int
|
|
||||||
// Views
|
// Views
|
||||||
mut sync.RWMutex
|
mut sync.RWMutex
|
||||||
viewDataMut sync.RWMutex
|
viewDataMut sync.RWMutex
|
||||||
viewData recorderViewData
|
viewData recorderViewData
|
||||||
viewCh chan *recorderViewData
|
viewCh chan *recorderViewData
|
||||||
recorderCh chan []byte
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewRecorder(fp string, address string, port int) (*Recorder, error) {
|
func NewRecorder(fp string, address string, port int) (*Recorder, error) {
|
||||||
// var recorder Recorder
|
var recorder Recorder
|
||||||
// var err error
|
var err error
|
||||||
//
|
|
||||||
// recorder.SDK, err = bngsdk.Init(address, port)
|
|
||||||
// if err != nil {
|
|
||||||
// return &Recorder{}, err
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// recorder.OutputFile, err = os.Create(fp)
|
|
||||||
// if err != nil {
|
|
||||||
// return &Recorder{}, err
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// recorder.viewData = recorderViewData{}
|
|
||||||
// recorder.viewCh = make(chan *recorderViewData, 1)
|
|
||||||
// recorder.recorderCh = make(chan []byte, 1)
|
|
||||||
|
|
||||||
// return &recorder, nil
|
recorder.SDK, err = bngsdk.NewBngSDK(bngsdk.Options{
|
||||||
return nil, errors.New("not implemented")
|
Logger: slog.Default().With("SDK", "BeamNG"),
|
||||||
|
SourceType: bngsdk.UDPData,
|
||||||
|
ImportUDPAddress: address,
|
||||||
|
ImportUDPPort: port,
|
||||||
|
ExportData: true,
|
||||||
|
ExportDataType: bngsdk.BinaryFile,
|
||||||
|
ExportDataPath: fp,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
recorder.viewData = recorderViewData{}
|
||||||
|
recorder.viewCh = make(chan *recorderViewData, 1)
|
||||||
|
|
||||||
|
return &recorder, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (r *Recorder) Close() {
|
func (r *Recorder) Close() {
|
||||||
r.SDK.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) {
|
func (r *Recorder) view(ctx context.Context) {
|
||||||
var buf bytes.Buffer
|
var buf bytes.Buffer
|
||||||
var nBytes int
|
|
||||||
|
|
||||||
buf.Grow(2048)
|
buf.Grow(2048)
|
||||||
|
|
||||||
@@ -90,16 +66,15 @@ func (r *Recorder) view(ctx context.Context) {
|
|||||||
case viewData := <-r.viewCh:
|
case viewData := <-r.viewCh:
|
||||||
buf.Reset()
|
buf.Reset()
|
||||||
buf.WriteString("\x1b[2J\x1b[H")
|
buf.WriteString("\x1b[2J\x1b[H")
|
||||||
fmt.Fprintf(&buf, "\x1b]0;%s - Recording ", ProgramName)
|
fmt.Fprintf(&buf, "\x1b]0;%s - Recording ", constants.ProgramName)
|
||||||
stringifyRecordingProgress(&buf, nBytes)
|
stringifyRecordingProgress(&buf, viewData.TotalBytes)
|
||||||
fmt.Fprintf(&buf, "\x07")
|
fmt.Fprintf(&buf, "\x07")
|
||||||
|
|
||||||
stringifyRecordingProgress(&buf, nBytes)
|
stringifyRecordingProgress(&buf, viewData.TotalBytes)
|
||||||
fmt.Fprintf(&buf, "\n\n")
|
fmt.Fprintf(&buf, "\n\n")
|
||||||
|
|
||||||
r.viewDataMut.RLock()
|
r.viewDataMut.RLock()
|
||||||
nBytes = viewData.TotalBytes
|
stringifyOutgaugeData(&buf, &viewData.Og)
|
||||||
// stringifyOutgaugeData(&buf, &r.SDK)
|
|
||||||
r.viewDataMut.RUnlock()
|
r.viewDataMut.RUnlock()
|
||||||
|
|
||||||
_, _ = buf.WriteTo(os.Stdout)
|
_, _ = buf.WriteTo(os.Stdout)
|
||||||
@@ -112,7 +87,6 @@ func (r *Recorder) Record(ctx context.Context) error {
|
|||||||
ticker := time.NewTicker(time.Second / 60)
|
ticker := time.NewTicker(time.Second / 60)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
|
|
||||||
go r.record(ctx)
|
|
||||||
go r.view(ctx)
|
go r.view(ctx)
|
||||||
|
|
||||||
for {
|
for {
|
||||||
@@ -120,28 +94,21 @@ func (r *Recorder) Record(ctx context.Context) error {
|
|||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
return nil
|
return nil
|
||||||
case <-ticker.C:
|
case <-ticker.C:
|
||||||
// err := r.SDK.ReadData(100 * time.Millisecond)
|
og, err := r.SDK.Update()
|
||||||
// if err != nil {
|
if err != nil {
|
||||||
// return err
|
return err
|
||||||
// }
|
}
|
||||||
//
|
|
||||||
// r.viewData.TotalBytes = r.TotalBytes
|
r.viewData.TotalBytes = r.SDK.GetTotalWritten()
|
||||||
//
|
r.viewData.Og = *og
|
||||||
// // Send the data to the view
|
|
||||||
// select {
|
// Send the data to the view
|
||||||
// case r.viewCh <- &r.viewData:
|
select {
|
||||||
// // Sent the data
|
case r.viewCh <- &r.viewData:
|
||||||
// default:
|
// Sent the data
|
||||||
// // Dropped the frame!
|
default:
|
||||||
// }
|
// Dropped the frame!
|
||||||
//
|
}
|
||||||
// // Write the data to the file
|
|
||||||
// select {
|
|
||||||
// case r.recorderCh <- r.SDK.GetBuffer():
|
|
||||||
// // Sent the data
|
|
||||||
// default:
|
|
||||||
// // Dropped the frame!
|
|
||||||
// }
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ import (
|
|||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/ESilva15/TelemetryMockserver/constants"
|
||||||
bngsdk "github.com/ESilva15/gobngsdk"
|
bngsdk "github.com/ESilva15/gobngsdk"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -41,6 +42,8 @@ func NewReplayer(address string, port int, fp string) (*Replayer, error) {
|
|||||||
Logger: slog.Default().With("SDK", "BeamNG"),
|
Logger: slog.Default().With("SDK", "BeamNG"),
|
||||||
SourceType: bngsdk.BinaryFile,
|
SourceType: bngsdk.BinaryFile,
|
||||||
BinSourcePath: fp,
|
BinSourcePath: fp,
|
||||||
|
ExportData: true,
|
||||||
|
ExportDataType: bngsdk.UDPData,
|
||||||
ExportUDPAddress: address,
|
ExportUDPAddress: address,
|
||||||
ExportUDPPort: port,
|
ExportUDPPort: port,
|
||||||
Loop: true,
|
Loop: true,
|
||||||
@@ -73,7 +76,7 @@ func (r *Replayer) renderToTerminal(ctx context.Context) {
|
|||||||
|
|
||||||
buf.Reset()
|
buf.Reset()
|
||||||
buf.WriteString("\x1b[2J\x1b[H")
|
buf.WriteString("\x1b[2J\x1b[H")
|
||||||
fmt.Fprintf(&buf, "\x1b]0;%s - Replaying %d%%\x07", ProgramName, percent)
|
fmt.Fprintf(&buf, "\x1b]0;%s - Replaying %d%%\x07", constants.ProgramName, percent)
|
||||||
|
|
||||||
fmt.Fprintf(&buf, "Replayed: %d%%\n", percent)
|
fmt.Fprintf(&buf, "Replayed: %d%%\n", percent)
|
||||||
|
|
||||||
@@ -88,7 +91,7 @@ func (r *Replayer) renderToTerminal(ctx context.Context) {
|
|||||||
func (r *Replayer) Replay(ctx context.Context, loop bool) error {
|
func (r *Replayer) Replay(ctx context.Context, loop bool) error {
|
||||||
go r.renderToTerminal(ctx)
|
go r.renderToTerminal(ctx)
|
||||||
|
|
||||||
ticker := time.NewTicker(time.Second / 240)
|
ticker := time.NewTicker(time.Second / 60)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
|
|||||||
@@ -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()
|
|
||||||
}
|
|
||||||
@@ -14,7 +14,7 @@ const (
|
|||||||
GiB
|
GiB
|
||||||
)
|
)
|
||||||
|
|
||||||
func stringifyRecordingProgress(s *bytes.Buffer, nBytes int) {
|
func stringifyRecordingProgress(s *bytes.Buffer, nBytes int64) {
|
||||||
if nBytes < KiB {
|
if nBytes < KiB {
|
||||||
fmt.Fprintf(s, "%d B", nBytes)
|
fmt.Fprintf(s, "%d B", nBytes)
|
||||||
} else if nBytes < MiB {
|
} else if nBytes < MiB {
|
||||||
|
|||||||
Reference in New Issue
Block a user