Shows progress and also added a loop flag so that it can go around and around
This commit is contained in:
+5
-1
@@ -1,6 +1,7 @@
|
|||||||
package cmd
|
package cmd
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
|
||||||
"github.com/ESilva15/BeamNGMockOg/mockserver"
|
"github.com/ESilva15/BeamNGMockOg/mockserver"
|
||||||
@@ -12,8 +13,10 @@ func replayAction(cmd *cobra.Command, args []string) {
|
|||||||
inputFile, _ := cmd.Flags().GetString("input")
|
inputFile, _ := cmd.Flags().GetString("input")
|
||||||
address, _ := cmd.Flags().GetString("address")
|
address, _ := cmd.Flags().GetString("address")
|
||||||
port, _ := cmd.Flags().GetInt("port")
|
port, _ := cmd.Flags().GetInt("port")
|
||||||
|
loop, _ := cmd.Flags().GetBool("loop")
|
||||||
|
|
||||||
if err := mockserver.Replay(address, port, inputFile); err != nil {
|
ctx := context.Background()
|
||||||
|
if err := mockserver.Replay(ctx, address, port, loop, inputFile); err != nil {
|
||||||
fmt.Printf("Something went wrong while playing the file: %v", err)
|
fmt.Printf("Something went wrong while playing the file: %v", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -31,4 +34,5 @@ func init() {
|
|||||||
rootCmd.AddCommand(replayCmd)
|
rootCmd.AddCommand(replayCmd)
|
||||||
|
|
||||||
replayCmd.PersistentFlags().StringP("input", "i", "nofile.bin", "input file for serving")
|
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
|
module github.com/ESilva15/BeamNGMockOg
|
||||||
|
|
||||||
|
replace github.com/ESilva15/gobngsdk => ../pkg/bngsdk
|
||||||
|
|
||||||
go 1.23.2
|
go 1.23.2
|
||||||
|
|
||||||
require (
|
require (
|
||||||
github.com/ESilva15/gobngsdk v0.0.2
|
github.com/ESilva15/gobngsdk v0.0.0-00010101000000-000000000000
|
||||||
github.com/spf13/cobra v1.10.1
|
github.com/spf13/cobra v1.10.1
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,55 @@
|
|||||||
|
package mockserver
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"encoding/binary"
|
||||||
|
"encoding/gob"
|
||||||
|
"io"
|
||||||
|
|
||||||
|
sdk "github.com/ESilva15/gobngsdk"
|
||||||
|
)
|
||||||
|
|
||||||
|
type GobReader struct {
|
||||||
|
TotalRead int
|
||||||
|
File io.ReadSeeker
|
||||||
|
Dec *gob.Decoder
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewGobReader(r io.ReadSeeker) *GobReader {
|
||||||
|
return &GobReader{
|
||||||
|
TotalRead: 0,
|
||||||
|
File: r,
|
||||||
|
Dec: gob.NewDecoder(r),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (g *GobReader) Reset() error {
|
||||||
|
_, err := g.File.Seek(0, io.SeekStart)
|
||||||
|
if err != nil {
|
||||||
|
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)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
buf := new(bytes.Buffer)
|
||||||
|
err = binary.Write(buf, binary.LittleEndian, &og)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
g.TotalRead += len(buf.Bytes())
|
||||||
|
|
||||||
|
return buf.Bytes(), nil
|
||||||
|
}
|
||||||
+40
-28
@@ -3,24 +3,27 @@
|
|||||||
package mockserver
|
package mockserver
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"context"
|
||||||
"encoding/binary"
|
"fmt"
|
||||||
"encoding/gob"
|
"io"
|
||||||
"log"
|
|
||||||
"os"
|
"os"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
sdk "github.com/ESilva15/gobngsdk"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
// Replay replays a given file <fp> in a UDP server <addr>:<port>
|
// Replay replays a given file <fp> in a UDP server <addr>:<port>
|
||||||
func Replay(address string, port int, fp string) error {
|
func Replay(ctx context.Context, address string, port int, loop bool, fp string) error {
|
||||||
bin, err := os.Open(fp)
|
fileInfo, err := os.Stat(fp)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Fatal("error opening file:", err)
|
return fmt.Errorf("error stating file: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
udp, err := newUDPServer(address, port)
|
bin, err := os.Open(fp)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("error opening file: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
reader := NewGobReader(bin)
|
||||||
|
udp, err := NewUDPTransport(address, port)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -28,27 +31,36 @@ func Replay(address string, port int, fp string) error {
|
|||||||
ticker := time.NewTicker(time.Second / 60)
|
ticker := time.NewTicker(time.Second / 60)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
|
|
||||||
dec := gob.NewDecoder(bin)
|
|
||||||
for {
|
for {
|
||||||
var og sdk.Outgauge
|
select {
|
||||||
if err := dec.Decode(&og); err != nil {
|
case <-ctx.Done():
|
||||||
log.Fatal("failed to read more data:", err)
|
return ctx.Err()
|
||||||
break
|
|
||||||
}
|
|
||||||
|
|
||||||
buf := new(bytes.Buffer)
|
case <-ticker.C:
|
||||||
if err := binary.Write(buf, binary.LittleEndian, &og); err != nil {
|
data, err := reader.Next()
|
||||||
log.Fatal("serialization failed:", err)
|
if err == io.EOF {
|
||||||
break
|
if !loop {
|
||||||
}
|
return udp.Close()
|
||||||
|
}
|
||||||
|
|
||||||
if _, err := udp.Conn.WriteToUDP(buf.Bytes(), udp.Addr); err != nil {
|
err = reader.Reset()
|
||||||
log.Fatal("send failed:", err)
|
if err != nil {
|
||||||
break
|
return err
|
||||||
}
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
<-ticker.C
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
_, err = udp.Send(data)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
percent := int(float64(reader.TotalRead) / float64(fileInfo.Size()) * 100)
|
||||||
|
fmt.Printf("\rReplayed: %d%%", percent)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
return udp.Close()
|
|
||||||
}
|
}
|
||||||
|
|||||||
+19
-9
@@ -5,25 +5,35 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
)
|
)
|
||||||
|
|
||||||
type udpServer struct {
|
type Transport interface {
|
||||||
Addr *net.UDPAddr
|
Send(data []byte) error
|
||||||
Conn *net.UDPConn
|
Close() error
|
||||||
}
|
}
|
||||||
|
|
||||||
func newUDPServer(addr string, port int) (udpServer, error) {
|
type UDPTransport struct {
|
||||||
udpAddr, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", addr, port))
|
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 {
|
if err != nil {
|
||||||
return udpServer{}, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
conn, err := net.ListenUDP("udp", nil)
|
conn, err := net.ListenUDP("udp", nil)
|
||||||
if err != 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()
|
return u.Conn.Close()
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user