partial subscriptions

not really partial. they just happen per peripheral now
This commit is contained in:
2026-09-23 17:41:06 +01:00
parent a519cf473e
commit 2952528a60
12 changed files with 172 additions and 110 deletions
+4 -22
View File
@@ -26,7 +26,7 @@ type DeviceService struct {
// Callbacks
// Telemetry service data fetchers
telemetryProvider func() (string, error)
triggerFieldSubscription func() []telemetry.FieldID
triggerFieldSubscription func([]telemetry.FieldID) []string
}
func NewDeviceService(logger *slog.Logger, msg chan string) *DeviceService {
@@ -80,24 +80,6 @@ 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] -------------------------------------------------------------
@@ -168,10 +150,10 @@ func (ds *DeviceService) transmit(ctx context.Context) {
}
// Callbacks [START] -----------------------------------------------------------
func (ds *DeviceService) peripheralConfigured(pname string) {
func (ds *DeviceService) peripheralConfigured(pname string, fields []telemetry.FieldID) {
// We need to retrigger field subscription here
fields := ds.triggerFieldSubscription()
ds.Messages <- fmt.Sprintf("subscribed to fields: %+v\n", fields)
subscribedTo := ds.triggerFieldSubscription(fields)
ds.Messages <- fmt.Sprintf("subscribed to fields: %q\n", subscribedTo)
}
// Callbacks [END] -------------------------------------------------------------
+12 -3
View File
@@ -10,6 +10,7 @@ import (
"esdi/devices"
"esdi/peripheral"
"esdi/telemetry"
)
var (
@@ -86,7 +87,7 @@ type PeripheralStateStore struct {
Messages chan string
// Callbacks
telemetryProvider func() (string, error)
onPeripheralConfigured func(string)
onPeripheralConfigured func(string, []telemetry.FieldID)
// Internal State
isStreaming bool
}
@@ -184,6 +185,15 @@ func (pss *PeripheralStateStore) GetStreamingState() bool {
return pss.isStreaming
}
func (pss *PeripheralStateStore) GetPeripheralFields(pname string) []telemetry.FieldID {
per, err := pss.GetState(pname)
if err != nil {
return nil
}
return per.Peripheral.RequiredFields()
}
// "Events" [START] ------------------------------------------------------------
func (ds *DeviceService) onDeviceTimedOut(pname string) {
ds.PSS.setDeviceTimedOut(pname)
@@ -280,7 +290,7 @@ func (pss *PeripheralStateStore) setDeviceConfigured(pname string) {
pss.store[pname].State = DeviceIsConfigured
pss.mu.Unlock()
pss.onPeripheralConfigured(pname)
pss.onPeripheralConfigured(pname, pss.GetPeripheralFields(pname))
}
func (pss *PeripheralStateStore) setDeviceIsStreaming(pname string) {
@@ -347,7 +357,6 @@ func (pss *PeripheralStateStore) configurePeripheral(
}
onSuccess(pname)
pss.setDeviceConfigured(pname)
}()
}
+2 -2
View File
@@ -26,10 +26,10 @@ func NewOrchestrator(logger *slog.Logger) (*Orchestrator, error) {
// Setup device service callbacks
devService.telemetryProvider = telemService.GetTelemetryProviderName
devService.triggerFieldSubscription = telemService.SubscribeToAllFields
devService.triggerFieldSubscription = telemService.SubscribeToFields
// Setup telemetry service callbacks
telemService.getRequiredFields = devService.GetRequiredFields
// telemService.getRequiredFields = devService.GetRequiredFields
return &Orchestrator{
DeviceService: devService,
+4 -6
View File
@@ -37,7 +37,7 @@ type TelemetryService struct {
// Callbacks
// Devices data request
// peripheralProvider func() []peripheral.Peripheral
getRequiredFields func() []telemetry.FieldID
// getRequiredFields func() []telemetry.FieldID
}
func NewTelemetryService(
@@ -114,13 +114,11 @@ func (t *TelemetryService) UnsubscribeListener(id string) {
}
}
func (t *TelemetryService) SubscribeToAllFields() []telem.FieldID {
fields := t.getRequiredFields()
func (t *TelemetryService) SubscribeToFields(fields []telemetry.FieldID) []string {
t.logger.Debug("requested fields", "fields", fields)
t.activeProvider.Subscribe(fields)
subscribed := t.activeProvider.Subscribe(fields)
return fields
return subscribed
}
// Listener Control [END] ------------------------------------------------------