1 Commits
5 changed files with 122 additions and 39 deletions
+5 -1
View File
@@ -1,6 +1,7 @@
package cmd
import (
"context"
"fmt"
"github.com/ESilva15/BeamNGMockOg/mockserver"
@@ -12,8 +13,10 @@ 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 {
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)
}
}
@@ -31,4 +34,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
)
+55
View File
@@ -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
View File
@@ -3,24 +3,27 @@
package mockserver
import (
"bytes"
"encoding/binary"
"encoding/gob"
"log"
"context"
"fmt"
"io"
"os"
"time"
sdk "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)
func Replay(ctx context.Context, address string, port int, loop bool, fp string) error {
fileInfo, err := os.Stat(fp)
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 {
return err
}
@@ -28,27 +31,36 @@ func Replay(address string, port int, fp string) error {
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:
data, err := reader.Next()
if err == io.EOF {
if !loop {
return udp.Close()
}
if _, err := udp.Conn.WriteToUDP(buf.Bytes(), udp.Addr); err != nil {
log.Fatal("send failed:", err)
break
}
err = reader.Reset()
if err != nil {
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
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()
}