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
This commit is contained in:
2026-09-17 22:55:17 +01:00
parent 20025dae1b
commit 65f03fd4c9
17 changed files with 327 additions and 111 deletions
+1
View File
@@ -0,0 +1 @@
# Peripheral Discovery
+20 -2
View File
@@ -1,9 +1,27 @@
package cdashdisplay package cdashdisplay
import "esdi/peripheral/devices" import (
"esdi/peripheral/devices"
"esdi/telemetry"
)
// TODO: I believe we don't need the #esdi/peripheral/devices thing anymore // TODO: I believe we don't need the #esdi/peripheral/devices thing anymore
const ( const (
ID = devices.CDashDisplayDevID ID = devices.CDashDisplayDevID
Name = devices.CDashDisplayDevName NAME = devices.CDashDisplayDevName
) )
func (cds *CDashDisplay) Name() string {
return NAME
}
func (cds *CDashDisplay) RequiredFields() map[int16]telemetry.FieldID {
fields := make(map[int16]telemetry.FieldID, len(cds.State.Layout.Windows))
for _, w := range cds.State.Layout.Windows {
fieldID, _ := telemetry.GetFieldID(w.UIData.TelemetryField)
fields[w.UIData.IDX] = fieldID
}
return fields
}
+56 -5
View File
@@ -1,9 +1,60 @@
// Package devices simply defines what a device should have // Package devices is a peripheral factory
// and some utils if necessary
package devices package devices
import "esdi/telemetry" import (
"errors"
type Device interface { "esdi/devices/cdashdisplay"
SendData(*telemetry.TelemetryData) "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) {
// // Find CDashDisplay
// {
// ds.Messages <- "looking for " + cdashdisplay.Name + "...\n"
// ds.Logger.Info("Looking for " + cdashdisplay.Name)
//
// cdashdisplay.SetLogger(ds.Logger.With("[device]", cdashdisplay.Name))
//
// // Create a cdashdisplay
// display, err := cdashdisplay.Discover()
// if err == nil {
// ds.Devices[cdashdisplay.Name] = display
// ds.Logger.Info("found " + cdashdisplay.Name + " on: " + display.WT.Cfg.Name)
// ds.Messages <- "found " + cdashdisplay.Name + " on: " + display.WT.Cfg.Name + "\n"
// return
// }
//
// ds.Logger.Info("didn't find " + cdashdisplay.Name)
// ds.Messages <- "didn't find " + cdashdisplay.Name + "\n"
// // No CDashDisplay available for one reason or another, so we don't set the
// // key
// }
return nil, errors.New("not implemented yet")
} }
+53
View File
@@ -0,0 +1,53 @@
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) SendData(data *telemetry.TelemetryData) {
if data == nil {
return
}
select {
case uid.dataChan <- *data:
default:
// Drop frame if buffer is full
}
}
func (uid *UIDevice) Name() string {
return NAME
}
func (uid *UIDevice) DataChannel() <-chan telemetry.TelemetryData {
return uid.dataChan
}
func (uid *UIDevice) RequiredFields() map[int16]telemetry.FieldID {
subscribeTo := []telemetry.FieldID{
telemetry.Speed,
telemetry.Gear,
telemetry.RPM,
}
fields := make(map[int16]telemetry.FieldID, len(subscribeTo))
for k, field := range subscribeTo {
fields[int16(k)] = field
}
return fields
}
+9 -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,12 @@ const (
DisplayPeripheral PeripheralType = "display" DisplayPeripheral PeripheralType = "display"
) )
type Peripheral interface {
Name() string
SendData(*telemetry.TelemetryData)
RequiredFields() map[int16]telemetry.FieldID
}
type PeripheralDeviceClerk struct { type PeripheralDeviceClerk struct {
// mu sync.RWMutex // mu sync.RWMutex
Devices map[uint8]*PeripheralDevice Devices map[uint8]*PeripheralDevice
+4
View File
@@ -88,6 +88,10 @@ func (b *BeamNG) Close() {
b.SDK.Close() b.SDK.Close()
} }
func (b *BeamNG) Name() string {
return NAME
}
func (b *BeamNG) IsAlive(timeout time.Duration) bool { func (b *BeamNG) IsAlive(timeout time.Duration) bool {
_, err := b.SDK.Update() _, err := b.SDK.Update()
return err == nil return err == nil
+4
View File
@@ -264,6 +264,10 @@ func (i *IRacing) Subscribe(requestFields map[int16]telemetry.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 { func (i *IRacing) IsAlive(timeout time.Duration) bool {
if !i.SDK.CheckForDataEvent(timeout) { if !i.SDK.CheckForDataEvent(timeout) {
return false return false
+16 -8
View File
@@ -16,7 +16,7 @@ import (
// the selected provider by its name // the selected provider by its name
type Provider struct { type Provider struct {
Name string Name string
NewProvider func(*slog.Logger) telemetry.TelemetryProvider NewProvider func(*slog.Logger) (telemetry.TelemetryProvider, error)
IsRunning func() bool // To check if this provider is up and running IsRunning func() bool // To check if this provider is up and running
} }
@@ -33,23 +33,31 @@ var Providers = map[string]Provider{
}, },
} }
func NewLiveIRacingProvider(logger *slog.Logger) telemetry.TelemetryProvider { func NewLiveIRacingProvider(logger *slog.Logger) (telemetry.TelemetryProvider, error) {
provider, _ := iracing.NewIRacingProvider(logger, goirsdk.Options{ provider, err := iracing.NewIRacingProvider(logger, goirsdk.Options{
Logger: logger, Logger: logger,
SourceType: goirsdk.SharedMemoryFile, SourceType: goirsdk.SharedMemoryFile,
}) })
if err != nil {
return provider logger.Error("failed to create iRacing provider", "err", err)
return nil, err
} }
func NewBeamNGProvider(logger *slog.Logger) telemetry.TelemetryProvider { return provider, nil
}
func NewBeamNGProvider(logger *slog.Logger) (telemetry.TelemetryProvider, error) {
// TODO: these should come from some kind of config // TODO: these should come from some kind of config
provider, _ := beamng.NewBeamNGProvider(logger, &bngsdk.Options{ provider, err := beamng.NewBeamNGProvider(logger, &bngsdk.Options{
Logger: logger.With("TelemetryProvider", beamng.NAME), Logger: logger.With("TelemetryProvider", beamng.NAME),
SourceType: bngsdk.UDPData, SourceType: bngsdk.UDPData,
ImportUDPAddress: "127.0.0.1", ImportUDPAddress: "127.0.0.1",
ImportUDPPort: 4444, ImportUDPPort: 4444,
}) })
if err != nil {
logger.Error("failed to create BeamNG provider", "err", err)
return nil, err
}
return provider return provider, nil
} }
+68 -22
View File
@@ -2,15 +2,20 @@ package services
import ( import (
"context" "context"
"errors"
"fmt" "fmt"
"log/slog" "log/slog"
"sync"
"sync/atomic" "sync/atomic"
"time"
"esdi/devices" "esdi/devices"
"esdi/devices/cdashdisplay" "esdi/peripheral"
"esdi/telemetry" "esdi/telemetry"
) )
var ErrPeripheralAlreadyRegistered = errors.New("peripheral is already registered")
// DeviceService will handle sending the data from the telemetry service to the // DeviceService will handle sending the data from the telemetry service to the
// actual devices // actual devices
// NOTE: create a virtual device and make it be the output window or something so // NOTE: create a virtual device and make it be the output window or something so
@@ -18,7 +23,11 @@ import (
// that would be pretty cool I think // that would be pretty cool I think
type DeviceService struct { type DeviceService struct {
Logger *slog.Logger Logger *slog.Logger
Devices map[string]devices.Device // Device discovery
mu sync.RWMutex
ctxDiscovery context.Context
ctxDiscoveryCancel context.CancelFunc
Devices map[string]peripheral.Peripheral
// Strem handling // Strem handling
streamCancel context.CancelFunc streamCancel context.CancelFunc
TelemCh <-chan telemetry.TelemetryData TelemCh <-chan telemetry.TelemetryData
@@ -28,41 +37,75 @@ type DeviceService struct {
func NewDeviceService(logger *slog.Logger) *DeviceService { func NewDeviceService(logger *slog.Logger) *DeviceService {
sharedChannel := make(chan string, 10) sharedChannel := make(chan string, 10)
return &DeviceService{
Devices: make(map[string]devices.Device), dev := &DeviceService{
Devices: make(map[string]peripheral.Peripheral),
Logger: logger, Logger: logger,
Messages: sharedChannel, 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
} }
func (ds *DeviceService) FindDevices() { func (ds *DeviceService) FindDevices() {
// Need to define a list of devices to search for // Need to define a list of devices to search for
// For now lets just try to find our cdashdisplay - will think about the rest later // For now lets just try to find our cdashdisplay - will think about the rest later
ticker := time.NewTicker(2 * time.Second)
defer ticker.Stop()
// Find CDashDisplay for {
{ select {
ds.Messages <- "looking for " + cdashdisplay.Name + "...\n" case <-ds.ctxDiscovery.Done():
ds.Logger.Info("Looking for " + cdashdisplay.Name) // If requested to cancel we cancel background discovery
cdashdisplay.SetLogger(ds.Logger.With("[device]", cdashdisplay.Name))
// Create a cdashdisplay
display, err := cdashdisplay.Discover()
if err == nil {
ds.Devices[cdashdisplay.Name] = display
ds.Logger.Info("found " + cdashdisplay.Name + " on: " + display.WT.Cfg.Name)
ds.Messages <- "found " + cdashdisplay.Name + " on: " + display.WT.Cfg.Name + "\n"
return return
case <-ticker.C:
for pName, peripheral := range devices.List {
if ds.DeviceExists(pName) {
// We already discovered this device
continue
} }
ds.Logger.Info("didn't find " + cdashdisplay.Name) ds.Logger.Debug("looking for device", "name", pName)
ds.Messages <- "didn't find " + cdashdisplay.Name + "\n" dev, err := peripheral.Discover()
// No CDashDisplay available for one reason or another, so we don't set the if err != nil {
// key ds.Logger.Debug("didn't find device", "name", pName)
continue
}
// Register the device we just found
ds.RegisterDevice(dev)
}
}
} }
} }
func (ds *DeviceService) GetDevice(name string) (devices.Device, error) { // func (ds *DeviceService) SubscribeFields() error {
// for _, dev := range ds.Devices {
// fields := dev.RequiredFields()
// }
//
// return nil
// }
func (ds *DeviceService) RegisterDevice(dev peripheral.Peripheral) error {
ds.mu.Lock()
defer ds.mu.Unlock()
if ds.DeviceExists(dev.Name()) {
return ErrPeripheralAlreadyRegistered
}
ds.Devices[dev.Name()] = dev
return nil
}
func (ds *DeviceService) GetDevice(name string) (peripheral.Peripheral, error) {
val, ok := ds.Devices[name] val, ok := ds.Devices[name]
if !ok { if !ok {
return nil, fmt.Errorf("device `%s` couldn't be found", name) return nil, fmt.Errorf("device `%s` couldn't be found", name)
@@ -122,9 +165,12 @@ func (ds *DeviceService) transmit(ctx context.Context) {
// TODO: make a copy of the data and send that copy instead of keeping // TODO: make a copy of the data and send that copy instead of keeping
// the data locked // the data locked
ds.mu.RLock()
for _, dev := range ds.Devices { for _, dev := range ds.Devices {
dev.SendData(&data) dev.SendData(&data)
} }
ds.mu.RUnlock()
isSending.Store(false) isSending.Store(false)
} }
} }
+26 -9
View File
@@ -54,9 +54,9 @@ func (t *TelemetryService) ProviderMonitor(ctx context.Context) {
case <-ctx.Done(): case <-ctx.Done():
return return
case <-ticker.C: case <-ticker.C:
slog.Debug("checking if provider is still running") slog.Info("checking if provider is still running")
if !t.activeProvider.IsAlive(500 * time.Millisecond) { if !t.activeProvider.IsAlive(500 * time.Millisecond) {
slog.Debug("provider healthcheck failed") slog.Warn("provider healthcheck failed")
t.dropActiveProvider() t.dropActiveProvider()
t.onProviderHealthCheckFailed() t.onProviderHealthCheckFailed()
return return
@@ -70,10 +70,14 @@ func (t *TelemetryService) onProviderHealthCheckFailed() {
go t.FindProvider(t.CtxMonitor) go t.FindProvider(t.CtxMonitor)
} }
func (t *TelemetryService) onFindProvider(prov providers.Provider) { func (t *TelemetryService) onFindProvider(prov telem.TelemetryProvider) {
// Attach to the provider // Attach to the provider
t.logger.Info("found provider for " + prov.Name) t.logger.Info("found provider for " + prov.Name())
t.SwitchProvider(prov.NewProvider(t.logger)) 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 // 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()) t.CtxHealthcheck, t.healthCheckCancel = context.WithCancel(context.Background())
@@ -102,9 +106,15 @@ func (t *TelemetryService) FindProvider(ctx context.Context) {
return return
case <-ticker.C: case <-ticker.C:
for _, prov := range providers.Providers { for _, prov := range providers.Providers {
t.logger.Debug("checking provider: " + prov.Name) // t.logger.Debug("checking provider: " + prov.Name)
if prov.IsRunning() { provider, err := prov.NewProvider(t.logger)
t.onFindProvider(prov) if err != nil {
// Its not running
continue
}
if provider.IsAlive(500 * time.Millisecond) {
t.onFindProvider(provider)
return return
} }
} }
@@ -141,6 +151,7 @@ func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan tele
t.mut.RLock() t.mut.RLock()
for _, ch := range t.listeners { for _, ch := range t.listeners {
// t.logger.Debug("sending data to listener", "listener", key, "data", data)
select { select {
case ch <- data: case ch <- data:
// Sends data to the subscriber // Sends data to the subscriber
@@ -191,9 +202,15 @@ func (t *TelemetryService) UnsubscribeListener(id string) {
} }
} }
func (t *TelemetryService) SubscribeToFields(fields map[int16]telem.FieldID) { func (t *TelemetryService) SubscribeToFields() {
// _ = t.devService.SubscribeFields()
// TODO: we need to find a way of requesting devices to send all subscribed fields
// instead of going through the devices on DevService here
for _, dev := range t.devService.Devices {
fields := dev.RequiredFields()
t.activeProvider.Subscribe(fields) t.activeProvider.Subscribe(fields)
} }
}
func (t *TelemetryService) StartStream() { func (t *TelemetryService) StartStream() {
slog.Debug("Stream started") slog.Debug("Stream started")
+1 -1
View File
@@ -135,7 +135,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
+1
View File
@@ -8,5 +8,6 @@ type TelemetryProvider interface {
Stream() (<-chan TelemetryData, error) Stream() (<-chan TelemetryData, error)
Subscribe(map[int16]FieldID) Subscribe(map[int16]FieldID)
IsAlive(time.Duration) bool IsAlive(time.Duration) bool
Name() string
Close() Close()
} }
+21 -21
View File
@@ -211,14 +211,14 @@ func (lc *LayoutController) createWindow() {
} }
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.Name lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
} }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.Name lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return return
} }
// --- // ---
@@ -321,14 +321,14 @@ func (lc *LayoutController) newWindowAction() {
func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) { func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) {
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.Name lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
} }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.Name lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return return
} }
// --- // ---
@@ -347,14 +347,14 @@ func (lc *LayoutController) displayLoadedLayouts() {
lc.Logger.Debug("We want to view our layout!") lc.Logger.Debug("We want to view our layout!")
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.Name lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
} }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.Name lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return return
} }
// --- // ---
@@ -392,14 +392,14 @@ func (lc *LayoutController) getCurrentTreeNodeModel() (*tview.TreeNode, int16, e
func (lc *LayoutController) loadLayout() { func (lc *LayoutController) loadLayout() {
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.Name lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
} }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.Name lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return return
} }
// --- // ---
@@ -416,14 +416,14 @@ func (lc *LayoutController) loadLayout() {
func (lc *LayoutController) unloadLayout() { func (lc *LayoutController) unloadLayout() {
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.Name lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
} }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.Name lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return return
} }
// --- // ---
@@ -437,14 +437,14 @@ func (lc *LayoutController) unloadLayout() {
func (lc *LayoutController) saveLayout() { func (lc *LayoutController) saveLayout() {
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.Name lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
} }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.Name lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return return
} }
// --- // ---
@@ -470,14 +470,14 @@ func (lc *LayoutController) deleteWindow() {
wID := curNode.GetReference().(int16) wID := curNode.GetReference().(int16)
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.Name lc.Messages <- "failed to get " + cdashdisplay.NAME
return return
} }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.Name lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return return
} }
// --- // ---
@@ -52,14 +52,14 @@ func (lc *LayoutController) handleMovementCapture(idx int16,
} }
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.Name lc.Messages <- "failed to get " + cdashdisplay.NAME
return nil return nil
} }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.Name lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return nil return nil
} }
// --- // ---
@@ -86,14 +86,14 @@ func (lc *LayoutController) handleResizeCapture(idx int16,
} }
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
if err != nil { if err != nil {
lc.Messages <- "failed to get " + cdashdisplay.Name lc.Messages <- "failed to get " + cdashdisplay.NAME
return nil return nil
} }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { if !ok {
lc.Messages <- "failed to acquire " + cdashdisplay.Name lc.Messages <- "failed to acquire " + cdashdisplay.NAME
return nil return nil
} }
// --- // ---
+1 -1
View File
@@ -67,7 +67,7 @@ func (mc *DeviceController) AddDeviceAPIListItems() {
mc.DeviceAPIView.DevAPIList. mc.DeviceAPIView.DevAPIList.
AddItem("layout", "build a layout for CDashDisplay", func() { AddItem("layout", "build a layout for CDashDisplay", func() {
// This CDashDisplay specific, only load if we have a CDashDisplay // This CDashDisplay specific, only load if we have a CDashDisplay
if !mc.DevService.DeviceExists(cdashdisplay.Name) { if !mc.DevService.DeviceExists(cdashdisplay.NAME) {
mc.DevService.Messages <- "CDashDisplay it not loaded yet\n" mc.DevService.Messages <- "CDashDisplay it not loaded yet\n"
return return
} }
+37 -32
View File
@@ -6,7 +6,7 @@ import (
"sync/atomic" "sync/atomic"
"esdi/config" "esdi/config"
"esdi/devices/cdashdisplay" "esdi/devices/uidevice"
"esdi/providers" "esdi/providers"
"esdi/services" "esdi/services"
"esdi/telemetry" "esdi/telemetry"
@@ -56,15 +56,15 @@ func NewStreamingCtrl(
} }
ctrl.registerHooks() ctrl.registerHooks()
ctrl.subscribeListeners() // 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 {
@@ -109,16 +109,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.Service.SetTelemetryChannel(sc.TelemServ.SubscribeListener("DeviceService", 1))
slog.Debug("starting to stream data again") dev, err := sc.Service.GetDevice(uidevice.NAME)
if err == nil {
if uiDev, ok := dev.(*uidevice.UIDevice); ok {
sc.TelemetryCh = uiDev.DataChannel()
go sc.listenToUIStream()
}
}
slog.Debug("starting services")
sc.Service.StartStream() sc.Service.StartStream()
slog.Debug("starting the stream")
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")
@@ -157,28 +162,23 @@ func (sc *StreamingCtrl) updateStream() {
// so we can get away with using a map for convenience here // so we can get away with using a map for convenience here
func (sc *StreamingCtrl) SetInternalState() { func (sc *StreamingCtrl) SetInternalState() {
// Acquire the cdashdisplay // Acquire the cdashdisplay
displayIF, err := sc.Service.GetDevice(cdashdisplay.Name) // displayIF, err := sc.Service.GetDevice(cdashdisplay.NAME)
if err != nil { // if err != nil {
sc.Messages <- "failed to get " + cdashdisplay.Name // sc.Messages <- "failed to get " + cdashdisplay.NAME
return // return
} // }
display, ok := displayIF.(*cdashdisplay.CDashDisplay) // display, ok := displayIF.(*cdashdisplay.CDashDisplay)
if !ok { // if !ok {
sc.Messages <- "failed to acquire " + cdashdisplay.Name // sc.Messages <- "failed to acquire " + cdashdisplay.NAME
return // return
} // }
// --- // ---
fields := make(map[int16]telemetry.FieldID, len(display.State.Layout.Windows)) sc.TelemServ.SubscribeToFields()
for _, w := range display.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\n")
}
sc.TelemServ.SubscribeToFields(fields)
sc.Messages <- fmt.Sprintf("Subscribed Fields: %+v [%d]\n", fields, len(fields))
} }
func (sc *StreamingCtrl) listenToUIStream() { func (sc *StreamingCtrl) listenToUIStream() {
@@ -190,8 +190,13 @@ func (sc *StreamingCtrl) listenToUIStream() {
} }
isDrawing.Store(true) isDrawing.Store(true)
// sc.Logger.Debug("got data", "data", msg)
// 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
@@ -24,7 +24,7 @@ func NewControlPanel(logger *slog.Logger) *ControlPanel {
} }
// NOTE: create our device service here // NOTE: create our device service here
devService := services.NewDeviceService(logger) devService := services.NewDeviceService(logger.With("service", "DeviceService"))
telemService := services.NewTelemetryService(logger, devService) telemService := services.NewTelemetryService(logger, devService)
if telemService == nil { if telemService == nil {