Mid stream peripheral connection #16

Merged
esilva merged 2 commits from mid-stream-peripheral-connection into auto-detect-devices 2026-09-23 17:42:19 +01:00
6 changed files with 53 additions and 25 deletions
Showing only changes of commit a519cf473e - Show all commits
+29 -1
View File
@@ -2,6 +2,7 @@ package services
import ( import (
"context" "context"
"fmt"
"log/slog" "log/slog"
"sync/atomic" "sync/atomic"
@@ -25,6 +26,7 @@ type DeviceService struct {
// Callbacks // Callbacks
// Telemetry service data fetchers // 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 { 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 // Set the callbacks for PSS
dev.PSS.telemetryProvider = dev.getTelemetryProvider dev.PSS.telemetryProvider = dev.getTelemetryProvider
dev.PSS.onPeripheralConfigured = dev.peripheralConfigured
return dev return dev
} }
@@ -59,7 +62,7 @@ func (ds *DeviceService) GetDevices() []peripheral.Peripheral {
peripherals := make([]peripheral.Peripheral, 0, len(snapshot)) peripherals := make([]peripheral.Peripheral, 0, len(snapshot))
for _, state := range snapshot { for _, state := range snapshot {
if state.State < DeviceIsConnected { if state.State < DeviceIsConfigured {
continue continue
} }
peripherals = append(peripherals, state.Peripheral) peripherals = append(peripherals, state.Peripheral)
@@ -77,6 +80,24 @@ func (ds *DeviceService) PeripheralExists(pname string) bool {
return err == nil 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] --------------------------------------------------------------- // Getters [END] ---------------------------------------------------------------
// Actions [START] ------------------------------------------------------------- // Actions [START] -------------------------------------------------------------
@@ -95,6 +116,8 @@ func (ds *DeviceService) StopStream() {
return return
} }
ds.PSS.OnStopStream()
ds.streamCancel() ds.streamCancel()
ds.streamCancel = nil ds.streamCancel = nil
} }
@@ -145,5 +168,10 @@ func (ds *DeviceService) transmit(ctx context.Context) {
} }
// Callbacks [START] ----------------------------------------------------------- // 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] ------------------------------------------------------------- // Callbacks [END] -------------------------------------------------------------
+10 -1
View File
@@ -86,6 +86,7 @@ type PeripheralStateStore struct {
Messages chan string Messages chan string
// Callbacks // Callbacks
telemetryProvider func() (string, error) telemetryProvider func() (string, error)
onPeripheralConfigured func(string)
// Internal State // Internal State
isStreaming bool isStreaming bool
} }
@@ -204,6 +205,12 @@ func (pss *PeripheralStateStore) OnStopStream() {
pss.mu.Lock() pss.mu.Lock()
pss.isStreaming = false pss.isStreaming = false
pss.mu.Unlock() pss.mu.Unlock()
for _, state := range pss.GetStates() {
if state.State == DeviceIsStreaming {
pss.setDeviceConfigured(state.device.Name)
}
}
} }
// "Events" [END] -------------------------------------------------------------- // "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.Logger.Info("device is configured and ready for data", "device", pname)
pss.mu.Lock() pss.mu.Lock()
defer pss.mu.Unlock()
pss.store[pname].State = DeviceIsConfigured pss.store[pname].State = DeviceIsConfigured
pss.mu.Unlock()
pss.onPeripheralConfigured(pname)
} }
func (pss *PeripheralStateStore) setDeviceIsStreaming(pname string) { func (pss *PeripheralStateStore) setDeviceIsStreaming(pname string) {
+2 -1
View File
@@ -26,9 +26,10 @@ func NewOrchestrator(logger *slog.Logger) (*Orchestrator, error) {
// Setup device service callbacks // Setup device service callbacks
devService.telemetryProvider = telemService.GetTelemetryProviderName devService.telemetryProvider = telemService.GetTelemetryProviderName
devService.triggerFieldSubscription = telemService.SubscribeToAllFields
// Setup telemetry service callbacks // Setup telemetry service callbacks
telemService.peripheralProvider = devService.GetDevices telemService.getRequiredFields = devService.GetRequiredFields
return &Orchestrator{ return &Orchestrator{
DeviceService: devService, DeviceService: devService,
+7 -17
View File
@@ -7,7 +7,6 @@ import (
"sync" "sync"
"time" "time"
"esdi/peripheral"
"esdi/providers" "esdi/providers"
"esdi/telemetry" "esdi/telemetry"
telem "esdi/telemetry" telem "esdi/telemetry"
@@ -37,7 +36,8 @@ type TelemetryService struct {
healthCheckCancel context.CancelFunc healthCheckCancel context.CancelFunc
// Callbacks // Callbacks
// Devices data request // Devices data request
peripheralProvider func() []peripheral.Peripheral // peripheralProvider func() []peripheral.Peripheral
getRequiredFields func() []telemetry.FieldID
} }
func NewTelemetryService( func NewTelemetryService(
@@ -114,23 +114,13 @@ func (t *TelemetryService) UnsubscribeListener(id string) {
} }
} }
func (t *TelemetryService) SubscribeToFields() []telem.FieldID { func (t *TelemetryService) SubscribeToAllFields() []telem.FieldID {
seen := make(map[telemetry.FieldID]struct{}) fields := t.getRequiredFields()
var allFields []telemetry.FieldID
for _, dev := range t.peripheralProvider() { t.logger.Debug("requested fields", "fields", fields)
for _, field := range dev.RequiredFields() { t.activeProvider.Subscribe(fields)
if _, exists := seen[field]; !exists {
seen[field] = struct{}{}
allFields = append(allFields, field)
}
}
}
t.logger.Debug("requested fields", "fields", allFields) return fields
t.activeProvider.Subscribe(allFields)
return allFields
} }
// Listener Control [END] ------------------------------------------------------ // Listener Control [END] ------------------------------------------------------
+1 -1
View File
@@ -93,7 +93,7 @@ func (mc *DeviceController) AddDeviceAPIListItems() {
mc.StreamCtrl.StreamView.Flex, mc.StreamCtrl.StreamView.Flex,
) )
mc.StreamCtrl.SetInternalState() // mc.StreamCtrl.SetInternalState()
mc.App.SetFocus(mc.StreamCtrl.StreamView.Options.Form) mc.App.SetFocus(mc.StreamCtrl.StreamView.Options.Form)
}) })
+2 -2
View File
@@ -154,8 +154,8 @@ func (sc *StreamingCtrl) updateStream() {
// Performance reasoning: this is not used during the high frequency data transmission // Performance reasoning: this is not used during the high frequency data transmission
// so we can get away with using a map for convenience here // so we can get away with using a map for convenience here
func (sc *StreamingCtrl) SetInternalState() { func (sc *StreamingCtrl) SetInternalState() {
fields := sc.TelemServ.SubscribeToFields() // fields := sc.TelemServ.SubscribeToAllFields()
sc.Messages <- fmt.Sprintf("Subscribed to fields: %+v\n", fields) // sc.Messages <- fmt.Sprintf("Subscribed to fields: %+v\n", fields)
} }
func (sc *StreamingCtrl) listenToUIStream() { func (sc *StreamingCtrl) listenToUIStream() {