From a519cf473ecd25f469c453bbbce6ca5aca25edbf Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Wed, 23 Sep 2026 13:40:31 +0100 Subject: [PATCH] peripheral can connect mid stream and start go into streaming immidiately --- services/devices.go | 32 +++++++++++++++++++++++++-- services/devices_lookup.go | 13 +++++++++-- services/services.go | 3 ++- services/telemetry.go | 24 ++++++-------------- tui/internal/controllers/device.go | 2 +- tui/internal/controllers/streaming.go | 4 ++-- 6 files changed, 53 insertions(+), 25 deletions(-) diff --git a/services/devices.go b/services/devices.go index c688821..8448d90 100644 --- a/services/devices.go +++ b/services/devices.go @@ -2,6 +2,7 @@ package services import ( "context" + "fmt" "log/slog" "sync/atomic" @@ -24,7 +25,8 @@ type DeviceService struct { Messages chan string // Callbacks // Telemetry service data fetchers - telemetryProvider func() (string, error) + telemetryProvider func() (string, error) + triggerFieldSubscription func() []telemetry.FieldID } func NewDeviceService(logger *slog.Logger, msg chan string) *DeviceService { @@ -43,6 +45,7 @@ func NewDeviceService(logger *slog.Logger, msg chan string) *DeviceService { // Set the callbacks for PSS dev.PSS.telemetryProvider = dev.getTelemetryProvider + dev.PSS.onPeripheralConfigured = dev.peripheralConfigured return dev } @@ -59,7 +62,7 @@ func (ds *DeviceService) GetDevices() []peripheral.Peripheral { peripherals := make([]peripheral.Peripheral, 0, len(snapshot)) for _, state := range snapshot { - if state.State < DeviceIsConnected { + if state.State < DeviceIsConfigured { continue } peripherals = append(peripherals, state.Peripheral) @@ -77,6 +80,24 @@ func (ds *DeviceService) PeripheralExists(pname string) bool { return err == nil } +func (ds *DeviceService) GetRequiredFields() []telemetry.FieldID { + // NOTE: this can be optimized, not that it matters at this stage, but if + // it runs while telemetry is running we want it optimized I guess + seen := make(map[telemetry.FieldID]struct{}) + var allFields []telemetry.FieldID + + for _, dev := range ds.GetDevices() { + for _, field := range dev.RequiredFields() { + if _, exists := seen[field]; !exists { + seen[field] = struct{}{} + allFields = append(allFields, field) + } + } + } + + return allFields +} + // Getters [END] --------------------------------------------------------------- // Actions [START] ------------------------------------------------------------- @@ -95,6 +116,8 @@ func (ds *DeviceService) StopStream() { return } + ds.PSS.OnStopStream() + ds.streamCancel() ds.streamCancel = nil } @@ -145,5 +168,10 @@ func (ds *DeviceService) transmit(ctx context.Context) { } // Callbacks [START] ----------------------------------------------------------- +func (ds *DeviceService) peripheralConfigured(pname string) { + // We need to retrigger field subscription here + fields := ds.triggerFieldSubscription() + ds.Messages <- fmt.Sprintf("subscribed to fields: %+v\n", fields) +} // Callbacks [END] ------------------------------------------------------------- diff --git a/services/devices_lookup.go b/services/devices_lookup.go index 499b8ea..3c46eeb 100644 --- a/services/devices_lookup.go +++ b/services/devices_lookup.go @@ -85,7 +85,8 @@ type PeripheralStateStore struct { // Messaging for UI and stuff Messages chan string // Callbacks - telemetryProvider func() (string, error) + telemetryProvider func() (string, error) + onPeripheralConfigured func(string) // Internal State isStreaming bool } @@ -204,6 +205,12 @@ func (pss *PeripheralStateStore) OnStopStream() { pss.mu.Lock() pss.isStreaming = false pss.mu.Unlock() + + for _, state := range pss.GetStates() { + if state.State == DeviceIsStreaming { + pss.setDeviceConfigured(state.device.Name) + } + } } // "Events" [END] -------------------------------------------------------------- @@ -270,8 +277,10 @@ func (pss *PeripheralStateStore) setDeviceConfigured(pname string) { pss.Logger.Info("device is configured and ready for data", "device", pname) pss.mu.Lock() - defer pss.mu.Unlock() pss.store[pname].State = DeviceIsConfigured + pss.mu.Unlock() + + pss.onPeripheralConfigured(pname) } func (pss *PeripheralStateStore) setDeviceIsStreaming(pname string) { diff --git a/services/services.go b/services/services.go index 18aff69..d8f1ec7 100644 --- a/services/services.go +++ b/services/services.go @@ -26,9 +26,10 @@ func NewOrchestrator(logger *slog.Logger) (*Orchestrator, error) { // Setup device service callbacks devService.telemetryProvider = telemService.GetTelemetryProviderName + devService.triggerFieldSubscription = telemService.SubscribeToAllFields // Setup telemetry service callbacks - telemService.peripheralProvider = devService.GetDevices + telemService.getRequiredFields = devService.GetRequiredFields return &Orchestrator{ DeviceService: devService, diff --git a/services/telemetry.go b/services/telemetry.go index 7ca2112..2f9983d 100644 --- a/services/telemetry.go +++ b/services/telemetry.go @@ -7,7 +7,6 @@ import ( "sync" "time" - "esdi/peripheral" "esdi/providers" "esdi/telemetry" telem "esdi/telemetry" @@ -37,7 +36,8 @@ type TelemetryService struct { healthCheckCancel context.CancelFunc // Callbacks // Devices data request - peripheralProvider func() []peripheral.Peripheral + // peripheralProvider func() []peripheral.Peripheral + getRequiredFields func() []telemetry.FieldID } func NewTelemetryService( @@ -114,23 +114,13 @@ func (t *TelemetryService) UnsubscribeListener(id string) { } } -func (t *TelemetryService) SubscribeToFields() []telem.FieldID { - seen := make(map[telemetry.FieldID]struct{}) - var allFields []telemetry.FieldID +func (t *TelemetryService) SubscribeToAllFields() []telem.FieldID { + fields := t.getRequiredFields() - for _, dev := range t.peripheralProvider() { - for _, field := range dev.RequiredFields() { - if _, exists := seen[field]; !exists { - seen[field] = struct{}{} - allFields = append(allFields, field) - } - } - } + t.logger.Debug("requested fields", "fields", fields) + t.activeProvider.Subscribe(fields) - t.logger.Debug("requested fields", "fields", allFields) - t.activeProvider.Subscribe(allFields) - - return allFields + return fields } // Listener Control [END] ------------------------------------------------------ diff --git a/tui/internal/controllers/device.go b/tui/internal/controllers/device.go index 9d3d653..86f67c7 100644 --- a/tui/internal/controllers/device.go +++ b/tui/internal/controllers/device.go @@ -93,7 +93,7 @@ func (mc *DeviceController) AddDeviceAPIListItems() { mc.StreamCtrl.StreamView.Flex, ) - mc.StreamCtrl.SetInternalState() + // mc.StreamCtrl.SetInternalState() mc.App.SetFocus(mc.StreamCtrl.StreamView.Options.Form) }) diff --git a/tui/internal/controllers/streaming.go b/tui/internal/controllers/streaming.go index 4fd57c1..8f40005 100644 --- a/tui/internal/controllers/streaming.go +++ b/tui/internal/controllers/streaming.go @@ -154,8 +154,8 @@ func (sc *StreamingCtrl) updateStream() { // Performance reasoning: this is not used during the high frequency data transmission // so we can get away with using a map for convenience here func (sc *StreamingCtrl) SetInternalState() { - fields := sc.TelemServ.SubscribeToFields() - sc.Messages <- fmt.Sprintf("Subscribed to fields: %+v\n", fields) + // fields := sc.TelemServ.SubscribeToAllFields() + // sc.Messages <- fmt.Sprintf("Subscribed to fields: %+v\n", fields) } func (sc *StreamingCtrl) listenToUIStream() {