Author SHA1 Message Date
esilva a519cf473e peripheral can connect mid stream and start go into streaming immidiately 2026-09-23 13:40:31 +01:00
esilva 76d758e264 Merge pull request 'Device reconnection handling' (#15) from device-reconnection-handling into auto-detect-devices
Reviewed-on: #15
2026-09-23 10:27:52 +01:00
esilva df1e3c26fa reconnects are slightly improved
I don't really like the reset thing with the delay. I need to implement
something better later on
2026-09-22 17:07:53 +01:00
esilva 636f20df77 delete this dead code 2026-09-22 16:27:44 +01:00
esilva 82972c9665 Added healthchecks, and maybe got a bit lost 2026-09-22 16:21:50 +01:00
esilva 009609da93 fix after implementing new changes to PSS. streaming was broken 2026-09-21 16:50:55 +01:00
esilva 91f33c79ed refactored the peripheral configuration loop 2026-09-21 16:07:00 +01:00
esilva d10c6997f4 clean up a bit the logic of how we get to know whether the telemetry started on devices lookup 2026-09-21 15:55:31 +01:00
esilva b364931c80 automatic setup for devices
Its not close to being done and still need to work out the flows due to
who arrives first: telemetry or the peripheral?
2026-09-21 15:27:42 +01:00
esilva 6ff4cc3014 added a central constants package so we don't have to pull a provider if all we need is its name 2026-09-21 15:26:57 +01:00
esilva 263d3c29a3 services decoupling 2026-09-21 11:45:10 +01:00
esilva 5d5b7bfdc3 working on the peripheral state handling 2026-09-20 21:44:15 +01:00
esilva dcba0774c2 Merge pull request 'Start stop behaviour' (#12) from start-stop-behaviour into auto-detect-devices
Reviewed-on: #12
2026-09-20 15:22:18 +01:00
esilva a5c465a203 has way too many things
In short: added the ability to detect when a device has disconnected.
From this we have to achieve a way of:
- reconnecting the device
- re-subscribing to the fields it needs (it should already be subscribed
since the ESDI didn't stop tho - but for the sake of it, or in case a
device joins later)
- start sending data to it (which should be automatic given how we are
handling the devices)
2026-09-20 15:06:48 +01:00
esilva 174b27004f added error returns to the SendData method on the peripheral interface 2026-09-18 16:37:30 +01:00
esilva a2e14ba6a9 applied the same start-stop scheme to the BeamNG provider 2026-09-18 15:42:34 +01:00
esilva a6683df10b fixed the start and stop behaviour. it was crashing 2026-09-18 15:38:17 +01:00
esilva 2d4875fbd0 removed log that was just polluting everything 2026-09-18 11:07:29 +01:00
esilva 0a8f34423c Merge pull request 'decoupled the winIDs from the telemetry data' (#11) from decouple-cdash-from-field-subscription into auto-detect-devices
Reviewed-on: #11
2026-09-18 00:07:39 +01:00
esilva bd10d6b4d6 subscription fix. should subscribe too fields at once 2026-09-18 00:07:27 +01:00
esilva c888dc5436 remove commented code and made the subscription be based on a list 2026-09-18 00:05:50 +01:00
esilva dda300ce6b decoupled the winIDs from the telemetry data
hell yeah, brother!
2026-09-17 23:55:54 +01:00
34 changed files with 1302 additions and 568 deletions
+1
View File
@@ -0,0 +1 @@
# Streaming Flow
+87 -93
View File
@@ -1,103 +1,97 @@
package cmd package cmd
import ( import (
"esdi/peripheral"
"fmt"
"strconv"
repl "github.com/ESilva15/ESgoRepl"
"github.com/spf13/cobra" "github.com/spf13/cobra"
) )
func replCmdAction(cmd *cobra.Command, args []string) { func replCmdAction(cmd *cobra.Command, args []string) {
r := repl.NewREPL(repl.REPLCfg{ // r := repl.NewREPL(repl.REPLCfg{
PS1: "\rESDI > ", // PS1: "\rESDI > ",
}) // })
//
perClerk := peripheral.NewPeripheralDeviceClerk() // perClerk := peripheral.NewPeripheralDeviceClerk()
//
discoverDevicesREPLCmd := repl.Command{ // discoverDevicesREPLCmd := repl.Command{
Name: "discover", // Name: "discover",
Usage: "discovers connected devices", // Usage: "discovers connected devices",
Action: func(r *repl.REPL, args []string) error { // Action: func(r *repl.REPL, args []string) error {
err := perClerk.FindDevices() // err := perClerk.FindDevices()
if err != nil { // if err != nil {
return err // return err
} // }
//
return nil // return nil
}, // },
} // }
//
listDevicesREPLCmd := repl.Command{ // listDevicesREPLCmd := repl.Command{
Name: "list", // Name: "list",
Usage: "lists connected devices", // Usage: "lists connected devices",
Action: func(r *repl.REPL, args []string) error { // Action: func(r *repl.REPL, args []string) error {
_ = perClerk.ListDevices() // _ = perClerk.ListDevices()
//
return nil // return nil
}, // },
} // }
//
listDeviceAPIREPLCmd := repl.Command{ // listDeviceAPIREPLCmd := repl.Command{
Name: "v-api", // Name: "v-api",
Usage: "shows API of a device - pass its ID", // Usage: "shows API of a device - pass its ID",
Action: func(r *repl.REPL, args []string) error { // Action: func(r *repl.REPL, args []string) error {
// We should add this to the REPL instead // // We should add this to the REPL instead
if len(args) < 1 { // if len(args) < 1 {
return fmt.Errorf("requires at least on argument") // return fmt.Errorf("requires at least on argument")
} // }
//
// First and only argument should be the ID of the device we want to use // // First and only argument should be the ID of the device we want to use
targetID, err := strconv.ParseInt(args[0], 10, 0) // targetID, err := strconv.ParseInt(args[0], 10, 0)
if err != nil { // if err != nil {
return err // return err
} // }
//
err = perClerk.ListDeviceAPI(uint8(targetID)) // err = perClerk.ListDeviceAPI(uint8(targetID))
if err != nil { // if err != nil {
fmt.Println("failed to view device API: ", err.Error()) // fmt.Println("failed to view device API: ", err.Error())
} // }
//
return nil // return nil
}, // },
} // }
//
runDeviceAPIREPLCmd := repl.Command{ // runDeviceAPIREPLCmd := repl.Command{
Name: "v-run", // Name: "v-run",
Usage: "runs a funcion of a device - pass its ID and function name", // Usage: "runs a funcion of a device - pass its ID and function name",
Action: func(r *repl.REPL, args []string) error { // Action: func(r *repl.REPL, args []string) error {
// We should add this to the REPL instead // // We should add this to the REPL instead
if len(args) < 3 { // if len(args) < 3 {
return fmt.Errorf("requires at least on argument") // return fmt.Errorf("requires at least on argument")
} // }
//
// First and only argument should be the ID of the device we want to use // // First and only argument should be the ID of the device we want to use
targetID, err := strconv.ParseInt(args[0], 10, 0) // targetID, err := strconv.ParseInt(args[0], 10, 0)
if err != nil { // if err != nil {
return err // return err
} // }
//
fnName := args[1] // fnName := args[1]
fnArgs := args[2:] // fnArgs := args[2:]
//
err = perClerk.RunDeviceFunction(uint8(targetID), fnName, fnArgs) // err = perClerk.RunDeviceFunction(uint8(targetID), fnName, fnArgs)
if err != nil { // if err != nil {
return err // return err
} // }
//
return nil // return nil
}, // },
} // }
//
r.RegisterCMD(discoverDevicesREPLCmd) // r.RegisterCMD(discoverDevicesREPLCmd)
r.RegisterCMD(listDevicesREPLCmd) // r.RegisterCMD(listDevicesREPLCmd)
r.RegisterCMD(listDeviceAPIREPLCmd) // r.RegisterCMD(listDeviceAPIREPLCmd)
r.RegisterCMD(runDeviceAPIREPLCmd) // r.RegisterCMD(runDeviceAPIREPLCmd)
//
r.Start() // r.Start()
r.Close() // r.Close()
} }
// removeLabelCmd represents the removeLabel command // removeLabelCmd represents the removeLabel command
+6
View File
@@ -0,0 +1,6 @@
package constants
const (
IRacingProviderName = "iRacing"
BeamNGProviderName = "BeamNG.drive"
)
+1 -1
View File
@@ -73,7 +73,7 @@ func findDisplayPort() (*communication.WalkieTalkie, error) {
select { select {
case err = <-probeResult: case err = <-probeResult:
// Probe completed normally (could be success or error) // Probe completed normally (could be success or error)
case <-time.After(2 * time.Second): case <-time.After(1000 * time.Millisecond):
// Hard timeout reached // Hard timeout reached
err = fmt.Errorf("probe completely hung/timed out: %s", port) err = fmt.Errorf("probe completely hung/timed out: %s", port)
} }
+76 -12
View File
@@ -8,9 +8,11 @@ import (
"log/slog" "log/slog"
"os" "os"
"path" "path"
"sync"
"time" "time"
helper "esdi/helpers" helper "esdi/helpers"
"esdi/peripheral"
"esdi/peripheral/communication" "esdi/peripheral/communication"
"esdi/peripheral/communication/packets" "esdi/peripheral/communication/packets"
"esdi/peripheral/types" "esdi/peripheral/types"
@@ -31,6 +33,8 @@ const (
updateWindowCMDID types.Command = 6 // Change this to a move cmd instead updateWindowCMDID types.Command = 6 // Change this to a move cmd instead
sendDataCMDID types.Command = 7 sendDataCMDID types.Command = 7
newLayoutCMDID types.Command = 8 newLayoutCMDID types.Command = 8
healthCheckCMDID types.Command = 9
resetCMDID types.Command = 10
) )
const ( const (
@@ -106,8 +110,12 @@ func NewCDashState() *CDashState {
} }
type CDashDisplay struct { type CDashDisplay struct {
WT *communication.WalkieTalkie WT *communication.WalkieTalkie
State *CDashState State *CDashState
fieldToWindows map[telemetry.FieldID][]int16
bufPool sync.Pool
failedSends int
FailedSendsConsecutiveLimit int
} }
// Connect will try to find and connect to the CDashDisplay // Connect will try to find and connect to the CDashDisplay
@@ -115,17 +123,52 @@ func NewCDashDisplay() (*CDashDisplay, error) {
// Look for the port // Look for the port
p, err := findDisplayPort() p, err := findDisplayPort()
if err != nil { if err != nil {
slog.Info("failed to find cdashdisplay port: %s", err.Error()) slog.Info("failed to find cdashdisplay port", "reason", err.Error())
return nil, err return nil, err
} }
return &CDashDisplay{ return &CDashDisplay{
WT: p, WT: p,
State: NewCDashState(), 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,
}, nil }, nil
} }
func (d *CDashDisplay) SendCommand() { 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) CreateWindow(win *DesktopUIWindow) (*DesktopUIWindow, error) { func (d *CDashDisplay) CreateWindow(win *DesktopUIWindow) (*DesktopUIWindow, error) {
@@ -146,6 +189,9 @@ func (d *CDashDisplay) CreateWindow(win *DesktopUIWindow) (*DesktopUIWindow, err
slog.Info(fmt.Sprintf("Recived ID message: %v", wID)) slog.Info(fmt.Sprintf("Recived ID message: %v", wID))
d.State.Layout.AddWindow(win) d.State.Layout.AddWindow(win)
if fieldID, ok := telemetry.GetFieldID(win.UIData.TelemetryField); ok {
d.RegisterFieldMapping(fieldID, win.UIData.IDX)
}
return win, nil return win, nil
} }
@@ -177,6 +223,8 @@ func (d *CDashDisplay) UpdateWindow(win *DesktopUIWindow) error {
// Yeah, same address as suspected // Yeah, same address as suspected
// I can't think about it right now. I'll think about that tomorrow // I can't think about it right now. I'll think about that tomorrow
// TODO: need to update the field mappings here!
return nil return nil
} }
@@ -199,7 +247,8 @@ func (d *CDashDisplay) DestroyWindow(wID int16) error {
return err return err
} }
// NODE: add this // NOTE: add this
d.UnregisterFieldMapping(wID)
d.State.Layout.RemoveWindow(wID) d.State.Layout.RemoveWindow(wID)
return nil return nil
@@ -356,12 +405,23 @@ func (d *CDashDisplay) UnloadLayout() error {
return nil return nil
} }
func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) { func (cds *CDashDisplay) reset() error {
packet := data.Pack() 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)
bytes, err := helper.StructToBytes(packet) bytes, err := helper.StructToBytes(packet)
if err != nil { if err != nil {
return return peripheral.ErrFailureToPackData
} }
curStr := "" curStr := ""
@@ -371,7 +431,6 @@ func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) {
curStr += fmt.Sprintf("%02x ", byte) curStr += fmt.Sprintf("%02x ", byte)
if byteCount == 8 { if byteCount == 8 {
// slog.Debug(curStr)
curStr = "" curStr = ""
byteCount = 0 byteCount = 0
} }
@@ -380,6 +439,11 @@ func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) {
// var ack packets.AckPacket // var ack packets.AckPacket
err = d.WT.SendCommand(sendDataCMDID, bytes, nil) err = d.WT.SendCommand(sendDataCMDID, bytes, nil)
if err != nil && err != io.EOF { if err != nil && err != io.EOF {
return if d.failedSends == d.FailedSendsConsecutiveLimit {
return peripheral.ErrDeviceTimedOut
}
d.failedSends++
} }
return nil
} }
+9
View File
@@ -0,0 +1,9 @@
package cdashdisplay
func (cds *CDashDisplay) OnLoad() error {
return nil
}
func (cds *CDashDisplay) OnTelemetryProviderFound() error {
return nil
}
+17 -4
View File
@@ -1,6 +1,7 @@
package cdashdisplay package cdashdisplay
import ( import (
"esdi/peripheral/communication/packets"
"esdi/peripheral/devices" "esdi/peripheral/devices"
"esdi/telemetry" "esdi/telemetry"
) )
@@ -15,13 +16,25 @@ func (cds *CDashDisplay) Name() string {
return NAME return NAME
} }
func (cds *CDashDisplay) RequiredFields() map[int16]telemetry.FieldID { func (cds *CDashDisplay) RequiredFields() []telemetry.FieldID {
fields := make(map[int16]telemetry.FieldID, len(cds.State.Layout.Windows)) fields := make([]telemetry.FieldID, 0, len(cds.State.Layout.Windows))
for _, w := range cds.State.Layout.Windows { for _, w := range cds.State.Layout.Windows {
fieldID, _ := telemetry.GetFieldID(w.UIData.TelemetryField) if fieldID, ok := telemetry.GetFieldID(w.UIData.TelemetryField); ok {
fields[w.UIData.IDX] = fieldID fields = append(fields, fieldID)
}
} }
return fields 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
}
+44
View File
@@ -0,0 +1,44 @@
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)
}
}
+74 -3
View File
@@ -1,5 +1,11 @@
package cdashdisplay package cdashdisplay
import (
"math"
"esdi/telemetry"
)
// In this file we will place all structs that are 1:1 representation of the // In this file we will place all structs that are 1:1 representation of the
// types in the device transport layer -> // types in the device transport layer ->
@@ -65,6 +71,71 @@ type UIWindowUpdatePacket struct {
Window UIWindow Window UIWindow
} }
// func getUIWindowDTO(w *DesktopWindowData) *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
}
+28 -5
View File
@@ -2,28 +2,33 @@
package devices package devices
import ( import (
"errors"
"esdi/devices/cdashdisplay" "esdi/devices/cdashdisplay"
"esdi/devices/uidevice" "esdi/devices/uidevice"
"esdi/peripheral" "esdi/peripheral"
) )
var ErrInvalidDevice = errors.New("invalid device")
type Device struct { type Device struct {
Name string Name string
Discover func() (peripheral.Peripheral, error) 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: { uidevice.NAME: {
Name: uidevice.NAME, Name: uidevice.NAME,
Discover: DiscoverUIDevice, Discover: UIDeviceDiscover,
}, },
cdashdisplay.NAME: { cdashdisplay.NAME: {
Name: cdashdisplay.NAME, Name: cdashdisplay.NAME,
Discover: DiscoverCDashDisplay, Discover: CDashDisplayDiscover,
}, },
} }
func DiscoverUIDevice() (peripheral.Peripheral, error) { func UIDeviceDiscover() (peripheral.Peripheral, error) {
uidev, err := uidevice.NewUIDevice() uidev, err := uidevice.NewUIDevice()
if err != nil { if err != nil {
return nil, err return nil, err
@@ -32,7 +37,11 @@ func DiscoverUIDevice() (peripheral.Peripheral, error) {
return uidev, nil return uidev, nil
} }
func DiscoverCDashDisplay() (peripheral.Peripheral, error) { // func UIDeviceSetup(peripheral peripheral.Peripheral) error {
// return nil
// }
func CDashDisplayDiscover() (peripheral.Peripheral, error) {
// Create a cdashdisplay // Create a cdashdisplay
display, err := cdashdisplay.NewCDashDisplay() display, err := cdashdisplay.NewCDashDisplay()
if err != nil { if err != nil {
@@ -41,3 +50,17 @@ func DiscoverCDashDisplay() (peripheral.Peripheral, error) {
return display, nil return display, nil
} }
// 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
// }
+21 -10
View File
@@ -17,9 +17,21 @@ func NewUIDevice() (peripheral.Peripheral, error) {
}, nil }, nil
} }
func (uid *UIDevice) SendData(data *telemetry.TelemetryData) { 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 {
if data == nil { if data == nil {
return return peripheral.ErrInvalidData
} }
select { select {
@@ -27,6 +39,8 @@ func (uid *UIDevice) SendData(data *telemetry.TelemetryData) {
default: default:
// Drop frame if buffer is full // Drop frame if buffer is full
} }
return nil
} }
func (uid *UIDevice) Name() string { func (uid *UIDevice) Name() string {
@@ -37,17 +51,14 @@ func (uid *UIDevice) DataChannel() <-chan telemetry.TelemetryData {
return uid.dataChan return uid.dataChan
} }
func (uid *UIDevice) RequiredFields() map[int16]telemetry.FieldID { func (uid *UIDevice) RequiredFields() []telemetry.FieldID {
subscribeTo := []telemetry.FieldID{ return []telemetry.FieldID{
telemetry.Speed, telemetry.Speed,
telemetry.Gear, telemetry.Gear,
telemetry.RPM, telemetry.RPM,
} }
}
fields := make(map[int16]telemetry.FieldID, len(subscribeTo)) func (uid *UIDevice) HealthCheck() bool {
for k, field := range subscribeTo { return true
fields[int16(k)] = field
}
return fields
} }
+9
View File
@@ -0,0 +1,9 @@
package uidevice
func (uid *UIDevice) OnLoad() error {
return nil
}
func (uid *UIDevice) OnTelemetryProviderFound() error {
return nil
}
-4
View File
@@ -10,7 +10,6 @@ import (
"esdi/cmd" "esdi/cmd"
"esdi/config" "esdi/config"
"esdi/telemetry"
"github.com/arl/statsviz" "github.com/arl/statsviz"
) )
@@ -42,9 +41,6 @@ func initApplication() {
slog.Error(fmt.Sprintf("failed to setup metrics server: %+v", err)) slog.Error(fmt.Sprintf("failed to setup metrics server: %+v", err))
} }
} }
// Setting up some internal data structures
telemetry.Init()
} }
func setupLogger() error { func setupLogger() error {
@@ -17,6 +17,7 @@ const (
CmdAckID types.Command = 2 CmdAckID types.Command = 2
CmdCreateWindow types.Command = 3 CmdCreateWindow types.Command = 3
CmdDestroyWindow types.Command = 4 CmdDestroyWindow types.Command = 4
CmdHealthCheck types.Command = 9
) )
var crc8Table = [256]byte{ var crc8Table = [256]byte{
@@ -18,3 +18,22 @@ func (pkt *NewWindowID) Validate() bool {
return true 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
}
+23 -48
View File
@@ -26,8 +26,7 @@ func (wt *WalkieTalkie) ReadFramedData(size int, packet any) error {
b := make([]byte, 1) b := make([]byte, 1)
_, err := wt.Serial.Read(b) _, err := wt.Serial.Read(b)
if err != nil { if err != nil {
// fmt.Fprintf(os.Stderr, "dev read: %s\n", err.Error()) return fmt.Errorf("error reading incoming: %+v, err:", b, err)
return err
} }
if b[0] == constvar.StartOfText { if b[0] == constvar.StartOfText {
@@ -44,9 +43,11 @@ func (wt *WalkieTalkie) ReadFramedData(size int, packet any) error {
reader := bytes.NewReader(buf) reader := bytes.NewReader(buf)
err = binary.Read(reader, binary.LittleEndian, packet) err = binary.Read(reader, binary.LittleEndian, packet)
if err != nil { if err != nil {
return err return fmt.Errorf("error parsing incoming: %+v, err:", buf, err)
} }
wt.Serial.Flush()
return nil return nil
} }
@@ -108,6 +109,15 @@ func (wt *WalkieTalkie) sendPacket(cmd types.Command, data any) error {
return err 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{ packet := CMDDataPacket{
StartMarker: constvar.StartOfText, StartMarker: constvar.StartOfText,
CMD: cmd, CMD: cmd,
@@ -119,13 +129,17 @@ func (wt *WalkieTalkie) sendPacket(cmd types.Command, data any) error {
// Send the payload // Send the payload
serializedPacket := packet.Serialize() serializedPacket := packet.Serialize()
// fmt.Fprintf(os.Stderr, "%+v", serializedPacket)
// slog.Debug("Serialized packet", "packet", packet)
// slog.Debug("# END ##########################################################")
_, err = wt.Serial.Write(serializedPacket) _, err = wt.Serial.Write(serializedPacket)
if err != nil { if err != nil {
return err return err
} }
wt.Serial.Flush()
return nil return nil
} }
@@ -147,50 +161,11 @@ func (wt *WalkieTalkie) readPacket(resp packets.Packet) error {
return nil return nil
} }
// func (wt *WalkieTalkie) sendHeader(h *header) error { func (wt *WalkieTalkie) SendCommand(
// err := wt.sendPacket(h) cmd types.Command,
// if err != nil { payload any,
// return err responseBody packets.Packet,
// } ) error {
//
// // 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 // Send the body
err := wt.sendPacket(cmd, payload) err := wt.sendPacket(cmd, payload)
if err != nil { if err != nil {
+9
View File
@@ -0,0 +1,9 @@
package peripheral
import "errors"
var (
ErrInvalidData = errors.New("invalid data")
ErrDeviceTimedOut = errors.New("device timed out")
ErrFailureToPackData = errors.New("failed to pack received data")
)
+7 -2
View File
@@ -17,8 +17,13 @@ const (
type Peripheral interface { type Peripheral interface {
Name() string Name() string
SendData(*telemetry.TelemetryData) Setup(string) error
RequiredFields() map[int16]telemetry.FieldID HealthCheck() bool
SendData(*telemetry.TelemetryData) error
RequiredFields() []telemetry.FieldID
OnLoad() error
OnTelemetryProviderFound() error
Close() error
} }
type PeripheralDeviceClerk struct { type PeripheralDeviceClerk struct {
+21 -19
View File
@@ -8,6 +8,7 @@ import (
"sync" "sync"
"time" "time"
"esdi/constants"
"esdi/telemetry" "esdi/telemetry"
bngsdk "github.com/ESilva15/gobngsdk" bngsdk "github.com/ESilva15/gobngsdk"
@@ -28,7 +29,7 @@ type BeamNG struct {
updaters [telemetry.MaxFields]func(*telemetry.TelemetryField) updaters [telemetry.MaxFields]func(*telemetry.TelemetryField)
// stream control // stream control
streamCh chan telemetry.TelemetryData wg sync.WaitGroup
streamCancel context.CancelFunc streamCancel context.CancelFunc
// timing // timing
@@ -36,7 +37,7 @@ type BeamNG struct {
} }
const ( const (
NAME = "BeamNG.drive" NAME = constants.BeamNGProviderName
) )
func NewBeamNGProvider(logger *slog.Logger, opts *bngsdk.Options) (*BeamNG, error) { func NewBeamNGProvider(logger *slog.Logger, opts *bngsdk.Options) (*BeamNG, error) {
@@ -46,12 +47,11 @@ func NewBeamNGProvider(logger *slog.Logger, opts *bngsdk.Options) (*BeamNG, erro
} }
provider := &BeamNG{ provider := &BeamNG{
logger: logger.With("TelemetryProvider", NAME), logger: logger.With("TelemetryProvider", NAME),
streamCh: make(chan telemetry.TelemetryData, 1), data: telemetry.NewTelemetryData(),
data: telemetry.NewTelemetryData(), SDK: beam,
SDK: beam, og: &bngsdk.Outgauge{},
og: &bngsdk.Outgauge{}, ticker: time.NewTicker(time.Second / 60),
ticker: time.NewTicker(time.Second / 60),
} }
provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){ 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()) ctx, b.streamCancel = context.WithCancel(context.Background())
// Start the stream // Start the stream
b.stream(ctx) ch := b.stream(ctx)
return b.streamCh, nil return ch, nil
} }
func (b *BeamNG) Subscribe(requestFields map[int16]telemetry.FieldID) { func (b *BeamNG) Subscribe(requestFields []telemetry.FieldID) {
// NOTE: document how the Subscribe funtion works // NOTE: document how the Subscribe funtion works
slog.Debug(fmt.Sprintf("Len Req: %d\n", len(requestFields))) slog.Debug(fmt.Sprintf("Len Req: %d\n", len(requestFields)))
@@ -125,9 +125,7 @@ func (b *BeamNG) Subscribe(requestFields map[int16]telemetry.FieldID) {
// we will add their dependencies and the primitives to a slice // we will add their dependencies and the primitives to a slice
pendingBinds := make([]telemetry.FieldID, telemetry.MaxFields) pendingBinds := make([]telemetry.FieldID, telemetry.MaxFields)
for winID, id := range requestFields { for _, id := range requestFields {
b.data.Values[id].IDs = append(b.data.Values[id].IDs, winID)
switch id { switch id {
case telemetry.RPMStateColour: case telemetry.RPMStateColour:
b.data.VirtualBinds = append(b.data.VirtualBinds, telemetry.NewRPMLights()) b.data.VirtualBinds = append(b.data.VirtualBinds, telemetry.NewRPMLights())
@@ -164,7 +162,6 @@ func (b *BeamNG) Subscribe(requestFields map[int16]telemetry.FieldID) {
func (b *BeamNG) readData() { func (b *BeamNG) readData() {
slog.Debug("READING THIS DATA") slog.Debug("READING THIS DATA")
// BUG: getting stuck in here
ogSnapshot, err := b.SDK.Update() ogSnapshot, err := b.SDK.Update()
slog.Debug("THE DATA WAS READ") slog.Debug("THE DATA WAS READ")
if err != nil { if err != nil {
@@ -195,10 +192,15 @@ func (b *BeamNG) readData() {
b.data.LastDataPoll = time.Now() b.data.LastDataPoll = time.Now()
} }
func (b *BeamNG) stream(ctx context.Context) { func (b *BeamNG) stream(ctx context.Context) <-chan telemetry.TelemetryData {
b.data.InitialTime = time.Now() b.data.InitialTime = time.Now()
outCh := make(chan telemetry.TelemetryData)
b.wg.Add(1)
go func() { go func() {
defer b.wg.Done()
defer close(outCh)
for { for {
// Explicitly intercept cancellation // Explicitly intercept cancellation
select { select {
@@ -207,8 +209,6 @@ func (b *BeamNG) stream(ctx context.Context) {
default: default:
} }
// NOTE: add a method to check if there's data available, or make this happen
select { select {
case <-ctx.Done(): case <-ctx.Done():
return return
@@ -219,7 +219,7 @@ func (b *BeamNG) stream(ctx context.Context) {
// Publish data // Publish data
select { select {
case b.streamCh <- *b.data: case outCh <- *b.data:
slog.Debug("PUBLISHED DATA") slog.Debug("PUBLISHED DATA")
default: default:
// skip this data, don't allow publishers to lag behind // skip this data, don't allow publishers to lag behind
@@ -227,4 +227,6 @@ func (b *BeamNG) stream(ctx context.Context) {
} }
} }
}() }()
return outCh
} }
+21 -18
View File
@@ -10,13 +10,14 @@ import (
"sync" "sync"
"time" "time"
"esdi/constants"
"esdi/telemetry" "esdi/telemetry"
"github.com/ESilva15/goirsdk" "github.com/ESilva15/goirsdk"
) )
const ( const (
NAME = "iRacing" NAME = constants.IRacingProviderName
) )
// IRacing is our iRacing telemetry data provider - its a TelemetryProvider interface // IRacing is our iRacing telemetry data provider - its a TelemetryProvider interface
@@ -33,7 +34,8 @@ type IRacing struct {
ticker *time.Ticker // ticker will keep polling intervals constant ticker *time.Ticker // ticker will keep polling intervals constant
// Stream // Stream
streamCh chan telemetry.TelemetryData wg sync.WaitGroup
// streamCh chan telemetry.TelemetryData
streamCancel context.CancelFunc streamCancel context.CancelFunc
} }
@@ -50,13 +52,11 @@ func NewIRacingProvider(
} }
provider := &IRacing{ provider := &IRacing{
logger: logger, logger: logger,
SDK: sdk, SDK: sdk,
data: telemetry.NewTelemetryData(), 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 // TODO: make this configurable from the user side
ticker: time.NewTicker(time.Second / 240), ticker: time.NewTicker(time.Second / 60),
} }
provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){ provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){
@@ -129,11 +129,14 @@ func (i *IRacing) isDataAvailable() bool {
return true return true
} }
func (i *IRacing) stream(ctx context.Context) { func (i *IRacing) stream(ctx context.Context) <-chan telemetry.TelemetryData {
i.data.InitialTime = time.Now() i.data.InitialTime = time.Now()
outCh := make(chan telemetry.TelemetryData)
i.wg.Add(1)
go func() { go func() {
defer close(i.streamCh) defer i.wg.Done()
defer close(outCh)
// Put this into the configuration file // Put this into the configuration file
consecutiveTimeouts := 0 consecutiveTimeouts := 0
@@ -149,12 +152,11 @@ func (i *IRacing) stream(ctx context.Context) {
if i.SDK.CheckForDataEvent(time.Duration(dataEvTimeout) * time.Millisecond) { if i.SDK.CheckForDataEvent(time.Duration(dataEvTimeout) * time.Millisecond) {
consecutiveTimeouts = 0 consecutiveTimeouts = 0
i.logger.Debug("sending data", "timeouts", consecutiveTimeouts)
i.readData() i.readData()
// Publish data // Publish data
select { select {
case i.streamCh <- *i.data: case outCh <- *i.data:
default: default:
// skip this data, don't allow publishers to lag behind // skip this data, don't allow publishers to lag behind
} }
@@ -170,6 +172,8 @@ func (i *IRacing) stream(ctx context.Context) {
} }
} }
}() }()
return outCh
} }
func (i *IRacing) readData() { func (i *IRacing) readData() {
@@ -206,9 +210,9 @@ func (i *IRacing) Stream() (<-chan telemetry.TelemetryData, error) {
ctx, i.streamCancel = context.WithCancel(context.Background()) ctx, i.streamCancel = context.WithCancel(context.Background())
// Start the stream // Start the stream
i.stream(ctx) ch := i.stream(ctx)
return i.streamCh, nil return ch, nil
} }
func (i *IRacing) StopStream() { func (i *IRacing) StopStream() {
@@ -217,10 +221,11 @@ func (i *IRacing) StopStream() {
} }
i.streamCancel() i.streamCancel()
i.wg.Wait()
i.streamCancel = nil i.streamCancel = nil
} }
func (i *IRacing) Subscribe(requestFields map[int16]telemetry.FieldID) { func (i *IRacing) Subscribe(requestFields []telemetry.FieldID) {
i.logger.Debug(fmt.Sprintf("Len Req: %d\n", len(requestFields))) i.logger.Debug(fmt.Sprintf("Len Req: %d\n", len(requestFields)))
i.data.ActiveBinds = make([]telemetry.BoundField, 0, len(requestFields)) i.data.ActiveBinds = make([]telemetry.BoundField, 0, len(requestFields))
@@ -229,9 +234,7 @@ func (i *IRacing) Subscribe(requestFields map[int16]telemetry.FieldID) {
// we will add their dependencies and the primitives to a slice // we will add their dependencies and the primitives to a slice
pendingBinds := make([]telemetry.FieldID, 0, telemetry.MaxFields) pendingBinds := make([]telemetry.FieldID, 0, telemetry.MaxFields)
for winID, id := range requestFields { for _, id := range requestFields {
i.data.Values[id].IDs = append(i.data.Values[id].IDs, winID)
switch id { switch id {
case telemetry.RPMStateColour: case telemetry.RPMStateColour:
i.data.VirtualBinds = append(i.data.VirtualBinds, telemetry.NewRPMLights()) i.data.VirtualBinds = append(i.data.VirtualBinds, telemetry.NewRPMLights())
+82 -82
View File
@@ -2,46 +2,40 @@ package services
import ( import (
"context" "context"
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"sync"
"sync/atomic" "sync/atomic"
"time"
"esdi/devices" "esdi/devices"
"esdi/peripheral" "esdi/peripheral"
"esdi/telemetry" "esdi/telemetry"
) )
var ErrPeripheralAlreadyRegistered = errors.New("peripheral is already registered") // DeviceService is the API for the peripherals
// 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 { type DeviceService struct {
Logger *slog.Logger Logger *slog.Logger
// Device discovery // Device discovery
mu sync.RWMutex PSS *PeripheralStateStore // Store to track peripheral state
ctxDiscovery context.Context ctxDiscovery context.Context
ctxDiscoveryCancel context.CancelFunc ctxDiscoveryCancel context.CancelFunc
Devices map[string]peripheral.Peripheral
// Strem handling // Strem handling
streamCancel context.CancelFunc streamCancel context.CancelFunc
TelemCh <-chan telemetry.TelemetryData TelemCh <-chan telemetry.TelemetryData
// Output // Output
Messages chan string Messages chan string
// Callbacks
// Telemetry service data fetchers
telemetryProvider func() (string, error)
triggerFieldSubscription func() []telemetry.FieldID
} }
func NewDeviceService(logger *slog.Logger) *DeviceService { func NewDeviceService(logger *slog.Logger, msg chan string) *DeviceService {
sharedChannel := make(chan string, 10)
dev := &DeviceService{ dev := &DeviceService{
Devices: make(map[string]peripheral.Peripheral), PSS: NewPeripheralStateStore(
logger.With("Service", "PeripheralStateStore"), devices.List, msg,
),
Logger: logger, Logger: logger,
Messages: sharedChannel, Messages: msg,
} }
// Start the routine that looks for devices - should always be running in the background // Start the routine that looks for devices - should always be running in the background
@@ -49,84 +43,71 @@ func NewDeviceService(logger *slog.Logger) *DeviceService {
dev.ctxDiscovery, dev.ctxDiscoveryCancel = context.WithCancel(context.Background()) dev.ctxDiscovery, dev.ctxDiscoveryCancel = context.WithCancel(context.Background())
go dev.FindDevices() go dev.FindDevices()
// Set the callbacks for PSS
dev.PSS.telemetryProvider = dev.getTelemetryProvider
dev.PSS.onPeripheralConfigured = dev.peripheralConfigured
return dev return dev
} }
func (ds *DeviceService) FindDevices() { // Getters [START] -------------------------------------------------------------
// Need to define a list of devices to search for // This function is currently only being used by PSS, we may have to find a better
// For now lets just try to find our cdashdisplay - will think about the rest later // pattern for this
ticker := time.NewTicker(2 * time.Second) func (ds *DeviceService) getTelemetryProvider() (string, error) {
defer ticker.Stop() return ds.telemetryProvider()
}
for { func (ds *DeviceService) GetDevices() []peripheral.Peripheral {
select { snapshot := ds.PSS.GetStates()
case <-ds.ctxDiscovery.Done(): peripherals := make([]peripheral.Peripheral, 0, len(snapshot))
// 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) for _, state := range snapshot {
dev, err := peripheral.Discover() if state.State < DeviceIsConfigured {
if err != nil { continue
ds.Logger.Debug("didn't find device", "name", pName) }
continue peripherals = append(peripherals, state.Peripheral)
} }
// Register the device we just found return peripherals
ds.RegisterDevice(dev) }
func (ds *DeviceService) GetPeripheral(pname string) (peripheral.Peripheral, error) {
return ds.PSS.GetPeripheral(pname)
}
func (ds *DeviceService) PeripheralExists(pname string) bool {
_, err := ds.PSS.GetPeripheral(pname)
return err == nil
}
func (ds *DeviceService) GetRequiredFields() []telemetry.FieldID {
// NOTE: this can be optimized, not that it matters at this stage, but if
// it runs while telemetry is running we want it optimized I guess
seen := make(map[telemetry.FieldID]struct{})
var allFields []telemetry.FieldID
for _, dev := range ds.GetDevices() {
for _, field := range dev.RequiredFields() {
if _, exists := seen[field]; !exists {
seen[field] = struct{}{}
allFields = append(allFields, field)
} }
} }
} }
return allFields
} }
// func (ds *DeviceService) SubscribeFields() error { // Getters [END] ---------------------------------------------------------------
// for _, dev := range ds.Devices {
// fields := dev.RequiredFields()
// }
//
// return nil
// }
func (ds *DeviceService) RegisterDevice(dev peripheral.Peripheral) error {
ds.mu.Lock()
defer ds.mu.Unlock()
if ds.DeviceExists(dev.Name()) {
return ErrPeripheralAlreadyRegistered
}
ds.Devices[dev.Name()] = dev
return nil
}
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) DeviceExists(name string) bool {
if _, ok := ds.Devices[name]; !ok {
return false
}
return true
}
// Actions [START] -------------------------------------------------------------
func (ds *DeviceService) StartStream() { func (ds *DeviceService) StartStream() {
// NOTE: i'm using this pattern a whole lot. Maybe I can create a struct to handle this // NOTE: i'm using this pattern a whole lot. Maybe I can create a struct to handle this
var ctx context.Context var ctx context.Context
ctx, ds.streamCancel = context.WithCancel(context.Background()) ctx, ds.streamCancel = context.WithCancel(context.Background())
ds.PSS.OnStartStream()
go ds.transmit(ctx) go ds.transmit(ctx)
} }
@@ -135,6 +116,8 @@ func (ds *DeviceService) StopStream() {
return return
} }
ds.PSS.OnStopStream()
ds.streamCancel() ds.streamCancel()
ds.streamCancel = nil ds.streamCancel = nil
} }
@@ -144,6 +127,8 @@ func (ds *DeviceService) SetTelemetryChannel(ch <-chan telemetry.TelemetryData)
ds.TelemCh = ch ds.TelemCh = ch
} }
// Actions [END] ---------------------------------------------------------------
// transmit will send the data to the devices themselves // transmit will send the data to the devices themselves
func (ds *DeviceService) transmit(ctx context.Context) { func (ds *DeviceService) transmit(ctx context.Context) {
var isSending atomic.Bool var isSending atomic.Bool
@@ -165,13 +150,28 @@ func (ds *DeviceService) transmit(ctx context.Context) {
// TODO: make a copy of the data and send that copy instead of keeping // TODO: make a copy of the data and send that copy instead of keeping
// the data locked // the data locked
ds.mu.RLock() for _, dev := range ds.PSS.GetStates() {
for _, dev := range ds.Devices { if dev.State != DeviceIsStreaming {
dev.SendData(&data) continue
}
err := dev.Peripheral.SendData(&data)
if err == peripheral.ErrDeviceTimedOut {
ds.onDeviceTimedOut(dev.device.Name)
}
} }
ds.mu.RUnlock()
isSending.Store(false) isSending.Store(false)
} }
} }
} }
// Callbacks [START] -----------------------------------------------------------
func (ds *DeviceService) peripheralConfigured(pname string) {
// We need to retrigger field subscription here
fields := ds.triggerFieldSubscription()
ds.Messages <- fmt.Sprintf("subscribed to fields: %+v\n", fields)
}
// Callbacks [END] -------------------------------------------------------------
+1
View File
@@ -0,0 +1 @@
package services
+489
View File
@@ -0,0 +1,489 @@
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)
onPeripheralConfigured func(string)
// 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()
for _, state := range pss.GetStates() {
if state.State == DeviceIsStreaming {
pss.setDeviceConfigured(state.device.Name)
}
}
}
// "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()
pss.store[pname].State = DeviceIsConfigured
pss.mu.Unlock()
pss.onPeripheralConfigured(pname)
}
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()
}
}
}
+37
View File
@@ -1,2 +1,39 @@
// Package services interacts with the other libraries required for this UI // Package services interacts with the other libraries required for this UI
package services 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
devService.triggerFieldSubscription = telemService.SubscribeToAllFields
// Setup telemetry service callbacks
telemService.getRequiredFields = devService.GetRequiredFields
return &Orchestrator{
DeviceService: devService,
TelemetryService: telemService,
Messages: msg,
}, nil
}
+134 -122
View File
@@ -2,19 +2,24 @@ package services
import ( import (
"context" "context"
"errors"
"log/slog" "log/slog"
"sync" "sync"
"time" "time"
"esdi/providers" "esdi/providers"
"esdi/telemetry"
telem "esdi/telemetry" telem "esdi/telemetry"
) )
var ErrNoActiveProviderAvailable = errors.New("no active provider available")
// TelemetryService will be our base struct to handle telemetry data // 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 // It should hook to a data sink and handle it like iRacing, BeamNG, AC and so on
type TelemetryService struct { type TelemetryService struct {
logger *slog.Logger logger *slog.Logger
devService *DeviceService // Streaming
isStreaming bool
// Concurrency protection // Concurrency protection
mut sync.RWMutex mut sync.RWMutex
activeProvider telem.TelemetryProvider activeProvider telem.TelemetryProvider
@@ -29,16 +34,21 @@ type TelemetryService struct {
cancelMonitor context.CancelFunc cancelMonitor context.CancelFunc
CtxHealthcheck context.Context CtxHealthcheck context.Context
healthCheckCancel context.CancelFunc healthCheckCancel context.CancelFunc
// Callbacks
// Devices data request
// peripheralProvider func() []peripheral.Peripheral
getRequiredFields func() []telemetry.FieldID
} }
func NewTelemetryService(logger *slog.Logger, devServo *DeviceService) *TelemetryService { func NewTelemetryService(
sharedChannel := make(chan string, 10) logger *slog.Logger,
msg chan string,
) *TelemetryService {
newService := &TelemetryService{ newService := &TelemetryService{
logger: logger, logger: logger,
isConnected: false, isConnected: false,
devService: devServo,
listeners: make(map[string]chan telem.TelemetryData), listeners: make(map[string]chan telem.TelemetryData),
Messages: sharedChannel, Messages: msg,
} }
newService.CtxMonitor, newService.cancelMonitor = context.WithCancel(context.Background()) newService.CtxMonitor, newService.cancelMonitor = context.WithCancel(context.Background())
@@ -56,6 +66,7 @@ func (t *TelemetryService) ProviderMonitor(ctx context.Context) {
case <-ticker.C: case <-ticker.C:
slog.Info("checking if provider is still running") slog.Info("checking if provider is still running")
if !t.activeProvider.IsAlive(500 * time.Millisecond) { if !t.activeProvider.IsAlive(500 * time.Millisecond) {
t.Messages <- "Healthcheck on provider failing. Dropping provider.\n"
slog.Warn("provider healthcheck failed") slog.Warn("provider healthcheck failed")
t.dropActiveProvider() t.dropActiveProvider()
t.onProviderHealthCheckFailed() t.onProviderHealthCheckFailed()
@@ -65,32 +76,88 @@ func (t *TelemetryService) ProviderMonitor(ctx context.Context) {
} }
} }
func (t *TelemetryService) onProviderHealthCheckFailed() { func (t *TelemetryService) GetTelemetryProviderName() (string, error) {
// Just restart the whole lookup process if t.activeProvider == nil {
go t.FindProvider(t.CtxMonitor) return "", ErrNoActiveProviderAvailable
}
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
} }
// Create a routine to poll this provider while we wait to start the stream or pause it return t.activeProvider.Name(), nil
t.CtxHealthcheck, t.healthCheckCancel = context.WithCancel(context.Background())
go t.ProviderMonitor(t.CtxHealthcheck)
} }
func (t *TelemetryService) onProviderStopsMidStream() { // Listener Control [START] ----------------------------------------------------
// clear the current provider
// TODO: now we need to also clear the devices to restart everything, func (t *TelemetryService) SubscribeListener(id string, bufferSize int) <-chan telem.TelemetryData {
// if the stream stopped we have to restart the devices and everything t.mut.Lock()
t.logger.Info("cleaning dropped provider and restarting lookup service") defer t.mut.Unlock()
t.dropActiveProvider()
go t.FindProvider(t.CtxMonitor) // NOTE: is this truly necessary?
// return the channel if it already exists
if ch, exists := t.listeners[id]; exists {
return ch
}
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) SubscribeToAllFields() []telem.FieldID {
fields := t.getRequiredFields()
t.logger.Debug("requested fields", "fields", fields)
t.activeProvider.Subscribe(fields)
return fields
}
// 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
} }
// TODO: add some way of retriggering this. Currently it should: // TODO: add some way of retriggering this. Currently it should:
@@ -122,19 +189,46 @@ func (t *TelemetryService) FindProvider(ctx context.Context) {
} }
} }
func (t *TelemetryService) SwitchProvider(newProvider telem.TelemetryProvider) error { // Provider Control [END] ------------------------------------------------------
// Streaming Control [START] ---------------------------------------------------
func (t *TelemetryService) StopStream() {
t.mut.Lock() t.mut.Lock()
defer t.mut.Unlock() defer t.mut.Unlock()
// Clean up the current to be old provider if t.cancelForward != nil {
if t.activeProvider != nil { t.cancelForward()
t.dropActiveProvider() t.cancelForward = nil
} }
// Assign the new provider t.activeProvider.StopStream()
t.activeProvider = newProvider t.isStreaming = false
}
return nil 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
} }
func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan telem.TelemetryData) { func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan telem.TelemetryData) {
@@ -164,94 +258,12 @@ func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan tele
} }
} }
func (t *TelemetryService) dropActiveProvider() { func (t *TelemetryService) IsStreaming() bool {
if t.cancelForward != nil { return t.isStreaming
t.cancelForward()
}
t.activeProvider.StopStream()
t.activeProvider.Close()
t.activeProvider = nil
} }
func (t *TelemetryService) SubscribeListener(id string, bufferSize int) <-chan telem.TelemetryData { // Streaming Control [END] -----------------------------------------------------
t.mut.Lock()
defer t.mut.Unlock()
// NOTE: is this truly necessary? // Callbacks [START] -----------------------------------------------------------
// return the channel if it already exists
if ch, exists := t.listeners[id]; exists {
return ch
}
ch := make(chan telem.TelemetryData, bufferSize) // Callbacks [END] -------------------------------------------------------------
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
}
+39
View File
@@ -0,0 +1,39 @@
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)
}
+8 -75
View File
@@ -1,12 +1,18 @@
package telemetry package telemetry
import ( import (
"math"
"strconv" "strconv"
"sync" "sync"
"time" "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 // 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. // 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. // allow him to make something that allows some data transformation to occur.
@@ -64,7 +70,7 @@ const (
// NOTE: we can optimize this via a special command that says a given piece of data // NOTE: we can optimize this via a special command that says a given piece of data
// is for multiple targets // is for multiple targets
type TelemetryField struct { type TelemetryField struct {
IDs []int16 // Identification for the serial device // IDs []int16 // Identification for the serial device
Type DataType Type DataType
Raw uint64 Raw uint64
Str string // Only to be used with DataTypeSTRING Str string // Only to be used with DataTypeSTRING
@@ -75,49 +81,6 @@ func (tf *TelemetryField) Unused() {
tf.Raw = uint64('-') 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 { func (tf *TelemetryField) String() string {
switch tf.Type { switch tf.Type {
case DataTypeSTRING: case DataTypeSTRING:
@@ -259,13 +222,6 @@ func GetFieldName(id FieldID) string {
var fieldNameToID map[string]FieldID 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) { func GetFieldID(name string) (FieldID, bool) {
id, ok := fieldNameToID[name] id, ok := fieldNameToID[name]
return id, ok return id, ok
@@ -285,26 +241,3 @@ type TelemetryData struct {
func NewTelemetryData() *TelemetryData { func NewTelemetryData() *TelemetryData {
return &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
}
+1 -1
View File
@@ -6,7 +6,7 @@ import "time"
type TelemetryProvider interface { type TelemetryProvider interface {
StopStream() StopStream()
Stream() (<-chan TelemetryData, error) Stream() (<-chan TelemetryData, error)
Subscribe(map[int16]FieldID) Subscribe([]FieldID)
IsAlive(time.Duration) bool IsAlive(time.Duration) bool
Name() string Name() string
Close() Close()
-4
View File
@@ -1,5 +1 @@
package telemetry package telemetry
func Init() {
initFieldNamesMap()
}
@@ -35,6 +35,7 @@ func NewLayoutController(base *Controller, service *services.DeviceService) *Lay
DevService: service, DevService: service,
MoveToolState: &windowManipState{Mode: moveMode}, MoveToolState: &windowManipState{Mode: moveMode},
// SelectedLayout: "beamng.yaml", // SelectedLayout: "beamng.yaml",
// TODO: this can't be here - the service/peripheral needs to know about it
SelectedLayout: "layout.yaml", SelectedLayout: "layout.yaml",
} }
@@ -211,7 +212,7 @@ func (lc *LayoutController) createWindow() {
} }
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
@@ -321,7 +322,7 @@ func (lc *LayoutController) newWindowAction() {
func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) { func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) {
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
@@ -347,7 +348,7 @@ func (lc *LayoutController) displayLoadedLayouts() {
lc.Logger.Debug("We want to view our layout!") lc.Logger.Debug("We want to view our layout!")
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
@@ -392,7 +393,7 @@ func (lc *LayoutController) getCurrentTreeNodeModel() (*tview.TreeNode, int16, e
func (lc *LayoutController) loadLayout() { func (lc *LayoutController) loadLayout() {
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
@@ -404,7 +405,8 @@ func (lc *LayoutController) loadLayout() {
} }
// --- // ---
// We would get the layout path from somewhere but for nots its layout.yaml // TODO: This can happen here, but we need to address how the layout is gotten
// THE UI SHOULD SET STATE IN THE SERVICES ONLY
err = display.LoadLayout(lc.SelectedLayout) err = display.LoadLayout(lc.SelectedLayout)
if err != nil { if err != nil {
lc.Messages <- "failed to load layout: " + err.Error() lc.Messages <- "failed to load layout: " + err.Error()
@@ -416,7 +418,7 @@ func (lc *LayoutController) loadLayout() {
func (lc *LayoutController) unloadLayout() { func (lc *LayoutController) unloadLayout() {
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
@@ -437,7 +439,7 @@ func (lc *LayoutController) unloadLayout() {
func (lc *LayoutController) saveLayout() { func (lc *LayoutController) saveLayout() {
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
@@ -470,7 +472,7 @@ func (lc *LayoutController) deleteWindow() {
wID := curNode.GetReference().(int16) wID := curNode.GetReference().(int16)
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
@@ -52,7 +52,7 @@ func (lc *LayoutController) handleMovementCapture(idx int16,
} }
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME lc.Messages <- "failed to get " + cdashdisplay.NAME
return nil return nil
@@ -86,7 +86,7 @@ func (lc *LayoutController) handleResizeCapture(idx int16,
} }
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME lc.Messages <- "failed to get " + cdashdisplay.NAME
return nil return nil
+11 -12
View File
@@ -16,19 +16,18 @@ type DeviceController struct {
DeviceAPIView *views.DeviceAPIView DeviceAPIView *views.DeviceAPIView
LayoutCtrl *LayoutController LayoutCtrl *LayoutController
StreamCtrl *StreamingCtrl StreamCtrl *StreamingCtrl
DevService *serv.DeviceService Orchestrator *serv.Orchestrator
} }
func NewDeviceController( func NewDeviceController(
base *Controller, base *Controller,
devService *serv.DeviceService, orchestrator *serv.Orchestrator,
telemService *serv.TelemetryService,
) *DeviceController { ) *DeviceController {
mc := &DeviceController{ mc := &DeviceController{
Controller: base, Controller: base,
LayoutCtrl: NewLayoutController(base, devService), LayoutCtrl: NewLayoutController(base, orchestrator.DeviceService),
DevService: devService, Orchestrator: orchestrator,
StreamCtrl: NewStreamingCtrl(base, devService, telemService), StreamCtrl: NewStreamingCtrl(base, orchestrator.DeviceService, orchestrator.TelemetryService),
} }
return mc return mc
@@ -57,7 +56,7 @@ func (mc *DeviceController) setDeviceAPIViewEvents() {
SetInputCapture(func(ev *tcell.EventKey) *tcell.EventKey { SetInputCapture(func(ev *tcell.EventKey) *tcell.EventKey {
switch ev.Rune() { switch ev.Rune() {
case 'r': case 'r':
go mc.DevService.FindDevices() go mc.Orchestrator.DeviceService.FindDevices()
} }
return ev return ev
}) })
@@ -67,8 +66,8 @@ func (mc *DeviceController) AddDeviceAPIListItems() {
mc.DeviceAPIView.DevAPIList. mc.DeviceAPIView.DevAPIList.
AddItem("layout", "build a layout for CDashDisplay", func() { AddItem("layout", "build a layout for CDashDisplay", func() {
// This CDashDisplay specific, only load if we have a CDashDisplay // This CDashDisplay specific, only load if we have a CDashDisplay
if !mc.DevService.DeviceExists(cdashdisplay.NAME) { if !mc.Orchestrator.DeviceService.PeripheralExists(cdashdisplay.NAME) {
mc.DevService.Messages <- "CDashDisplay it not loaded yet\n" mc.Orchestrator.DeviceService.Messages <- "CDashDisplay it not loaded yet\n"
return return
} }
@@ -94,7 +93,7 @@ func (mc *DeviceController) AddDeviceAPIListItems() {
mc.StreamCtrl.StreamView.Flex, mc.StreamCtrl.StreamView.Flex,
) )
mc.StreamCtrl.SetInternalState() // mc.StreamCtrl.SetInternalState()
mc.App.SetFocus(mc.StreamCtrl.StreamView.Options.Form) mc.App.SetFocus(mc.StreamCtrl.StreamView.Options.Form)
}) })
@@ -116,7 +115,7 @@ func (mc *DeviceController) injectControllerCallbacks() {
func (mc *DeviceController) injectChannels() { func (mc *DeviceController) injectChannels() {
go func() { go func() {
for msg := range mc.DevService.Messages { for msg := range mc.Orchestrator.Messages {
mc.PrintToOutputWindow(msg) mc.PrintToOutputWindow(msg)
} }
}() }()
+9 -34
View File
@@ -18,7 +18,7 @@ import (
type StreamingCtrl struct { type StreamingCtrl struct {
*Controller *Controller
Service *services.DeviceService DevService *services.DeviceService
StreamView *views.StreamToolView StreamView *views.StreamToolView
Messages chan string Messages chan string
Internal chan string Internal chan string
@@ -45,27 +45,20 @@ func NewStreamingCtrl(
ctrl := &StreamingCtrl{ ctrl := &StreamingCtrl{
Controller: base, Controller: base,
Service: devService, DevService: devService,
TelemServ: serTelem, TelemServ: serTelem,
Messages: make(chan string, 10), Messages: make(chan string, 10),
Internal: make(chan string, 10), Internal: make(chan string, 10),
TelemetryCh: make(chan telemetry.TelemetryData, 1), TelemetryCh: make(chan telemetry.TelemetryData, 1),
Run: false, Run: false,
StreamView: streamView, StreamView: streamView,
isRunning: false,
} }
ctrl.registerHooks() ctrl.registerHooks()
// ctrl.subscribeListeners()
return ctrl return ctrl
} }
// func (sc *StreamingCtrl) subscribeListeners() {
// // Here I will set a UIDevice
// sc.TelemetryCh = sc.TelemServ.SubscribeListener("UI", 1)
// }
func (sc *StreamingCtrl) registerHooks() { func (sc *StreamingCtrl) registerHooks() {
sc.StreamView.Options.Form.SetInputCapture(func(ev *tcell.EventKey) *tcell.EventKey { sc.StreamView.Options.Form.SetInputCapture(func(ev *tcell.EventKey) *tcell.EventKey {
switch ev.Key() { switch ev.Key() {
@@ -96,11 +89,11 @@ func (sc *StreamingCtrl) registerHooks() {
} }
func (sc *StreamingCtrl) StartStop() { func (sc *StreamingCtrl) StartStop() {
if sc.isRunning { if sc.TelemServ.IsStreaming() {
slog.Info("stopping stream") slog.Info("stopping stream")
sc.TelemServ.StopStream() sc.TelemServ.StopStream()
sc.Service.StopStream() sc.DevService.StopStream()
sc.isRunning = false sc.isRunning = false
return return
@@ -110,9 +103,9 @@ func (sc *StreamingCtrl) StartStop() {
// NOTE: // NOTE:
// Subscribe the only existing device - needs to be discovered by now // Subscribe the only existing device - needs to be discovered by now
slog.Debug("setting the data stream for device servie") slog.Debug("setting the data stream for device servie")
sc.Service.SetTelemetryChannel(sc.TelemServ.SubscribeListener("DeviceService", 1)) sc.DevService.SetTelemetryChannel(sc.TelemServ.SubscribeListener("DeviceService", 1))
dev, err := sc.Service.GetDevice(uidevice.NAME) dev, err := sc.DevService.GetPeripheral(uidevice.NAME)
if err == nil { if err == nil {
if uiDev, ok := dev.(*uidevice.UIDevice); ok { if uiDev, ok := dev.(*uidevice.UIDevice); ok {
sc.TelemetryCh = uiDev.DataChannel() sc.TelemetryCh = uiDev.DataChannel()
@@ -121,7 +114,7 @@ func (sc *StreamingCtrl) StartStop() {
} }
slog.Debug("starting services") slog.Debug("starting services")
sc.Service.StartStream() sc.DevService.StartStream()
sc.TelemServ.StartStream() sc.TelemServ.StartStream()
sc.isRunning = true sc.isRunning = true
@@ -161,24 +154,8 @@ func (sc *StreamingCtrl) updateStream() {
// Performance reasoning: this is not used during the high frequency data transmission // Performance reasoning: this is not used during the high frequency data transmission
// so we can get away with using a map for convenience here // so we can get away with using a map for convenience here
func (sc *StreamingCtrl) SetInternalState() { func (sc *StreamingCtrl) SetInternalState() {
// Acquire the cdashdisplay // fields := sc.TelemServ.SubscribeToAllFields()
// displayIF, err := sc.Service.GetDevice(cdashdisplay.NAME) // sc.Messages <- fmt.Sprintf("Subscribed to fields: %+v\n", fields)
// 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() { func (sc *StreamingCtrl) listenToUIStream() {
@@ -190,8 +167,6 @@ func (sc *StreamingCtrl) listenToUIStream() {
} }
isDrawing.Store(true) isDrawing.Store(true)
// sc.Logger.Debug("got data", "data", msg)
// Capture locally // Capture locally
telemetryMsg := msg telemetryMsg := msg
+5 -9
View File
@@ -23,19 +23,15 @@ func NewControlPanel(logger *slog.Logger) *ControlPanel {
App: tview.NewApplication(), App: tview.NewApplication(),
} }
// NOTE: create our device service here orchestrator, err := services.NewOrchestrator(logger)
devService := services.NewDeviceService(logger.With("service", "DeviceService")) if err != nil {
// TODO: no panic here
telemService := services.NewTelemetryService(logger, devService) panic("failed to create services orchestrator")
if telemService == nil {
panic("failed to create the telemetry service")
} }
go telemService.FindProvider(telemService.CtxMonitor)
return &ControlPanel{ return &ControlPanel{
Controller: baseController, Controller: baseController,
DeviceController: controllers.NewDeviceController(baseController, devService, telemService), DeviceController: controllers.NewDeviceController(baseController, orchestrator),
} }
} }