6 Commits
8 changed files with 328 additions and 46 deletions
+35 -1
View File
@@ -12,7 +12,41 @@ without having to be playing the game while doing it (fans are noisy).
# Usage
To start replaying:
`BeamNGMockOg replay -a 127.0.0.1 -p 4443 -i sunburstManual.bin`
`BeamNGMockOg replay [--loop] -a 127.0.0.1 -p 4443 -i sunburstManual.bin`
- `loop` allows the replay functionality to keep replaying the same data file
To start recording:
`BeamNGMockOg record -a 127.0.0.1 -p 4443 -o sunburstDCT.bin`
## Development
Use `tcpdump` to listen to the socket and check if data is coming through:
`tcpdump -i any udp port <port> -X`
### Benchmark
Run with:
`go test ./mockserver -bench=BenchmarkReplayAsync -benchmem -benchtime=1x -memprofile=mem.pprof`
Analyze the output with:
`go tool pprof -sample_index=alloc_objects mem.pprof` -> `top`
- `ignore=net` will ignore the net package for example
- `focus=mockserver` will show only the results from this code we are testing
or
`go tool pprof -http=:8080 mem.pprof`
#### Previous results:
goos: linux
goarch: amd64
pkg: github.com/ESilva15/BeamNGMockOg/mockserver
cpu: AMD Ryzen 5 5600G with Radeon Graphics
BenchmarkReplayAsync-12 1 109433574741 ns/op 4533008 B/op 66059 allocs/op
PASS
ok github.com/ESilva15/BeamNGMockOg/mockserver 109.438s
#### Current results
goos: linux
goarch: amd64
pkg: github.com/ESilva15/BeamNGMockOg/mockserver
cpu: AMD Ryzen 5 5600G with Radeon Graphics
BenchmarkReplayAsync-12 1 109433931841 ns/op 356088 B/op 13272 allocs/op
PASS
ok github.com/ESilva15/BeamNGMockOg/mockserver 109.440s
+11 -1
View File
@@ -1,6 +1,7 @@
package cmd
import (
"context"
"fmt"
"github.com/ESilva15/BeamNGMockOg/mockserver"
@@ -12,8 +13,16 @@ func replayAction(cmd *cobra.Command, args []string) {
inputFile, _ := cmd.Flags().GetString("input")
address, _ := cmd.Flags().GetString("address")
port, _ := cmd.Flags().GetInt("port")
loop, _ := cmd.Flags().GetBool("loop")
if err := mockserver.Replay(address, port, inputFile); err != nil {
replayer, err := mockserver.NewReplayer(address, port, inputFile)
if err != nil {
fmt.Printf("Something went wrong setting up the player: %+v", err)
return
}
ctx := context.Background()
if err := replayer.Replay(ctx, loop); err != nil {
fmt.Printf("Something went wrong while playing the file: %v", err)
}
}
@@ -31,4 +40,5 @@ func init() {
rootCmd.AddCommand(replayCmd)
replayCmd.PersistentFlags().StringP("input", "i", "nofile.bin", "input file for serving")
replayCmd.PersistentFlags().BoolP("loop", "l", false, "loop replay after finishing")
}
+3 -1
View File
@@ -1,9 +1,11 @@
module github.com/ESilva15/BeamNGMockOg
replace github.com/ESilva15/gobngsdk => ../pkg/bngsdk
go 1.23.2
require (
github.com/ESilva15/gobngsdk v0.0.2
github.com/ESilva15/gobngsdk v0.0.0-00010101000000-000000000000
github.com/spf13/cobra v1.10.1
)
+45
View File
@@ -0,0 +1,45 @@
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
}
+9 -3
View File
@@ -1,7 +1,7 @@
package mockserver
import (
"encoding/gob"
"encoding/binary"
"log"
"os"
"time"
@@ -9,6 +9,9 @@ import (
bngsdk "github.com/ESilva15/gobngsdk"
)
// NOTE: add some visual feedback of whats happening.
// Maybe reuse the replay view function
// Record records data from the UDP connection created by address and port
func Record(address string, port int, filePath string) error {
ticker := time.NewTicker(time.Second / 60)
@@ -19,6 +22,7 @@ func Record(address string, port int, filePath string) error {
if err != nil {
return err
}
defer bin.Close()
// Create the BeamNGSDK instance
beam, err := bngsdk.Init(address, port)
@@ -27,15 +31,17 @@ func Record(address string, port int, filePath string) error {
}
defer beam.Close()
enc := gob.NewEncoder(bin)
for {
err := beam.ReadData()
if err != nil {
return err
}
if err := enc.Encode(beam.Data); err != nil {
err = binary.Write(bin, binary.LittleEndian, beam.Data)
if err != nil {
log.Fatal(err)
}
<-ticker.C
}
}
+142 -31
View File
@@ -3,52 +3,163 @@
package mockserver
import (
"bytes"
"encoding/binary"
"encoding/gob"
"log"
"context"
"fmt"
"io"
"os"
"sync"
"time"
"unsafe"
sdk "github.com/ESilva15/gobngsdk"
bngsdk "github.com/ESilva15/gobngsdk"
)
// Replay replays a given file <fp> in a UDP server <addr>:<port>
func Replay(address string, port int, fp string) error {
bin, err := os.Open(fp)
type ViewData struct {
Data []byte
SizeRead int64
}
// Replayer does the replaying
// Should we make a "player" struct that can record and replay?
type Replayer struct {
DataSourcePath string
Socket *UDPTransport
// Streams
dataViewCh chan ViewData
socketCh chan []byte
// Mut
mut sync.RWMutex
data []byte
// View
viewData ViewData
}
func NewReplayer(address string, port int, fp string) (*Replayer, error) {
udp, err := NewUDPTransport(address, port)
if err != nil {
log.Fatal("error opening file:", err)
return nil, err
}
udp, err := newUDPServer(address, port)
if err != nil {
return err
replayer := &Replayer{
DataSourcePath: fp,
Socket: udp,
data: make([]byte, unsafe.Sizeof(bngsdk.Outgauge{})),
viewData: ViewData{},
dataViewCh: make(chan ViewData, 1),
socketCh: make(chan []byte, 1),
}
return replayer, nil
}
// 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)
}
// NOTE: temporary until I make a better view
lastPercent := -1
for {
select {
case <-ctx.Done():
return
case data := <-r.dataViewCh:
percent := int(float64(data.SizeRead) / float64(fileInfo.Size()) * 100)
if percent != lastPercent {
fmt.Printf("\rReplayed: %d%%", percent)
lastPercent = percent
}
}
}
}
// 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()
dec := gob.NewDecoder(bin)
for {
var og sdk.Outgauge
if err := dec.Decode(&og); err != nil {
log.Fatal("failed to read more data:", err)
break
}
select {
case <-ctx.Done():
return ctx.Err()
buf := new(bytes.Buffer)
if err := binary.Write(buf, binary.LittleEndian, &og); err != nil {
log.Fatal("serialization failed:", err)
break
}
case <-ticker.C:
r.mut.Lock()
err := reader.Next(r.data)
r.mut.Unlock()
if _, err := udp.Conn.WriteToUDP(buf.Bytes(), udp.Addr); err != nil {
log.Fatal("send failed:", err)
break
}
if err == io.EOF {
if !loop {
return r.Socket.Close()
}
<-ticker.C
err = reader.Reset()
if err != nil {
return err
}
continue
}
if err != nil {
return err
}
r.viewData.Data = r.data
r.viewData.SizeRead = reader.TotalRead
// Send the data to the view
select {
case r.dataViewCh <- r.viewData:
// Sent the data
default:
// Dropped the frame!
}
// Send the data to the UDP socket
select {
case r.socketCh <- r.data:
// Sent the data
default:
// Dropped the frame!
}
}
}
return udp.Close()
}
+64
View File
@@ -0,0 +1,64 @@
package mockserver
import (
"context"
"net"
"os"
"testing"
)
func BenchmarkReplayAsync(b *testing.B) {
// Dummy socket
addr, err := net.ResolveUDPAddr("udp", "127.0.0.1:0")
if err != nil {
b.Fatalf("failed to resolve UDP address: %v", err)
}
listener, err := net.ListenUDP("udp", addr)
if err != nil {
b.Fatalf("failed to start background UDP listener: %v", err)
}
defer listener.Close()
// Drain the socket
go func() {
buf := make([]byte, 65535)
for {
_, _, err := listener.ReadFrom(buf)
if err != nil {
return
}
}
}()
// Extract the random port assigned by the OS
assignedAddr := listener.LocalAddr().(*net.UDPAddr)
// Get the file path
filePath := "sunburstManual.bin"
if _, err := os.Stat(filePath); os.IsNotExist(err) {
filePath = "../" + filePath
}
// Reset the timer
b.ResetTimer()
// Benchmark loop
for i := 0; i < b.N; i++ {
ctx, cancel := context.WithCancel(context.Background())
replayer, err := NewReplayer(assignedAddr.IP.String(), assignedAddr.Port, filePath)
if err != nil {
cancel()
b.Fatalf("Error setting up replayer: %+v", err)
}
err = replayer.Replay(ctx, false)
if err != nil {
cancel()
b.Fatalf("Replay failed during benchmark run: %+v", err)
}
cancel()
}
}
+19 -9
View File
@@ -5,25 +5,35 @@ import (
"net"
)
type udpServer struct {
Addr *net.UDPAddr
Conn *net.UDPConn
type Transport interface {
Send(data []byte) error
Close() error
}
func newUDPServer(addr string, port int) (udpServer, error) {
udpAddr, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", addr, port))
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 udpServer{}, err
return nil, err
}
conn, err := net.ListenUDP("udp", nil)
if err != nil {
return udpServer{}, err
return nil, err
}
return udpServer{Addr: udpAddr, Conn: conn}, nil
return &UDPTransport{Conn: conn, Addr: addr}, nil
}
func (u *udpServer) Close() error {
// 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()
}