From a6ab64117c5c29b7410afa72fc7266844aa96566 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Thu, 17 Sep 2026 22:55:17 +0100 Subject: [PATCH] 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 --- Docs/Peripheral_Discovery.md | 1 + devices/cdashdisplay/info.go | 22 ++++- devices/device.go | 61 +++++++++++- devices/uidevice/device.go | 53 +++++++++++ peripheral/peripheral.go | 10 +- providers/beamng/beamng.go | 4 + providers/iracing/iracing.go | 4 + providers/providers.go | 22 +++-- services/devices.go | 94 ++++++++++++++----- services/telemetry.go | 37 ++++++-- telemetry/data.go | 2 +- telemetry/provider.go | 1 + .../controllers/cdashdisplay_layout.go | 42 ++++----- .../cdashdisplay_layout_moveTool.go | 12 +-- tui/internal/controllers/device.go | 2 +- tui/internal/controllers/streaming.go | 69 +++++++------- tui/tui.go | 2 +- 17 files changed, 327 insertions(+), 111 deletions(-) create mode 100644 Docs/Peripheral_Discovery.md create mode 100644 devices/uidevice/device.go diff --git a/Docs/Peripheral_Discovery.md b/Docs/Peripheral_Discovery.md new file mode 100644 index 0000000..e3d041b --- /dev/null +++ b/Docs/Peripheral_Discovery.md @@ -0,0 +1 @@ +# Peripheral Discovery diff --git a/devices/cdashdisplay/info.go b/devices/cdashdisplay/info.go index ecd930a..14cbbe4 100644 --- a/devices/cdashdisplay/info.go +++ b/devices/cdashdisplay/info.go @@ -1,9 +1,27 @@ 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 const ( 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 +} diff --git a/devices/device.go b/devices/device.go index c7baab1..6c85b4f 100644 --- a/devices/device.go +++ b/devices/device.go @@ -1,9 +1,60 @@ -// Package devices simply defines what a device should have -// and some utils if necessary +// Package devices is a peripheral factory package devices -import "esdi/telemetry" +import ( + "errors" -type Device interface { - SendData(*telemetry.TelemetryData) + "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) { + // // 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") } diff --git a/devices/uidevice/device.go b/devices/uidevice/device.go new file mode 100644 index 0000000..d228982 --- /dev/null +++ b/devices/uidevice/device.go @@ -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 +} diff --git a/peripheral/peripheral.go b/peripheral/peripheral.go index e9d6f25..bdb035d 100644 --- a/peripheral/peripheral.go +++ b/peripheral/peripheral.go @@ -2,9 +2,11 @@ package peripheral import ( - "esdi/peripheral/devices" "fmt" "path/filepath" + + "esdi/peripheral/devices" + "esdi/telemetry" ) type PeripheralType string @@ -13,6 +15,12 @@ const ( DisplayPeripheral PeripheralType = "display" ) +type Peripheral interface { + Name() string + SendData(*telemetry.TelemetryData) + RequiredFields() map[int16]telemetry.FieldID +} + type PeripheralDeviceClerk struct { // mu sync.RWMutex Devices map[uint8]*PeripheralDevice diff --git a/providers/beamng/beamng.go b/providers/beamng/beamng.go index 58e245c..70e90ed 100644 --- a/providers/beamng/beamng.go +++ b/providers/beamng/beamng.go @@ -88,6 +88,10 @@ func (b *BeamNG) Close() { b.SDK.Close() } +func (b *BeamNG) Name() string { + return NAME +} + func (b *BeamNG) IsAlive(timeout time.Duration) bool { _, err := b.SDK.Update() return err == nil diff --git a/providers/iracing/iracing.go b/providers/iracing/iracing.go index 333364a..2c57f6e 100644 --- a/providers/iracing/iracing.go +++ b/providers/iracing/iracing.go @@ -264,6 +264,10 @@ func (i *IRacing) Subscribe(requestFields map[int16]telemetry.FieldID) { 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 diff --git a/providers/providers.go b/providers/providers.go index 00f0f08..a10e23d 100644 --- a/providers/providers.go +++ b/providers/providers.go @@ -16,7 +16,7 @@ import ( // the selected provider by its name type Provider struct { 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 } @@ -33,23 +33,31 @@ var Providers = map[string]Provider{ }, } -func NewLiveIRacingProvider(logger *slog.Logger) telemetry.TelemetryProvider { - provider, _ := iracing.NewIRacingProvider(logger, goirsdk.Options{ +func NewLiveIRacingProvider(logger *slog.Logger) (telemetry.TelemetryProvider, error) { + provider, err := iracing.NewIRacingProvider(logger, goirsdk.Options{ Logger: logger, 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 { +func NewBeamNGProvider(logger *slog.Logger) (telemetry.TelemetryProvider, error) { // 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), 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 + return provider, nil } diff --git a/services/devices.go b/services/devices.go index 9716930..c1119ff 100644 --- a/services/devices.go +++ b/services/devices.go @@ -2,23 +2,32 @@ package services import ( "context" + "errors" "fmt" "log/slog" + "sync" "sync/atomic" + "time" "esdi/devices" - "esdi/devices/cdashdisplay" + "esdi/peripheral" "esdi/telemetry" ) +var ErrPeripheralAlreadyRegistered = errors.New("peripheral is already registered") + // DeviceService will handle sending the data from the telemetry service to the // actual devices // NOTE: create a virtual device and make it be the output window or something so // we can just add it as a device or whatever instead of being a custom made thing // that would be pretty cool I think type DeviceService struct { - Logger *slog.Logger - Devices map[string]devices.Device + Logger *slog.Logger + // Device discovery + mu sync.RWMutex + ctxDiscovery context.Context + ctxDiscoveryCancel context.CancelFunc + Devices map[string]peripheral.Peripheral // Strem handling streamCancel context.CancelFunc TelemCh <-chan telemetry.TelemetryData @@ -28,41 +37,75 @@ type DeviceService struct { func NewDeviceService(logger *slog.Logger) *DeviceService { sharedChannel := make(chan string, 10) - return &DeviceService{ - Devices: make(map[string]devices.Device), + + dev := &DeviceService{ + Devices: make(map[string]peripheral.Peripheral), 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 } func (ds *DeviceService) FindDevices() { // Need to define a list of devices to search for // For now lets just try to find our cdashdisplay - will think about the rest later + ticker := time.NewTicker(2 * time.Second) + defer ticker.Stop() - // 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" + for { + select { + case <-ds.ctxDiscovery.Done(): + // If requested to cancel we cancel background discovery return - } + case <-ticker.C: + for pName, peripheral := range devices.List { + if ds.DeviceExists(pName) { + // We already discovered this device + continue + } - ds.Logger.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 + ds.Logger.Debug("looking for device", "name", pName) + dev, err := peripheral.Discover() + if err != nil { + ds.Logger.Debug("didn't find device", "name", pName) + continue + } + + // Register the device we just found + ds.RegisterDevice(dev) + } + } } } -func (ds *DeviceService) 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] if !ok { 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 // the data locked + ds.mu.RLock() for _, dev := range ds.Devices { dev.SendData(&data) } + ds.mu.RUnlock() + isSending.Store(false) } } diff --git a/services/telemetry.go b/services/telemetry.go index f35ae24..9dd7c53 100644 --- a/services/telemetry.go +++ b/services/telemetry.go @@ -54,9 +54,9 @@ func (t *TelemetryService) ProviderMonitor(ctx context.Context) { case <-ctx.Done(): return 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) { - slog.Debug("provider healthcheck failed") + slog.Warn("provider healthcheck failed") t.dropActiveProvider() t.onProviderHealthCheckFailed() return @@ -70,10 +70,14 @@ func (t *TelemetryService) onProviderHealthCheckFailed() { go t.FindProvider(t.CtxMonitor) } -func (t *TelemetryService) onFindProvider(prov providers.Provider) { +func (t *TelemetryService) onFindProvider(prov telem.TelemetryProvider) { // Attach to the provider - t.logger.Info("found provider for " + prov.Name) - t.SwitchProvider(prov.NewProvider(t.logger)) + 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()) @@ -102,9 +106,15 @@ func (t *TelemetryService) FindProvider(ctx context.Context) { return case <-ticker.C: for _, prov := range providers.Providers { - t.logger.Debug("checking provider: " + prov.Name) - if prov.IsRunning() { - t.onFindProvider(prov) + // 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 } } @@ -141,6 +151,7 @@ func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan tele 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 @@ -191,8 +202,14 @@ func (t *TelemetryService) UnsubscribeListener(id string) { } } -func (t *TelemetryService) SubscribeToFields(fields map[int16]telem.FieldID) { - t.activeProvider.Subscribe(fields) +func (t *TelemetryService) SubscribeToFields() { + // _ = t.devService.SubscribeFields() + // TODO: we need to find a way of requesting devices to send all subscribed fields + // instead of going through the devices on DevService here + for _, dev := range t.devService.Devices { + fields := dev.RequiredFields() + t.activeProvider.Subscribe(fields) + } } func (t *TelemetryService) StartStream() { diff --git a/telemetry/data.go b/telemetry/data.go index 158e815..c4f5ae0 100644 --- a/telemetry/data.go +++ b/telemetry/data.go @@ -135,7 +135,7 @@ func (tf *TelemetryField) String() string { return "NaN" } -type FieldID uint16 +type FieldID = uint16 // We use FirstField to start the count on the fields the user can select // the first three will be for internal use diff --git a/telemetry/provider.go b/telemetry/provider.go index e2cfeb5..93a405a 100644 --- a/telemetry/provider.go +++ b/telemetry/provider.go @@ -8,5 +8,6 @@ type TelemetryProvider interface { Stream() (<-chan TelemetryData, error) Subscribe(map[int16]FieldID) IsAlive(time.Duration) bool + Name() string Close() } diff --git a/tui/internal/controllers/cdashdisplay_layout.go b/tui/internal/controllers/cdashdisplay_layout.go index da21306..6400b47 100644 --- a/tui/internal/controllers/cdashdisplay_layout.go +++ b/tui/internal/controllers/cdashdisplay_layout.go @@ -211,14 +211,14 @@ func (lc *LayoutController) createWindow() { } // Acquire the cdashdisplay - displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) + displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) if err != nil { - lc.Messages <- "failed to get " + cdashdisplay.Name + lc.Messages <- "failed to get " + cdashdisplay.NAME return } display, ok := displayIF.(*cdashdisplay.CDashDisplay) if !ok { - lc.Messages <- "failed to acquire " + cdashdisplay.Name + lc.Messages <- "failed to acquire " + cdashdisplay.NAME return } // --- @@ -321,14 +321,14 @@ func (lc *LayoutController) newWindowAction() { func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) { // Acquire the cdashdisplay - displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) + displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) if err != nil { - lc.Messages <- "failed to get " + cdashdisplay.Name + lc.Messages <- "failed to get " + cdashdisplay.NAME return } display, ok := displayIF.(*cdashdisplay.CDashDisplay) if !ok { - lc.Messages <- "failed to acquire " + cdashdisplay.Name + lc.Messages <- "failed to acquire " + cdashdisplay.NAME return } // --- @@ -347,14 +347,14 @@ func (lc *LayoutController) displayLoadedLayouts() { lc.Logger.Debug("We want to view our layout!") // Acquire the cdashdisplay - displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) + displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) if err != nil { - lc.Messages <- "failed to get " + cdashdisplay.Name + lc.Messages <- "failed to get " + cdashdisplay.NAME return } display, ok := displayIF.(*cdashdisplay.CDashDisplay) if !ok { - lc.Messages <- "failed to acquire " + cdashdisplay.Name + lc.Messages <- "failed to acquire " + cdashdisplay.NAME return } // --- @@ -392,14 +392,14 @@ func (lc *LayoutController) getCurrentTreeNodeModel() (*tview.TreeNode, int16, e func (lc *LayoutController) loadLayout() { // Acquire the cdashdisplay - displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) + displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) if err != nil { - lc.Messages <- "failed to get " + cdashdisplay.Name + lc.Messages <- "failed to get " + cdashdisplay.NAME return } display, ok := displayIF.(*cdashdisplay.CDashDisplay) if !ok { - lc.Messages <- "failed to acquire " + cdashdisplay.Name + lc.Messages <- "failed to acquire " + cdashdisplay.NAME return } // --- @@ -416,14 +416,14 @@ func (lc *LayoutController) loadLayout() { func (lc *LayoutController) unloadLayout() { // Acquire the cdashdisplay - displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) + displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) if err != nil { - lc.Messages <- "failed to get " + cdashdisplay.Name + lc.Messages <- "failed to get " + cdashdisplay.NAME return } display, ok := displayIF.(*cdashdisplay.CDashDisplay) if !ok { - lc.Messages <- "failed to acquire " + cdashdisplay.Name + lc.Messages <- "failed to acquire " + cdashdisplay.NAME return } // --- @@ -437,14 +437,14 @@ func (lc *LayoutController) unloadLayout() { func (lc *LayoutController) saveLayout() { // Acquire the cdashdisplay - displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) + displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) if err != nil { - lc.Messages <- "failed to get " + cdashdisplay.Name + lc.Messages <- "failed to get " + cdashdisplay.NAME return } display, ok := displayIF.(*cdashdisplay.CDashDisplay) if !ok { - lc.Messages <- "failed to acquire " + cdashdisplay.Name + lc.Messages <- "failed to acquire " + cdashdisplay.NAME return } // --- @@ -470,14 +470,14 @@ func (lc *LayoutController) deleteWindow() { wID := curNode.GetReference().(int16) // Acquire the cdashdisplay - displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) + displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) if err != nil { - lc.Messages <- "failed to get " + cdashdisplay.Name + lc.Messages <- "failed to get " + cdashdisplay.NAME return } display, ok := displayIF.(*cdashdisplay.CDashDisplay) if !ok { - lc.Messages <- "failed to acquire " + cdashdisplay.Name + lc.Messages <- "failed to acquire " + cdashdisplay.NAME return } // --- diff --git a/tui/internal/controllers/cdashdisplay_layout_moveTool.go b/tui/internal/controllers/cdashdisplay_layout_moveTool.go index d363a58..ba31690 100644 --- a/tui/internal/controllers/cdashdisplay_layout_moveTool.go +++ b/tui/internal/controllers/cdashdisplay_layout_moveTool.go @@ -52,14 +52,14 @@ func (lc *LayoutController) handleMovementCapture(idx int16, } // Acquire the cdashdisplay - displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) + displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) if err != nil { - lc.Messages <- "failed to get " + cdashdisplay.Name + lc.Messages <- "failed to get " + cdashdisplay.NAME return nil } display, ok := displayIF.(*cdashdisplay.CDashDisplay) if !ok { - lc.Messages <- "failed to acquire " + cdashdisplay.Name + lc.Messages <- "failed to acquire " + cdashdisplay.NAME return nil } // --- @@ -86,14 +86,14 @@ func (lc *LayoutController) handleResizeCapture(idx int16, } // Acquire the cdashdisplay - displayIF, err := lc.DevService.GetDevice(cdashdisplay.Name) + displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME) if err != nil { - lc.Messages <- "failed to get " + cdashdisplay.Name + lc.Messages <- "failed to get " + cdashdisplay.NAME return nil } display, ok := displayIF.(*cdashdisplay.CDashDisplay) if !ok { - lc.Messages <- "failed to acquire " + cdashdisplay.Name + lc.Messages <- "failed to acquire " + cdashdisplay.NAME return nil } // --- diff --git a/tui/internal/controllers/device.go b/tui/internal/controllers/device.go index be1fa58..22c3e73 100644 --- a/tui/internal/controllers/device.go +++ b/tui/internal/controllers/device.go @@ -67,7 +67,7 @@ func (mc *DeviceController) AddDeviceAPIListItems() { mc.DeviceAPIView.DevAPIList. AddItem("layout", "build a layout for CDashDisplay", func() { // This CDashDisplay specific, only load if we have a CDashDisplay - if !mc.DevService.DeviceExists(cdashdisplay.Name) { + if !mc.DevService.DeviceExists(cdashdisplay.NAME) { mc.DevService.Messages <- "CDashDisplay it not loaded yet\n" return } diff --git a/tui/internal/controllers/streaming.go b/tui/internal/controllers/streaming.go index 0edfd28..cad0299 100644 --- a/tui/internal/controllers/streaming.go +++ b/tui/internal/controllers/streaming.go @@ -6,7 +6,7 @@ import ( "sync/atomic" "esdi/config" - "esdi/devices/cdashdisplay" + "esdi/devices/uidevice" "esdi/providers" "esdi/services" "esdi/telemetry" @@ -56,15 +56,15 @@ func NewStreamingCtrl( } ctrl.registerHooks() - ctrl.subscribeListeners() - go ctrl.listenToUIStream() + // ctrl.subscribeListeners() return ctrl } -func (sc *StreamingCtrl) subscribeListeners() { - sc.TelemetryCh = sc.TelemServ.SubscribeListener("UI", 1) -} +// func (sc *StreamingCtrl) subscribeListeners() { +// // Here I will set a UIDevice +// sc.TelemetryCh = sc.TelemServ.SubscribeListener("UI", 1) +// } func (sc *StreamingCtrl) registerHooks() { sc.StreamView.Options.Form.SetInputCapture(func(ev *tcell.EventKey) *tcell.EventKey { @@ -109,16 +109,21 @@ func (sc *StreamingCtrl) StartStop() { // stream is not running, we have to start it now // NOTE: // Subscribe the only existing device - needs to be discovered by now - slog.Debug("setting the data stream for cdash") - sc.Service.SetTelemetryChannel(sc.TelemServ.SubscribeListener("cdash", 1)) + slog.Debug("setting the data stream for device servie") + 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() - - slog.Debug("starting the stream") sc.TelemServ.StartStream() - slog.Debug("setting local control variables") sc.isRunning = true slog.Debug("starting stream") @@ -157,28 +162,23 @@ func (sc *StreamingCtrl) updateStream() { // so we can get away with using a map for convenience here func (sc *StreamingCtrl) SetInternalState() { // Acquire the cdashdisplay - displayIF, err := sc.Service.GetDevice(cdashdisplay.Name) - if err != nil { - sc.Messages <- "failed to get " + cdashdisplay.Name - return - } - display, ok := displayIF.(*cdashdisplay.CDashDisplay) - if !ok { - sc.Messages <- "failed to acquire " + cdashdisplay.Name - return - } + // displayIF, err := sc.Service.GetDevice(cdashdisplay.NAME) + // if err != nil { + // sc.Messages <- "failed to get " + cdashdisplay.NAME + // return + // } + // display, ok := displayIF.(*cdashdisplay.CDashDisplay) + // if !ok { + // sc.Messages <- "failed to acquire " + cdashdisplay.NAME + // return + // } // --- - fields := make(map[int16]telemetry.FieldID, len(display.State.Layout.Windows)) + sc.TelemServ.SubscribeToFields() - for _, w := range display.State.Layout.Windows { - fieldID, _ := telemetry.GetFieldID(w.UIData.TelemetryField) - fields[w.UIData.IDX] = fieldID - } - - sc.TelemServ.SubscribeToFields(fields) - - sc.Messages <- fmt.Sprintf("Subscribed Fields: %+v [%d]\n", fields, len(fields)) + // sc.Messages <- fmt.Sprintf("Subscribed Fields: %+v [%d]\n", fields, len(fields)) + // Should I update this? + sc.Messages <- fmt.Sprintf("Subscribed to fields\n") } func (sc *StreamingCtrl) listenToUIStream() { @@ -190,8 +190,13 @@ func (sc *StreamingCtrl) listenToUIStream() { } isDrawing.Store(true) + // sc.Logger.Debug("got data", "data", msg) + + // Capture locally + telemetryMsg := msg + sc.App.QueueUpdateDraw(func() { - sc.StreamView.Visualizer.Update(&msg) + sc.StreamView.Visualizer.Update(&telemetryMsg) isDrawing.Store(false) }) } diff --git a/tui/tui.go b/tui/tui.go index 374d21e..4add4bb 100644 --- a/tui/tui.go +++ b/tui/tui.go @@ -24,7 +24,7 @@ func NewControlPanel(logger *slog.Logger) *ControlPanel { } // NOTE: create our device service here - devService := services.NewDeviceService(logger) + devService := services.NewDeviceService(logger.With("service", "DeviceService")) telemService := services.NewTelemetryService(logger, devService) if telemService == nil {