Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7c558eda9c | ||
|
|
c23c39d889 | ||
|
|
7939f07f34 | ||
|
|
4d05e636b6 | ||
|
|
540cb90f42 | ||
|
|
9d3e9942d7 |
@@ -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
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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
|
||||
)
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
@@ -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()
|
||||
}
|
||||
|
||||
@@ -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
@@ -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()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user