diff --git a/README.md b/README.md index 03837a3..b1de2b9 100644 --- a/README.md +++ b/README.md @@ -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 diff --git a/cmd/replay.go b/cmd/replay.go index 6daf248..ce297bf 100644 --- a/cmd/replay.go +++ b/cmd/replay.go @@ -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) } } diff --git a/mockserver/reader.go b/mockserver/reader.go index 16f83ae..36bd87c 100644 --- a/mockserver/reader.go +++ b/mockserver/reader.go @@ -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 } diff --git a/mockserver/replay.go b/mockserver/replay.go index 4d9d2fe..d280b91 100644 --- a/mockserver/replay.go +++ b/mockserver/replay.go @@ -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 in a UDP server : -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! diff --git a/mockserver/replay_test.go b/mockserver/replay_test.go index 4cfee00..be8d256 100644 --- a/mockserver/replay_test.go +++ b/mockserver/replay_test.go @@ -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()