Performance and API improvements for the replay function
This commit is contained in:
@@ -24,4 +24,29 @@ Use `tcpdump` to listen to the socket and check if data is coming through:
|
||||
|
||||
### Benchmark
|
||||
Run with:
|
||||
`go test -bench=. -benchmem`
|
||||
`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
|
||||
|
||||
+7
-1
@@ -15,8 +15,14 @@ func replayAction(cmd *cobra.Command, args []string) {
|
||||
port, _ := cmd.Flags().GetInt("port")
|
||||
loop, _ := cmd.Flags().GetBool("loop")
|
||||
|
||||
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 := mockserver.Replay(ctx, address, port, loop, inputFile); err != nil {
|
||||
if err := replayer.Replay(ctx, loop); err != nil {
|
||||
fmt.Printf("Something went wrong while playing the file: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
+10
-20
@@ -1,25 +1,23 @@
|
||||
package mockserver
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/binary"
|
||||
"encoding/gob"
|
||||
"io"
|
||||
"unsafe"
|
||||
|
||||
sdk "github.com/ESilva15/gobngsdk"
|
||||
)
|
||||
|
||||
type GobReader struct {
|
||||
TotalRead int
|
||||
TotalRead int64
|
||||
File io.ReadSeeker
|
||||
Dec *gob.Decoder
|
||||
Buf []byte
|
||||
}
|
||||
|
||||
func NewGobReader(r io.ReadSeeker) *GobReader {
|
||||
return &GobReader{
|
||||
TotalRead: 0,
|
||||
File: r,
|
||||
Dec: gob.NewDecoder(r),
|
||||
Buf: make([]byte, unsafe.Sizeof(sdk.Outgauge{})),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -29,27 +27,19 @@ func (g *GobReader) Reset() error {
|
||||
return err
|
||||
}
|
||||
|
||||
g.Dec = gob.NewDecoder(g.File)
|
||||
g.TotalRead = 0
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (g *GobReader) Next() ([]byte, error) {
|
||||
var og sdk.Outgauge
|
||||
|
||||
err := g.Dec.Decode(&og)
|
||||
func (g *GobReader) Next(buffer []byte) error {
|
||||
_, err := io.ReadFull(g.File, buffer)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return err
|
||||
}
|
||||
|
||||
buf := new(bytes.Buffer)
|
||||
err = binary.Write(buf, binary.LittleEndian, &og)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
pos, _ := g.File.Seek(0, io.SeekCurrent)
|
||||
g.TotalRead = pos
|
||||
|
||||
g.TotalRead += len(buf.Bytes())
|
||||
|
||||
return buf.Bytes(), nil
|
||||
return nil
|
||||
}
|
||||
|
||||
+72
-25
@@ -7,42 +7,91 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
"unsafe"
|
||||
|
||||
bngsdk "github.com/ESilva15/gobngsdk"
|
||||
)
|
||||
|
||||
type ViewData struct {
|
||||
Data []byte
|
||||
SizeRead int
|
||||
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 {
|
||||
return nil, 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 renderToTerminal(ctx context.Context, stream chan ViewData, fp string) {
|
||||
fileInfo, err := os.Stat(fp)
|
||||
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 := <-stream:
|
||||
case data := <-r.dataViewCh:
|
||||
percent := int(float64(data.SizeRead) / float64(fileInfo.Size()) * 100)
|
||||
fmt.Printf("\rReplayed: %d%%", percent)
|
||||
if percent != lastPercent {
|
||||
fmt.Printf("\rReplayed: %d%%", percent)
|
||||
lastPercent = percent
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// writeToUDPSocket will write the telemetry data to the UDP socket
|
||||
func writeToUDPSocket(ctx context.Context, stream chan []byte, socket *UDPTransport) {
|
||||
func (r *Replayer) writeToUDPSocket(ctx context.Context) {
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case data := <-stream:
|
||||
_, err := socket.Send(data)
|
||||
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
|
||||
@@ -52,37 +101,32 @@ func writeToUDPSocket(ctx context.Context, stream chan []byte, socket *UDPTransp
|
||||
}
|
||||
|
||||
// Replay replays a given file <fp> in a UDP server <addr>:<port>
|
||||
func Replay(ctx context.Context, address string, port int, loop bool, fp string) error {
|
||||
bin, err := os.Open(fp)
|
||||
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)
|
||||
udp, err := NewUDPTransport(address, port)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
go r.renderToTerminal(ctx)
|
||||
go r.writeToUDPSocket(ctx)
|
||||
|
||||
ticker := time.NewTicker(time.Second / 60)
|
||||
defer ticker.Stop()
|
||||
|
||||
uiChannel := make(chan ViewData, 1)
|
||||
socketChannel := make(chan []byte, 1)
|
||||
|
||||
go renderToTerminal(ctx, uiChannel, fp)
|
||||
go writeToUDPSocket(ctx, socketChannel, udp)
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
|
||||
case <-ticker.C:
|
||||
data, err := reader.Next()
|
||||
r.mut.Lock()
|
||||
err := reader.Next(r.data)
|
||||
r.mut.Unlock()
|
||||
|
||||
if err == io.EOF {
|
||||
if !loop {
|
||||
return udp.Close()
|
||||
return r.Socket.Close()
|
||||
}
|
||||
|
||||
err = reader.Reset()
|
||||
@@ -97,9 +141,12 @@ func Replay(ctx context.Context, address string, port int, loop bool, fp string)
|
||||
return err
|
||||
}
|
||||
|
||||
r.viewData.Data = r.data
|
||||
r.viewData.SizeRead = reader.TotalRead
|
||||
|
||||
// Send the data to the view
|
||||
select {
|
||||
case uiChannel <- ViewData{Data: data, SizeRead: reader.TotalRead}:
|
||||
case r.dataViewCh <- r.viewData:
|
||||
// Sent the data
|
||||
default:
|
||||
// Dropped the frame!
|
||||
@@ -107,7 +154,7 @@ func Replay(ctx context.Context, address string, port int, loop bool, fp string)
|
||||
|
||||
// Send the data to the UDP socket
|
||||
select {
|
||||
case socketChannel <- data:
|
||||
case r.socketCh <- r.data:
|
||||
// Sent the data
|
||||
default:
|
||||
// Dropped the frame!
|
||||
|
||||
@@ -47,10 +47,16 @@ func BenchmarkReplayAsync(b *testing.B) {
|
||||
for i := 0; i < b.N; i++ {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
|
||||
err := Replay(ctx, assignedAddr.IP.String(), assignedAddr.Port, false, filePath)
|
||||
replayer, err := NewReplayer(assignedAddr.IP.String(), assignedAddr.Port, filePath)
|
||||
if err != nil {
|
||||
cancel()
|
||||
b.Fatalf("Replay failed during benchmark run: %v", err)
|
||||
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()
|
||||
|
||||
Reference in New Issue
Block a user