Author SHA1 Message Date
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
esilva 1ac4ab1865 logger fixes on cdashdisplay device implementation 2026-09-17 23:08:36 +01:00
esilva a6ab64117c this is a doozy - device auto discovery and handling
The main goal was to decouple more things
We got devices being looked for in the background and the stream view
is like a device now too
2026-09-17 23:01:50 +01:00
esilva 1321977f18 updated beamng provider for the new SDK 2026-09-17 23:01:50 +01:00
esilva ffc03e9e49 added the IsAlive check for beamng 2026-09-17 23:01:50 +01:00
esilva 15b997c22f fixed some little UI bugs
- We could open the CDashDisplay specific layout before it existed -
Crash!
- We could open the stream UI before there was a provider- Crash!
2026-09-17 23:01:50 +01:00
esilva 59fa825da1 removing some unecessary logs 2026-09-17 23:01:50 +01:00
esilva dbd4086bf9 Provider loop discovery almost finished
Provider lookup will lookup on startup, do the healthcheck until the
stream starts and also restart on stream closure or provider stall
2026-09-17 23:01:50 +01:00
esilva c1366b010f remove this very verbose log 2026-09-17 23:01:50 +01:00
esilva d6362eeeb5 remove the fatals from here, but I do need to do error handling there 2026-09-17 23:01:50 +01:00
esilva 1b9a5173b6 Updated the function that checks for the sim running and how we stream data 2026-09-17 23:01:50 +01:00
esilva 7675c7e34a fixed the iracing Gear transform 2026-09-17 23:01:50 +01:00
esilva 922c2df88d Don't allow program to crash if there's no active provider on StartStream action 2026-09-17 23:01:50 +01:00
esilva 6b09eb5be0 some fixes regarding iracing provider transforms and working the auto detection 2026-09-17 23:01:50 +01:00
esilva 1d022a4fa8 update the iracing provider to be on the same page as the BeamNG one
Deprecated the Fetch and Transform methods
2026-09-17 23:01:50 +01:00
esilva a251f45859 basic loop to discover providers 2026-09-17 23:01:50 +01:00
esilva afc4b794bb Created a devices service and moved the CDashDisplay service to the sahdow realm 2026-09-17 23:01:50 +01:00
esilva 66012b6904 Merge pull request 'remove old logger. We use actual loggers now' (#10) from remove-old-logger into master
Reviewed-on: #10
2026-09-17 23:00:17 +01:00
esilva a807eb3782 remove old logger. We use actual loggers now 2026-09-17 23:00:02 +01:00
esilva fe86993c6d update the irsdk version 2026-08-24 16:54:36 +01:00
esilva 5ae0542149 Another README edit 2026-07-23 10:47:27 +01:00
esilva 4088e4d8f8 Edited the README 2026-07-23 10:45:25 +01:00
esilva ce7aabfa09 Instead of panicing on unset updaters, set them as unused
This might have impact on different types of data other than a window
set for STRING, but I will worry about that later on
2026-07-08 09:47:45 +01:00
esilva ca57cbd9f0 Small validation for the updaters array in BeamNG 2026-07-07 18:16:52 +01:00
esilva 0887fbfd75 simplified a bit how the TelemetryField updaters work
- now we only use one function to update the TelemetryField instead
of having to fetch data and then update it.

Still need to update the iRacing provider to use this pattern. Way
better and creates better units.

Got rid of that stupid long ass switch statement too
2026-07-07 17:52:47 +01:00
esilva 94ba0e2539 window title reporting 2026-07-07 00:06:59 +01:00
esilva c431c389e3 added a TODO so I won't forget it 2026-07-06 23:47:12 +01:00
esilva 16e6a336b0 Added fields for the BeamNG provider and am working on some changes for the telem
I need to simplify the way telemetry is fetched and updated on the app
side.
Reduce the use of any for example and so on
2026-07-06 23:31:41 +01:00
esilva 8def8cac68 Bumped the BeamNG SDK version 2026-07-06 23:30:40 +01:00
esilva 4e91edc036 small change to the gear transform for BeamNG 2026-07-06 14:43:45 +01:00
esilva abc60a81d3 added an important TODO 2026-06-26 17:47:38 +01:00
esilva 20dbf15adf Support for BeamNG as a data provider 2026-06-25 15:35:57 +01:00
50 changed files with 2347 additions and 1024 deletions
+1
View File
@@ -0,0 +1 @@
# Peripheral Discovery
+31
View File
@@ -0,0 +1,31 @@
# Provider Discovery
```
NewControlPanel()
go telemService.FindProvider() - has a callback for onFind
|
|iterates over the known providers
|
onFind
telemService.SwitchProvider() - activates the found provider
|
|-→ Create a background job while the stream hasn't initiated to listened
| for data, otherwise the connection might die before we start streaming
| ↓
| go telemService.ProviderMonitor()
↓ | |
/-→waits--\ if the provider stops we this healthcheck is stopped
\_________/ we clear the provider and once the stream starts
go back to the
FindProvider()
```
The provider discovery routine starts on `tui/tui.go`.
`FindProvider` is called here and it starts a background job.
On successful discovery the background job calls the `SwitchProvider` method
and dies.
## Provider stalls
A provider stalls once there is no new data.
After stalling
+1
View File
@@ -0,0 +1 @@
# Streaming Flow
+25 -13
View File
@@ -21,30 +21,42 @@ or to view the live data:
Games implemented so far: Games implemented so far:
- [iRacing](https://www.iracing.com/) using the [goirsdk](https://github.com/ESilva15/goirsdk) - [iRacing](https://www.iracing.com/) using the [goirsdk](https://github.com/ESilva15/goirsdk)
Games being implemented: <!-- Games being implemented: -->
- [BeamNG.drive](https://www.beamng.com/game/) using the [gobngsdk](https://github.com/ESilva15/gobngsdk) <!-- - [BeamNG.drive](https://www.beamng.com/game/) using the [gobngsdk](https://github.com/ESilva15/gobngsdk) -->
<!---->
Games to be implemented: <!-- Games to be implemented: -->
- [Assetto Corsa](https://assettocorsa.gg/) <!-- - [Assetto Corsa](https://assettocorsa.gg/) -->
## Roadmap ## Roadmap
- [ ] Implement the interface for a data source - [X] Implement the interface for a data source
- [ ] Finish implementing BeamNG
- [ ] Configuration of the peripherals via ESDI - [ ] Configuration of the peripherals via ESDI
- [ ] Detection of the display - [ ] Detection of the display
- [ ] Fuel Calculator
- [ ] LapTime Calculator - [ ] LapTime Calculator
- [ ] A very long list useful stuff like flags, position, more info about - [ ] A very long list useful stuff like flags, position, more info about
other drivers, track conditions and so on so forth other drivers, track conditions and so on so forth
- [ ] Dynamic data packets
- [ ] More roadmap entries
- [ ] Telemetry analysis tool - [ ] Telemetry analysis tool
- [ ] Better user interface - [ ] Better user interface
## Debugging ## Development
### Freezes: ### Mockservers
To mock [BeamNG.drive](https://www.beamng.com/game/), I built
[BeaMNGMockOg](https://github.com/ESilva15/BeamNGMockOg) (should have though longer
about the name).
You create a recording by launching the BeamNG and then launching the mockserver with:
`BeamNGMockOg record -a 127.0.0.1 -p 4443 -o output.bin`.
To replay the recording:
`BeamNGMockOg replay [--loop] -a 127.0.0.1 -p 4443 -i input.bin`
### Troubleshooting
#### Can't see BeamNG data coming through
Use, for example, `tcpdump -i any udp port <port> -X` to check if any data is available.
### Debugging
#### Freezes:
Using delve: Using delve:
- Launch Terminal1 with `dlv debug . --headless --listen=:2345 -- tui` - Launch Terminal1 with `dlv debug . --headless --listen=:2345 -- tui`
- Launch Terminal2 and connect to with with `dlv connect :2345` - Launch Terminal2 and connect to with with `dlv connect :2345`
@@ -59,7 +71,7 @@ pprof:
- type `web` to view a graph version - type `web` to view a graph version
### Shameless begging ## Shameless begging
Hey, doesn't hurt to try, its free either way: Hey, doesn't hurt to try, its free either way:
[Buy me a coffee!](buymeacoffee.com/ESilva_15) [Buy me a coffee!](buymeacoffee.com/ESilva_15)
-70
View File
@@ -1,70 +0,0 @@
package cdashdisplay
// In this file we will place all structs that are 1:1 representation of the
// types in the device transport layer ->
const (
ShowIDFalse uint8 = 0
ShowIDTrue uint8 = 1
)
const (
WinTypeBASE uint8 = iota
WinTypeBAR
WinTypeSTRING
WinTypeTABLE
)
var WinTYPES = []string{"BASE", "BAR", "STRING", "TABLE"}
type UIDimensions struct {
X0 uint16 `yaml:"X0"`
Y0 uint16 `yaml:"Y0"`
Width uint16 `yaml:"Width"`
Height uint16 `yaml:"Height"`
}
type UIDecorations struct {
BGColour uint16 `yaml:"BGColour"`
FGColour uint16 `yaml:"FGColour"`
TitleColour uint16 `yaml:"TitleColour"`
BorderColour uint16 `yaml:"BorderColour"`
TitleSize uint8 `yaml:"TitleSize"`
TextSize uint8 `yaml:"TextSize"`
HasBorder uint8 `yaml:"HasBorder"`
Padding uint8 `yaml:"Padding"`
}
// NOTE: we can't use FString32 for this - too many bytes
// NOTE: use a bit flags for this options instead
type UIWindowOpts struct {
ShowID uint8 `yaml:"ShowID"`
WinType uint8 `yaml:"WinType"`
PreviewValue FString32 `yaml:"PreviewValue"`
}
type UIWindow struct {
Dims UIDimensions `yaml:"Dims"`
Decor UIDecorations `yaml:"Decor"`
Opts UIWindowOpts `yaml:"Opts"`
Title FString32 `yaml:"Title"`
}
type DesktopUIWindow struct {
UIWindow
UIData DesktopUIData
}
type DesktopUIData struct {
IDX int16 `yaml:"WID"`
TelemetryField string `yaml:"TelemetryField"`
}
type UIWindowUpdatePacket struct {
WinID int16
Window UIWindow
}
// func getUIWindowDTO(w *DesktopWindowData) *UIWindow {
//
// }
+1 -4
View File
@@ -1,20 +1,17 @@
package cmd package cmd
import ( import (
"esdi/logger"
esdi "esdi/oldEsdi" esdi "esdi/oldEsdi"
"github.com/spf13/cobra" "github.com/spf13/cobra"
) )
func liveTelemetryCmdAction(cmd *cobra.Command, args []string) { func liveTelemetryCmdAction(cmd *cobra.Command, args []string) {
log := logger.GetInstance()
ddPort, _ := cmd.Flags().GetString("port") ddPort, _ := cmd.Flags().GetString("port")
outputFile, _ := cmd.Flags().GetString("out") outputFile, _ := cmd.Flags().GetString("out")
sessionFile, _ := cmd.Flags().GetString("session") 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) esdi.RunLiveTelemetry(ddPort, outputFile, sessionFile)
} }
+2 -3
View File
@@ -1,7 +1,6 @@
package cmd package cmd
import ( import (
"esdi/logger"
esdi "esdi/oldEsdi" esdi "esdi/oldEsdi"
// "github.com/ESilva15/goirsdk" // "github.com/ESilva15/goirsdk"
@@ -10,14 +9,14 @@ import (
) )
func offlineTelemetryCmdAction(cmd *cobra.Command, args []string) { func offlineTelemetryCmdAction(cmd *cobra.Command, args []string) {
log := logger.GetInstance() // log := logger.GetInstance()
ddPort, _ := cmd.Flags().GetString("port") ddPort, _ := cmd.Flags().GetString("port")
inFile, _ := cmd.Flags().GetString("in") inFile, _ := cmd.Flags().GetString("in")
outFile, _ := cmd.Flags().GetString("out") outFile, _ := cmd.Flags().GetString("out")
sessionFile, _ := cmd.Flags().GetString("session") 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) esdi.RunOfflineTelemetry(ddPort, inFile, outFile, sessionFile)
} }
+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
@@ -1,10 +1,12 @@
package cdashdisplay package cdashdisplay
import ( import (
"fmt"
"log/slog"
"time"
"esdi/peripheral/communication" "esdi/peripheral/communication"
"esdi/peripheral/communication/packets" "esdi/peripheral/communication/packets"
"fmt"
"time"
"github.com/tarm/serial" "github.com/tarm/serial"
portp "go.bug.st/serial" portp "go.bug.st/serial"
@@ -46,11 +48,11 @@ func findDisplayPort() (*communication.WalkieTalkie, error) {
return nil, err return nil, err
} }
pLogger.Info(fmt.Sprintf("Looking into %v", ports)) slog.Info(fmt.Sprintf("Looking into %v", ports))
var wt *communication.WalkieTalkie var wt *communication.WalkieTalkie
for _, port := range ports { for _, port := range ports {
pLogger.Info(fmt.Sprintf("Trying port %s", port)) slog.Info(fmt.Sprintf("Trying port %s", port))
wt = &communication.WalkieTalkie{ wt = &communication.WalkieTalkie{
Cfg: &serial.Config{ Cfg: &serial.Config{
@@ -60,7 +62,7 @@ func findDisplayPort() (*communication.WalkieTalkie, error) {
}, },
} }
pLogger.Info(fmt.Sprintf("Started probing port %s", port)) slog.Info(fmt.Sprintf("Started probing port %s", port))
probeResult := make(chan error, 1) probeResult := make(chan error, 1)
@@ -76,14 +78,14 @@ func findDisplayPort() (*communication.WalkieTalkie, error) {
err = fmt.Errorf("probe completely hung/timed out: %s", port) err = fmt.Errorf("probe completely hung/timed out: %s", port)
} }
pLogger.Info(fmt.Sprintf("Finished probing port %s", port)) slog.Info(fmt.Sprintf("Finished probing port %s", port))
if err == nil { if err == nil {
pLogger.Info(fmt.Sprintf("Success probing port %s: %+v", port, err)) slog.Info(fmt.Sprintf("Success probing port %s: %+v", port, err))
break break
} }
pLogger.Info(fmt.Sprintf("wasn't port %s", port)) slog.Info(fmt.Sprintf("wasn't port %s", port))
wt = nil wt = nil
} }
@@ -91,6 +93,6 @@ func findDisplayPort() (*communication.WalkieTalkie, error) {
return nil, fmt.Errorf("couldn't find cdashdisplay") return nil, fmt.Errorf("couldn't find cdashdisplay")
} }
pLogger.Info(fmt.Sprintf("found cdashdisplay on port: %s", wt.Cfg.Name)) slog.Info(fmt.Sprintf("found cdashdisplay on port: %s", wt.Cfg.Name))
return wt, nil return wt, nil
} }
@@ -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"
@@ -19,12 +21,6 @@ import (
"gopkg.in/yaml.v3" "gopkg.in/yaml.v3"
) )
var pLogger *slog.Logger
func SetLogger(l *slog.Logger) {
pLogger = l
}
// I have to move this to some kind of configuration place // I have to move this to some kind of configuration place
const ( const (
layoutsDir = "./layouts/" layoutsDir = "./layouts/"
@@ -112,27 +108,70 @@ 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
func NewCDashDisplay() (*CDashDisplay, error) { func NewCDashDisplay() (*CDashDisplay, error) {
// Look for the port // Look for the port
p, err := findDisplayPort() p, err := findDisplayPort()
if err != nil { if err != nil {
pLogger.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 (cds *CDashDisplay) Close() error {
// if cds.WT != nil {
// cds.Close()
// }
return nil
}
func (d *CDashDisplay) SendCommand() { func (d *CDashDisplay) SendCommand() {
} }
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) {
bytes, err := helper.StructToBytes(win.UIWindow) bytes, err := helper.StructToBytes(win.UIWindow)
if err != nil { if err != nil {
@@ -148,9 +187,12 @@ func (d *CDashDisplay) CreateWindow(win *DesktopUIWindow) (*DesktopUIWindow, err
win.UIData.IDX = wID.ID win.UIData.IDX = wID.ID
pLogger.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
} }
@@ -176,12 +218,14 @@ func (d *CDashDisplay) UpdateWindow(win *DesktopUIWindow) error {
// I get it and update it in the controller // I get it and update it in the controller
// I send the pointer here // I send the pointer here
// -> it should be the same pointer then right? // -> it should be the same pointer then right?
pLogger.Debug(fmt.Sprintf("PreUpdate ID: %p", win)) slog.Debug(fmt.Sprintf("PreUpdate ID: %p", win))
d.State.Layout.Windows[win.UIData.IDX] = win d.State.Layout.Windows[win.UIData.IDX] = win
pLogger.Debug(fmt.Sprintf("PostUpdate ID: %p", win)) slog.Debug(fmt.Sprintf("PostUpdate ID: %p", win))
// 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
} }
@@ -204,14 +248,15 @@ 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
} }
func (d *CDashDisplay) updateWindowDimensions(win *UIWindow, packet UpdateDimsPacket) error { func (d *CDashDisplay) updateWindowDimensions(win *UIWindow, packet UpdateDimsPacket) error {
pLogger.Debug(fmt.Sprintf("UPDATE: %v", packet)) slog.Debug(fmt.Sprintf("UPDATE: %v", packet))
bytes, err := helper.StructToBytes(packet) bytes, err := helper.StructToBytes(packet)
if err != nil { if err != nil {
@@ -225,9 +270,9 @@ func (d *CDashDisplay) updateWindowDimensions(win *UIWindow, packet UpdateDimsPa
} }
// Nothing bad happened afaik // Nothing bad happened afaik
pLogger.Debug(fmt.Sprintf("cur dims: %v", win.Dims)) slog.Debug(fmt.Sprintf("cur dims: %v", win.Dims))
win.Dims = packet.Dims win.Dims = packet.Dims
pLogger.Debug(fmt.Sprintf("new dims: %v", win.Dims)) slog.Debug(fmt.Sprintf("new dims: %v", win.Dims))
return nil return nil
} }
@@ -343,30 +388,30 @@ func (d *CDashDisplay) LoadLayout(layoutName string) error {
func (d *CDashDisplay) UnloadLayout() error { func (d *CDashDisplay) UnloadLayout() error {
var err error var err error
for _, w := range d.State.Layout.Windows { for _, w := range d.State.Layout.Windows {
pLogger.Debug(fmt.Sprintf("= Removing %d ==============================================", slog.Debug(fmt.Sprintf("= Removing %d ==============================================",
w.UIData.IDX)) w.UIData.IDX))
err = d.DestroyWindow(w.UIData.IDX) err = d.DestroyWindow(w.UIData.IDX)
time.Sleep(75 * time.Millisecond) time.Sleep(75 * time.Millisecond)
if err != nil { if err != nil {
pLogger.Error(fmt.Sprintf("failed to destroy window: %+v", err)) slog.Error(fmt.Sprintf("failed to destroy window: %+v", err))
// NOTE: Add a way to handle multiple errors ? // NOTE: Add a way to handle multiple errors ?
return err return err
} }
pLogger.Debug(fmt.Sprintf("= Removing %d ==============================================", slog.Debug(fmt.Sprintf("= Removing %d ==============================================",
w.UIData.IDX)) w.UIData.IDX))
} }
return nil return nil
} }
func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) { func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) error {
packet := data.Pack() 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 := ""
@@ -376,7 +421,6 @@ func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) {
curStr += fmt.Sprintf("%02x ", byte) curStr += fmt.Sprintf("%02x ", byte)
if byteCount == 8 { if byteCount == 8 {
pLogger.Debug(curStr)
curStr = "" curStr = ""
byteCount = 0 byteCount = 0
} }
@@ -385,6 +429,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
} }
+28
View File
@@ -0,0 +1,28 @@
package cdashdisplay
import (
"esdi/peripheral/devices"
"esdi/telemetry"
)
// TODO: I believe we don't need the #esdi/peripheral/devices thing anymore
const (
ID = devices.CDashDisplayDevID
NAME = devices.CDashDisplayDevName
)
func (cds *CDashDisplay) Name() string {
return NAME
}
func (cds *CDashDisplay) RequiredFields() []telemetry.FieldID {
fields := make([]telemetry.FieldID, 0, 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)
}
}
return fields
}
@@ -1,6 +1,9 @@
package cdashdisplay package cdashdisplay
import "fmt" import (
"fmt"
"log/slog"
)
type LayoutTree struct { type LayoutTree struct {
Windows map[int16]*DesktopUIWindow `yaml:"Windows"` Windows map[int16]*DesktopUIWindow `yaml:"Windows"`
@@ -13,13 +16,13 @@ func NewLayoutTree() *LayoutTree {
} }
func (l *LayoutTree) AddWindow(w *DesktopUIWindow) { func (l *LayoutTree) AddWindow(w *DesktopUIWindow) {
pLogger.Debug(fmt.Sprintf("adding window '%d' - %v", w.UIData.IDX)) slog.Debug(fmt.Sprintf("adding window '%d' - %v", w.UIData.IDX))
l.Windows[w.UIData.IDX] = w l.Windows[w.UIData.IDX] = w
pLogger.Debug(fmt.Sprintf("new map - %v", l.Windows)) slog.Debug(fmt.Sprintf("new map - %v", l.Windows))
} }
func (l *LayoutTree) RemoveWindow(idx int16) { func (l *LayoutTree) RemoveWindow(idx int16) {
pLogger.Debug(fmt.Sprintf("removing window '%d'", idx)) slog.Debug(fmt.Sprintf("removing window '%d'", idx))
delete(l.Windows, idx) delete(l.Windows, idx)
pLogger.Debug(fmt.Sprintf("new map - %v", l.Windows)) slog.Debug(fmt.Sprintf("new map - %v", l.Windows))
} }
+141
View File
@@ -0,0 +1,141 @@
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 ->
const (
ShowIDFalse uint8 = 0
ShowIDTrue uint8 = 1
)
const (
WinTypeBASE uint8 = iota
WinTypeBAR
WinTypeSTRING
WinTypeTABLE
)
var WinTYPES = []string{"BASE", "BAR", "STRING", "TABLE"}
type UIDimensions struct {
X0 uint16 `yaml:"X0"`
Y0 uint16 `yaml:"Y0"`
Width uint16 `yaml:"Width"`
Height uint16 `yaml:"Height"`
}
type UIDecorations struct {
BGColour uint16 `yaml:"BGColour"`
FGColour uint16 `yaml:"FGColour"`
TitleColour uint16 `yaml:"TitleColour"`
BorderColour uint16 `yaml:"BorderColour"`
TitleSize uint8 `yaml:"TitleSize"`
TextSize uint8 `yaml:"TextSize"`
HasBorder uint8 `yaml:"HasBorder"`
Padding uint8 `yaml:"Padding"`
}
// NOTE: we can't use FString32 for this - too many bytes
// NOTE: use a bit flags for this options instead
type UIWindowOpts struct {
ShowID uint8 `yaml:"ShowID"`
WinType uint8 `yaml:"WinType"`
PreviewValue FString32 `yaml:"PreviewValue"`
}
type UIWindow struct {
Dims UIDimensions `yaml:"Dims"`
Decor UIDecorations `yaml:"Decor"`
Opts UIWindowOpts `yaml:"Opts"`
Title FString32 `yaml:"Title"`
}
type DesktopUIWindow struct {
UIWindow
UIData DesktopUIData
}
type DesktopUIData struct {
IDX int16 `yaml:"WID"`
TelemetryField string `yaml:"TelemetryField"`
}
type UIWindowUpdatePacket struct {
WinID int16
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
}
+43
View File
@@ -0,0 +1,43 @@
// Package devices is a peripheral factory
package devices
import (
"esdi/devices/cdashdisplay"
"esdi/devices/uidevice"
"esdi/peripheral"
)
type Device struct {
Name string
Discover func() (peripheral.Peripheral, error)
}
var List map[string]*Device = map[string]*Device{
uidevice.NAME: {
Name: uidevice.NAME,
Discover: DiscoverUIDevice,
},
cdashdisplay.NAME: {
Name: cdashdisplay.NAME,
Discover: DiscoverCDashDisplay,
},
}
func DiscoverUIDevice() (peripheral.Peripheral, error) {
uidev, err := uidevice.NewUIDevice()
if err != nil {
return nil, err
}
return uidev, nil
}
func DiscoverCDashDisplay() (peripheral.Peripheral, error) {
// Create a cdashdisplay
display, err := cdashdisplay.NewCDashDisplay()
if err != nil {
return nil, err
}
return display, nil
}
+56
View File
@@ -0,0 +1,56 @@
package uidevice
import (
"esdi/peripheral"
"esdi/telemetry"
)
type UIDevice struct {
dataChan chan telemetry.TelemetryData
}
const NAME = "UIView"
func NewUIDevice() (peripheral.Peripheral, error) {
return &UIDevice{
dataChan: make(chan telemetry.TelemetryData, 1),
}, nil
}
func (uid *UIDevice) Close() error {
if uid.dataChan != nil {
close(uid.dataChan)
}
return nil
}
func (uid *UIDevice) SendData(data *telemetry.TelemetryData) error {
if data == nil {
return peripheral.ErrInvalidData
}
select {
case uid.dataChan <- *data:
default:
// Drop frame if buffer is full
}
return nil
}
func (uid *UIDevice) Name() string {
return NAME
}
func (uid *UIDevice) DataChannel() <-chan telemetry.TelemetryData {
return uid.dataChan
}
func (uid *UIDevice) RequiredFields() []telemetry.FieldID {
return []telemetry.FieldID{
telemetry.Speed,
telemetry.Gear,
telemetry.RPM,
}
}
+3 -3
View File
@@ -1,11 +1,11 @@
module esdi module esdi
go 1.25.5 go 1.27.0
require ( require (
github.com/ESilva15/ESgoRepl v0.1.0 github.com/ESilva15/ESgoRepl v0.1.0
github.com/ESilva15/gobngsdk v0.0.2 github.com/ESilva15/gobngsdk v1.1.3
github.com/ESilva15/goirsdk v0.2.11 github.com/ESilva15/goirsdk v0.3.0
github.com/arl/statsviz v0.8.0 github.com/arl/statsviz v0.8.0
github.com/gdamore/tcell/v2 v2.8.1 github.com/gdamore/tcell/v2 v2.8.1
github.com/rivo/tview v0.42.0 github.com/rivo/tview v0.42.0
+4 -4
View File
@@ -1,9 +1,9 @@
github.com/ESilva15/ESgoRepl v0.1.0 h1:hOVcRBfdP8Yw1kLr2tFoM3SetV7vxFswKjzfAtq8UgM= github.com/ESilva15/ESgoRepl v0.1.0 h1:hOVcRBfdP8Yw1kLr2tFoM3SetV7vxFswKjzfAtq8UgM=
github.com/ESilva15/ESgoRepl v0.1.0/go.mod h1:O2JQMyEgHy0zf9UvzvLNvp2w95T2uazc3torrzRmG9Y= github.com/ESilva15/ESgoRepl v0.1.0/go.mod h1:O2JQMyEgHy0zf9UvzvLNvp2w95T2uazc3torrzRmG9Y=
github.com/ESilva15/gobngsdk v0.0.2 h1:N6stNOE14Wg80BAZMFu/Qd9Qa/J0DmExCfdPoDf7lco= github.com/ESilva15/gobngsdk v1.1.3 h1:CFaz2KjHKnBLrQuEN1QHOd1761cWM/rmTwLpSLrgbe4=
github.com/ESilva15/gobngsdk v0.0.2/go.mod h1:cKLaZRgM0tGGXDvosaSjHnxeyC5Y6rg7zCLqEUJnOzw= github.com/ESilva15/gobngsdk v1.1.3/go.mod h1:cKLaZRgM0tGGXDvosaSjHnxeyC5Y6rg7zCLqEUJnOzw=
github.com/ESilva15/goirsdk v0.2.11 h1:C/HuO9xGmRJ01vMR0KRjrhcMEfq1Mea1rQK0CICaOU4= github.com/ESilva15/goirsdk v0.3.0 h1:6D95Avq7chikGHwEUKh7Y/HaLHOWQ75Qxaz8oMDSfrw=
github.com/ESilva15/goirsdk v0.2.11/go.mod h1:5borQbw+L4fe9b58JFHRHZwrzz0sMVAteuYbRnFOc0Q= github.com/ESilva15/goirsdk v0.3.0/go.mod h1:5borQbw+L4fe9b58JFHRHZwrzz0sMVAteuYbRnFOc0Q=
github.com/arl/statsviz v0.8.0 h1:O6GjjVxEDxcByAucOSl29HaGYLXsuwA3ujJw8H9E7/U= github.com/arl/statsviz v0.8.0 h1:O6GjjVxEDxcByAucOSl29HaGYLXsuwA3ujJw8H9E7/U=
github.com/arl/statsviz v0.8.0/go.mod h1:XlrbiT7xYT03xaW9JMMfD8KFUhBOESJwfyNJu83PbB0= github.com/arl/statsviz v0.8.0/go.mod h1:XlrbiT7xYT03xaW9JMMfD8KFUhBOESJwfyNJu83PbB0=
github.com/cpuguy83/go-md2man/v2 v2.0.4/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46tRHOmNcaadrF8o= github.com/cpuguy83/go-md2man/v2 v2.0.4/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46tRHOmNcaadrF8o=
+232 -88
View File
@@ -1,5 +1,53 @@
Windows: Windows:
1:
uiwindow:
Dims:
X0: 170
Y0: 136
Width: 40
Height: 55
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 4
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: '>'
Title: ->
uidata:
WID: 1
TelemetryField: Right Indicator
2: 2:
uiwindow:
Dims:
X0: 341
Y0: 417
Width: 40
Height: 55
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 4
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: A
Title: ABS
uidata:
WID: 2
TelemetryField: ABS Dash Light
3:
uiwindow: uiwindow:
Dims: Dims:
X0: 470 X0: 470
@@ -21,81 +69,9 @@ Windows:
PreviewValue: "9999" PreviewValue: "9999"
Title: Dig Tacho Title: Dig Tacho
uidata: uidata:
WID: 2 WID: 3
TelemetryField: RPM TelemetryField: RPM
4: 4:
uiwindow:
Dims:
X0: 321
Y0: 195
Width: 134
Height: 80
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 7
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: "234"
Title: SPEEDO
uidata:
WID: 4
TelemetryField: Speed
5:
uiwindow:
Dims:
X0: 590
Y0: 367
Width: 100
Height: 50
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 3
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: "101.5"
Title: Water T
uidata:
WID: 5
TelemetryField: Water Temperature
10:
uiwindow:
Dims:
X0: 470
Y0: 132
Width: 169
Height: 50
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 3
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: "123.4"
Title: Fuel Level
uidata:
WID: 10
TelemetryField: Fuel Level
11:
uiwindow: uiwindow:
Dims: Dims:
X0: 321 X0: 321
@@ -117,33 +93,33 @@ Windows:
PreviewValue: R PreviewValue: R
Title: Box Title: Box
uidata: uidata:
WID: 11 WID: 4
TelemetryField: Gear TelemetryField: Gear
16: 5:
uiwindow: uiwindow:
Dims: Dims:
X0: 694 X0: 80
Y0: 423 Y0: 136
Width: 100 Width: 40
Height: 50 Height: 55
Decor: Decor:
BGColour: 4161 BGColour: 4161
FGColour: 65535 FGColour: 65535
TitleColour: 65535 TitleColour: 65535
BorderColour: 63488 BorderColour: 63488
TitleSize: 2 TitleSize: 2
TextSize: 3 TextSize: 4
HasBorder: 1 HasBorder: 1
Padding: 0 Padding: 0
Opts: Opts:
ShowID: 0 ShowID: 0
WinType: 2 WinType: 2
PreviewValue: "4.56" PreviewValue: <
Title: Oil P Title: <-
uidata: uidata:
WID: 16 WID: 5
TelemetryField: Oil Pressure TelemetryField: Left Indicator
17: 6:
uiwindow: uiwindow:
Dims: Dims:
X0: 694 X0: 694
@@ -165,9 +141,177 @@ Windows:
PreviewValue: "101.4" PreviewValue: "101.4"
Title: Oil T Title: Oil T
uidata: uidata:
WID: 17 WID: 6
TelemetryField: Oil Temperature TelemetryField: Oil Temperature
20: 7:
uiwindow:
Dims:
X0: 386
Y0: 417
Width: 40
Height: 55
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 4
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: T
Title: TC
uidata:
WID: 7
TelemetryField: Traction Control Light
8:
uiwindow:
Dims:
X0: 431
Y0: 417
Width: 40
Height: 55
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 4
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: B
Title: Bat
uidata:
WID: 8
TelemetryField: Battery Light
9:
uiwindow:
Dims:
X0: 321
Y0: 195
Width: 134
Height: 80
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 7
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: "234"
Title: SPEEDO
uidata:
WID: 9
TelemetryField: Speed
10:
uiwindow:
Dims:
X0: 590
Y0: 367
Width: 100
Height: 50
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 3
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: "101.5"
Title: Water T
uidata:
WID: 10
TelemetryField: Water Temperature
11:
uiwindow:
Dims:
X0: 470
Y0: 132
Width: 169
Height: 50
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 3
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: "123.4"
Title: Fuel Level
uidata:
WID: 11
TelemetryField: Fuel Level
12:
uiwindow:
Dims:
X0: 694
Y0: 423
Width: 100
Height: 50
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 3
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: "4.56"
Title: Oil P
uidata:
WID: 12
TelemetryField: Oil Pressure
13:
uiwindow:
Dims:
X0: 296
Y0: 417
Width: 40
Height: 55
Decor:
BGColour: 4161
FGColour: 65535
TitleColour: 65535
BorderColour: 63488
TitleSize: 2
TextSize: 4
HasBorder: 1
Padding: 0
Opts:
ShowID: 0
WinType: 2
PreviewValue: P
Title: HB
uidata:
WID: 13
TelemetryField: Parking Brake Dash Light
14:
uiwindow: uiwindow:
Dims: Dims:
X0: 142 X0: 142
@@ -189,5 +333,5 @@ Windows:
PreviewValue: "5678" PreviewValue: "5678"
Title: TACHO Title: TACHO
uidata: uidata:
WID: 20 WID: 14
TelemetryField: RPM TelemetryField: RPM
-22
View File
@@ -1,22 +0,0 @@
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
}
+5 -4
View File
@@ -10,12 +10,16 @@ import (
"esdi/cmd" "esdi/cmd"
"esdi/config" "esdi/config"
"esdi/telemetry"
"github.com/arl/statsviz" "github.com/arl/statsviz"
) )
// NOTE: in go we have the init() function. Its a function that runs before
// everything else in a package. Make use of that
func initApplication() { func initApplication() {
fmt.Fprint(os.Stdout, "\x1b]0;ESDI\x07")
// 1. This is the first initialization setup we do so we can log // 1. This is the first initialization setup we do so we can log
err := setupLogger() err := setupLogger()
if err != nil { if err != nil {
@@ -37,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 {
+1 -1
View File
@@ -54,7 +54,7 @@ func RunLiveTelemetry(port string, output string, session string) {
// log.Fatalf("Failed to create iRacing interface: %v", err) // log.Fatalf("Failed to create iRacing interface: %v", err)
// } // }
irsdk, err := goirsdk.Init(nil, output, session) irsdk, err := goirsdk.Init(goirsdk.Options{})
if err != nil { if err != nil {
log.Fatalf("Failed to create irsdk instance: %v\n", err) log.Fatalf("Failed to create irsdk instance: %v\n", err)
} }
+17 -16
View File
@@ -108,6 +108,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,7 +128,9 @@ 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 {
@@ -176,21 +187,11 @@ func (wt *WalkieTalkie) readPacket(resp packets.Packet) error {
// return nil // return nil
// } // }
func (wt *WalkieTalkie) SendCommand(cmd types.Command, payload any, func (wt *WalkieTalkie) SendCommand(
responseBody packets.Packet) error { cmd types.Command,
// Prepare the header payload any,
// header := header{ responseBody packets.Packet,
// StartByte: constvar.StartOfText, ) error {
// 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 {
+3 -2
View File
@@ -6,8 +6,9 @@ import "esdi/peripheral/types"
// IDs for our devices. They need to be correctly mapped on the devices // IDs for our devices. They need to be correctly mapped on the devices
// themselves so we can discover them // themselves so we can discover them
const ( const (
CDashDisplayDevID = 0x01 CDashDisplayDevID = 0x01
ESBtnBoxDevID = 0x02 CDashDisplayDevName = "CDashDisplay"
ESBtnBoxDevID = 0x02
) )
// DeviceMap maps the implemented devices // DeviceMap maps the implemented devices
+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")
)
+10 -1
View File
@@ -2,9 +2,11 @@
package peripheral package peripheral
import ( import (
"esdi/peripheral/devices"
"fmt" "fmt"
"path/filepath" "path/filepath"
"esdi/peripheral/devices"
"esdi/telemetry"
) )
type PeripheralType string type PeripheralType string
@@ -13,6 +15,13 @@ const (
DisplayPeripheral PeripheralType = "display" DisplayPeripheral PeripheralType = "display"
) )
type Peripheral interface {
Name() string
SendData(*telemetry.TelemetryData) error
RequiredFields() []telemetry.FieldID
Close() error
}
type PeripheralDeviceClerk struct { type PeripheralDeviceClerk struct {
// mu sync.RWMutex // mu sync.RWMutex
Devices map[uint8]*PeripheralDevice Devices map[uint8]*PeripheralDevice
+206 -18
View File
@@ -1,43 +1,231 @@
// Package beamng is the BeamNG.drive data provider
package beamng package beamng
import ( import (
"context"
"fmt" "fmt"
"log/slog"
"sync"
"time"
"esdi/telemetry"
bngsdk "github.com/ESilva15/gobngsdk" bngsdk "github.com/ESilva15/gobngsdk"
) )
// BeamNG is the concrete implementation of the TelemetryProvider interface
// for BeamNG.drive
// NOTE: document this please. What is a TelemetryData????
type BeamNG struct { type BeamNG struct {
SDK *bngsdk.BNGSDK logger *slog.Logger
SDK *bngsdk.BeamNGSDK
og *bngsdk.Outgauge
// data handling
mut sync.Mutex
data *telemetry.TelemetryData
updaters [telemetry.MaxFields]func(*telemetry.TelemetryField)
// stream control
wg sync.WaitGroup
streamCancel context.CancelFunc
// timing
ticker *time.Ticker
} }
const ( const (
NAME = "BeamNG.drive" NAME = "BeamNG.drive"
) )
func Init(ip string, port int) (BeamNG, error) { func NewBeamNGProvider(logger *slog.Logger, opts *bngsdk.Options) (*BeamNG, error) {
var err error beam, err := bngsdk.NewBngSDK(*opts)
sdk, err := bngsdk.Init(ip, port)
if err != nil { if err != nil {
return BeamNG{}, err return &BeamNG{}, err
} }
return BeamNG{SDK: &sdk}, nil provider := &BeamNG{
} logger: logger.With("TelemetryProvider", NAME),
data: telemetry.NewTelemetryData(),
// GetData will retrieve a given field by its name from the OutGauge data SDK: beam,
func (b *BeamNG) GetData(fieldName string) (interface{}, error) { og: &bngsdk.Outgauge{},
if val, ok := b.SDK.DataDict[fieldName]; ok { ticker: time.NewTicker(time.Second / 60),
return val, nil
} }
return nil, fmt.Errorf("key `%s` doesn't exist", fieldName) provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){
telemetry.Speed: provider.updateSpeed,
telemetry.Gear: provider.updateGear,
telemetry.RPM: provider.updateRPM,
telemetry.FuelLevel: provider.fuelLevel,
// Engine Data
telemetry.OilPress: provider.oilPressure,
telemetry.OilTemp: provider.oilTemp,
telemetry.WaterTemp: provider.engTemp,
// Electrics (dash lights and so on)
telemetry.PitSpeedLimiter: provider.pitSpeedLimiter,
telemetry.LeftIndicator: provider.leftIndicator,
telemetry.RightIndicator: provider.rightIndicator,
telemetry.Hazards: provider.unused,
telemetry.ABSWarningLight: provider.absLight,
telemetry.ParkingBrakeLight: provider.handbrakeLight,
telemetry.TCLight: provider.tcLight,
telemetry.BatteryLight: provider.batteryLight,
}
// Set the unset telemetry fields on the updaters as unused fields
for k := range int(telemetry.MaxFields) {
if provider.updaters[k] == nil {
provider.updaters[k] = provider.unused
}
}
return provider, nil
} }
func (b *BeamNG) UpdateData() error { func (b *BeamNG) Close() {
return b.SDK.ReadData() b.SDK.Close()
} }
func (b *BeamNG) GetSessionInfo() (interface{}, error) { func (b *BeamNG) Name() string {
return nil, fmt.Errorf("Not implemented") return NAME
}
func (b *BeamNG) IsAlive(timeout time.Duration) bool {
_, err := b.SDK.Update()
return err == nil
}
func (b *BeamNG) StopStream() {
if b.streamCancel == nil {
return
}
b.streamCancel()
}
func (b *BeamNG) Stream() (<-chan telemetry.TelemetryData, error) {
var ctx context.Context
ctx, b.streamCancel = context.WithCancel(context.Background())
// Start the stream
ch := b.stream(ctx)
return ch, nil
}
func (b *BeamNG) Subscribe(requestFields []telemetry.FieldID) {
// NOTE: document how the Subscribe funtion works
slog.Debug(fmt.Sprintf("Len Req: %d\n", len(requestFields)))
b.data.ActiveBinds = make([]telemetry.BoundField, 0, len(requestFields))
// First we must add the virtual fields
// we will add their dependencies and the primitives to a slice
pendingBinds := make([]telemetry.FieldID, telemetry.MaxFields)
for _, id := range requestFields {
switch id {
case telemetry.RPMStateColour:
b.data.VirtualBinds = append(b.data.VirtualBinds, telemetry.NewRPMLights())
case telemetry.FCCurrentLap:
b.data.VirtualBinds = append(b.data.VirtualBinds,
telemetry.NewFuelCalculator(slog.Default().WithGroup("FUEL CALC")))
default:
// primitive telemetry field
pendingBinds = append(pendingBinds, id)
}
}
boundCheck := make(map[telemetry.FieldID]bool)
// Now that we know all the fields we need to bind we follow the binding procedure
for _, id := range pendingBinds {
// Check if we already bound this FieldID
if boundCheck[id] {
continue
}
binding := telemetry.BoundField{
ID: id,
}
b.data.ActiveBinds = append(b.data.ActiveBinds, binding)
boundCheck[id] = true
}
slog.Debug(fmt.Sprintf("Subscribed: %+v\n", b.data.ActiveBinds))
}
// Internal
func (b *BeamNG) readData() {
slog.Debug("READING THIS DATA")
ogSnapshot, err := b.SDK.Update()
slog.Debug("THE DATA WAS READ")
if err != nil {
slog.Error("Error getting data", "error", err)
return
}
b.mut.Lock()
b.og = ogSnapshot
defer b.mut.Unlock()
// Read 1 to 1 data
slog.Debug("Reading normal data binds")
for _, bind := range b.data.ActiveBinds {
slog.Debug("Current bind: ", "id", bind.ID)
b.updaters[bind.ID](&b.data.Values[bind.ID])
}
// Set up virtual binds
slog.Debug("Entering virtual binds loop")
for _, vBind := range b.data.VirtualBinds {
// NOTE: delete the logs here, they are really bad
slog.Debug("Processing virtual binds")
vBind.Process(b.data)
}
b.data.PenultimateDataPoll = b.data.LastDataPoll
b.data.LastDataPoll = time.Now()
}
func (b *BeamNG) stream(ctx context.Context) <-chan telemetry.TelemetryData {
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 {
case <-ctx.Done():
return
default:
}
select {
case <-ctx.Done():
return
case <-b.ticker.C:
slog.Debug("READING DATA")
b.readData()
slog.Debug("READ DATA")
// Publish data
select {
case outCh <- *b.data:
slog.Debug("PUBLISHED DATA")
default:
// skip this data, don't allow publishers to lag behind
}
}
}
}()
return outCh
} }
+113
View File
@@ -0,0 +1,113 @@
package beamng
import (
"strconv"
conv "esdi/conversions"
"esdi/telemetry"
)
const (
LapTimeFormatStr = "04:05.000"
)
func (b *BeamNG) unused(out *telemetry.TelemetryField) {
out.Unused()
}
func (b *BeamNG) updateSpeed(out *telemetry.TelemetryField) {
out.Type = telemetry.DataTypeUINT16
out.Raw = uint64(conv.MsToKph(b.og.Speed))
}
func (b *BeamNG) updateGear(out *telemetry.TelemetryField) {
out.Type = telemetry.DataTypeSTRING
// NOTE: stupid idea but we can cache these values
out.Str = strconv.Itoa(int(b.og.Gear))
}
func (b *BeamNG) updateRPM(out *telemetry.TelemetryField) {
out.Type = telemetry.DataTypeUINT16
out.Raw = uint64(uint16(b.og.RPM))
}
func (b *BeamNG) fuelLevel(out *telemetry.TelemetryField) {
telemetry.FloatToStringTransform(b.og.Fuel, out)
}
func (b *BeamNG) oilPressure(out *telemetry.TelemetryField) {
telemetry.FloatToStringTransform(b.og.OilPressure, out)
}
func (b *BeamNG) oilTemp(out *telemetry.TelemetryField) {
telemetry.FloatToStringTransform(b.og.OilTemp, out)
}
func (b *BeamNG) engTemp(out *telemetry.TelemetryField) {
telemetry.FloatToStringTransform(b.og.EngTemp, out)
}
// NOTE: find how to empty this
func (b *BeamNG) pitSpeedLimiter(out *telemetry.TelemetryField) {
}
func (b *BeamNG) leftIndicator(out *telemetry.TelemetryField) {
chr := ' '
if b.og.LeftIndicator() {
chr = '<'
}
out.Type = telemetry.DataTypeCHAR
out.Raw = uint64(chr)
}
func (b *BeamNG) rightIndicator(out *telemetry.TelemetryField) {
chr := ' '
if b.og.RightIndicator() {
chr = '>'
}
out.Type = telemetry.DataTypeCHAR
out.Raw = uint64(chr)
}
func (b *BeamNG) absLight(out *telemetry.TelemetryField) {
chr := ' '
if b.og.ABS() {
chr = 'A'
}
out.Type = telemetry.DataTypeCHAR
out.Raw = uint64(chr)
}
func (b *BeamNG) handbrakeLight(out *telemetry.TelemetryField) {
chr := ' '
if b.og.Handbrake() {
chr = 'P'
}
out.Type = telemetry.DataTypeCHAR
out.Raw = uint64(chr)
}
func (b *BeamNG) tcLight(out *telemetry.TelemetryField) {
chr := ' '
if b.og.TractionControl() {
chr = 'T'
}
out.Type = telemetry.DataTypeCHAR
out.Raw = uint64(chr)
}
func (b *BeamNG) batteryLight(out *telemetry.TelemetryField) {
chr := ' '
if b.og.BatteryLight() {
chr = 'B'
}
out.Type = telemetry.DataTypeCHAR
out.Raw = uint64(chr)
}
+33
View File
@@ -0,0 +1,33 @@
package beamng
import (
"log/slog"
"net"
"time"
)
func IsRunning() bool {
// TODO: The address should be loaded from some type of configuration
addr, err := net.ResolveUDPAddr("udp", "127.0.0.1:4444")
if err != nil {
return false
}
conn, err := net.ListenUDP("udp", addr)
defer conn.Close()
if err != nil {
slog.Debug("failed to listen to udp socket: " + err.Error())
return false
}
buf := make([]byte, 512)
conn.SetReadDeadline(time.Now().Add(50 * time.Millisecond))
n, _, err := conn.ReadFromUDP(buf)
if err != nil {
return false
}
// Think of a better number or something
return n >= 80
}
-72
View File
@@ -1,72 +0,0 @@
package iracing
import (
telem "esdi/telemetry"
"strconv"
"time"
)
const (
LapTimeFormatStr = "04:05.000"
)
func LapTimeTransform(v any, out *telem.TelemetryField) {
lapTimeInSeconds := v.(float32)
if lapTimeInSeconds < 0 {
lapTimeInSeconds = 0
}
wholeSeconds := int64(lapTimeInSeconds)
lapTime := time.Unix(wholeSeconds, int64((lapTimeInSeconds-float32(wholeSeconds))*1e9))
out.Type = telem.DataTypeSTRING
out.Str = lapTime.Format(LapTimeFormatStr)
}
func EmptyTransform(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeCHAR
out.Raw = uint64('-')
}
func PitSpeedLimiterTransform(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeSTRING
if v.(bool) {
out.Str = "PIT"
} else {
out.Str = " "
}
}
func UInt8Transform(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeUINT8
if v == nil {
out.Raw = 0
return
}
out.Raw = uint64(v.(int))
}
func FloatToStringTransform(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeSTRING
if v == nil {
out.Str = "inv"
return
}
out.Str = strconv.FormatFloat(float64(v.(float32)), 'f', 1, 32)
}
func FloatToUInt8Transform(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeUINT8
if v == nil {
out.Raw = uint64(0)
return
}
out.Raw = uint64(v.(float32))
}
+20 -41
View File
@@ -2,44 +2,23 @@ package iracing
// This file maps the data from the desktop provider data structure to iRacing // This file maps the data from the desktop provider data structure to iRacing
import ( // TODO: pass these to the new Updaters system
"esdi/telemetry" // var internalToSDKFieldNames = map[telemetry.FieldID]string{
) // // Tire data
// telemetry.LFtempL: "LFtempCL",
var internalToSDKFieldNames = map[telemetry.FieldID]string{ // telemetry.LFtempM: "LFtempCM",
telemetry.Speed: "Speed", // telemetry.LFtempR: "LFtempCR",
telemetry.Gear: "Gear", // telemetry.RFtempL: "RFtempCL",
telemetry.RPM: "RPM", // telemetry.RFtempM: "RFtempCM",
telemetry.FuelLevel: "FuelLevel", // telemetry.RFtempR: "RFtempCR",
// Engine Data // telemetry.LRtempL: "LRtempCL",
telemetry.OilPress: "OilPress", // telemetry.LRtempM: "LRtempCM",
telemetry.OilTemp: "OilTemp", // telemetry.LRtempR: "LRtempCR",
telemetry.WaterTemp: "WaterTemp", // telemetry.RRtempL: "RRtempCL",
// Engine Warnings // telemetry.RRtempM: "RRtempCM",
telemetry.PitSpeedLimiter: "irsdk_pitSpeedLimiter", // telemetry.RRtempR: "RRtempCR",
// Ajudstements // // Session Data
telemetry.BrakeBias: "dcBrakeBias", // telemetry.SessionTime: "SessionTime",
telemetry.ABSSetting: "dcABS", // telemetry.ReplaySessionTime: "ReplaySessionTime",
telemetry.TCSetting: "dcTractionControl", // telemetry.Empty: "empty",
telemetry.ThrottleSetting: "dcThrottleShape", // }
// Lap Data
telemetry.LapLastLapTime: "LapLastLapTime",
telemetry.LapNumber: "Lap",
// Tire data
telemetry.LFtempL: "LFtempCL",
telemetry.LFtempM: "LFtempCM",
telemetry.LFtempR: "LFtempCR",
telemetry.RFtempL: "RFtempCL",
telemetry.RFtempM: "RFtempCM",
telemetry.RFtempR: "RFtempCR",
telemetry.LRtempL: "LRtempCL",
telemetry.LRtempM: "LRtempCM",
telemetry.LRtempR: "LRtempCR",
telemetry.RRtempL: "RRtempCL",
telemetry.RRtempM: "RRtempCM",
telemetry.RRtempR: "RRtempCR",
// Session Data
telemetry.SessionTime: "SessionTime",
telemetry.ReplaySessionTime: "ReplaySessionTime",
telemetry.Empty: "empty",
}
+130 -161
View File
@@ -7,13 +7,10 @@ import (
"context" "context"
"fmt" "fmt"
"log/slog" "log/slog"
"os"
"strconv"
"sync" "sync"
"time" "time"
conv "esdi/conversions" "esdi/telemetry"
telem "esdi/telemetry"
"github.com/ESilva15/goirsdk" "github.com/ESilva15/goirsdk"
) )
@@ -28,99 +25,154 @@ type IRacing struct {
SDK *goirsdk.IBT SDK *goirsdk.IBT
// Data Handling // Data Handling
mut sync.Mutex mut sync.Mutex
data *telem.TelemetryData data *telemetry.TelemetryData
updaters [telemetry.MaxFields]func(*telemetry.TelemetryField)
// Timing information // Timing information
ticker *time.Ticker // ticker will keep polling intervals constant ticker *time.Ticker // ticker will keep polling intervals constant
// Stream // Stream
streamCh chan telem.TelemetryData wg sync.WaitGroup
// streamCh chan telemetry.TelemetryData
streamCancel context.CancelFunc streamCancel context.CancelFunc
} }
func NewIRacingProvider( func NewIRacingProvider(
logger *slog.Logger, logger *slog.Logger,
source string, opts goirsdk.Options,
telemOut string,
yamlOut string,
) (*IRacing, error) { ) (*IRacing, error) {
var err error var err error
// Open the input file if provided - otherwise live telemetry was requested sdk, err := goirsdk.Init(opts)
// Maybe this can be changed so we don't have to run it with these ifs but by configuring our
// provider
var file goirsdk.Reader = nil
if source != "" {
file, err = os.Open(source)
if err != nil {
return &IRacing{}, err
// log.Fatalf("Failed to open IBT file: %v", err)
}
}
sdk, err := goirsdk.Init(file, telemOut, yamlOut)
if err != nil { if err != nil {
logger.Error("failed to open the IRSDK instance") logger.Error("failed to open the IRSDK instance")
return &IRacing{}, err return &IRacing{}, err
} }
return &IRacing{ provider := &IRacing{
logger: logger, logger: logger,
SDK: sdk, SDK: sdk,
data: telem.NewTelemetryData(), data: telemetry.NewTelemetryData(),
streamCh: make(chan telem.TelemetryData, 1), // TODO: make this configurable from the user side
// NOTE: This is because I stupidly recorded a test IBT file in 240 ticker: time.NewTicker(time.Second / 60),
ticker: time.NewTicker(time.Second / 240), }
}, nil
provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){
telemetry.Speed: provider.speed,
telemetry.Gear: provider.gear,
telemetry.RPM: provider.rpm,
telemetry.FuelLevel: provider.fuelLevel,
// Engine Data
telemetry.OilPress: provider.oilPress,
telemetry.OilTemp: provider.oilTemp,
telemetry.WaterTemp: provider.waterTemp,
// EngineWarnings
telemetry.PitSpeedLimiter: provider.pitSpeedLimiter,
// Adjustements
telemetry.BrakeBias: provider.brakeBias,
telemetry.ABSSetting: provider.absSetting,
telemetry.TCSetting: provider.tcSetting,
telemetry.ThrottleSetting: provider.throttleSetting,
// Lap Data
telemetry.LapLastLapTime: provider.lapTime,
telemetry.LapNumber: provider.lapNumber,
// case telemetry.LFtempM:
// binding.Transform = func(v any, out *telemetry.TelemetryField) {
// out.Type = telemetry.DataTypeSTRING
// out.Str = strconv.FormatFloat(float64(v.(float32)), 'f', 1, 32)
// }
// case telemetry.SessionTime:
// binding.Transform = func(v any, out *telemetry.TelemetryField) {
// out.Type = telemetry.DataTypeSTRING
// out.Str = strconv.FormatFloat(v.(float64), 'f', 1, 32)
// }
// case telemetry.ReplaySessionTime:
// binding.Transform = func(v any, out *telemetry.TelemetryField) {
// out.Type = telemetry.DataTypeSTRING
// out.Str = strconv.FormatFloat(v.(float64), 'f', 1, 32)
// }
// case telemetry.Empty:
// binding.Transform = telemetry.EmptyTransform
// }
}
// Set the unset telemetry fields on the updaters as unused fields
for k := range int(telemetry.MaxFields) {
if provider.updaters[k] == nil {
provider.updaters[k] = provider.unused
}
}
return provider, nil
}
func (i *IRacing) Close() {
// Need to find a way of gracefully closing the channel
// close(i.streamCh)
i.SDK.Close()
i.ticker.Stop()
} }
func (i *IRacing) isDataAvailable() bool { func (i *IRacing) isDataAvailable() bool {
// Its offline telemetry, data must be available // Its offline telemetry, data must be available
if i.SDK.File != nil { if i.SDK.File == nil {
return true return false
} }
// Check if live telemetry is on // Check if live telemetry is on
if i.SDK.IsConnected() { if !i.SDK.IsConnected() {
return true return false
} }
return false 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 i.wg.Done()
defer close(outCh)
// Put this into the configuration file
consecutiveTimeouts := 0
maxTimeouts := 30
dataEvTimeout := 100
for { for {
// Explicitly intercpt cancellation // Explicitly intercept cancellation
select { select {
case <-ctx.Done(): case <-ctx.Done():
return return
default: default:
} }
// We start by checking if we do or do not have data available if i.SDK.CheckForDataEvent(time.Duration(dataEvTimeout) * time.Millisecond) {
if !i.isDataAvailable() { consecutiveTimeouts = 0
continue
}
select {
case <-ctx.Done():
return
case <-i.ticker.C:
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
} }
} else {
consecutiveTimeouts++
if consecutiveTimeouts >= maxTimeouts {
i.logger.Info("Telemetry stream stalled")
if i.streamCancel != nil {
i.streamCancel()
}
return
}
} }
} }
}() }()
return outCh
} }
func (i *IRacing) readData() { func (i *IRacing) readData() {
@@ -134,18 +186,13 @@ func (i *IRacing) readData() {
return return
} }
// Read 1 to 1 data // Read 1 to 1 data using our updaters
for _, b := range i.data.ActiveBinds { for _, bind := range i.data.ActiveBinds {
v := i.SDK.Vars.Vars[b.Key].Value i.updaters[bind.ID](&i.data.Values[bind.ID])
b.Transform(v, &i.data.Values[b.ID])
} }
// Set up virtual binds // Set up virtual binds
i.logger.Debug("Entering virtual binds loop")
for _, vBind := range i.data.VirtualBinds { for _, vBind := range i.data.VirtualBinds {
// NOTE: delete the logs here, they are really bad
i.logger.Debug("Processing virtual binds")
vBind.Process(i.data) vBind.Process(i.data)
} }
@@ -157,14 +204,14 @@ func (i *IRacing) readData() {
// Stream returns a channel that we will use to funnel the telemetry data back to the // Stream returns a channel that we will use to funnel the telemetry data back to the
// UI, which then should broadcast it to the devices // UI, which then should broadcast it to the devices
func (i *IRacing) Stream() (<-chan telem.TelemetryData, error) { func (i *IRacing) Stream() (<-chan telemetry.TelemetryData, error) {
var ctx context.Context var ctx context.Context
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() {
@@ -173,34 +220,33 @@ func (i *IRacing) StopStream() {
} }
i.streamCancel() i.streamCancel()
i.wg.Wait()
i.streamCancel = nil i.streamCancel = nil
} }
func (i *IRacing) Subscribe(requestFields map[int16]telem.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([]telem.BoundField, 0, len(requestFields)) i.data.ActiveBinds = make([]telemetry.BoundField, 0, len(requestFields))
// First we must add the virtual fields // First we must add the virtual fields
// we will add their dependencies and the primitives to a slice // we will add their dependencies and the primitives to a slice
pendingBinds := make([]telem.FieldID, telem.MaxFields) pendingBinds := make([]telemetry.FieldID, 0, telemetry.MaxFields)
for winID, id := range requestFields {
i.data.Values[id].IDs = append(i.data.Values[id].IDs, winID)
for _, id := range requestFields {
switch id { switch id {
case telem.RPMStateColour: case telemetry.RPMStateColour:
i.data.VirtualBinds = append(i.data.VirtualBinds, telem.NewRPMLights()) i.data.VirtualBinds = append(i.data.VirtualBinds, telemetry.NewRPMLights())
case telem.FCCurrentLap: case telemetry.FCCurrentLap:
i.data.VirtualBinds = append(i.data.VirtualBinds, i.data.VirtualBinds = append(i.data.VirtualBinds,
telem.NewFuelCalculator(i.logger.WithGroup("FUEL CALC"))) telemetry.NewFuelCalculator(i.logger.WithGroup("FUEL CALC")))
default: default:
// primitive telemetry field // primitive telemetry field
pendingBinds = append(pendingBinds, id) pendingBinds = append(pendingBinds, id)
} }
} }
boundCheck := make(map[telem.FieldID]bool) boundCheck := make(map[telemetry.FieldID]bool)
// Now that we know all the fields we need to bind we follow the binding procedure // Now that we know all the fields we need to bind we follow the binding procedure
for _, id := range pendingBinds { for _, id := range pendingBinds {
@@ -209,97 +255,8 @@ func (i *IRacing) Subscribe(requestFields map[int16]telem.FieldID) {
continue continue
} }
// Translate the UI FieldIDs to this provider's field names binding := telemetry.BoundField{
sdkKey, ok := internalToSDKFieldNames[id] ID: id,
if !ok {
i.logger.Debug("failed to get internal id")
// Need to find a way to pass a message saying something wasn't right
continue
}
binding := telem.BoundField{
Key: sdkKey,
ID: id,
}
switch id {
case telem.Speed:
binding.Transform = func(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeUINT16
out.Raw = uint64(conv.MsToKph(v.(float32)))
}
case telem.Gear:
binding.Transform = func(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeCHAR
// out.Raw = uint64(v.(int))
gear := 0
if val, ok := v.(int32); ok {
gear = int(val)
} else if val, ok := v.(int); ok {
gear = val
}
switch {
case gear == 0:
out.Raw = uint64('N') // ASCII 78
case gear < 0:
out.Raw = uint64('R') // ASCII 82
case gear > 0 && gear < 10:
// Quickest way to turn 1 into '1', 2 into '2', etc.
// ASCII '0' is 48, so 48 + 1 = 49 ('1')
out.Raw = uint64('0' + gear)
default:
out.Raw = uint64('?') // Fallback
}
}
case telem.RPM:
binding.Transform = func(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeUINT16
out.Raw = uint64(uint16(v.(float32)))
}
case telem.FuelLevel:
binding.Transform = FloatToStringTransform
// Engine Data
case telem.OilPress:
binding.Transform = FloatToStringTransform
case telem.OilTemp:
binding.Transform = FloatToStringTransform
case telem.WaterTemp:
binding.Transform = FloatToStringTransform
// Something else
case telem.PitSpeedLimiter:
binding.Transform = PitSpeedLimiterTransform
// Adjustements
case telem.BrakeBias:
binding.Transform = FloatToStringTransform
case telem.ABSSetting:
binding.Transform = FloatToUInt8Transform
case telem.TCSetting:
binding.Transform = FloatToUInt8Transform
case telem.ThrottleSetting:
binding.Transform = FloatToUInt8Transform
case telem.LFtempM:
binding.Transform = func(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeSTRING
out.Str = strconv.FormatFloat(float64(v.(float32)), 'f', 1, 32)
}
case telem.SessionTime:
binding.Transform = func(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeSTRING
out.Str = strconv.FormatFloat(v.(float64), 'f', 1, 32)
}
case telem.ReplaySessionTime:
binding.Transform = func(v any, out *telem.TelemetryField) {
out.Type = telem.DataTypeSTRING
out.Str = strconv.FormatFloat(v.(float64), 'f', 1, 32)
}
case telem.Empty:
binding.Transform = EmptyTransform
case telem.LapLastLapTime:
binding.Transform = LapTimeTransform
case telem.LapNumber:
binding.Transform = UInt8Transform
} }
i.data.ActiveBinds = append(i.data.ActiveBinds, binding) i.data.ActiveBinds = append(i.data.ActiveBinds, binding)
@@ -308,3 +265,15 @@ func (i *IRacing) Subscribe(requestFields map[int16]telem.FieldID) {
i.logger.Debug(fmt.Sprintf("Subscribed: %+v\n", i.data.ActiveBinds)) i.logger.Debug(fmt.Sprintf("Subscribed: %+v\n", i.data.ActiveBinds))
} }
func (i *IRacing) Name() string {
return NAME
}
func (i *IRacing) IsAlive(timeout time.Duration) bool {
if !i.SDK.CheckForDataEvent(timeout) {
return false
}
return true
}
+149
View File
@@ -0,0 +1,149 @@
package iracing
import (
"log/slog"
"time"
conv "esdi/conversions"
"esdi/telemetry"
"github.com/ESilva15/goirsdk"
)
const (
LapTimeFormatStr = "04:05.000"
)
func getVar[T any](
vars map[string]goirsdk.Var,
key string,
fallback T,
logger *slog.Logger,
) T {
raw, ok := vars[key]
if !ok {
logger.Error("key not present in map", "key", key)
return fallback
}
val, ok := raw.Value.(T)
if !ok {
logger.Error("could not cast value", "key", key, "val", raw)
return fallback
}
return val
}
func (i *IRacing) unused(out *telemetry.TelemetryField) {
out.Unused()
}
func (i *IRacing) speed(out *telemetry.TelemetryField) {
speed := getVar(i.SDK.Vars.Vars, "Speed", float32(0), i.logger)
out.Type = telemetry.DataTypeUINT16
out.Raw = uint64(conv.MsToKph(speed))
}
func (i *IRacing) gear(out *telemetry.TelemetryField) {
gear := getVar(i.SDK.Vars.Vars, "Gear", int(0), i.logger)
out.Type = telemetry.DataTypeCHAR
switch {
case gear == 0:
out.Raw = uint64('N')
case gear == -1:
out.Raw = uint64('R')
case gear > 0:
out.Raw = uint64('0' + gear)
default:
out.Raw = uint64('?') // Fallback
}
}
func (i *IRacing) rpm(out *telemetry.TelemetryField) {
rpm := getVar(i.SDK.Vars.Vars, "RPM", float32(0), i.logger)
out.Type = telemetry.DataTypeUINT16
out.Raw = uint64(uint16(rpm))
}
func (i *IRacing) fuelLevel(out *telemetry.TelemetryField) {
fuelLevel := getVar(i.SDK.Vars.Vars, "FuelLevel", float32(-1.0), i.logger)
telemetry.FloatToStringTransform(fuelLevel, out)
}
// EngineWarnings
func (i *IRacing) pitSpeedLimiter(out *telemetry.TelemetryField) {
out.Type = telemetry.DataTypeSTRING
out.Str = " "
if i.SDK.PitSpeedLimiter() {
out.Str = "PIT"
}
}
// EngineData ↓
func (i *IRacing) oilPress(out *telemetry.TelemetryField) {
oilPress := getVar(i.SDK.Vars.Vars, "OilPress", float32(-1.0), i.logger)
telemetry.FloatToStringTransform(oilPress, out)
}
func (i *IRacing) oilTemp(out *telemetry.TelemetryField) {
oilTemp := getVar(i.SDK.Vars.Vars, "OilTemp", float32(-1.0), i.logger)
telemetry.FloatToStringTransform(oilTemp, out)
}
func (i *IRacing) waterTemp(out *telemetry.TelemetryField) {
waterTemp := getVar(i.SDK.Vars.Vars, "WaterTemp", float32(-1.0), i.logger)
telemetry.FloatToStringTransform(waterTemp, out)
}
// EngineData ↑
// Adjustments ↓
func (i *IRacing) brakeBias(out *telemetry.TelemetryField) {
bb := getVar(i.SDK.Vars.Vars, "dcBrakeBias", float32(-1.0), i.logger)
telemetry.FloatToStringTransform(bb, out)
}
func (i *IRacing) absSetting(out *telemetry.TelemetryField) {
abs := getVar(i.SDK.Vars.Vars, "dcABS", float32(-1.0), i.logger)
telemetry.FloatToUInt8Transform(abs, out)
}
func (i *IRacing) tcSetting(out *telemetry.TelemetryField) {
tc := getVar(i.SDK.Vars.Vars, "dcTractionControl", float32(-1.0), i.logger)
telemetry.FloatToUInt8Transform(tc, out)
}
func (i *IRacing) throttleSetting(out *telemetry.TelemetryField) {
throttle := getVar(i.SDK.Vars.Vars, "dcThrottleShape", float32(-1.0), i.logger)
telemetry.FloatToUInt8Transform(throttle, out)
}
// Adjustments ↑
// Laps ↓
func (i *IRacing) lapTime(out *telemetry.TelemetryField) {
lapTimeInSeconds := getVar(i.SDK.Vars.Vars, "LapLastLapTime", float32(0), i.logger)
if lapTimeInSeconds < 0 {
lapTimeInSeconds = 0
}
wholeSeconds := int64(lapTimeInSeconds)
lapTime := time.Unix(wholeSeconds, int64((lapTimeInSeconds-float32(wholeSeconds))*1e9))
out.Type = telemetry.DataTypeSTRING
out.Str = lapTime.Format(LapTimeFormatStr)
}
func (i *IRacing) lapNumber(out *telemetry.TelemetryField) {
lapNumber := getVar(i.SDK.Vars.Vars, "Lap", int(0), i.logger)
telemetry.UInt8Transform(lapNumber, out)
}
// Laps ↑
+32
View File
@@ -0,0 +1,32 @@
package iracing
import (
"time"
"github.com/ESilva15/goirsdk"
eventutils "github.com/ESilva15/goirsdk/eventutils"
)
func IsRunning() bool {
irUtils, err := eventutils.Init()
if err != nil {
// we need error validation or something here
return false
}
defer irUtils.Close()
if err := irUtils.OpenEvent(goirsdk.IRSDK_DATAVALIDEVENTNAME); err != nil {
// we need error validation or something here
return false
}
// We now check for some consecutive data events
for range 3 {
if !irUtils.CheckValidDataEvent(1 * time.Second) {
// slog.Debug("Timed out waiting for DataValidEvent")
return false
}
}
return true
}
+38 -9
View File
@@ -7,28 +7,57 @@ import (
"esdi/providers/beamng" "esdi/providers/beamng"
"esdi/providers/iracing" "esdi/providers/iracing"
"esdi/telemetry" "esdi/telemetry"
bngsdk "github.com/ESilva15/gobngsdk"
"github.com/ESilva15/goirsdk"
) )
// Make this be some kind of struct where we can access a function that returns // Make this be some kind of struct where we can access a function that returns
// the selected provider by its name // the selected provider by its name
type Provider struct { type Provider struct {
Name string Name string
Provider telemetry.TelemetryProvider NewProvider func(*slog.Logger) (telemetry.TelemetryProvider, error)
IsRunning func() bool // To check if this provider is up and running
} }
var Providers = map[string]Provider{ var Providers = map[string]Provider{
beamng.NAME: { beamng.NAME: {
Name: beamng.NAME, Name: beamng.NAME,
NewProvider: NewBeamNGProvider,
IsRunning: beamng.IsRunning,
}, },
iracing.NAME: { iracing.NAME: {
Name: iracing.NAME, Name: iracing.NAME,
NewProvider: NewLiveIRacingProvider,
IsRunning: iracing.IsRunning,
}, },
} }
func NewIRacingProvider(logger *slog.Logger, source string, func NewLiveIRacingProvider(logger *slog.Logger) (telemetry.TelemetryProvider, error) {
telemOut string, yamlOut string, provider, err := iracing.NewIRacingProvider(logger, goirsdk.Options{
) telemetry.TelemetryProvider { Logger: logger,
provider, _ := iracing.NewIRacingProvider(logger, source, "", "") SourceType: goirsdk.SharedMemoryFile,
})
if err != nil {
logger.Error("failed to create iRacing provider", "err", err)
return nil, err
}
return provider return provider, nil
}
func NewBeamNGProvider(logger *slog.Logger) (telemetry.TelemetryProvider, error) {
// TODO: these should come from some kind of config
provider, err := beamng.NewBeamNGProvider(logger, &bngsdk.Options{
Logger: logger.With("TelemetryProvider", beamng.NAME),
SourceType: bngsdk.UDPData,
ImportUDPAddress: "127.0.0.1",
ImportUDPPort: 4444,
})
if err != nil {
logger.Error("failed to create BeamNG provider", "err", err)
return nil, err
}
return provider, nil
} }
-148
View File
@@ -1,148 +0,0 @@
package services
import (
"context"
"fmt"
"log/slog"
"sync/atomic"
"esdi/cdashdisplay"
helper "esdi/helpers"
"esdi/peripheral"
"esdi/telemetry"
)
type CDashService struct {
Logger *slog.Logger
CDash *cdashdisplay.CDashDisplay
DevClerk *peripheral.PeripheralDeviceClerk
Messages chan string
// Telemetry Channel
streamCancel context.CancelFunc
TelemCh <-chan telemetry.TelemetryData
}
func NewCDashService(logger *slog.Logger) *CDashService {
sharedChannel := make(chan string, 10)
return &CDashService{
Logger: logger,
CDash: nil,
DevClerk: peripheral.NewPeripheralDeviceClerk(),
Messages: sharedChannel,
}
}
func (cds *CDashService) FindDevice() {
cds.Messages <- "looking for cdash display...\n"
cds.Logger.Info("Looking for CDashDisplay")
cdashdisplay.SetLogger(cds.Logger.With("[device]", "cdashdisplay"))
display, err := cdashdisplay.NewCDashDisplay()
if err != nil {
cds.Logger.Info("didn't find cdashdisplay")
cds.Messages <- "didn't find cdash display\n"
return
}
cds.CDash = display
cds.Logger.Info("found cdashdisplay on: " + display.WT.Cfg.Name)
cds.Messages <- "found cdashdisplay on: " + display.WT.Cfg.Name + "\n"
}
func (cds *CDashService) CreateWindow(
win *cdashdisplay.DesktopUIWindow,
) (*cdashdisplay.DesktopUIWindow, error) {
updatedWindow, err := cds.CDash.CreateWindow(win)
if err != nil {
return nil, err
}
return updatedWindow, nil
}
func (cds *CDashService) LoadLayout(layoutPath string) error {
return cds.CDash.LoadLayout(layoutPath)
}
func (cds *CDashService) SaveLayout(layoutPath string) error {
return cds.CDash.SaveLayout(layoutPath)
}
func (cds *CDashService) UnloadLayout() error {
return cds.CDash.UnloadLayout()
}
func (cds *CDashService) UpdateWindow(win *cdashdisplay.DesktopUIWindow) error {
cds.Messages <- fmt.Sprintf("Updating a window:\n%+v\n", win)
return cds.CDash.UpdateWindow(win)
}
func (cds *CDashService) DeleteWindow(idx int16) error {
return cds.CDash.DestroyWindow(idx)
}
func (cds *CDashService) ResizeWindow(idx int16, vec *helper.Vector) error {
err := cds.CDash.ResizeWindow(idx, vec)
if err != nil {
return err
}
return nil
}
func (cds *CDashService) MoveWindow(idx int16, vec *helper.Vector) error {
err := cds.CDash.MoveWindow(idx, vec)
if err != nil {
return err
}
return nil
}
func (cds *CDashService) SetTelemetryChannel(ch <-chan telemetry.TelemetryData) {
cds.TelemCh = ch
}
func (cds *CDashService) StartStream() {
// NOTE: i'm using this pattern a whole lot. Maybe I can create a struct to handle this
var ctx context.Context
ctx, cds.streamCancel = context.WithCancel(context.Background())
go cds.transmit(ctx)
}
func (cds *CDashService) StopStream() {
if cds.streamCancel == nil {
return
}
cds.streamCancel()
cds.streamCancel = nil
}
// INTERNAL
func (cds *CDashService) transmit(ctx context.Context) {
var isSending atomic.Bool
for {
select {
case <-ctx.Done():
return
case data, ok := <-cds.TelemCh:
if !ok {
return
}
if isSending.Load() {
continue
}
isSending.Store(true)
cds.CDash.SendData(&data)
isSending.Store(false)
}
}
}
-9
View File
@@ -1,9 +0,0 @@
package services
// 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 {
}
+131
View File
@@ -0,0 +1,131 @@
package services
import (
"context"
"log/slog"
"sync/atomic"
"esdi/devices"
"esdi/peripheral"
"esdi/telemetry"
)
// DeviceService is the API for the peripherals
type DeviceService struct {
Logger *slog.Logger
// Device discovery
PSS *PeripheralStateStore // Store to track peripheral state
ctxDiscovery context.Context
ctxDiscoveryCancel context.CancelFunc
// Strem handling
streamCancel context.CancelFunc
TelemCh <-chan telemetry.TelemetryData
// Output
Messages chan string
}
func NewDeviceService(logger *slog.Logger) *DeviceService {
sharedChannel := make(chan string, 10)
dev := &DeviceService{
PSS: NewPeripheralStateStore(logger.With("Service", "PeripheralStateStore"), devices.List),
Logger: logger,
Messages: sharedChannel,
}
// Start the routine that looks for devices - should always be running in the background
// Create a routine to poll this provider while we wait to start the stream or pause it
dev.ctxDiscovery, dev.ctxDiscoveryCancel = context.WithCancel(context.Background())
go dev.FindDevices()
return dev
}
// Getters [START] -------------------------------------------------------------
func (ds *DeviceService) GetDevices() []peripheral.Peripheral {
snapshot := ds.PSS.GetStates()
peripherals := make([]peripheral.Peripheral, 0, len(snapshot))
for _, state := range snapshot {
if state.State != DeviceIsConnected {
continue
}
peripherals = append(peripherals, state.Peripheral)
}
return peripherals
}
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
}
// 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())
go ds.transmit(ctx)
}
func (ds *DeviceService) StopStream() {
if ds.streamCancel == nil {
return
}
ds.streamCancel()
ds.streamCancel = nil
}
// SetTelemetryChannel sets the TelemCh to the passed channel
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
for {
select {
case <-ctx.Done():
return
case data, ok := <-ds.TelemCh:
if !ok {
return
}
if isSending.Load() {
continue
}
isSending.Store(true)
// 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 != DeviceIsConnected {
continue
}
err := dev.Peripheral.SendData(&data)
if err == peripheral.ErrDeviceTimedOut {
ds.onDeviceTimedOut(dev.device.Name)
}
}
isSending.Store(false)
}
}
}
+222
View File
@@ -0,0 +1,222 @@
package services
import (
"errors"
"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")
)
type DeviceState = uint8
const (
DeviceTimedOut uint8 = iota
DeviceIsConnected
DeviceIsDisconnected
DeviceReconnected
)
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
}
func NewPeripheralStateStore(
nLogger *slog.Logger,
devList map[string]*devices.Device,
) *PeripheralStateStore {
store := PeripheralStateStore{
Logger: nLogger,
store: make(map[string]*PeripheralState),
}
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
}
// "Events" [START] ------------------------------------------------------------
func (ds *DeviceService) onDeviceTimedOut(pname string) {
// We need to deregister the device
ds.PSS.setDeviceTimedOut(pname)
}
// "Events" [END] --------------------------------------------------------------
// Device State Handling [START] -----------------------------------------------
func (pss *PeripheralStateStore) setDeviceConnected(pname string, per peripheral.Peripheral) {
pss.Logger.Info("found device", "device", pname)
pss.mu.Lock()
defer pss.mu.Unlock()
pss.store[pname].Peripheral = per
pss.store[pname].State = DeviceIsConnected
}
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
}
// Device State Handling [END] -------------------------------------------------
// Device Handling [START] -----------------------------------------------------
func (pss *PeripheralStateStore) HandleDeviceState() {
peripherals := pss.GetStates()
for pName, pState := range peripherals {
switch pState.State {
case DeviceIsDisconnected:
pss.Logger.Debug("looking for device", "name", pName)
dev, err := pState.device.Discover()
if err != nil {
continue
}
// Register the device we just found
pss.setDeviceConnected(pName, dev)
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)
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)
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)
}
}
}
// 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(2 * time.Second)
defer ticker.Stop()
for {
select {
case <-ds.ctxDiscovery.Done():
return
case <-ticker.C:
ds.PSS.HandleDeviceState()
}
}
}
+204 -69
View File
@@ -4,91 +4,71 @@ import (
"context" "context"
"log/slog" "log/slog"
"sync" "sync"
"time"
"esdi/providers" "esdi/providers"
"esdi/telemetry"
telem "esdi/telemetry" telem "esdi/telemetry"
) )
// 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
cdash *CDashService devService *DeviceService
// Streaming
isStreaming bool
// Concurrency protection // Concurrency protection
mut sync.RWMutex mut sync.RWMutex
ativeProvider telem.TelemetryProvider activeProvider telem.TelemetryProvider
isConnected bool
// Channel for the UI // Channel for the UI
listeners map[string]chan telem.TelemetryData listeners map[string]chan telem.TelemetryData
uiOutCh chan telem.TelemetryData
cancelForward context.CancelFunc cancelForward context.CancelFunc
// Output window
Messages chan string
// Cancel looking for providers
CtxMonitor context.Context
cancelMonitor context.CancelFunc
CtxHealthcheck context.Context
healthCheckCancel context.CancelFunc
} }
func NewTelemetryService(logger *slog.Logger, cdash *CDashService) *TelemetryService { func NewTelemetryService(logger *slog.Logger, devServo *DeviceService) *TelemetryService {
sharedChannel := make(chan string, 10)
newService := &TelemetryService{ newService := &TelemetryService{
logger: logger, logger: logger,
cdash: cdash, isConnected: false,
uiOutCh: make(chan telem.TelemetryData, 5), devService: devServo,
listeners: make(map[string]chan telem.TelemetryData), listeners: make(map[string]chan telem.TelemetryData),
Messages: sharedChannel,
} }
newService.CtxMonitor, newService.cancelMonitor = context.WithCancel(context.Background())
// Need to instantiate a default provider here
source := "/home/esilva/Desktop/projetos/simracing_peripherals/testTelemetry/gt3_mustang_bathurst.ibt"
firstProvider := providers.NewIRacingProvider(slog.Default(), source, "", "")
newService.SwitchProvider(firstProvider)
return newService return newService
} }
// func (t *TelemetryService) GetUIStream() <-chan telem.TelemetryData { func (t *TelemetryService) ProviderMonitor(ctx context.Context) {
// return t.uiOutCh ticker := time.NewTicker(2 * time.Second)
// } defer ticker.Stop()
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.ativeProvider != nil {
t.dropActiveProvider()
}
// Assign the new provider
t.ativeProvider = newProvider
return nil
}
func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan telem.TelemetryData) {
for { for {
select { select {
case <-ctx.Done(): case <-ctx.Done():
return return
case data, ok := <-dataCh: case <-ticker.C:
if !ok { slog.Info("checking if provider is still running")
if !t.activeProvider.IsAlive(500 * time.Millisecond) {
slog.Warn("provider healthcheck failed")
t.dropActiveProvider()
t.onProviderHealthCheckFailed()
return return
} }
t.mut.RLock()
for _, ch := range t.listeners {
select {
case ch <- data:
// Sends data to the subscriber
default:
// Subscriber is full, we just skip ahead. Maybe find a how to add metrics here
}
}
t.mut.RUnlock()
} }
} }
} }
func (t *TelemetryService) dropActiveProvider() { // Listener Control [START] ----------------------------------------------------
if t.cancelForward != nil {
t.cancelForward()
}
t.ativeProvider.StopStream()
}
func (t *TelemetryService) SubscribeListener(id string, bufferSize int) <-chan telem.TelemetryData { func (t *TelemetryService) SubscribeListener(id string, bufferSize int) <-chan telem.TelemetryData {
t.mut.Lock() t.mut.Lock()
@@ -118,27 +98,123 @@ func (t *TelemetryService) UnsubscribeListener(id string) {
} }
} }
func (t *TelemetryService) SubscribeToFields(fields map[int16]telem.FieldID) { func (t *TelemetryService) SubscribeToFields() []telem.FieldID {
t.ativeProvider.Subscribe(fields) seen := make(map[telemetry.FieldID]struct{})
var allFields []telemetry.FieldID
for _, dev := range t.devService.GetDevices() {
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
} }
func (t *TelemetryService) StartStream() { // Listener Control [END] ------------------------------------------------------
slog.Debug("Stream started")
// Start the new stream // Provider Control [START] ----------------------------------------------------
simInCh, _ := t.ativeProvider.Stream()
// if err != nil {
// // NOTE
// }
// Create the context so we can control the lifecycle func (t *TelemetryService) HasActiveProvider() bool {
ctx, cancel := context.WithCancel(context.Background()) if t.activeProvider == nil {
t.cancelForward = cancel return false
}
// Multiplex this data return true
go t.multiplexData(ctx, simInCh)
} }
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) 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())
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
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.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:
// start monitoring on startup -> find provider -> stop monitoring (when game closes for example)
func (t *TelemetryService) FindProvider(ctx context.Context) {
ticker := time.NewTicker(2 * time.Second)
defer ticker.Stop()
t.logger.Debug("monitoring for providers")
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
for _, prov := range providers.Providers {
// t.logger.Debug("checking provider: " + prov.Name)
provider, err := prov.NewProvider(t.logger)
if err != nil {
// Its not running
continue
}
if provider.IsAlive(500 * time.Millisecond) {
t.onFindProvider(provider)
return
}
}
}
}
}
// Provider Control [END] ------------------------------------------------------
// Streaming Control [START] ---------------------------------------------------
func (t *TelemetryService) StopStream() { func (t *TelemetryService) StopStream() {
t.mut.Lock() t.mut.Lock()
defer t.mut.Unlock() defer t.mut.Unlock()
@@ -148,5 +224,64 @@ func (t *TelemetryService) StopStream() {
t.cancelForward = nil t.cancelForward = nil
} }
t.ativeProvider.StopStream() t.activeProvider.StopStream()
t.isStreaming = false
} }
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) {
for {
select {
case <-ctx.Done():
return
case data, ok := <-dataCh:
if !ok {
t.logger.Debug("something happened on the provider - stream closed")
t.onProviderStopsMidStream()
return
}
t.mut.RLock()
for _, ch := range t.listeners {
// t.logger.Debug("sending data to listener", "listener", key, "data", data)
select {
case ch <- data:
// Sends data to the subscriber
default:
// Subscriber is full, we just skip ahead. Maybe find a how to add metrics here
}
}
t.mut.RUnlock()
}
}
}
func (t *TelemetryService) IsStreaming() bool {
return t.isStreaming
}
// Streaming Control [END] -----------------------------------------------------
+41
View File
@@ -0,0 +1,41 @@
package telemetry
import "strconv"
func EmptyTransform(v any, out *TelemetryField) {
out.Type = DataTypeCHAR
out.Raw = uint64('-')
}
func UInt8Transform(v int, out *TelemetryField) {
out.Type = DataTypeUINT8
if v < 0 {
out.Raw = uint64(0)
return
}
out.Raw = uint64(v)
}
func FloatToStringTransform(v float32, out *TelemetryField) {
out.Type = DataTypeSTRING
if v < 0.0 {
out.Str = "inv"
return
}
out.Str = strconv.FormatFloat(float64(v), 'f', 1, 32)
}
func FloatToUInt8Transform(v float32, out *TelemetryField) {
out.Type = DataTypeUINT8
if v < 0.0 {
out.Raw = uint64(0)
return
}
out.Raw = uint64(v)
}
+32 -76
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.
@@ -27,10 +33,11 @@ type VirtualField interface {
// NOTE: Update the iracing SDK to write data to the same map ALWAYS, then // NOTE: Update the iracing SDK to write data to the same map ALWAYS, then
// I can bind that address and read directly from there on the transform // I can bind that address and read directly from there on the transform
// BoundField is the data structure we use to bind the telemetry provider's data
// to our internal telemetry fields
type BoundField struct { type BoundField struct {
Key string ID FieldID
ID FieldID Update func(out *TelemetryField)
Transform func(any, *TelemetryField)
} }
var bufferPool = sync.Pool{ var bufferPool = sync.Pool{
@@ -63,52 +70,15 @@ 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
} }
// Pack will pack this current TelemetryField into bytes to send over the wire func (tf *TelemetryField) Unused() {
// Format: tf.Type = DataTypeCHAR
// 0x00 - Field ID tf.Raw = uint64('-')
// 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 {
@@ -128,7 +98,7 @@ func (tf *TelemetryField) String() string {
return "NaN" return "NaN"
} }
type FieldID uint16 type FieldID = uint16
// We use FirstField to start the count on the fields the user can select // We use FirstField to start the count on the fields the user can select
// the first three will be for internal use // the first three will be for internal use
@@ -143,6 +113,14 @@ const (
WaterTemp WaterTemp
// Engine Warnings // Engine Warnings
PitSpeedLimiter PitSpeedLimiter
// Electrics (dash lights and so on)
LeftIndicator
RightIndicator
Hazards
ABSWarningLight
ParkingBrakeLight
TCLight
BatteryLight
// Adjustements // Adjustements
BrakeBias BrakeBias
ABSSetting ABSSetting
@@ -191,6 +169,14 @@ var FieldNames = [MaxFields]string{
WaterTemp: "Water Temperature", WaterTemp: "Water Temperature",
// Engine Warnings // Engine Warnings
PitSpeedLimiter: "Pit Speed Limiter", PitSpeedLimiter: "Pit Speed Limiter",
// Electrics (dash lights and so on)
LeftIndicator: "Left Indicator",
RightIndicator: "Right Indicator",
Hazards: "Hazards",
ABSWarningLight: "ABS Dash Light",
ParkingBrakeLight: "Parking Brake Dash Light",
TCLight: "Traction Control Light",
BatteryLight: "Battery Light",
// Ajustments // Ajustments
BrakeBias: "BrakeBias", BrakeBias: "BrakeBias",
ABSSetting: "ABS Control", ABSSetting: "ABS Control",
@@ -236,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
@@ -262,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
}
+6 -1
View File
@@ -1,8 +1,13 @@
// Package telemetry is our interface with our data sources // Package telemetry is our interface with our data sources
package telemetry package telemetry
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
Name() string
Close()
} }
-4
View File
@@ -1,5 +1 @@
package telemetry package telemetry
func Init() {
initFieldNamesMap()
}
@@ -5,11 +5,12 @@ import (
"os" "os"
"strconv" "strconv"
"esdi/cdashdisplay" "esdi/devices/cdashdisplay"
helper "esdi/helpers"
"esdi/services" "esdi/services"
"esdi/tui/internal/views" "esdi/tui/internal/views"
helper "esdi/helpers"
"github.com/gdamore/tcell/v2" "github.com/gdamore/tcell/v2"
"github.com/rivo/tview" "github.com/rivo/tview"
) )
@@ -19,18 +20,21 @@ type LayoutController struct {
OnExit func() OnExit func()
LayoutToolView *views.LayoutToolView LayoutToolView *views.LayoutToolView
Messages chan string Messages chan string
DevService *services.CDashService DevService *services.DeviceService
MoveToolState *windowManipState MoveToolState *windowManipState
SelectedLayout string // NOTE: This should be a struct to handle its own things SelectedLayout string // NOTE: This should be a struct to handle its own things
// TODO: create a filter to select the layout
// filter telemetry provider, then filter vehicle in use and so on
} }
func NewLayoutController(base *Controller, service *services.CDashService) *LayoutController { func NewLayoutController(base *Controller, service *services.DeviceService) *LayoutController {
lc := &LayoutController{ lc := &LayoutController{
Controller: base, Controller: base,
LayoutToolView: views.NewLayoutToolView(), LayoutToolView: views.NewLayoutToolView(),
Messages: make(chan string, 10), Messages: make(chan string, 10),
DevService: service, DevService: service,
MoveToolState: &windowManipState{Mode: moveMode}, MoveToolState: &windowManipState{Mode: moveMode},
// SelectedLayout: "beamng.yaml",
SelectedLayout: "layout.yaml", SelectedLayout: "layout.yaml",
} }
@@ -206,7 +210,20 @@ func (lc *LayoutController) createWindow() {
return return
} }
window, err = lc.DevService.CreateWindow(window) // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME
return
}
display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return
}
// ---
window, err = display.CreateWindow(window)
if err != nil { if err != nil {
lc.Messages <- "failed to create window\n" lc.Messages <- "failed to create window\n"
return return
@@ -303,7 +320,20 @@ func (lc *LayoutController) newWindowAction() {
} }
func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) { func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) {
err := lc.DevService.UpdateWindow(win) // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME
return
}
display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return
}
// ---
err = display.UpdateWindow(win)
lc.Messages <- fmt.Sprintf("Window: %v\n", win) lc.Messages <- fmt.Sprintf("Window: %v\n", win)
@@ -316,7 +346,20 @@ func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow
func (lc *LayoutController) displayLoadedLayouts() { func (lc *LayoutController) displayLoadedLayouts() {
lc.Logger.Debug("We want to view our layout!") lc.Logger.Debug("We want to view our layout!")
for _, w := range lc.DevService.CDash.State.Layout.Windows { // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME
return
}
display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return
}
// ---
for _, w := range display.State.Layout.Windows {
lc.Logger.Debug("==========================================================================") lc.Logger.Debug("==========================================================================")
lc.Logger.Debug(fmt.Sprintf("updating form view for a layout: %+v", w.UIData.TelemetryField)) lc.Logger.Debug(fmt.Sprintf("updating form view for a layout: %+v", w.UIData.TelemetryField))
err := lc.updateFormView(w) err := lc.updateFormView(w)
@@ -348,8 +391,21 @@ func (lc *LayoutController) getCurrentTreeNodeModel() (*tview.TreeNode, int16, e
} }
func (lc *LayoutController) loadLayout() { func (lc *LayoutController) loadLayout() {
// Acquire the cdashdisplay
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME
return
}
display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return
}
// ---
// We would get the layout path from somewhere but for nots its layout.yaml // We would get the layout path from somewhere but for nots its layout.yaml
err := lc.DevService.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()
return return
@@ -359,7 +415,20 @@ func (lc *LayoutController) loadLayout() {
} }
func (lc *LayoutController) unloadLayout() { func (lc *LayoutController) unloadLayout() {
err := lc.DevService.UnloadLayout() // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME
return
}
display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return
}
// ---
err = display.UnloadLayout()
if err != nil { if err != nil {
lc.Logger.Error(fmt.Sprintf("Failed to unload layout: %+v", err)) lc.Logger.Error(fmt.Sprintf("Failed to unload layout: %+v", err))
return return
@@ -367,7 +436,20 @@ func (lc *LayoutController) unloadLayout() {
} }
func (lc *LayoutController) saveLayout() { func (lc *LayoutController) saveLayout() {
err := lc.DevService.SaveLayout(lc.SelectedLayout) // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME
return
}
display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return
}
// ---
err = display.SaveLayout(lc.SelectedLayout)
if err != nil { if err != nil {
lc.Messages <- "failed to save layout: " + err.Error() lc.Messages <- "failed to save layout: " + err.Error()
return return
@@ -387,8 +469,21 @@ func (lc *LayoutController) deleteWindow() {
} }
wID := curNode.GetReference().(int16) wID := curNode.GetReference().(int16)
// Acquire the cdashdisplay
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME
return
}
display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return
}
// ---
// Delete it // Delete it
err := lc.DevService.DeleteWindow(wID) err = display.DestroyWindow(wID)
if err != nil { if err != nil {
lc.Messages <- "failed to delete window: " + err.Error() + "\n" lc.Messages <- "failed to delete window: " + err.Error() + "\n"
return return
@@ -1,6 +1,7 @@
package controllers package controllers
import ( import (
"esdi/devices/cdashdisplay"
helper "esdi/helpers" helper "esdi/helpers"
"github.com/gdamore/tcell/v2" "github.com/gdamore/tcell/v2"
@@ -43,40 +44,68 @@ func keyToVector(r rune) (helper.Vector, bool) {
} }
func (lc *LayoutController) handleMovementCapture(idx int16, func (lc *LayoutController) handleMovementCapture(idx int16,
ev *tcell.EventKey) *tcell.EventKey { ev *tcell.EventKey,
) *tcell.EventKey {
vec, ok := keyToVector(ev.Rune()) vec, ok := keyToVector(ev.Rune())
if !ok { if !ok {
return nil return nil
} }
err := lc.DevService.MoveWindow(idx, &vec) // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME
return nil
}
display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return nil
}
// ---
err = display.MoveWindow(idx, &vec)
if err != nil { if err != nil {
lc.Messages <- "failed to move window: " + err.Error() + "\n" lc.Messages <- "failed to move window: " + err.Error() + "\n"
return nil return nil
} }
// Success - update the form // Success - update the form
window := lc.DevService.CDash.State.Layout.Windows[idx] window := display.State.Layout.Windows[idx]
lc.LayoutToolView.UpdateFormView(idx, window) lc.LayoutToolView.UpdateFormView(idx, window)
return nil return nil
} }
func (lc *LayoutController) handleResizeCapture(idx int16, func (lc *LayoutController) handleResizeCapture(idx int16,
ev *tcell.EventKey) *tcell.EventKey { ev *tcell.EventKey,
) *tcell.EventKey {
vec, ok := keyToVector(ev.Rune()) vec, ok := keyToVector(ev.Rune())
if !ok { if !ok {
return nil return nil
} }
err := lc.DevService.ResizeWindow(idx, &vec) // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.NAME
return nil
}
display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return nil
}
// ---
err = display.ResizeWindow(idx, &vec)
if err != nil { if err != nil {
lc.Messages <- "failed to resize window: " + err.Error() + "\n" lc.Messages <- "failed to resize window: " + err.Error() + "\n"
return nil return nil
} }
// Success - update the form // Success - update the form
window := lc.DevService.CDash.State.Layout.Windows[idx] window := display.State.Layout.Windows[idx]
lc.LayoutToolView.UpdateFormView(idx, window) lc.LayoutToolView.UpdateFormView(idx, window)
return nil return nil
+18 -4
View File
@@ -4,6 +4,7 @@ package controllers
import ( import (
"fmt" "fmt"
"esdi/devices/cdashdisplay"
serv "esdi/services" serv "esdi/services"
"esdi/tui/internal/views" "esdi/tui/internal/views"
@@ -15,12 +16,12 @@ type DeviceController struct {
DeviceAPIView *views.DeviceAPIView DeviceAPIView *views.DeviceAPIView
LayoutCtrl *LayoutController LayoutCtrl *LayoutController
StreamCtrl *StreamingCtrl StreamCtrl *StreamingCtrl
DevService *serv.CDashService DevService *serv.DeviceService
} }
func NewDeviceController( func NewDeviceController(
base *Controller, base *Controller,
devService *serv.CDashService, devService *serv.DeviceService,
telemService *serv.TelemetryService, telemService *serv.TelemetryService,
) *DeviceController { ) *DeviceController {
mc := &DeviceController{ mc := &DeviceController{
@@ -56,7 +57,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.FindDevice() go mc.DevService.FindDevices()
} }
return ev return ev
}) })
@@ -65,6 +66,12 @@ func (mc *DeviceController) setDeviceAPIViewEvents() {
func (mc *DeviceController) AddDeviceAPIListItems() { 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
if !mc.DevService.PeripheralExists(cdashdisplay.NAME) {
mc.DevService.Messages <- "CDashDisplay it not loaded yet\n"
return
}
// Get the api pages // Get the api pages
views.AddAndShowPage( views.AddAndShowPage(
mc.DeviceAPIView.DevAPIToolView.Pages, mc.DeviceAPIView.DevAPIToolView.Pages,
@@ -75,7 +82,14 @@ func (mc *DeviceController) AddDeviceAPIListItems() {
}) })
mc.DeviceAPIView.DevAPIList. mc.DeviceAPIView.DevAPIList.
AddItem("stream", "stream data to the display", func() { AddItem("stream", "stream data to the display", func() {
views.AddAndShowPage(mc.DeviceAPIView.DevAPIToolView.Pages, // If we don't have a CDashDisplay or data source, this should be blocked
if !mc.StreamCtrl.TelemServ.HasActiveProvider() {
mc.StreamCtrl.Messages <- "no active provider present\n"
return
}
views.AddAndShowPage(
mc.DeviceAPIView.DevAPIToolView.Pages,
"streaming-tool", "streaming-tool",
mc.StreamCtrl.StreamView.Flex, mc.StreamCtrl.StreamView.Flex,
) )
+30 -28
View File
@@ -6,6 +6,7 @@ import (
"sync/atomic" "sync/atomic"
"esdi/config" "esdi/config"
"esdi/devices/uidevice"
"esdi/providers" "esdi/providers"
"esdi/services" "esdi/services"
"esdi/telemetry" "esdi/telemetry"
@@ -17,7 +18,7 @@ import (
type StreamingCtrl struct { type StreamingCtrl struct {
*Controller *Controller
Service *services.CDashService DevService *services.DeviceService
StreamView *views.StreamToolView StreamView *views.StreamToolView
Messages chan string Messages chan string
Internal chan string Internal chan string
@@ -32,10 +33,10 @@ type StreamingCtrl struct {
func NewStreamingCtrl( func NewStreamingCtrl(
base *Controller, base *Controller,
serCDash *services.CDashService, devService *services.DeviceService,
serTelem *services.TelemetryService, serTelem *services.TelemetryService,
) *StreamingCtrl { ) *StreamingCtrl {
// NOTE: looks sus // NOTE: looks sus, put this somewhere also. Not very good in here
providerList := []providers.Provider{} providerList := []providers.Provider{}
for _, item := range providers.Providers { for _, item := range providers.Providers {
providerList = append(providerList, item) providerList = append(providerList, item)
@@ -44,26 +45,24 @@ func NewStreamingCtrl(
ctrl := &StreamingCtrl{ ctrl := &StreamingCtrl{
Controller: base, Controller: base,
Service: serCDash, 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()
go ctrl.listenToUIStream()
return ctrl return ctrl
} }
func (sc *StreamingCtrl) subscribeListeners() { // func (sc *StreamingCtrl) subscribeListeners() {
sc.TelemetryCh = sc.TelemServ.SubscribeListener("UI", 1) // // 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 {
@@ -95,11 +94,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
@@ -108,16 +107,21 @@ func (sc *StreamingCtrl) StartStop() {
// stream is not running, we have to start it now // stream is not running, we have to start it now
// 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 cdash") slog.Debug("setting the data stream for device servie")
sc.Service.SetTelemetryChannel(sc.TelemServ.SubscribeListener("cdash", 1)) sc.DevService.SetTelemetryChannel(sc.TelemServ.SubscribeListener("DeviceService", 1))
slog.Debug("starting to stream data again") dev, err := sc.DevService.GetPeripheral(uidevice.NAME)
sc.Service.StartStream() if err == nil {
if uiDev, ok := dev.(*uidevice.UIDevice); ok {
sc.TelemetryCh = uiDev.DataChannel()
go sc.listenToUIStream()
}
}
slog.Debug("starting the stream") slog.Debug("starting services")
sc.DevService.StartStream()
sc.TelemServ.StartStream() sc.TelemServ.StartStream()
slog.Debug("setting local control variables")
sc.isRunning = true sc.isRunning = true
slog.Debug("starting stream") slog.Debug("starting stream")
@@ -155,16 +159,11 @@ 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() {
fields := make(map[int16]telemetry.FieldID, len(sc.Service.CDash.State.Layout.Windows)) fields := sc.TelemServ.SubscribeToFields()
for _, w := range sc.Service.CDash.State.Layout.Windows { // sc.Messages <- fmt.Sprintf("Subscribed Fields: %+v [%d]\n", fields, len(fields))
fieldID, _ := telemetry.GetFieldID(w.UIData.TelemetryField) // Should I update this?
fields[w.UIData.IDX] = fieldID sc.Messages <- fmt.Sprintf("Subscribed to fields: %+v", fields)
}
sc.TelemServ.SubscribeToFields(fields)
sc.Messages <- fmt.Sprintf("Subscribed Fields: %+v [%d]\n", fields, len(fields))
} }
func (sc *StreamingCtrl) listenToUIStream() { func (sc *StreamingCtrl) listenToUIStream() {
@@ -176,8 +175,11 @@ func (sc *StreamingCtrl) listenToUIStream() {
} }
isDrawing.Store(true) isDrawing.Store(true)
// Capture locally
telemetryMsg := msg
sc.App.QueueUpdateDraw(func() { sc.App.QueueUpdateDraw(func() {
sc.StreamView.Visualizer.Update(&msg) sc.StreamView.Visualizer.Update(&telemetryMsg)
isDrawing.Store(false) isDrawing.Store(false)
}) })
} }
+1 -1
View File
@@ -4,7 +4,7 @@ import (
"fmt" "fmt"
"log/slog" "log/slog"
"esdi/cdashdisplay" "esdi/devices/cdashdisplay"
tviewh "esdi/tui/internal/tview_helpers" tviewh "esdi/tui/internal/tview_helpers"
"github.com/gdamore/tcell/v2" "github.com/gdamore/tcell/v2"
+1 -1
View File
@@ -3,7 +3,7 @@ package views
import ( import (
"fmt" "fmt"
"esdi/cdashdisplay" "esdi/devices/cdashdisplay"
"esdi/telemetry" "esdi/telemetry"
"github.com/rivo/tview" "github.com/rivo/tview"
+5 -1
View File
@@ -23,12 +23,16 @@ func NewControlPanel(logger *slog.Logger) *ControlPanel {
App: tview.NewApplication(), App: tview.NewApplication(),
} }
devService := services.NewCDashService(logger) // NOTE: create our device service here
devService := services.NewDeviceService(logger.With("service", "DeviceService"))
telemService := services.NewTelemetryService(logger, devService) telemService := services.NewTelemetryService(logger, devService)
if telemService == nil { if telemService == nil {
panic("failed to create the telemetry service") 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, devService, telemService),