Compare commits
15
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
65f03fd4c9 | ||
|
|
20025dae1b | ||
|
|
03d0f8aefc | ||
|
|
df99994a75 | ||
|
|
80651c21cf | ||
|
|
528c62288a | ||
|
|
196e1d8916 | ||
|
|
decd91c406 | ||
|
|
b8e632d3eb | ||
|
|
d8589394f6 | ||
|
|
76f2e265d7 | ||
|
|
7c03088470 | ||
|
|
0006b649f3 | ||
|
|
3493a73a82 | ||
|
|
f4bde52f40 |
@@ -1 +0,0 @@
|
||||
# Streaming Flow
|
||||
@@ -1,17 +1,20 @@
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"esdi/logger"
|
||||
esdi "esdi/oldEsdi"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
)
|
||||
|
||||
func liveTelemetryCmdAction(cmd *cobra.Command, args []string) {
|
||||
log := logger.GetInstance()
|
||||
|
||||
ddPort, _ := cmd.Flags().GetString("port")
|
||||
outputFile, _ := cmd.Flags().GetString("out")
|
||||
sessionFile, _ := cmd.Flags().GetString("session")
|
||||
|
||||
// log.Printf("Called `live`:\nPort: '%s'\nOutFile: '%s'\n", ddPort, outputFile)
|
||||
log.Printf("Called `live`:\nPort: '%s'\nOutFile: '%s'\n", ddPort, outputFile)
|
||||
|
||||
esdi.RunLiveTelemetry(ddPort, outputFile, sessionFile)
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"esdi/logger"
|
||||
esdi "esdi/oldEsdi"
|
||||
|
||||
// "github.com/ESilva15/goirsdk"
|
||||
@@ -9,14 +10,14 @@ import (
|
||||
)
|
||||
|
||||
func offlineTelemetryCmdAction(cmd *cobra.Command, args []string) {
|
||||
// log := logger.GetInstance()
|
||||
log := logger.GetInstance()
|
||||
|
||||
ddPort, _ := cmd.Flags().GetString("port")
|
||||
inFile, _ := cmd.Flags().GetString("in")
|
||||
outFile, _ := cmd.Flags().GetString("out")
|
||||
sessionFile, _ := cmd.Flags().GetString("session")
|
||||
|
||||
// log.Printf("Called `offline`:\nPort: '%s'\nSource: '%s'\nOutFile: '%s'\n", ddPort, inFile, outFile)
|
||||
log.Printf("Called `offline`:\nPort: '%s'\nSource: '%s'\nOutFile: '%s'\n", ddPort, inFile, outFile)
|
||||
|
||||
esdi.RunOfflineTelemetry(ddPort, inFile, outFile, sessionFile)
|
||||
}
|
||||
|
||||
+93
-87
@@ -1,97 +1,103 @@
|
||||
package cmd
|
||||
|
||||
import (
|
||||
"esdi/peripheral"
|
||||
"fmt"
|
||||
"strconv"
|
||||
|
||||
repl "github.com/ESilva15/ESgoRepl"
|
||||
|
||||
"github.com/spf13/cobra"
|
||||
)
|
||||
|
||||
func replCmdAction(cmd *cobra.Command, args []string) {
|
||||
// r := repl.NewREPL(repl.REPLCfg{
|
||||
// PS1: "\rESDI > ",
|
||||
// })
|
||||
//
|
||||
// perClerk := peripheral.NewPeripheralDeviceClerk()
|
||||
//
|
||||
// discoverDevicesREPLCmd := repl.Command{
|
||||
// Name: "discover",
|
||||
// Usage: "discovers connected devices",
|
||||
// Action: func(r *repl.REPL, args []string) error {
|
||||
// err := perClerk.FindDevices()
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
//
|
||||
// return nil
|
||||
// },
|
||||
// }
|
||||
//
|
||||
// listDevicesREPLCmd := repl.Command{
|
||||
// Name: "list",
|
||||
// Usage: "lists connected devices",
|
||||
// Action: func(r *repl.REPL, args []string) error {
|
||||
// _ = perClerk.ListDevices()
|
||||
//
|
||||
// return nil
|
||||
// },
|
||||
// }
|
||||
//
|
||||
// listDeviceAPIREPLCmd := repl.Command{
|
||||
// Name: "v-api",
|
||||
// Usage: "shows API of a device - pass its ID",
|
||||
// Action: func(r *repl.REPL, args []string) error {
|
||||
// // We should add this to the REPL instead
|
||||
// if len(args) < 1 {
|
||||
// return fmt.Errorf("requires at least on argument")
|
||||
// }
|
||||
//
|
||||
// // First and only argument should be the ID of the device we want to use
|
||||
// targetID, err := strconv.ParseInt(args[0], 10, 0)
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
//
|
||||
// err = perClerk.ListDeviceAPI(uint8(targetID))
|
||||
// if err != nil {
|
||||
// fmt.Println("failed to view device API: ", err.Error())
|
||||
// }
|
||||
//
|
||||
// return nil
|
||||
// },
|
||||
// }
|
||||
//
|
||||
// runDeviceAPIREPLCmd := repl.Command{
|
||||
// Name: "v-run",
|
||||
// Usage: "runs a funcion of a device - pass its ID and function name",
|
||||
// Action: func(r *repl.REPL, args []string) error {
|
||||
// // We should add this to the REPL instead
|
||||
// if len(args) < 3 {
|
||||
// return fmt.Errorf("requires at least on argument")
|
||||
// }
|
||||
//
|
||||
// // First and only argument should be the ID of the device we want to use
|
||||
// targetID, err := strconv.ParseInt(args[0], 10, 0)
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
//
|
||||
// fnName := args[1]
|
||||
// fnArgs := args[2:]
|
||||
//
|
||||
// err = perClerk.RunDeviceFunction(uint8(targetID), fnName, fnArgs)
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
//
|
||||
// return nil
|
||||
// },
|
||||
// }
|
||||
//
|
||||
// r.RegisterCMD(discoverDevicesREPLCmd)
|
||||
// r.RegisterCMD(listDevicesREPLCmd)
|
||||
// r.RegisterCMD(listDeviceAPIREPLCmd)
|
||||
// r.RegisterCMD(runDeviceAPIREPLCmd)
|
||||
//
|
||||
// r.Start()
|
||||
// r.Close()
|
||||
r := repl.NewREPL(repl.REPLCfg{
|
||||
PS1: "\rESDI > ",
|
||||
})
|
||||
|
||||
perClerk := peripheral.NewPeripheralDeviceClerk()
|
||||
|
||||
discoverDevicesREPLCmd := repl.Command{
|
||||
Name: "discover",
|
||||
Usage: "discovers connected devices",
|
||||
Action: func(r *repl.REPL, args []string) error {
|
||||
err := perClerk.FindDevices()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
listDevicesREPLCmd := repl.Command{
|
||||
Name: "list",
|
||||
Usage: "lists connected devices",
|
||||
Action: func(r *repl.REPL, args []string) error {
|
||||
_ = perClerk.ListDevices()
|
||||
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
listDeviceAPIREPLCmd := repl.Command{
|
||||
Name: "v-api",
|
||||
Usage: "shows API of a device - pass its ID",
|
||||
Action: func(r *repl.REPL, args []string) error {
|
||||
// We should add this to the REPL instead
|
||||
if len(args) < 1 {
|
||||
return fmt.Errorf("requires at least on argument")
|
||||
}
|
||||
|
||||
// First and only argument should be the ID of the device we want to use
|
||||
targetID, err := strconv.ParseInt(args[0], 10, 0)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
err = perClerk.ListDeviceAPI(uint8(targetID))
|
||||
if err != nil {
|
||||
fmt.Println("failed to view device API: ", err.Error())
|
||||
}
|
||||
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
runDeviceAPIREPLCmd := repl.Command{
|
||||
Name: "v-run",
|
||||
Usage: "runs a funcion of a device - pass its ID and function name",
|
||||
Action: func(r *repl.REPL, args []string) error {
|
||||
// We should add this to the REPL instead
|
||||
if len(args) < 3 {
|
||||
return fmt.Errorf("requires at least on argument")
|
||||
}
|
||||
|
||||
// First and only argument should be the ID of the device we want to use
|
||||
targetID, err := strconv.ParseInt(args[0], 10, 0)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
fnName := args[1]
|
||||
fnArgs := args[2:]
|
||||
|
||||
err = perClerk.RunDeviceFunction(uint8(targetID), fnName, fnArgs)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
},
|
||||
}
|
||||
|
||||
r.RegisterCMD(discoverDevicesREPLCmd)
|
||||
r.RegisterCMD(listDevicesREPLCmd)
|
||||
r.RegisterCMD(listDeviceAPIREPLCmd)
|
||||
r.RegisterCMD(runDeviceAPIREPLCmd)
|
||||
|
||||
r.Start()
|
||||
r.Close()
|
||||
}
|
||||
|
||||
// removeLabelCmd represents the removeLabel command
|
||||
|
||||
@@ -1,6 +0,0 @@
|
||||
package constants
|
||||
|
||||
const (
|
||||
IRacingProviderName = "iRacing"
|
||||
BeamNGProviderName = "BeamNG.drive"
|
||||
)
|
||||
@@ -1,12 +1,10 @@
|
||||
package cdashdisplay
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"time"
|
||||
|
||||
"esdi/peripheral/communication"
|
||||
"esdi/peripheral/communication/packets"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/tarm/serial"
|
||||
portp "go.bug.st/serial"
|
||||
@@ -48,11 +46,11 @@ func findDisplayPort() (*communication.WalkieTalkie, error) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
slog.Info(fmt.Sprintf("Looking into %v", ports))
|
||||
pLogger.Info(fmt.Sprintf("Looking into %v", ports))
|
||||
|
||||
var wt *communication.WalkieTalkie
|
||||
for _, port := range ports {
|
||||
slog.Info(fmt.Sprintf("Trying port %s", port))
|
||||
pLogger.Info(fmt.Sprintf("Trying port %s", port))
|
||||
|
||||
wt = &communication.WalkieTalkie{
|
||||
Cfg: &serial.Config{
|
||||
@@ -62,7 +60,7 @@ func findDisplayPort() (*communication.WalkieTalkie, error) {
|
||||
},
|
||||
}
|
||||
|
||||
slog.Info(fmt.Sprintf("Started probing port %s", port))
|
||||
pLogger.Info(fmt.Sprintf("Started probing port %s", port))
|
||||
|
||||
probeResult := make(chan error, 1)
|
||||
|
||||
@@ -73,19 +71,19 @@ func findDisplayPort() (*communication.WalkieTalkie, error) {
|
||||
select {
|
||||
case err = <-probeResult:
|
||||
// Probe completed normally (could be success or error)
|
||||
case <-time.After(1000 * time.Millisecond):
|
||||
case <-time.After(2 * time.Second):
|
||||
// Hard timeout reached
|
||||
err = fmt.Errorf("probe completely hung/timed out: %s", port)
|
||||
}
|
||||
|
||||
slog.Info(fmt.Sprintf("Finished probing port %s", port))
|
||||
pLogger.Info(fmt.Sprintf("Finished probing port %s", port))
|
||||
|
||||
if err == nil {
|
||||
slog.Info(fmt.Sprintf("Success probing port %s: %+v", port, err))
|
||||
pLogger.Info(fmt.Sprintf("Success probing port %s: %+v", port, err))
|
||||
break
|
||||
}
|
||||
|
||||
slog.Info(fmt.Sprintf("wasn't port %s", port))
|
||||
pLogger.Info(fmt.Sprintf("wasn't port %s", port))
|
||||
wt = nil
|
||||
}
|
||||
|
||||
@@ -93,6 +91,6 @@ func findDisplayPort() (*communication.WalkieTalkie, error) {
|
||||
return nil, fmt.Errorf("couldn't find cdashdisplay")
|
||||
}
|
||||
|
||||
slog.Info(fmt.Sprintf("found cdashdisplay on port: %s", wt.Cfg.Name))
|
||||
pLogger.Info(fmt.Sprintf("found cdashdisplay on port: %s", wt.Cfg.Name))
|
||||
return wt, nil
|
||||
}
|
||||
|
||||
@@ -8,11 +8,9 @@ import (
|
||||
"log/slog"
|
||||
"os"
|
||||
"path"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
helper "esdi/helpers"
|
||||
"esdi/peripheral"
|
||||
"esdi/peripheral/communication"
|
||||
"esdi/peripheral/communication/packets"
|
||||
"esdi/peripheral/types"
|
||||
@@ -21,6 +19,13 @@ import (
|
||||
"gopkg.in/yaml.v3"
|
||||
)
|
||||
|
||||
// TODO: remove this logger. It should be on the device struct
|
||||
var pLogger *slog.Logger
|
||||
|
||||
func SetLogger(l *slog.Logger) {
|
||||
pLogger = l
|
||||
}
|
||||
|
||||
// I have to move this to some kind of configuration place
|
||||
const (
|
||||
layoutsDir = "./layouts/"
|
||||
@@ -33,8 +38,6 @@ const (
|
||||
updateWindowCMDID types.Command = 6 // Change this to a move cmd instead
|
||||
sendDataCMDID types.Command = 7
|
||||
newLayoutCMDID types.Command = 8
|
||||
healthCheckCMDID types.Command = 9
|
||||
resetCMDID types.Command = 10
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -110,65 +113,26 @@ func NewCDashState() *CDashState {
|
||||
}
|
||||
|
||||
type CDashDisplay struct {
|
||||
WT *communication.WalkieTalkie
|
||||
State *CDashState
|
||||
fieldToWindows map[telemetry.FieldID][]int16
|
||||
bufPool sync.Pool
|
||||
failedSends int
|
||||
FailedSendsConsecutiveLimit int
|
||||
WT *communication.WalkieTalkie
|
||||
State *CDashState
|
||||
}
|
||||
|
||||
// Connect will try to find and connect to the CDashDisplay
|
||||
func NewCDashDisplay() (*CDashDisplay, error) {
|
||||
func Discover() (*CDashDisplay, error) {
|
||||
// Look for the port
|
||||
p, err := findDisplayPort()
|
||||
if err != nil {
|
||||
slog.Info("failed to find cdashdisplay port", "reason", err.Error())
|
||||
pLogger.Info("failed to find cdashdisplay port: %s", err.Error())
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &CDashDisplay{
|
||||
WT: p,
|
||||
State: NewCDashState(),
|
||||
fieldToWindows: make(map[telemetry.FieldID][]int16),
|
||||
bufPool: sync.Pool{
|
||||
New: func() any {
|
||||
b := make([]byte, 0, telemetry.MaxFields*8)
|
||||
return &b
|
||||
},
|
||||
},
|
||||
failedSends: 0,
|
||||
FailedSendsConsecutiveLimit: 5,
|
||||
WT: p,
|
||||
State: NewCDashState(),
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (cds *CDashDisplay) Close() error {
|
||||
// if cds.WT != nil {
|
||||
// cds.Close()
|
||||
// }
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) RegisterFieldMapping(fieldID telemetry.FieldID, winID int16) {
|
||||
d.fieldToWindows[fieldID] = append(d.fieldToWindows[fieldID], winID)
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) UnregisterFieldMapping(winID int16) {
|
||||
for fieldID, windows := range d.fieldToWindows {
|
||||
updated := windows[:0]
|
||||
for _, w := range windows {
|
||||
if w != winID {
|
||||
updated = append(updated, w)
|
||||
}
|
||||
|
||||
if len(updated) == 0 {
|
||||
delete(d.fieldToWindows, fieldID)
|
||||
} else {
|
||||
d.fieldToWindows[fieldID] = updated
|
||||
}
|
||||
}
|
||||
}
|
||||
func (d *CDashDisplay) SendCommand() {
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) CreateWindow(win *DesktopUIWindow) (*DesktopUIWindow, error) {
|
||||
@@ -186,12 +150,9 @@ func (d *CDashDisplay) CreateWindow(win *DesktopUIWindow) (*DesktopUIWindow, err
|
||||
|
||||
win.UIData.IDX = wID.ID
|
||||
|
||||
slog.Info(fmt.Sprintf("Recived ID message: %v", wID))
|
||||
pLogger.Info(fmt.Sprintf("Recived ID message: %v", wID))
|
||||
|
||||
d.State.Layout.AddWindow(win)
|
||||
if fieldID, ok := telemetry.GetFieldID(win.UIData.TelemetryField); ok {
|
||||
d.RegisterFieldMapping(fieldID, win.UIData.IDX)
|
||||
}
|
||||
|
||||
return win, nil
|
||||
}
|
||||
@@ -217,14 +178,12 @@ func (d *CDashDisplay) UpdateWindow(win *DesktopUIWindow) error {
|
||||
// I get it and update it in the controller
|
||||
// I send the pointer here
|
||||
// -> it should be the same pointer then right?
|
||||
slog.Debug(fmt.Sprintf("PreUpdate ID: %p", win))
|
||||
pLogger.Debug(fmt.Sprintf("PreUpdate ID: %p", win))
|
||||
d.State.Layout.Windows[win.UIData.IDX] = win
|
||||
slog.Debug(fmt.Sprintf("PostUpdate ID: %p", win))
|
||||
pLogger.Debug(fmt.Sprintf("PostUpdate ID: %p", win))
|
||||
// Yeah, same address as suspected
|
||||
// I can't think about it right now. I'll think about that tomorrow
|
||||
|
||||
// TODO: need to update the field mappings here!
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -247,15 +206,14 @@ func (d *CDashDisplay) DestroyWindow(wID int16) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// NOTE: add this
|
||||
d.UnregisterFieldMapping(wID)
|
||||
// NODE: add this
|
||||
d.State.Layout.RemoveWindow(wID)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) updateWindowDimensions(win *UIWindow, packet UpdateDimsPacket) error {
|
||||
slog.Debug(fmt.Sprintf("UPDATE: %v", packet))
|
||||
pLogger.Debug(fmt.Sprintf("UPDATE: %v", packet))
|
||||
|
||||
bytes, err := helper.StructToBytes(packet)
|
||||
if err != nil {
|
||||
@@ -269,9 +227,9 @@ func (d *CDashDisplay) updateWindowDimensions(win *UIWindow, packet UpdateDimsPa
|
||||
}
|
||||
|
||||
// Nothing bad happened afaik
|
||||
slog.Debug(fmt.Sprintf("cur dims: %v", win.Dims))
|
||||
pLogger.Debug(fmt.Sprintf("cur dims: %v", win.Dims))
|
||||
win.Dims = packet.Dims
|
||||
slog.Debug(fmt.Sprintf("new dims: %v", win.Dims))
|
||||
pLogger.Debug(fmt.Sprintf("new dims: %v", win.Dims))
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -387,41 +345,30 @@ func (d *CDashDisplay) LoadLayout(layoutName string) error {
|
||||
func (d *CDashDisplay) UnloadLayout() error {
|
||||
var err error
|
||||
for _, w := range d.State.Layout.Windows {
|
||||
slog.Debug(fmt.Sprintf("= Removing %d ==============================================",
|
||||
pLogger.Debug(fmt.Sprintf("= Removing %d ==============================================",
|
||||
w.UIData.IDX))
|
||||
|
||||
err = d.DestroyWindow(w.UIData.IDX)
|
||||
time.Sleep(75 * time.Millisecond)
|
||||
if err != nil {
|
||||
slog.Error(fmt.Sprintf("failed to destroy window: %+v", err))
|
||||
pLogger.Error(fmt.Sprintf("failed to destroy window: %+v", err))
|
||||
// NOTE: Add a way to handle multiple errors ?
|
||||
return err
|
||||
}
|
||||
|
||||
slog.Debug(fmt.Sprintf("= Removing %d ==============================================",
|
||||
pLogger.Debug(fmt.Sprintf("= Removing %d ==============================================",
|
||||
w.UIData.IDX))
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (cds *CDashDisplay) reset() error {
|
||||
err := cds.WT.SendCommand(resetCMDID, []byte{0x01, 0x02, 0x03, 0x04}, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
time.Sleep(3000 * time.Millisecond)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) error {
|
||||
packet := d.encodePacket(data)
|
||||
func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) {
|
||||
packet := data.Pack()
|
||||
|
||||
bytes, err := helper.StructToBytes(packet)
|
||||
if err != nil {
|
||||
return peripheral.ErrFailureToPackData
|
||||
return
|
||||
}
|
||||
|
||||
curStr := ""
|
||||
@@ -431,6 +378,7 @@ func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) error {
|
||||
curStr += fmt.Sprintf("%02x ", byte)
|
||||
|
||||
if byteCount == 8 {
|
||||
// pLogger.Debug(curStr)
|
||||
curStr = ""
|
||||
byteCount = 0
|
||||
}
|
||||
@@ -439,11 +387,6 @@ func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) error {
|
||||
// var ack packets.AckPacket
|
||||
err = d.WT.SendCommand(sendDataCMDID, bytes, nil)
|
||||
if err != nil && err != io.EOF {
|
||||
if d.failedSends == d.FailedSendsConsecutiveLimit {
|
||||
return peripheral.ErrDeviceTimedOut
|
||||
}
|
||||
d.failedSends++
|
||||
return
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1,9 +0,0 @@
|
||||
package cdashdisplay
|
||||
|
||||
func (cds *CDashDisplay) OnLoad() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (cds *CDashDisplay) OnTelemetryProviderFound() error {
|
||||
return nil
|
||||
}
|
||||
@@ -1,7 +1,6 @@
|
||||
package cdashdisplay
|
||||
|
||||
import (
|
||||
"esdi/peripheral/communication/packets"
|
||||
"esdi/peripheral/devices"
|
||||
"esdi/telemetry"
|
||||
)
|
||||
@@ -16,25 +15,13 @@ func (cds *CDashDisplay) Name() string {
|
||||
return NAME
|
||||
}
|
||||
|
||||
func (cds *CDashDisplay) RequiredFields() []telemetry.FieldID {
|
||||
fields := make([]telemetry.FieldID, 0, len(cds.State.Layout.Windows))
|
||||
func (cds *CDashDisplay) RequiredFields() map[int16]telemetry.FieldID {
|
||||
fields := make(map[int16]telemetry.FieldID, len(cds.State.Layout.Windows))
|
||||
|
||||
for _, w := range cds.State.Layout.Windows {
|
||||
if fieldID, ok := telemetry.GetFieldID(w.UIData.TelemetryField); ok {
|
||||
fields = append(fields, fieldID)
|
||||
}
|
||||
fieldID, _ := telemetry.GetFieldID(w.UIData.TelemetryField)
|
||||
fields[w.UIData.IDX] = fieldID
|
||||
}
|
||||
|
||||
return fields
|
||||
}
|
||||
|
||||
func (cds *CDashDisplay) HealthCheck() bool {
|
||||
// Send the command
|
||||
var health packets.HealthCheck
|
||||
err := cds.WT.SendCommand(healthCheckCMDID, []byte{0x01, 0x02, 0x03, 0x04}, &health)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
@@ -1,9 +1,6 @@
|
||||
package cdashdisplay
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
)
|
||||
import "fmt"
|
||||
|
||||
type LayoutTree struct {
|
||||
Windows map[int16]*DesktopUIWindow `yaml:"Windows"`
|
||||
@@ -16,13 +13,13 @@ func NewLayoutTree() *LayoutTree {
|
||||
}
|
||||
|
||||
func (l *LayoutTree) AddWindow(w *DesktopUIWindow) {
|
||||
slog.Debug(fmt.Sprintf("adding window '%d' - %v", w.UIData.IDX))
|
||||
pLogger.Debug(fmt.Sprintf("adding window '%d' - %v", w.UIData.IDX))
|
||||
l.Windows[w.UIData.IDX] = w
|
||||
slog.Debug(fmt.Sprintf("new map - %v", l.Windows))
|
||||
pLogger.Debug(fmt.Sprintf("new map - %v", l.Windows))
|
||||
}
|
||||
|
||||
func (l *LayoutTree) RemoveWindow(idx int16) {
|
||||
slog.Debug(fmt.Sprintf("removing window '%d'", idx))
|
||||
pLogger.Debug(fmt.Sprintf("removing window '%d'", idx))
|
||||
delete(l.Windows, idx)
|
||||
slog.Debug(fmt.Sprintf("new map - %v", l.Windows))
|
||||
pLogger.Debug(fmt.Sprintf("new map - %v", l.Windows))
|
||||
}
|
||||
|
||||
@@ -1,44 +0,0 @@
|
||||
package cdashdisplay
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"esdi/constants"
|
||||
)
|
||||
|
||||
func (cds *CDashDisplay) setupForIracing() error {
|
||||
err := cds.LoadLayout("layout.yaml")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (cds *CDashDisplay) setupForBeamNG() error {
|
||||
err := cds.LoadLayout("beamng.yaml")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (cds *CDashDisplay) Setup(provider string) error {
|
||||
// Doesn't matter which one we are picking, we need to reset the CDashDisplay
|
||||
// first - we will improve this setup behaviour later on with an Update to
|
||||
// change it without dropping connection or some shit
|
||||
err := cds.reset()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
switch provider {
|
||||
case constants.IRacingProviderName:
|
||||
return cds.setupForIracing()
|
||||
case constants.BeamNGProviderName:
|
||||
return cds.setupForBeamNG()
|
||||
default:
|
||||
return fmt.Errorf("unknown provider: %s", provider)
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,5 @@
|
||||
package cdashdisplay
|
||||
|
||||
import (
|
||||
"math"
|
||||
|
||||
"esdi/telemetry"
|
||||
)
|
||||
|
||||
// In this file we will place all structs that are 1:1 representation of the
|
||||
// types in the device transport layer ->
|
||||
|
||||
@@ -71,71 +65,6 @@ type UIWindowUpdatePacket struct {
|
||||
Window UIWindow
|
||||
}
|
||||
|
||||
// TODO: this is here because this is the only device I use like this
|
||||
// once I add another serial device I will move this somewhere else that can
|
||||
// be reused by all devices but still is decoupled from the telemetry package
|
||||
func (cds *CDashDisplay) encodePacket(td *telemetry.TelemetryData) []byte {
|
||||
bufPtr := cds.bufPool.Get().(*[]byte)
|
||||
buf := (*bufPtr)[:0]
|
||||
|
||||
for fieldID, windowIDs := range cds.fieldToWindows {
|
||||
if int(fieldID) >= len(td.Values) || len(windowIDs) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
tf := &td.Values[fieldID]
|
||||
|
||||
for _, winID := range windowIDs {
|
||||
buf = packField(winID, tf, buf)
|
||||
}
|
||||
}
|
||||
|
||||
// We have to copy here because we have to return the buffer
|
||||
result := make([]byte, len(buf))
|
||||
copy(result, buf)
|
||||
|
||||
cds.bufPool.Put(&buf)
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
// Pack will pack this current TelemetryField into bytes to send over the wire
|
||||
// Format:
|
||||
// 0x00 - Field ID
|
||||
// 0x00 |
|
||||
// 0x01 - DataType
|
||||
// 0x02 - if its a (u)int8
|
||||
// or
|
||||
// 0x02 - if its a (u)int16 - first byte
|
||||
// 0x02 - if its a (u)int16 - second byte
|
||||
// or
|
||||
// 0x02 - str len max is 255 chars
|
||||
// [0x02] - str
|
||||
func packField(winID int16, tf *telemetry.TelemetryField, dest []byte) []byte {
|
||||
// NOTE: maybe we can have a pool of these so we don't have to create them here
|
||||
// or whatever
|
||||
dest = append(dest, uint8(winID), uint8(winID>>8))
|
||||
dest = append(dest, uint8(tf.Type))
|
||||
|
||||
switch tf.Type {
|
||||
case telemetry.DataTypeINT8, telemetry.DataTypeUINT8, telemetry.DataTypeCHAR:
|
||||
dest = append(dest, uint8(tf.Raw))
|
||||
case telemetry.DataTypeINT16, telemetry.DataTypeUINT16:
|
||||
dest = append(dest, uint8(tf.Raw), uint8(tf.Raw>>8))
|
||||
case telemetry.DataTypeINT32, telemetry.DataTypeUINT32:
|
||||
dest = append(dest, uint8(tf.Raw), uint8(tf.Raw>>8), uint8(tf.Raw>>16), uint8(tf.Raw>>24))
|
||||
case telemetry.DataTypeINT64, telemetry.DataTypeUINT64:
|
||||
dest = append(
|
||||
dest, uint8(tf.Raw), uint8(tf.Raw>>8), uint8(tf.Raw>>16),
|
||||
uint8(tf.Raw>>24), uint8(tf.Raw>>32), uint8(tf.Raw>>40), uint8(tf.Raw>>48),
|
||||
uint8(tf.Raw>>56),
|
||||
)
|
||||
case telemetry.DataTypeSTRING:
|
||||
l := min(len(tf.Str), math.MaxUint8)
|
||||
|
||||
dest = append(dest, uint8(l))
|
||||
dest = append(dest, tf.Str[:l]...)
|
||||
}
|
||||
|
||||
return dest
|
||||
}
|
||||
// func getUIWindowDTO(w *DesktopWindowData) *UIWindow {
|
||||
//
|
||||
// }
|
||||
|
||||
+27
-33
@@ -9,26 +9,23 @@ import (
|
||||
"esdi/peripheral"
|
||||
)
|
||||
|
||||
var ErrInvalidDevice = errors.New("invalid device")
|
||||
|
||||
type Device struct {
|
||||
Name string
|
||||
Discover func() (peripheral.Peripheral, error)
|
||||
// DefaultSetup func(peripheral.Peripheral) error
|
||||
}
|
||||
|
||||
var List map[string]*Device = map[string]*Device{
|
||||
var List map[string]Device = map[string]Device{
|
||||
uidevice.NAME: {
|
||||
Name: uidevice.NAME,
|
||||
Discover: UIDeviceDiscover,
|
||||
Discover: DiscoverUIDevice,
|
||||
},
|
||||
cdashdisplay.NAME: {
|
||||
Name: cdashdisplay.NAME,
|
||||
Discover: CDashDisplayDiscover,
|
||||
Discover: DiscoverCDashDisplay,
|
||||
},
|
||||
}
|
||||
|
||||
func UIDeviceDiscover() (peripheral.Peripheral, error) {
|
||||
func DiscoverUIDevice() (peripheral.Peripheral, error) {
|
||||
uidev, err := uidevice.NewUIDevice()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -37,30 +34,27 @@ func UIDeviceDiscover() (peripheral.Peripheral, error) {
|
||||
return uidev, nil
|
||||
}
|
||||
|
||||
// func UIDeviceSetup(peripheral peripheral.Peripheral) error {
|
||||
// return nil
|
||||
// }
|
||||
|
||||
func CDashDisplayDiscover() (peripheral.Peripheral, error) {
|
||||
// Create a cdashdisplay
|
||||
display, err := cdashdisplay.NewCDashDisplay()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return display, nil
|
||||
func DiscoverCDashDisplay() (peripheral.Peripheral, error) {
|
||||
// // Find CDashDisplay
|
||||
// {
|
||||
// ds.Messages <- "looking for " + cdashdisplay.Name + "...\n"
|
||||
// ds.Logger.Info("Looking for " + cdashdisplay.Name)
|
||||
//
|
||||
// cdashdisplay.SetLogger(ds.Logger.With("[device]", cdashdisplay.Name))
|
||||
//
|
||||
// // Create a cdashdisplay
|
||||
// display, err := cdashdisplay.Discover()
|
||||
// if err == nil {
|
||||
// ds.Devices[cdashdisplay.Name] = display
|
||||
// ds.Logger.Info("found " + cdashdisplay.Name + " on: " + display.WT.Cfg.Name)
|
||||
// ds.Messages <- "found " + cdashdisplay.Name + " on: " + display.WT.Cfg.Name + "\n"
|
||||
// return
|
||||
// }
|
||||
//
|
||||
// ds.Logger.Info("didn't find " + cdashdisplay.Name)
|
||||
// ds.Messages <- "didn't find " + cdashdisplay.Name + "\n"
|
||||
// // No CDashDisplay available for one reason or another, so we don't set the
|
||||
// // key
|
||||
// }
|
||||
return nil, errors.New("not implemented yet")
|
||||
}
|
||||
|
||||
// func CDashDisplaySetup(peripheral peripheral.Peripheral) error {
|
||||
// cdash, ok := peripheral.(*cdashdisplay.CDashDisplay)
|
||||
// if !ok {
|
||||
// return ErrInvalidDevice
|
||||
// }
|
||||
//
|
||||
// err := cdash.LoadLayout("layout.yaml")
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
//
|
||||
// return nil
|
||||
// }
|
||||
|
||||
+10
-21
@@ -17,21 +17,9 @@ func NewUIDevice() (peripheral.Peripheral, error) {
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (uid *UIDevice) Close() error {
|
||||
if uid.dataChan != nil {
|
||||
close(uid.dataChan)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (uid *UIDevice) Setup(provider string) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (uid *UIDevice) SendData(data *telemetry.TelemetryData) error {
|
||||
func (uid *UIDevice) SendData(data *telemetry.TelemetryData) {
|
||||
if data == nil {
|
||||
return peripheral.ErrInvalidData
|
||||
return
|
||||
}
|
||||
|
||||
select {
|
||||
@@ -39,8 +27,6 @@ func (uid *UIDevice) SendData(data *telemetry.TelemetryData) error {
|
||||
default:
|
||||
// Drop frame if buffer is full
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (uid *UIDevice) Name() string {
|
||||
@@ -51,14 +37,17 @@ func (uid *UIDevice) DataChannel() <-chan telemetry.TelemetryData {
|
||||
return uid.dataChan
|
||||
}
|
||||
|
||||
func (uid *UIDevice) RequiredFields() []telemetry.FieldID {
|
||||
return []telemetry.FieldID{
|
||||
func (uid *UIDevice) RequiredFields() map[int16]telemetry.FieldID {
|
||||
subscribeTo := []telemetry.FieldID{
|
||||
telemetry.Speed,
|
||||
telemetry.Gear,
|
||||
telemetry.RPM,
|
||||
}
|
||||
}
|
||||
|
||||
func (uid *UIDevice) HealthCheck() bool {
|
||||
return true
|
||||
fields := make(map[int16]telemetry.FieldID, len(subscribeTo))
|
||||
for k, field := range subscribeTo {
|
||||
fields[int16(k)] = field
|
||||
}
|
||||
|
||||
return fields
|
||||
}
|
||||
|
||||
@@ -1,9 +0,0 @@
|
||||
package uidevice
|
||||
|
||||
func (uid *UIDevice) OnLoad() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (uid *UIDevice) OnTelemetryProviderFound() error {
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
package logger
|
||||
|
||||
import (
|
||||
"log"
|
||||
"os"
|
||||
"sync"
|
||||
)
|
||||
|
||||
var l *log.Logger
|
||||
var once sync.Once
|
||||
|
||||
func createLogger() {
|
||||
l = log.New(os.Stdout, "[esdi] ", log.LstdFlags | log.Lshortfile)
|
||||
}
|
||||
|
||||
func GetInstance() *log.Logger {
|
||||
once.Do(func() {
|
||||
createLogger()
|
||||
})
|
||||
|
||||
return l
|
||||
}
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
|
||||
"esdi/cmd"
|
||||
"esdi/config"
|
||||
"esdi/telemetry"
|
||||
|
||||
"github.com/arl/statsviz"
|
||||
)
|
||||
@@ -41,6 +42,9 @@ func initApplication() {
|
||||
slog.Error(fmt.Sprintf("failed to setup metrics server: %+v", err))
|
||||
}
|
||||
}
|
||||
|
||||
// Setting up some internal data structures
|
||||
telemetry.Init()
|
||||
}
|
||||
|
||||
func setupLogger() error {
|
||||
|
||||
@@ -17,7 +17,6 @@ const (
|
||||
CmdAckID types.Command = 2
|
||||
CmdCreateWindow types.Command = 3
|
||||
CmdDestroyWindow types.Command = 4
|
||||
CmdHealthCheck types.Command = 9
|
||||
)
|
||||
|
||||
var crc8Table = [256]byte{
|
||||
|
||||
@@ -18,22 +18,3 @@ func (pkt *NewWindowID) Validate() bool {
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
type HealthCheck struct {
|
||||
StartMarker byte
|
||||
Response byte
|
||||
EndMarker byte
|
||||
}
|
||||
|
||||
func (pkt *HealthCheck) Validate() bool {
|
||||
if pkt.StartMarker != constvar.StartOfText ||
|
||||
pkt.EndMarker != constvar.EndOfText {
|
||||
return false
|
||||
}
|
||||
|
||||
if pkt.Response != 0x06 {
|
||||
return false
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
@@ -26,7 +26,8 @@ func (wt *WalkieTalkie) ReadFramedData(size int, packet any) error {
|
||||
b := make([]byte, 1)
|
||||
_, err := wt.Serial.Read(b)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error reading incoming: %+v, err:", b, err)
|
||||
// fmt.Fprintf(os.Stderr, "dev read: %s\n", err.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
if b[0] == constvar.StartOfText {
|
||||
@@ -43,11 +44,9 @@ func (wt *WalkieTalkie) ReadFramedData(size int, packet any) error {
|
||||
reader := bytes.NewReader(buf)
|
||||
err = binary.Read(reader, binary.LittleEndian, packet)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error parsing incoming: %+v, err:", buf, err)
|
||||
return err
|
||||
}
|
||||
|
||||
wt.Serial.Flush()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -109,15 +108,6 @@ func (wt *WalkieTalkie) sendPacket(cmd types.Command, data any) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// slog.Debug("# START ########################################################")
|
||||
// slog.Debug(fmt.Sprintf("StartMarker: %02x", constvar.StartOfText))
|
||||
// slog.Debug(fmt.Sprintf("CMD: %02x", cmd))
|
||||
// slog.Debug(fmt.Sprintf("Len: %d", len(payload)))
|
||||
// slog.Debug(fmt.Sprintf("Payload: %v", payload))
|
||||
// slog.Debug(fmt.Sprintf("CRC: %v", CRC8(payload)))
|
||||
// slog.Debug(fmt.Sprintf("EndMarker: %02x", constvar.EndOfText))
|
||||
// slog.Debug("-")
|
||||
|
||||
packet := CMDDataPacket{
|
||||
StartMarker: constvar.StartOfText,
|
||||
CMD: cmd,
|
||||
@@ -129,17 +119,13 @@ func (wt *WalkieTalkie) sendPacket(cmd types.Command, data any) error {
|
||||
|
||||
// Send the payload
|
||||
serializedPacket := packet.Serialize()
|
||||
|
||||
// slog.Debug("Serialized packet", "packet", packet)
|
||||
// slog.Debug("# END ##########################################################")
|
||||
// fmt.Fprintf(os.Stderr, "%+v", serializedPacket)
|
||||
|
||||
_, err = wt.Serial.Write(serializedPacket)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
wt.Serial.Flush()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -161,11 +147,50 @@ func (wt *WalkieTalkie) readPacket(resp packets.Packet) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (wt *WalkieTalkie) SendCommand(
|
||||
cmd types.Command,
|
||||
payload any,
|
||||
responseBody packets.Packet,
|
||||
) error {
|
||||
// func (wt *WalkieTalkie) sendHeader(h *header) error {
|
||||
// err := wt.sendPacket(h)
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
//
|
||||
// // var ack packets.AckPacket
|
||||
// // err = wt.readPacket(&ack)
|
||||
// // if err != nil {
|
||||
// // return err
|
||||
// // }
|
||||
//
|
||||
// return nil
|
||||
// }
|
||||
|
||||
// func (wt *WalkieTalkie) sendBody(payload any, resp packets.Packet) error {
|
||||
// err := wt.sendPacket(payload)
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
//
|
||||
// // err = wt.readPacket(resp)
|
||||
// // if err != nil {
|
||||
// // return err
|
||||
// // }
|
||||
//
|
||||
// return nil
|
||||
// }
|
||||
|
||||
func (wt *WalkieTalkie) SendCommand(cmd types.Command, payload any,
|
||||
responseBody packets.Packet) error {
|
||||
// Prepare the header
|
||||
// header := header{
|
||||
// StartByte: constvar.StartOfText,
|
||||
// CMD: cmd,
|
||||
// EndByte: constvar.EndOfText,
|
||||
// }
|
||||
|
||||
// Send the header
|
||||
// err := wt.sendHeader(&header)
|
||||
// if err != nil {
|
||||
// return err
|
||||
// }
|
||||
|
||||
// Send the body
|
||||
err := wt.sendPacket(cmd, payload)
|
||||
if err != nil {
|
||||
|
||||
@@ -1,9 +0,0 @@
|
||||
package peripheral
|
||||
|
||||
import "errors"
|
||||
|
||||
var (
|
||||
ErrInvalidData = errors.New("invalid data")
|
||||
ErrDeviceTimedOut = errors.New("device timed out")
|
||||
ErrFailureToPackData = errors.New("failed to pack received data")
|
||||
)
|
||||
@@ -17,13 +17,8 @@ const (
|
||||
|
||||
type Peripheral interface {
|
||||
Name() string
|
||||
Setup(string) error
|
||||
HealthCheck() bool
|
||||
SendData(*telemetry.TelemetryData) error
|
||||
RequiredFields() []telemetry.FieldID
|
||||
OnLoad() error
|
||||
OnTelemetryProviderFound() error
|
||||
Close() error
|
||||
SendData(*telemetry.TelemetryData)
|
||||
RequiredFields() map[int16]telemetry.FieldID
|
||||
}
|
||||
|
||||
type PeripheralDeviceClerk struct {
|
||||
|
||||
+19
-21
@@ -8,7 +8,6 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"esdi/constants"
|
||||
"esdi/telemetry"
|
||||
|
||||
bngsdk "github.com/ESilva15/gobngsdk"
|
||||
@@ -29,7 +28,7 @@ type BeamNG struct {
|
||||
updaters [telemetry.MaxFields]func(*telemetry.TelemetryField)
|
||||
|
||||
// stream control
|
||||
wg sync.WaitGroup
|
||||
streamCh chan telemetry.TelemetryData
|
||||
streamCancel context.CancelFunc
|
||||
|
||||
// timing
|
||||
@@ -37,7 +36,7 @@ type BeamNG struct {
|
||||
}
|
||||
|
||||
const (
|
||||
NAME = constants.BeamNGProviderName
|
||||
NAME = "BeamNG.drive"
|
||||
)
|
||||
|
||||
func NewBeamNGProvider(logger *slog.Logger, opts *bngsdk.Options) (*BeamNG, error) {
|
||||
@@ -47,11 +46,12 @@ func NewBeamNGProvider(logger *slog.Logger, opts *bngsdk.Options) (*BeamNG, erro
|
||||
}
|
||||
|
||||
provider := &BeamNG{
|
||||
logger: logger.With("TelemetryProvider", NAME),
|
||||
data: telemetry.NewTelemetryData(),
|
||||
SDK: beam,
|
||||
og: &bngsdk.Outgauge{},
|
||||
ticker: time.NewTicker(time.Second / 60),
|
||||
logger: logger.With("TelemetryProvider", NAME),
|
||||
streamCh: make(chan telemetry.TelemetryData, 1),
|
||||
data: telemetry.NewTelemetryData(),
|
||||
SDK: beam,
|
||||
og: &bngsdk.Outgauge{},
|
||||
ticker: time.NewTicker(time.Second / 60),
|
||||
}
|
||||
|
||||
provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){
|
||||
@@ -110,12 +110,12 @@ func (b *BeamNG) Stream() (<-chan telemetry.TelemetryData, error) {
|
||||
ctx, b.streamCancel = context.WithCancel(context.Background())
|
||||
|
||||
// Start the stream
|
||||
ch := b.stream(ctx)
|
||||
b.stream(ctx)
|
||||
|
||||
return ch, nil
|
||||
return b.streamCh, nil
|
||||
}
|
||||
|
||||
func (b *BeamNG) Subscribe(requestFields []telemetry.FieldID) {
|
||||
func (b *BeamNG) Subscribe(requestFields map[int16]telemetry.FieldID) {
|
||||
// NOTE: document how the Subscribe funtion works
|
||||
slog.Debug(fmt.Sprintf("Len Req: %d\n", len(requestFields)))
|
||||
|
||||
@@ -125,7 +125,9 @@ func (b *BeamNG) Subscribe(requestFields []telemetry.FieldID) {
|
||||
// we will add their dependencies and the primitives to a slice
|
||||
pendingBinds := make([]telemetry.FieldID, telemetry.MaxFields)
|
||||
|
||||
for _, id := range requestFields {
|
||||
for winID, id := range requestFields {
|
||||
b.data.Values[id].IDs = append(b.data.Values[id].IDs, winID)
|
||||
|
||||
switch id {
|
||||
case telemetry.RPMStateColour:
|
||||
b.data.VirtualBinds = append(b.data.VirtualBinds, telemetry.NewRPMLights())
|
||||
@@ -162,6 +164,7 @@ func (b *BeamNG) Subscribe(requestFields []telemetry.FieldID) {
|
||||
|
||||
func (b *BeamNG) readData() {
|
||||
slog.Debug("READING THIS DATA")
|
||||
// BUG: getting stuck in here
|
||||
ogSnapshot, err := b.SDK.Update()
|
||||
slog.Debug("THE DATA WAS READ")
|
||||
if err != nil {
|
||||
@@ -192,15 +195,10 @@ func (b *BeamNG) readData() {
|
||||
b.data.LastDataPoll = time.Now()
|
||||
}
|
||||
|
||||
func (b *BeamNG) stream(ctx context.Context) <-chan telemetry.TelemetryData {
|
||||
func (b *BeamNG) stream(ctx context.Context) {
|
||||
b.data.InitialTime = time.Now()
|
||||
outCh := make(chan telemetry.TelemetryData)
|
||||
b.wg.Add(1)
|
||||
|
||||
go func() {
|
||||
defer b.wg.Done()
|
||||
defer close(outCh)
|
||||
|
||||
for {
|
||||
// Explicitly intercept cancellation
|
||||
select {
|
||||
@@ -209,6 +207,8 @@ func (b *BeamNG) stream(ctx context.Context) <-chan telemetry.TelemetryData {
|
||||
default:
|
||||
}
|
||||
|
||||
// NOTE: add a method to check if there's data available, or make this happen
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
@@ -219,7 +219,7 @@ func (b *BeamNG) stream(ctx context.Context) <-chan telemetry.TelemetryData {
|
||||
|
||||
// Publish data
|
||||
select {
|
||||
case outCh <- *b.data:
|
||||
case b.streamCh <- *b.data:
|
||||
slog.Debug("PUBLISHED DATA")
|
||||
default:
|
||||
// skip this data, don't allow publishers to lag behind
|
||||
@@ -227,6 +227,4 @@ func (b *BeamNG) stream(ctx context.Context) <-chan telemetry.TelemetryData {
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return outCh
|
||||
}
|
||||
|
||||
@@ -10,14 +10,13 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"esdi/constants"
|
||||
"esdi/telemetry"
|
||||
|
||||
"github.com/ESilva15/goirsdk"
|
||||
)
|
||||
|
||||
const (
|
||||
NAME = constants.IRacingProviderName
|
||||
NAME = "iRacing"
|
||||
)
|
||||
|
||||
// IRacing is our iRacing telemetry data provider - its a TelemetryProvider interface
|
||||
@@ -34,8 +33,7 @@ type IRacing struct {
|
||||
ticker *time.Ticker // ticker will keep polling intervals constant
|
||||
|
||||
// Stream
|
||||
wg sync.WaitGroup
|
||||
// streamCh chan telemetry.TelemetryData
|
||||
streamCh chan telemetry.TelemetryData
|
||||
streamCancel context.CancelFunc
|
||||
}
|
||||
|
||||
@@ -52,11 +50,13 @@ func NewIRacingProvider(
|
||||
}
|
||||
|
||||
provider := &IRacing{
|
||||
logger: logger,
|
||||
SDK: sdk,
|
||||
data: telemetry.NewTelemetryData(),
|
||||
logger: logger,
|
||||
SDK: sdk,
|
||||
data: telemetry.NewTelemetryData(),
|
||||
streamCh: make(chan telemetry.TelemetryData, 1),
|
||||
// NOTE: This is because I stupidly recorded a test IBT file in 240
|
||||
// TODO: make this configurable from the user side
|
||||
ticker: time.NewTicker(time.Second / 60),
|
||||
ticker: time.NewTicker(time.Second / 240),
|
||||
}
|
||||
|
||||
provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){
|
||||
@@ -129,14 +129,11 @@ func (i *IRacing) isDataAvailable() bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func (i *IRacing) stream(ctx context.Context) <-chan telemetry.TelemetryData {
|
||||
func (i *IRacing) stream(ctx context.Context) {
|
||||
i.data.InitialTime = time.Now()
|
||||
outCh := make(chan telemetry.TelemetryData)
|
||||
i.wg.Add(1)
|
||||
|
||||
go func() {
|
||||
defer i.wg.Done()
|
||||
defer close(outCh)
|
||||
defer close(i.streamCh)
|
||||
|
||||
// Put this into the configuration file
|
||||
consecutiveTimeouts := 0
|
||||
@@ -152,11 +149,12 @@ func (i *IRacing) stream(ctx context.Context) <-chan telemetry.TelemetryData {
|
||||
|
||||
if i.SDK.CheckForDataEvent(time.Duration(dataEvTimeout) * time.Millisecond) {
|
||||
consecutiveTimeouts = 0
|
||||
i.logger.Debug("sending data", "timeouts", consecutiveTimeouts)
|
||||
i.readData()
|
||||
|
||||
// Publish data
|
||||
select {
|
||||
case outCh <- *i.data:
|
||||
case i.streamCh <- *i.data:
|
||||
default:
|
||||
// skip this data, don't allow publishers to lag behind
|
||||
}
|
||||
@@ -172,8 +170,6 @@ func (i *IRacing) stream(ctx context.Context) <-chan telemetry.TelemetryData {
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return outCh
|
||||
}
|
||||
|
||||
func (i *IRacing) readData() {
|
||||
@@ -210,9 +206,9 @@ func (i *IRacing) Stream() (<-chan telemetry.TelemetryData, error) {
|
||||
ctx, i.streamCancel = context.WithCancel(context.Background())
|
||||
|
||||
// Start the stream
|
||||
ch := i.stream(ctx)
|
||||
i.stream(ctx)
|
||||
|
||||
return ch, nil
|
||||
return i.streamCh, nil
|
||||
}
|
||||
|
||||
func (i *IRacing) StopStream() {
|
||||
@@ -221,11 +217,10 @@ func (i *IRacing) StopStream() {
|
||||
}
|
||||
|
||||
i.streamCancel()
|
||||
i.wg.Wait()
|
||||
i.streamCancel = nil
|
||||
}
|
||||
|
||||
func (i *IRacing) Subscribe(requestFields []telemetry.FieldID) {
|
||||
func (i *IRacing) Subscribe(requestFields map[int16]telemetry.FieldID) {
|
||||
i.logger.Debug(fmt.Sprintf("Len Req: %d\n", len(requestFields)))
|
||||
|
||||
i.data.ActiveBinds = make([]telemetry.BoundField, 0, len(requestFields))
|
||||
@@ -234,7 +229,9 @@ func (i *IRacing) Subscribe(requestFields []telemetry.FieldID) {
|
||||
// we will add their dependencies and the primitives to a slice
|
||||
pendingBinds := make([]telemetry.FieldID, 0, telemetry.MaxFields)
|
||||
|
||||
for _, id := range requestFields {
|
||||
for winID, id := range requestFields {
|
||||
i.data.Values[id].IDs = append(i.data.Values[id].IDs, winID)
|
||||
|
||||
switch id {
|
||||
case telemetry.RPMStateColour:
|
||||
i.data.VirtualBinds = append(i.data.VirtualBinds, telemetry.NewRPMLights())
|
||||
|
||||
+81
-53
@@ -2,38 +2,46 @@ package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"esdi/devices"
|
||||
"esdi/peripheral"
|
||||
"esdi/telemetry"
|
||||
)
|
||||
|
||||
// DeviceService is the API for the peripherals
|
||||
var ErrPeripheralAlreadyRegistered = errors.New("peripheral is already registered")
|
||||
|
||||
// DeviceService will handle sending the data from the telemetry service to the
|
||||
// actual devices
|
||||
// NOTE: create a virtual device and make it be the output window or something so
|
||||
// we can just add it as a device or whatever instead of being a custom made thing
|
||||
// that would be pretty cool I think
|
||||
type DeviceService struct {
|
||||
Logger *slog.Logger
|
||||
// Device discovery
|
||||
PSS *PeripheralStateStore // Store to track peripheral state
|
||||
mu sync.RWMutex
|
||||
ctxDiscovery context.Context
|
||||
ctxDiscoveryCancel context.CancelFunc
|
||||
Devices map[string]peripheral.Peripheral
|
||||
// Strem handling
|
||||
streamCancel context.CancelFunc
|
||||
TelemCh <-chan telemetry.TelemetryData
|
||||
// Output
|
||||
Messages chan string
|
||||
// Callbacks
|
||||
// Telemetry service data fetchers
|
||||
telemetryProvider func() (string, error)
|
||||
}
|
||||
|
||||
func NewDeviceService(logger *slog.Logger, msg chan string) *DeviceService {
|
||||
func NewDeviceService(logger *slog.Logger) *DeviceService {
|
||||
sharedChannel := make(chan string, 10)
|
||||
|
||||
dev := &DeviceService{
|
||||
PSS: NewPeripheralStateStore(
|
||||
logger.With("Service", "PeripheralStateStore"), devices.List, msg,
|
||||
),
|
||||
Devices: make(map[string]peripheral.Peripheral),
|
||||
Logger: logger,
|
||||
Messages: msg,
|
||||
Messages: sharedChannel,
|
||||
}
|
||||
|
||||
// Start the routine that looks for devices - should always be running in the background
|
||||
@@ -41,52 +49,84 @@ func NewDeviceService(logger *slog.Logger, msg chan string) *DeviceService {
|
||||
dev.ctxDiscovery, dev.ctxDiscoveryCancel = context.WithCancel(context.Background())
|
||||
go dev.FindDevices()
|
||||
|
||||
// Set the callbacks for PSS
|
||||
dev.PSS.telemetryProvider = dev.getTelemetryProvider
|
||||
|
||||
return dev
|
||||
}
|
||||
|
||||
// Getters [START] -------------------------------------------------------------
|
||||
// This function is currently only being used by PSS, we may have to find a better
|
||||
// pattern for this
|
||||
func (ds *DeviceService) getTelemetryProvider() (string, error) {
|
||||
return ds.telemetryProvider()
|
||||
func (ds *DeviceService) FindDevices() {
|
||||
// Need to define a list of devices to search for
|
||||
// For now lets just try to find our cdashdisplay - will think about the rest later
|
||||
ticker := time.NewTicker(2 * time.Second)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ds.ctxDiscovery.Done():
|
||||
// If requested to cancel we cancel background discovery
|
||||
return
|
||||
case <-ticker.C:
|
||||
for pName, peripheral := range devices.List {
|
||||
if ds.DeviceExists(pName) {
|
||||
// We already discovered this device
|
||||
continue
|
||||
}
|
||||
|
||||
ds.Logger.Debug("looking for device", "name", pName)
|
||||
dev, err := peripheral.Discover()
|
||||
if err != nil {
|
||||
ds.Logger.Debug("didn't find device", "name", pName)
|
||||
continue
|
||||
}
|
||||
|
||||
// Register the device we just found
|
||||
ds.RegisterDevice(dev)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (ds *DeviceService) GetDevices() []peripheral.Peripheral {
|
||||
snapshot := ds.PSS.GetStates()
|
||||
peripherals := make([]peripheral.Peripheral, 0, len(snapshot))
|
||||
// func (ds *DeviceService) SubscribeFields() error {
|
||||
// for _, dev := range ds.Devices {
|
||||
// fields := dev.RequiredFields()
|
||||
// }
|
||||
//
|
||||
// return nil
|
||||
// }
|
||||
|
||||
for _, state := range snapshot {
|
||||
if state.State < DeviceIsConnected {
|
||||
continue
|
||||
}
|
||||
peripherals = append(peripherals, state.Peripheral)
|
||||
func (ds *DeviceService) RegisterDevice(dev peripheral.Peripheral) error {
|
||||
ds.mu.Lock()
|
||||
defer ds.mu.Unlock()
|
||||
|
||||
if ds.DeviceExists(dev.Name()) {
|
||||
return ErrPeripheralAlreadyRegistered
|
||||
}
|
||||
|
||||
return peripherals
|
||||
ds.Devices[dev.Name()] = dev
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (ds *DeviceService) GetPeripheral(pname string) (peripheral.Peripheral, error) {
|
||||
return ds.PSS.GetPeripheral(pname)
|
||||
func (ds *DeviceService) GetDevice(name string) (peripheral.Peripheral, error) {
|
||||
val, ok := ds.Devices[name]
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("device `%s` couldn't be found", name)
|
||||
}
|
||||
|
||||
return val, nil
|
||||
}
|
||||
|
||||
func (ds *DeviceService) PeripheralExists(pname string) bool {
|
||||
_, err := ds.PSS.GetPeripheral(pname)
|
||||
return err == nil
|
||||
func (ds *DeviceService) DeviceExists(name string) bool {
|
||||
if _, ok := ds.Devices[name]; !ok {
|
||||
return false
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
// Getters [END] ---------------------------------------------------------------
|
||||
|
||||
// Actions [START] -------------------------------------------------------------
|
||||
func (ds *DeviceService) StartStream() {
|
||||
// NOTE: i'm using this pattern a whole lot. Maybe I can create a struct to handle this
|
||||
var ctx context.Context
|
||||
ctx, ds.streamCancel = context.WithCancel(context.Background())
|
||||
|
||||
ds.PSS.OnStartStream()
|
||||
|
||||
go ds.transmit(ctx)
|
||||
}
|
||||
|
||||
@@ -104,8 +144,6 @@ func (ds *DeviceService) SetTelemetryChannel(ch <-chan telemetry.TelemetryData)
|
||||
ds.TelemCh = ch
|
||||
}
|
||||
|
||||
// Actions [END] ---------------------------------------------------------------
|
||||
|
||||
// transmit will send the data to the devices themselves
|
||||
func (ds *DeviceService) transmit(ctx context.Context) {
|
||||
var isSending atomic.Bool
|
||||
@@ -127,23 +165,13 @@ func (ds *DeviceService) transmit(ctx context.Context) {
|
||||
|
||||
// TODO: make a copy of the data and send that copy instead of keeping
|
||||
// the data locked
|
||||
for _, dev := range ds.PSS.GetStates() {
|
||||
if dev.State != DeviceIsStreaming {
|
||||
continue
|
||||
}
|
||||
|
||||
err := dev.Peripheral.SendData(&data)
|
||||
|
||||
if err == peripheral.ErrDeviceTimedOut {
|
||||
ds.onDeviceTimedOut(dev.device.Name)
|
||||
}
|
||||
ds.mu.RLock()
|
||||
for _, dev := range ds.Devices {
|
||||
dev.SendData(&data)
|
||||
}
|
||||
ds.mu.RUnlock()
|
||||
|
||||
isSending.Store(false)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Callbacks [START] -----------------------------------------------------------
|
||||
|
||||
// Callbacks [END] -------------------------------------------------------------
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
package services
|
||||
@@ -1,480 +0,0 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"maps"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"esdi/devices"
|
||||
"esdi/peripheral"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrPeripheralAlreadyRegistered = errors.New("peripheral is already registered")
|
||||
ErrNoSuchDevice = errors.New("device doesn't exist")
|
||||
ErrDeviceIsNotConnected = errors.New("device isn't connected")
|
||||
ErrFailedToSetupPeripheral = errors.New("peripheral setup failed")
|
||||
)
|
||||
|
||||
type DeviceState = uint8
|
||||
|
||||
const (
|
||||
DeviceTimedOut uint8 = iota
|
||||
DeviceIsDisconnected
|
||||
DeviceReconnected
|
||||
DeviceIsDiscovering
|
||||
DeviceIsConnected
|
||||
DeviceIsUnconfigured
|
||||
DeviceIsConfiguring
|
||||
DeviceIsConfigured
|
||||
DeviceIsStreaming
|
||||
)
|
||||
|
||||
func DeviceStateToStr(state DeviceState) string {
|
||||
switch state {
|
||||
case DeviceTimedOut:
|
||||
return "DeviceTimedOut"
|
||||
case DeviceIsDisconnected:
|
||||
return "DeviceIsDisconnected"
|
||||
case DeviceReconnected:
|
||||
return "DeviceReconnected"
|
||||
case DeviceIsDiscovering:
|
||||
return "DeviceIsDiscovering"
|
||||
case DeviceIsConnected:
|
||||
return "DeviceIsConnected"
|
||||
case DeviceIsUnconfigured:
|
||||
return "DeviceIsUnconfigured"
|
||||
case DeviceIsConfiguring:
|
||||
return "DeviceIsConfiguring"
|
||||
case DeviceIsConfigured:
|
||||
return "DeviceIsConfigured"
|
||||
case DeviceIsStreaming:
|
||||
return "DeviceIsStreaming"
|
||||
default:
|
||||
return "UnknownState"
|
||||
}
|
||||
}
|
||||
|
||||
type PeripheralState struct {
|
||||
device *devices.Device
|
||||
Peripheral peripheral.Peripheral
|
||||
State DeviceState
|
||||
}
|
||||
|
||||
func NewPeripheralState(
|
||||
dev *devices.Device,
|
||||
peripheral peripheral.Peripheral,
|
||||
state DeviceState,
|
||||
) *PeripheralState {
|
||||
perState := PeripheralState{
|
||||
device: dev,
|
||||
Peripheral: peripheral,
|
||||
State: state,
|
||||
}
|
||||
|
||||
return &perState
|
||||
}
|
||||
|
||||
type PeripheralStateStore struct {
|
||||
Logger *slog.Logger
|
||||
mu sync.RWMutex
|
||||
store map[string]*PeripheralState
|
||||
// Messaging for UI and stuff
|
||||
Messages chan string
|
||||
// Callbacks
|
||||
telemetryProvider func() (string, error)
|
||||
// Internal State
|
||||
isStreaming bool
|
||||
}
|
||||
|
||||
func NewPeripheralStateStore(
|
||||
nLogger *slog.Logger,
|
||||
devList map[string]*devices.Device,
|
||||
msg chan string,
|
||||
) *PeripheralStateStore {
|
||||
store := PeripheralStateStore{
|
||||
Logger: nLogger,
|
||||
store: make(map[string]*PeripheralState),
|
||||
Messages: msg,
|
||||
}
|
||||
|
||||
for _, dev := range devList {
|
||||
store.AddDevice(dev)
|
||||
}
|
||||
|
||||
return &store
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) GetStates() map[string]*PeripheralState {
|
||||
pss.mu.RLock()
|
||||
defer pss.mu.RUnlock()
|
||||
|
||||
return maps.Clone(pss.store)
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) GetState(pname string) (*PeripheralState, error) {
|
||||
if !pss.DeviceExists(pname) {
|
||||
return nil, ErrNoSuchDevice
|
||||
}
|
||||
|
||||
pss.mu.RLock()
|
||||
defer pss.mu.RUnlock()
|
||||
return pss.store[pname], nil
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) GetPeripheral(pname string) (peripheral.Peripheral, error) {
|
||||
state, err := pss.GetState(pname)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if state.State < DeviceIsConnected {
|
||||
return nil, ErrDeviceIsNotConnected
|
||||
}
|
||||
|
||||
return state.Peripheral, nil
|
||||
}
|
||||
|
||||
// AddDevice adds a new device for tracking
|
||||
func (pss *PeripheralStateStore) AddDevice(dev *devices.Device) error {
|
||||
if pss.DeviceExists(dev.Name) {
|
||||
return ErrPeripheralAlreadyRegistered
|
||||
}
|
||||
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
pss.store[dev.Name] = NewPeripheralState(dev, nil, DeviceIsDisconnected)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// DeviceExists returns whether the store is already tracking `pname`
|
||||
func (pss *PeripheralStateStore) DeviceExists(pname string) bool {
|
||||
pss.mu.RLock()
|
||||
defer pss.mu.RUnlock()
|
||||
|
||||
if _, ok := pss.store[pname]; ok {
|
||||
return true
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
// DeleteDevice deletes `pname` from tracking
|
||||
func (pss *PeripheralStateStore) DeleteDevice(pname string) error {
|
||||
if !pss.DeviceExists(pname) {
|
||||
return ErrNoSuchDevice
|
||||
}
|
||||
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
|
||||
delete(pss.store, pname)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) GetStreamingState() bool {
|
||||
pss.mu.RLock()
|
||||
defer pss.mu.RUnlock()
|
||||
return pss.isStreaming
|
||||
}
|
||||
|
||||
// "Events" [START] ------------------------------------------------------------
|
||||
func (ds *DeviceService) onDeviceTimedOut(pname string) {
|
||||
ds.PSS.setDeviceTimedOut(pname)
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) OnStartStream() {
|
||||
pss.mu.Lock()
|
||||
pss.isStreaming = true
|
||||
pss.mu.Unlock()
|
||||
|
||||
for _, state := range pss.GetStates() {
|
||||
if state.State == DeviceIsConfigured {
|
||||
pss.setDeviceIsStreaming(state.device.Name)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) OnStopStream() {
|
||||
pss.mu.Lock()
|
||||
pss.isStreaming = false
|
||||
pss.mu.Unlock()
|
||||
}
|
||||
|
||||
// "Events" [END] --------------------------------------------------------------
|
||||
|
||||
// Device State Handling [START] -----------------------------------------------
|
||||
func (pss *PeripheralStateStore) setDeviceDisconnected(pname string) {
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
pss.store[pname].State = DeviceIsDisconnected
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) setDeviceIsDiscovering(pname string) {
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
pss.store[pname].State = DeviceIsDiscovering
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) setDeviceConnected(pname string, per peripheral.Peripheral) {
|
||||
pss.Logger.Info("found device", "device", pname)
|
||||
|
||||
pss.mu.Lock()
|
||||
pss.store[pname].Peripheral = per
|
||||
pss.store[pname].State = DeviceIsConnected
|
||||
pss.mu.Unlock()
|
||||
|
||||
pss.Messages <- fmt.Sprintf("Device successfuly connected: %s\n", pname)
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) setDeviceTimedOut(pname string) {
|
||||
pss.Logger.Info("device timed out", "device", pname)
|
||||
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
pss.store[pname].Peripheral = nil
|
||||
pss.store[pname].State = DeviceTimedOut
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) setDeviceReconnected(pname string, per peripheral.Peripheral) {
|
||||
pss.Logger.Info("device reconnecting", "device", pname)
|
||||
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
pss.store[pname].Peripheral = per
|
||||
pss.store[pname].State = DeviceReconnected
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) setDeviceUnconfigured(pname string) {
|
||||
pss.Logger.Info("device is connected but not configured", "device", pname)
|
||||
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
pss.store[pname].State = DeviceIsUnconfigured
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) setDeviceIsConfiguring(pname string) {
|
||||
pss.Logger.Info("device is configuring for new provider", "device", pname)
|
||||
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
pss.store[pname].State = DeviceIsConfiguring
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) setDeviceConfigured(pname string) {
|
||||
pss.Logger.Info("device is configured and ready for data", "device", pname)
|
||||
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
pss.store[pname].State = DeviceIsConfigured
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) setDeviceIsStreaming(pname string) {
|
||||
pss.Logger.Info("device is configured and ready for data", "device", pname)
|
||||
|
||||
pss.mu.Lock()
|
||||
defer pss.mu.Unlock()
|
||||
pss.store[pname].State = DeviceIsStreaming
|
||||
}
|
||||
|
||||
// Device State Handling [END] -------------------------------------------------
|
||||
|
||||
// Device Handling [START] -----------------------------------------------------
|
||||
func (pss *PeripheralStateStore) discoverPeripheral(
|
||||
pname string,
|
||||
onDiscovery func(string, peripheral.Peripheral),
|
||||
onFailure func(string),
|
||||
) {
|
||||
pss.Logger.Debug("looking for device", "name", pname)
|
||||
pss.setDeviceIsDiscovering(pname)
|
||||
pss.Messages <- "Discovering " + pname + "\n"
|
||||
|
||||
go func() {
|
||||
state, err := pss.GetState(pname)
|
||||
if err != nil {
|
||||
pss.Logger.Error("Can't reconnect device", "device", pname, "error", err)
|
||||
onFailure(pname)
|
||||
return
|
||||
}
|
||||
|
||||
dev, err := state.device.Discover()
|
||||
if err != nil {
|
||||
onFailure(pname)
|
||||
return
|
||||
}
|
||||
|
||||
// Register the device we just found
|
||||
onDiscovery(pname, dev)
|
||||
}()
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) configurePeripheral(
|
||||
pname string,
|
||||
onSuccess func(string),
|
||||
onFailure func(string),
|
||||
) {
|
||||
go func() {
|
||||
state, err := pss.GetState(pname)
|
||||
if err != nil {
|
||||
onFailure(pname)
|
||||
return
|
||||
}
|
||||
|
||||
provider, err := pss.telemetryProvider()
|
||||
if err != nil {
|
||||
onFailure(pname)
|
||||
return
|
||||
}
|
||||
|
||||
err = state.Peripheral.Setup(provider)
|
||||
if err != nil {
|
||||
onFailure(pname)
|
||||
return
|
||||
}
|
||||
|
||||
onSuccess(pname)
|
||||
pss.setDeviceConfigured(pname)
|
||||
}()
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) handleDeviceTimedOut(pname string) error {
|
||||
pss.Messages <- "device " + pname + " timed out\n"
|
||||
pss.discoverPeripheral(pname, pss.setDeviceReconnected, pss.setDeviceTimedOut)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) handleDeviceReconnected(pname string) error {
|
||||
// The device has reconnected, but we must set its state again
|
||||
state, err := pss.GetState(pname)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Update the peripheral state
|
||||
pss.setDeviceConnected(pname, state.Peripheral)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// handleDeviceConnected will handle the device setup after it connects
|
||||
// NOTE: should this be a state after Connected?
|
||||
// Connected -> Unconfigured -> Configured I believe this would work nicely
|
||||
// THIS IS A TODO ↑↑↑↑↑↑
|
||||
func (pss *PeripheralStateStore) handleDeviceConnected(pname string) error {
|
||||
pss.setDeviceUnconfigured(pname)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) handleDeviceIsUnconfigured(pname string) error {
|
||||
// Here we need to configure our device. If no error occurs its configured!
|
||||
_, err := pss.telemetryProvider()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
pss.setDeviceIsConfiguring(pname)
|
||||
pss.configurePeripheral(pname, pss.setDeviceConfigured, pss.setDeviceUnconfigured)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) handleDeviceIsConfigured(pname string) error {
|
||||
// Here we have to check wheter we are streaming or not. If we aren't streaming
|
||||
// then we ought to do a healthcheck on the peripheral
|
||||
if pss.GetStreamingState() {
|
||||
pss.Messages <- "returning device " + pname + " into streaming\n"
|
||||
pss.setDeviceIsStreaming(pname)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) performHealthCheck(pname string, state *PeripheralState) bool {
|
||||
healthStatus := state.Peripheral.HealthCheck()
|
||||
if !healthStatus {
|
||||
pss.Messages <- "peripheral " + pname + " failed healthcheck"
|
||||
pss.setDeviceTimedOut(pname)
|
||||
return false
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
func (pss *PeripheralStateStore) HandleDeviceState() {
|
||||
peripherals := pss.GetStates()
|
||||
|
||||
for pName, pState := range peripherals {
|
||||
switch pState.State {
|
||||
case DeviceIsDisconnected:
|
||||
pss.discoverPeripheral(pName, pss.setDeviceConnected, pss.setDeviceDisconnected)
|
||||
case DeviceIsDiscovering:
|
||||
// We need to set a device into discovery mode so we won't retrigger discoveries
|
||||
// and pool them up
|
||||
case DeviceIsConnected:
|
||||
// Need to check if its streaming, if its not streaming than we have to do a healthcheck
|
||||
pss.Logger.Debug("Device is connected. Normal", "device", pName)
|
||||
pss.handleDeviceConnected(pName)
|
||||
case DeviceIsUnconfigured:
|
||||
pss.Logger.Debug("Device is still being configured.", "device", pName)
|
||||
pss.handleDeviceIsUnconfigured(pName)
|
||||
case DeviceIsConfiguring:
|
||||
// Do nothing configuration is happening in the background
|
||||
case DeviceIsConfigured:
|
||||
// Nothing to do here
|
||||
pss.handleDeviceIsConfigured(pName)
|
||||
case DeviceIsStreaming:
|
||||
//
|
||||
case DeviceReconnected:
|
||||
// If the device has reconnected we need to reset the device and then set it as connected
|
||||
pss.Logger.Debug("Device has reconnected. Clearing up state", "device", pName)
|
||||
err := pss.handleDeviceReconnected(pName)
|
||||
if err != nil {
|
||||
pss.Logger.Error("device reconnection handler failed", "error", err)
|
||||
continue
|
||||
}
|
||||
case DeviceTimedOut:
|
||||
// If the device has timed out we need to re-discover it or something
|
||||
pss.Logger.Debug("Device is timed out. Attempting to recconect", "device", pName)
|
||||
err := pss.handleDeviceTimedOut(pName)
|
||||
if err != nil {
|
||||
pss.Logger.Error("device timing out handler failed", "error", err)
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
updatedState, err := pss.GetState(pName)
|
||||
if err != nil {
|
||||
// TODO: log something useful here
|
||||
continue
|
||||
}
|
||||
if updatedState.State == DeviceIsConnected ||
|
||||
updatedState.State == DeviceIsUnconfigured ||
|
||||
updatedState.State == DeviceIsConfigured {
|
||||
pss.performHealthCheck(pName, updatedState)
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
// Device Handling [END] -------------------------------------------------------
|
||||
|
||||
// FindDevices is a routine that goes over the devices in the PeripheralStateStore
|
||||
// and handles their state accordingly
|
||||
func (ds *DeviceService) FindDevices() {
|
||||
ticker := time.NewTicker(1 * time.Second)
|
||||
defer ticker.Stop()
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ds.ctxDiscovery.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
ds.PSS.HandleDeviceState()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,38 +1,2 @@
|
||||
// Package services interacts with the other libraries required for this UI
|
||||
package services
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"log/slog"
|
||||
)
|
||||
|
||||
type Orchestrator struct {
|
||||
DeviceService *DeviceService
|
||||
TelemetryService *TelemetryService
|
||||
Messages chan string
|
||||
}
|
||||
|
||||
func NewOrchestrator(logger *slog.Logger) (*Orchestrator, error) {
|
||||
msg := make(chan string, 10)
|
||||
|
||||
devService := NewDeviceService(logger.With("service", "DeviceService"), msg)
|
||||
|
||||
telemService := NewTelemetryService(logger.With("service", "TelemetryService"), msg)
|
||||
if telemService == nil {
|
||||
return nil, errors.New("failed to create telemetry service")
|
||||
}
|
||||
|
||||
go telemService.FindProvider(telemService.CtxMonitor)
|
||||
|
||||
// Setup device service callbacks
|
||||
devService.telemetryProvider = telemService.GetTelemetryProviderName
|
||||
|
||||
// Setup telemetry service callbacks
|
||||
telemService.peripheralProvider = devService.GetDevices
|
||||
|
||||
return &Orchestrator{
|
||||
DeviceService: devService,
|
||||
TelemetryService: telemService,
|
||||
Messages: msg,
|
||||
}, nil
|
||||
}
|
||||
|
||||
+120
-142
@@ -2,25 +2,19 @@ package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"log/slog"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"esdi/peripheral"
|
||||
"esdi/providers"
|
||||
"esdi/telemetry"
|
||||
telem "esdi/telemetry"
|
||||
)
|
||||
|
||||
var ErrNoActiveProviderAvailable = errors.New("no active provider available")
|
||||
|
||||
// TelemetryService will be our base struct to handle telemetry data
|
||||
// It should hook to a data sink and handle it like iRacing, BeamNG, AC and so on
|
||||
type TelemetryService struct {
|
||||
logger *slog.Logger
|
||||
// Streaming
|
||||
isStreaming bool
|
||||
logger *slog.Logger
|
||||
devService *DeviceService
|
||||
// Concurrency protection
|
||||
mut sync.RWMutex
|
||||
activeProvider telem.TelemetryProvider
|
||||
@@ -35,20 +29,16 @@ type TelemetryService struct {
|
||||
cancelMonitor context.CancelFunc
|
||||
CtxHealthcheck context.Context
|
||||
healthCheckCancel context.CancelFunc
|
||||
// Callbacks
|
||||
// Devices data request
|
||||
peripheralProvider func() []peripheral.Peripheral
|
||||
}
|
||||
|
||||
func NewTelemetryService(
|
||||
logger *slog.Logger,
|
||||
msg chan string,
|
||||
) *TelemetryService {
|
||||
func NewTelemetryService(logger *slog.Logger, devServo *DeviceService) *TelemetryService {
|
||||
sharedChannel := make(chan string, 10)
|
||||
newService := &TelemetryService{
|
||||
logger: logger,
|
||||
isConnected: false,
|
||||
devService: devServo,
|
||||
listeners: make(map[string]chan telem.TelemetryData),
|
||||
Messages: msg,
|
||||
Messages: sharedChannel,
|
||||
}
|
||||
newService.CtxMonitor, newService.cancelMonitor = context.WithCancel(context.Background())
|
||||
|
||||
@@ -66,7 +56,6 @@ func (t *TelemetryService) ProviderMonitor(ctx context.Context) {
|
||||
case <-ticker.C:
|
||||
slog.Info("checking if provider is still running")
|
||||
if !t.activeProvider.IsAlive(500 * time.Millisecond) {
|
||||
t.Messages <- "Healthcheck on provider failing. Dropping provider.\n"
|
||||
slog.Warn("provider healthcheck failed")
|
||||
t.dropActiveProvider()
|
||||
t.onProviderHealthCheckFailed()
|
||||
@@ -76,98 +65,32 @@ func (t *TelemetryService) ProviderMonitor(ctx context.Context) {
|
||||
}
|
||||
}
|
||||
|
||||
func (t *TelemetryService) GetTelemetryProviderName() (string, error) {
|
||||
if t.activeProvider == nil {
|
||||
return "", ErrNoActiveProviderAvailable
|
||||
}
|
||||
|
||||
return t.activeProvider.Name(), nil
|
||||
func (t *TelemetryService) onProviderHealthCheckFailed() {
|
||||
// Just restart the whole lookup process
|
||||
go t.FindProvider(t.CtxMonitor)
|
||||
}
|
||||
|
||||
// Listener Control [START] ----------------------------------------------------
|
||||
|
||||
func (t *TelemetryService) SubscribeListener(id string, bufferSize int) <-chan telem.TelemetryData {
|
||||
t.mut.Lock()
|
||||
defer t.mut.Unlock()
|
||||
|
||||
// NOTE: is this truly necessary?
|
||||
// return the channel if it already exists
|
||||
if ch, exists := t.listeners[id]; exists {
|
||||
return ch
|
||||
func (t *TelemetryService) onFindProvider(prov telem.TelemetryProvider) {
|
||||
// Attach to the provider
|
||||
t.logger.Info("found provider for " + prov.Name())
|
||||
err := t.SwitchProvider(prov)
|
||||
if err != nil {
|
||||
t.logger.Error("failed to switch to provider onFindProvider", "err", err)
|
||||
return
|
||||
}
|
||||
|
||||
ch := make(chan telem.TelemetryData, bufferSize)
|
||||
t.listeners[id] = ch
|
||||
|
||||
t.logger.Info("New stream subscriber registered", "id", id)
|
||||
return ch
|
||||
// Create a routine to poll this provider while we wait to start the stream or pause it
|
||||
t.CtxHealthcheck, t.healthCheckCancel = context.WithCancel(context.Background())
|
||||
go t.ProviderMonitor(t.CtxHealthcheck)
|
||||
}
|
||||
|
||||
func (t *TelemetryService) UnsubscribeListener(id string) {
|
||||
t.mut.Lock()
|
||||
defer t.mut.Unlock()
|
||||
|
||||
if ch, exists := t.listeners[id]; exists {
|
||||
close(ch)
|
||||
delete(t.listeners, id)
|
||||
t.logger.Info("Stream subscriber removed", "id", id)
|
||||
}
|
||||
}
|
||||
|
||||
func (t *TelemetryService) SubscribeToFields() []telem.FieldID {
|
||||
seen := make(map[telemetry.FieldID]struct{})
|
||||
var allFields []telemetry.FieldID
|
||||
|
||||
for _, dev := range t.peripheralProvider() {
|
||||
for _, field := range dev.RequiredFields() {
|
||||
if _, exists := seen[field]; !exists {
|
||||
seen[field] = struct{}{}
|
||||
allFields = append(allFields, field)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
t.logger.Debug("requested fields", "fields", allFields)
|
||||
t.activeProvider.Subscribe(allFields)
|
||||
|
||||
return allFields
|
||||
}
|
||||
|
||||
// Listener Control [END] ------------------------------------------------------
|
||||
|
||||
// Provider Control [START] ----------------------------------------------------
|
||||
|
||||
func (t *TelemetryService) HasActiveProvider() bool {
|
||||
if t.activeProvider == nil {
|
||||
return false
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
func (t *TelemetryService) dropActiveProvider() {
|
||||
if t.cancelForward != nil {
|
||||
t.cancelForward()
|
||||
}
|
||||
|
||||
t.activeProvider.StopStream()
|
||||
t.activeProvider.Close()
|
||||
t.activeProvider = nil
|
||||
}
|
||||
|
||||
func (t *TelemetryService) SwitchProvider(newProvider telem.TelemetryProvider) error {
|
||||
t.mut.Lock()
|
||||
defer t.mut.Unlock()
|
||||
|
||||
// Clean up the current to be old provider
|
||||
if t.activeProvider != nil {
|
||||
t.dropActiveProvider()
|
||||
}
|
||||
|
||||
// Assign the new provider
|
||||
t.activeProvider = newProvider
|
||||
|
||||
return nil
|
||||
func (t *TelemetryService) onProviderStopsMidStream() {
|
||||
// clear the current provider
|
||||
// TODO: now we need to also clear the devices to restart everything,
|
||||
// if the stream stopped we have to restart the devices and everything
|
||||
t.logger.Info("cleaning dropped provider and restarting lookup service")
|
||||
t.dropActiveProvider()
|
||||
go t.FindProvider(t.CtxMonitor)
|
||||
}
|
||||
|
||||
// TODO: add some way of retriggering this. Currently it should:
|
||||
@@ -199,46 +122,19 @@ func (t *TelemetryService) FindProvider(ctx context.Context) {
|
||||
}
|
||||
}
|
||||
|
||||
// Provider Control [END] ------------------------------------------------------
|
||||
|
||||
// Streaming Control [START] ---------------------------------------------------
|
||||
|
||||
func (t *TelemetryService) StopStream() {
|
||||
func (t *TelemetryService) SwitchProvider(newProvider telem.TelemetryProvider) error {
|
||||
t.mut.Lock()
|
||||
defer t.mut.Unlock()
|
||||
|
||||
if t.cancelForward != nil {
|
||||
t.cancelForward()
|
||||
t.cancelForward = nil
|
||||
// Clean up the current to be old provider
|
||||
if t.activeProvider != nil {
|
||||
t.dropActiveProvider()
|
||||
}
|
||||
|
||||
t.activeProvider.StopStream()
|
||||
t.isStreaming = false
|
||||
}
|
||||
// Assign the new provider
|
||||
t.activeProvider = newProvider
|
||||
|
||||
func (t *TelemetryService) StartStream() {
|
||||
slog.Debug("Stream started")
|
||||
|
||||
// Start the new stream
|
||||
if t.activeProvider == nil {
|
||||
slog.Debug("there's no active provider. not starting the stream")
|
||||
return
|
||||
}
|
||||
|
||||
// Stop the provider healthcheck
|
||||
t.healthCheckCancel()
|
||||
|
||||
simInCh, _ := t.activeProvider.Stream()
|
||||
// TODO: the provider needs to be able to tell the data has stopped
|
||||
// so we can restart the provider lookup routine
|
||||
|
||||
// Create the context so we can control the lifecycle
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.cancelForward = cancel
|
||||
|
||||
// Multiplex this data
|
||||
go t.multiplexData(ctx, simInCh)
|
||||
t.isStreaming = true
|
||||
return nil
|
||||
}
|
||||
|
||||
func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan telem.TelemetryData) {
|
||||
@@ -268,12 +164,94 @@ func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan tele
|
||||
}
|
||||
}
|
||||
|
||||
func (t *TelemetryService) IsStreaming() bool {
|
||||
return t.isStreaming
|
||||
func (t *TelemetryService) dropActiveProvider() {
|
||||
if t.cancelForward != nil {
|
||||
t.cancelForward()
|
||||
}
|
||||
|
||||
t.activeProvider.StopStream()
|
||||
t.activeProvider.Close()
|
||||
t.activeProvider = nil
|
||||
}
|
||||
|
||||
// Streaming Control [END] -----------------------------------------------------
|
||||
func (t *TelemetryService) SubscribeListener(id string, bufferSize int) <-chan telem.TelemetryData {
|
||||
t.mut.Lock()
|
||||
defer t.mut.Unlock()
|
||||
|
||||
// Callbacks [START] -----------------------------------------------------------
|
||||
// NOTE: is this truly necessary?
|
||||
// return the channel if it already exists
|
||||
if ch, exists := t.listeners[id]; exists {
|
||||
return ch
|
||||
}
|
||||
|
||||
// Callbacks [END] -------------------------------------------------------------
|
||||
ch := make(chan telem.TelemetryData, bufferSize)
|
||||
t.listeners[id] = ch
|
||||
|
||||
t.logger.Info("New stream subscriber registered", "id", id)
|
||||
return ch
|
||||
}
|
||||
|
||||
func (t *TelemetryService) UnsubscribeListener(id string) {
|
||||
t.mut.Lock()
|
||||
defer t.mut.Unlock()
|
||||
|
||||
if ch, exists := t.listeners[id]; exists {
|
||||
close(ch)
|
||||
delete(t.listeners, id)
|
||||
t.logger.Info("Stream subscriber removed", "id", id)
|
||||
}
|
||||
}
|
||||
|
||||
func (t *TelemetryService) SubscribeToFields() {
|
||||
// _ = t.devService.SubscribeFields()
|
||||
// TODO: we need to find a way of requesting devices to send all subscribed fields
|
||||
// instead of going through the devices on DevService here
|
||||
for _, dev := range t.devService.Devices {
|
||||
fields := dev.RequiredFields()
|
||||
t.activeProvider.Subscribe(fields)
|
||||
}
|
||||
}
|
||||
|
||||
func (t *TelemetryService) StartStream() {
|
||||
slog.Debug("Stream started")
|
||||
|
||||
// Start the new stream
|
||||
if t.activeProvider == nil {
|
||||
slog.Debug("there's no active provider. not starting the stream")
|
||||
return
|
||||
}
|
||||
|
||||
// Stop the provider healthcheck
|
||||
t.healthCheckCancel()
|
||||
|
||||
simInCh, _ := t.activeProvider.Stream()
|
||||
// TODO: the provider needs to be able to tell the data has stopped
|
||||
// so we can restart the provider lookup routine
|
||||
|
||||
// Create the context so we can control the lifecycle
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
t.cancelForward = cancel
|
||||
|
||||
// Multiplex this data
|
||||
go t.multiplexData(ctx, simInCh)
|
||||
}
|
||||
|
||||
func (t *TelemetryService) StopStream() {
|
||||
t.mut.Lock()
|
||||
defer t.mut.Unlock()
|
||||
|
||||
if t.cancelForward != nil {
|
||||
t.cancelForward()
|
||||
t.cancelForward = nil
|
||||
}
|
||||
|
||||
t.activeProvider.StopStream()
|
||||
}
|
||||
|
||||
func (t *TelemetryService) HasActiveProvider() bool {
|
||||
if t.activeProvider == nil {
|
||||
return false
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
@@ -1,39 +0,0 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
telem "esdi/telemetry"
|
||||
)
|
||||
|
||||
func (t *TelemetryService) onProviderHealthCheckFailed() {
|
||||
// Just restart the whole lookup process
|
||||
go t.FindProvider(t.CtxMonitor)
|
||||
}
|
||||
|
||||
func (t *TelemetryService) onFindProvider(prov telem.TelemetryProvider) {
|
||||
// Attach to the provider
|
||||
t.logger.Info("found provider for " + prov.Name())
|
||||
t.Messages <- fmt.Sprintf("Found provider \"%s\"\n", prov.Name())
|
||||
err := t.SwitchProvider(prov)
|
||||
if err != nil {
|
||||
t.Messages <- fmt.Sprintf("Failed to switch to provider: %+v\n", err.Error())
|
||||
t.logger.Error("failed to switch to provider onFindProvider", "err", err)
|
||||
return
|
||||
}
|
||||
|
||||
// Start the healthcheck on our provider so we can drop it if it stops
|
||||
t.CtxHealthcheck, t.healthCheckCancel = context.WithCancel(context.Background())
|
||||
go t.ProviderMonitor(t.CtxHealthcheck)
|
||||
}
|
||||
|
||||
func (t *TelemetryService) onProviderStopsMidStream() {
|
||||
// clear the current provider
|
||||
// TODO: now we need to also clear the devices to restart everything,
|
||||
// if the stream stopped we have to restart the devices and everything
|
||||
t.Messages <- "Telemetry provider stopped mid stream\n"
|
||||
t.logger.Info("cleaning dropped provider and restarting lookup service")
|
||||
t.dropActiveProvider()
|
||||
go t.FindProvider(t.CtxMonitor)
|
||||
}
|
||||
+75
-8
@@ -1,18 +1,12 @@
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"math"
|
||||
"strconv"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
func init() {
|
||||
fieldNameToID = make(map[string]FieldID, MaxFields)
|
||||
for id, name := range FieldNames {
|
||||
fieldNameToID[name] = FieldID(id)
|
||||
}
|
||||
}
|
||||
|
||||
// NOTE: allow the user to create custom data things. For example, iRacing provides
|
||||
// multiple surface temps, but I guess the user doesn't want all of them at once.
|
||||
// allow him to make something that allows some data transformation to occur.
|
||||
@@ -70,7 +64,7 @@ const (
|
||||
// NOTE: we can optimize this via a special command that says a given piece of data
|
||||
// is for multiple targets
|
||||
type TelemetryField struct {
|
||||
// IDs []int16 // Identification for the serial device
|
||||
IDs []int16 // Identification for the serial device
|
||||
Type DataType
|
||||
Raw uint64
|
||||
Str string // Only to be used with DataTypeSTRING
|
||||
@@ -81,6 +75,49 @@ func (tf *TelemetryField) Unused() {
|
||||
tf.Raw = uint64('-')
|
||||
}
|
||||
|
||||
// Pack will pack this current TelemetryField into bytes to send over the wire
|
||||
// Format:
|
||||
// 0x00 - Field ID
|
||||
// 0x00 |
|
||||
// 0x01 - DataType
|
||||
// 0x02 - if its a (u)int8
|
||||
// or
|
||||
// 0x02 - if its a (u)int16 - first byte
|
||||
// 0x02 - if its a (u)int16 - second byte
|
||||
// or
|
||||
// 0x02 - str len max is 255 chars
|
||||
// [0x02] - str
|
||||
func (tf *TelemetryField) Pack(dest []byte) []byte {
|
||||
// NOTE: maybe we can have a pool of these so we don't have to create them here
|
||||
// or whatever
|
||||
for _, id := range tf.IDs {
|
||||
dest = append(dest, uint8(id), uint8(id>>8))
|
||||
dest = append(dest, uint8(tf.Type))
|
||||
|
||||
switch tf.Type {
|
||||
case DataTypeINT8, DataTypeUINT8, DataTypeCHAR:
|
||||
dest = append(dest, uint8(tf.Raw))
|
||||
case DataTypeINT16, DataTypeUINT16:
|
||||
dest = append(dest, uint8(tf.Raw), uint8(tf.Raw>>8))
|
||||
case DataTypeINT32, DataTypeUINT32:
|
||||
dest = append(dest, uint8(tf.Raw), uint8(tf.Raw>>8), uint8(tf.Raw>>16), uint8(tf.Raw>>24))
|
||||
case DataTypeINT64, DataTypeUINT64:
|
||||
dest = append(
|
||||
dest, uint8(tf.Raw), uint8(tf.Raw>>8), uint8(tf.Raw>>16),
|
||||
uint8(tf.Raw>>24), uint8(tf.Raw>>32), uint8(tf.Raw>>40), uint8(tf.Raw>>48),
|
||||
uint8(tf.Raw>>56),
|
||||
)
|
||||
case DataTypeSTRING:
|
||||
l := min(len(tf.Str), math.MaxUint8)
|
||||
|
||||
dest = append(dest, uint8(l))
|
||||
dest = append(dest, tf.Str[:l]...)
|
||||
}
|
||||
}
|
||||
|
||||
return dest
|
||||
}
|
||||
|
||||
func (tf *TelemetryField) String() string {
|
||||
switch tf.Type {
|
||||
case DataTypeSTRING:
|
||||
@@ -222,6 +259,13 @@ func GetFieldName(id FieldID) string {
|
||||
|
||||
var fieldNameToID map[string]FieldID
|
||||
|
||||
func initFieldNamesMap() {
|
||||
fieldNameToID = make(map[string]FieldID, MaxFields)
|
||||
for id, name := range FieldNames {
|
||||
fieldNameToID[name] = FieldID(id)
|
||||
}
|
||||
}
|
||||
|
||||
func GetFieldID(name string) (FieldID, bool) {
|
||||
id, ok := fieldNameToID[name]
|
||||
return id, ok
|
||||
@@ -241,3 +285,26 @@ type TelemetryData struct {
|
||||
func NewTelemetryData() *TelemetryData {
|
||||
return &TelemetryData{}
|
||||
}
|
||||
|
||||
func (td *TelemetryData) Pack() []byte {
|
||||
bufPtr := bufferPool.Get().(*[]byte)
|
||||
buf := (*bufPtr)[:0]
|
||||
|
||||
// for _, bind := range td.ActiveBinds {
|
||||
// buf = td.Values[bind.ID].Pack(buf)
|
||||
// }
|
||||
|
||||
for k := range td.Values {
|
||||
if len(td.Values[k].IDs) > 0 {
|
||||
buf = td.Values[k].Pack(buf)
|
||||
}
|
||||
}
|
||||
|
||||
// We have to copy here because we have to return the buffer
|
||||
result := make([]byte, len(buf))
|
||||
copy(result, buf)
|
||||
|
||||
bufferPool.Put(&buf)
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
@@ -6,7 +6,7 @@ import "time"
|
||||
type TelemetryProvider interface {
|
||||
StopStream()
|
||||
Stream() (<-chan TelemetryData, error)
|
||||
Subscribe([]FieldID)
|
||||
Subscribe(map[int16]FieldID)
|
||||
IsAlive(time.Duration) bool
|
||||
Name() string
|
||||
Close()
|
||||
|
||||
@@ -1 +1,5 @@
|
||||
package telemetry
|
||||
|
||||
func Init() {
|
||||
initFieldNamesMap()
|
||||
}
|
||||
|
||||
@@ -35,7 +35,6 @@ func NewLayoutController(base *Controller, service *services.DeviceService) *Lay
|
||||
DevService: service,
|
||||
MoveToolState: &windowManipState{Mode: moveMode},
|
||||
// SelectedLayout: "beamng.yaml",
|
||||
// TODO: this can't be here - the service/peripheral needs to know about it
|
||||
SelectedLayout: "layout.yaml",
|
||||
}
|
||||
|
||||
@@ -212,7 +211,7 @@ func (lc *LayoutController) createWindow() {
|
||||
}
|
||||
|
||||
// Acquire the cdashdisplay
|
||||
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
return
|
||||
@@ -322,7 +321,7 @@ func (lc *LayoutController) newWindowAction() {
|
||||
|
||||
func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) {
|
||||
// Acquire the cdashdisplay
|
||||
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
return
|
||||
@@ -348,7 +347,7 @@ func (lc *LayoutController) displayLoadedLayouts() {
|
||||
lc.Logger.Debug("We want to view our layout!")
|
||||
|
||||
// Acquire the cdashdisplay
|
||||
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
return
|
||||
@@ -393,7 +392,7 @@ func (lc *LayoutController) getCurrentTreeNodeModel() (*tview.TreeNode, int16, e
|
||||
|
||||
func (lc *LayoutController) loadLayout() {
|
||||
// Acquire the cdashdisplay
|
||||
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
return
|
||||
@@ -405,8 +404,7 @@ func (lc *LayoutController) loadLayout() {
|
||||
}
|
||||
// ---
|
||||
|
||||
// TODO: This can happen here, but we need to address how the layout is gotten
|
||||
// THE UI SHOULD SET STATE IN THE SERVICES ONLY
|
||||
// We would get the layout path from somewhere but for nots its layout.yaml
|
||||
err = display.LoadLayout(lc.SelectedLayout)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to load layout: " + err.Error()
|
||||
@@ -418,7 +416,7 @@ func (lc *LayoutController) loadLayout() {
|
||||
|
||||
func (lc *LayoutController) unloadLayout() {
|
||||
// Acquire the cdashdisplay
|
||||
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
return
|
||||
@@ -439,7 +437,7 @@ func (lc *LayoutController) unloadLayout() {
|
||||
|
||||
func (lc *LayoutController) saveLayout() {
|
||||
// Acquire the cdashdisplay
|
||||
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
return
|
||||
@@ -472,7 +470,7 @@ func (lc *LayoutController) deleteWindow() {
|
||||
wID := curNode.GetReference().(int16)
|
||||
|
||||
// Acquire the cdashdisplay
|
||||
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
return
|
||||
|
||||
@@ -52,7 +52,7 @@ func (lc *LayoutController) handleMovementCapture(idx int16,
|
||||
}
|
||||
|
||||
// Acquire the cdashdisplay
|
||||
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
return nil
|
||||
@@ -86,7 +86,7 @@ func (lc *LayoutController) handleResizeCapture(idx int16,
|
||||
}
|
||||
|
||||
// Acquire the cdashdisplay
|
||||
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
||||
if err != nil {
|
||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
return nil
|
||||
|
||||
@@ -16,18 +16,19 @@ type DeviceController struct {
|
||||
DeviceAPIView *views.DeviceAPIView
|
||||
LayoutCtrl *LayoutController
|
||||
StreamCtrl *StreamingCtrl
|
||||
Orchestrator *serv.Orchestrator
|
||||
DevService *serv.DeviceService
|
||||
}
|
||||
|
||||
func NewDeviceController(
|
||||
base *Controller,
|
||||
orchestrator *serv.Orchestrator,
|
||||
devService *serv.DeviceService,
|
||||
telemService *serv.TelemetryService,
|
||||
) *DeviceController {
|
||||
mc := &DeviceController{
|
||||
Controller: base,
|
||||
LayoutCtrl: NewLayoutController(base, orchestrator.DeviceService),
|
||||
Orchestrator: orchestrator,
|
||||
StreamCtrl: NewStreamingCtrl(base, orchestrator.DeviceService, orchestrator.TelemetryService),
|
||||
Controller: base,
|
||||
LayoutCtrl: NewLayoutController(base, devService),
|
||||
DevService: devService,
|
||||
StreamCtrl: NewStreamingCtrl(base, devService, telemService),
|
||||
}
|
||||
|
||||
return mc
|
||||
@@ -56,7 +57,7 @@ func (mc *DeviceController) setDeviceAPIViewEvents() {
|
||||
SetInputCapture(func(ev *tcell.EventKey) *tcell.EventKey {
|
||||
switch ev.Rune() {
|
||||
case 'r':
|
||||
go mc.Orchestrator.DeviceService.FindDevices()
|
||||
go mc.DevService.FindDevices()
|
||||
}
|
||||
return ev
|
||||
})
|
||||
@@ -66,8 +67,8 @@ func (mc *DeviceController) AddDeviceAPIListItems() {
|
||||
mc.DeviceAPIView.DevAPIList.
|
||||
AddItem("layout", "build a layout for CDashDisplay", func() {
|
||||
// This CDashDisplay specific, only load if we have a CDashDisplay
|
||||
if !mc.Orchestrator.DeviceService.PeripheralExists(cdashdisplay.NAME) {
|
||||
mc.Orchestrator.DeviceService.Messages <- "CDashDisplay it not loaded yet\n"
|
||||
if !mc.DevService.DeviceExists(cdashdisplay.NAME) {
|
||||
mc.DevService.Messages <- "CDashDisplay it not loaded yet\n"
|
||||
return
|
||||
}
|
||||
|
||||
@@ -115,7 +116,7 @@ func (mc *DeviceController) injectControllerCallbacks() {
|
||||
|
||||
func (mc *DeviceController) injectChannels() {
|
||||
go func() {
|
||||
for msg := range mc.Orchestrator.Messages {
|
||||
for msg := range mc.DevService.Messages {
|
||||
mc.PrintToOutputWindow(msg)
|
||||
}
|
||||
}()
|
||||
|
||||
@@ -18,7 +18,7 @@ import (
|
||||
|
||||
type StreamingCtrl struct {
|
||||
*Controller
|
||||
DevService *services.DeviceService
|
||||
Service *services.DeviceService
|
||||
StreamView *views.StreamToolView
|
||||
Messages chan string
|
||||
Internal chan string
|
||||
@@ -45,20 +45,27 @@ func NewStreamingCtrl(
|
||||
|
||||
ctrl := &StreamingCtrl{
|
||||
Controller: base,
|
||||
DevService: devService,
|
||||
Service: devService,
|
||||
TelemServ: serTelem,
|
||||
Messages: make(chan string, 10),
|
||||
Internal: make(chan string, 10),
|
||||
TelemetryCh: make(chan telemetry.TelemetryData, 1),
|
||||
Run: false,
|
||||
StreamView: streamView,
|
||||
isRunning: false,
|
||||
}
|
||||
|
||||
ctrl.registerHooks()
|
||||
// ctrl.subscribeListeners()
|
||||
|
||||
return ctrl
|
||||
}
|
||||
|
||||
// func (sc *StreamingCtrl) subscribeListeners() {
|
||||
// // Here I will set a UIDevice
|
||||
// sc.TelemetryCh = sc.TelemServ.SubscribeListener("UI", 1)
|
||||
// }
|
||||
|
||||
func (sc *StreamingCtrl) registerHooks() {
|
||||
sc.StreamView.Options.Form.SetInputCapture(func(ev *tcell.EventKey) *tcell.EventKey {
|
||||
switch ev.Key() {
|
||||
@@ -89,11 +96,11 @@ func (sc *StreamingCtrl) registerHooks() {
|
||||
}
|
||||
|
||||
func (sc *StreamingCtrl) StartStop() {
|
||||
if sc.TelemServ.IsStreaming() {
|
||||
if sc.isRunning {
|
||||
slog.Info("stopping stream")
|
||||
|
||||
sc.TelemServ.StopStream()
|
||||
sc.DevService.StopStream()
|
||||
sc.Service.StopStream()
|
||||
|
||||
sc.isRunning = false
|
||||
return
|
||||
@@ -103,9 +110,9 @@ func (sc *StreamingCtrl) StartStop() {
|
||||
// NOTE:
|
||||
// Subscribe the only existing device - needs to be discovered by now
|
||||
slog.Debug("setting the data stream for device servie")
|
||||
sc.DevService.SetTelemetryChannel(sc.TelemServ.SubscribeListener("DeviceService", 1))
|
||||
sc.Service.SetTelemetryChannel(sc.TelemServ.SubscribeListener("DeviceService", 1))
|
||||
|
||||
dev, err := sc.DevService.GetPeripheral(uidevice.NAME)
|
||||
dev, err := sc.Service.GetDevice(uidevice.NAME)
|
||||
if err == nil {
|
||||
if uiDev, ok := dev.(*uidevice.UIDevice); ok {
|
||||
sc.TelemetryCh = uiDev.DataChannel()
|
||||
@@ -114,7 +121,7 @@ func (sc *StreamingCtrl) StartStop() {
|
||||
}
|
||||
|
||||
slog.Debug("starting services")
|
||||
sc.DevService.StartStream()
|
||||
sc.Service.StartStream()
|
||||
sc.TelemServ.StartStream()
|
||||
|
||||
sc.isRunning = true
|
||||
@@ -154,8 +161,24 @@ func (sc *StreamingCtrl) updateStream() {
|
||||
// Performance reasoning: this is not used during the high frequency data transmission
|
||||
// so we can get away with using a map for convenience here
|
||||
func (sc *StreamingCtrl) SetInternalState() {
|
||||
fields := sc.TelemServ.SubscribeToFields()
|
||||
sc.Messages <- fmt.Sprintf("Subscribed to fields: %+v\n", fields)
|
||||
// Acquire the cdashdisplay
|
||||
// displayIF, err := sc.Service.GetDevice(cdashdisplay.NAME)
|
||||
// if err != nil {
|
||||
// sc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||
// return
|
||||
// }
|
||||
// display, ok := displayIF.(*cdashdisplay.CDashDisplay)
|
||||
// if !ok {
|
||||
// sc.Messages <- "failed to acquire " + cdashdisplay.NAME
|
||||
// return
|
||||
// }
|
||||
// ---
|
||||
|
||||
sc.TelemServ.SubscribeToFields()
|
||||
|
||||
// sc.Messages <- fmt.Sprintf("Subscribed Fields: %+v [%d]\n", fields, len(fields))
|
||||
// Should I update this?
|
||||
sc.Messages <- fmt.Sprintf("Subscribed to fields\n")
|
||||
}
|
||||
|
||||
func (sc *StreamingCtrl) listenToUIStream() {
|
||||
@@ -167,6 +190,8 @@ func (sc *StreamingCtrl) listenToUIStream() {
|
||||
}
|
||||
isDrawing.Store(true)
|
||||
|
||||
// sc.Logger.Debug("got data", "data", msg)
|
||||
|
||||
// Capture locally
|
||||
telemetryMsg := msg
|
||||
|
||||
|
||||
+9
-5
@@ -23,15 +23,19 @@ func NewControlPanel(logger *slog.Logger) *ControlPanel {
|
||||
App: tview.NewApplication(),
|
||||
}
|
||||
|
||||
orchestrator, err := services.NewOrchestrator(logger)
|
||||
if err != nil {
|
||||
// TODO: no panic here
|
||||
panic("failed to create services orchestrator")
|
||||
// NOTE: create our device service here
|
||||
devService := services.NewDeviceService(logger.With("service", "DeviceService"))
|
||||
|
||||
telemService := services.NewTelemetryService(logger, devService)
|
||||
if telemService == nil {
|
||||
panic("failed to create the telemetry service")
|
||||
}
|
||||
|
||||
go telemService.FindProvider(telemService.CtxMonitor)
|
||||
|
||||
return &ControlPanel{
|
||||
Controller: baseController,
|
||||
DeviceController: controllers.NewDeviceController(baseController, orchestrator),
|
||||
DeviceController: controllers.NewDeviceController(baseController, devService, telemService),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user