wshub: persist calibration alongside sources; add config reload

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
Martino Ferrari
2026-08-16 19:27:23 +02:00
co-authored by Claude Sonnet 4.6
parent 66efd74dd5
commit dfd257cfd9
3 changed files with 213 additions and 12 deletions
+5
View File
@@ -239,6 +239,10 @@ type Hub struct {
sm *SourceManager // set after construction; used for WS-initiated source changes
// cal holds the per-signal calibration table. It is metadata only: the
// rings, the history and the trigger comparator all keep raw samples.
cal *calTable
// Ring buffers for hi-res zoom data.
// ringsMu protects the map structure; each sigRing has its own RWMutex for data.
ringsMu sync.RWMutex
@@ -270,6 +274,7 @@ func NewHub() *Hub {
rings: make(map[string]*sigRing),
statsMap: make(map[string]*SourceStat),
trigger: newTriggerEngine(),
cal: newCalTable(),
}
}
+76 -11
View File
@@ -1,12 +1,12 @@
package wshub
import (
"encoding/json"
"fmt"
"io"
"log"
"net"
"os"
"sort"
"strconv"
"strings"
"sync"
@@ -95,11 +95,16 @@ func (sm *SourceManager) Remove(id string) {
}
}
// Save writes the current source list to filePath.
func (sm *SourceManager) Save() error {
if sm.filePath == "" {
return fmt.Errorf("no sources-file configured")
// Path returns the configured config-file path ("" when none).
func (sm *SourceManager) Path() string {
sm.mu.RLock()
defer sm.mu.RUnlock()
return sm.filePath
}
// snapshotSources returns the current sources sorted by label, so the written
// file is byte-stable across runs (the map iteration order is not).
func (sm *SourceManager) snapshotSources() []SourceConfig {
sm.mu.RLock()
cfgs := make([]SourceConfig, 0, len(sm.sources))
for _, ms := range sm.sources {
@@ -111,26 +116,86 @@ func (sm *SourceManager) Save() error {
})
}
sm.mu.RUnlock()
sort.Slice(cfgs, func(i, j int) bool {
if cfgs[i].Label != cfgs[j].Label {
return cfgs[i].Label < cfgs[j].Label
}
return cfgs[i].Addr < cfgs[j].Addr
})
return cfgs
}
data, err := json.MarshalIndent(cfgs, "", " ")
// Save writes the current source list and calibration table to filePath as one
// flat JSON array.
func (sm *SourceManager) Save() error {
path := sm.Path()
if path == "" {
return fmt.Errorf("no sources-file configured")
}
data, err := encodeConfigFile(sm.snapshotSources(), sm.hub.cal.List())
if err != nil {
return err
}
return os.WriteFile(sm.filePath, data, 0644)
return os.WriteFile(path, data, 0644)
}
// Load reads sources from path and adds them.
// Load reads the config file at path, replaces the calibration table with its
// contents and starts every source it lists.
func (sm *SourceManager) Load(path string) error {
data, err := os.ReadFile(path)
if err != nil {
return err
}
var cfgs []SourceConfig
if err := json.Unmarshal(data, &cfgs); err != nil {
srcs, cals, err := parseConfigFile(data)
if err != nil {
return err
}
sm.mu.Lock()
sm.filePath = path
for _, cfg := range cfgs {
sm.mu.Unlock()
sm.hub.cal.Replace(cals)
for _, cfg := range srcs {
sm.Add(cfg.Label, cfg.Addr, cfg.MulticastGroup, cfg.DataPort)
}
return nil
}
// Reload re-reads the config file. The calibration table is replaced wholesale
// and sources listed in the file that are not already running are started; no
// live source is ever stopped, restarted or reconnected, because a reload must
// not interrupt streaming. The asymmetry is deliberate: calibration is cheap
// to reapply, a source is a live UDP session.
func (sm *SourceManager) Reload() error {
path := sm.Path()
if path == "" {
return fmt.Errorf("no sources-file configured")
}
data, err := os.ReadFile(path)
if err != nil {
return err
}
srcs, cals, err := parseConfigFile(data)
if err != nil {
return err
}
sm.hub.cal.Replace(cals)
sm.mu.RLock()
live := make(map[string]bool, len(sm.sources))
for _, ms := range sm.sources {
live[ms.label+"\x00"+ms.addr] = true
}
sm.mu.RUnlock()
for _, cfg := range srcs {
label := cfg.Label
if label == "" {
label = cfg.Addr // Add() applies the same default
}
if live[label+"\x00"+cfg.Addr] {
continue
}
sm.Add(cfg.Label, cfg.Addr, cfg.MulticastGroup, cfg.DataPort)
}
return nil
+131
View File
@@ -0,0 +1,131 @@
package wshub
import (
"os"
"path/filepath"
"testing"
)
// newTestManager builds a hub + manager pair with no goroutines running.
func newTestManager(t *testing.T) (*Hub, *SourceManager, string) {
t.Helper()
path := filepath.Join(t.TempDir(), "sources.json")
h := NewHub()
sm := NewSourceManager(h, path)
h.SetSourceManager(sm)
return h, sm, path
}
func TestSaveWritesSourcesAndCalibration(t *testing.T) {
h, sm, path := newTestManager(t)
// Register two sources without starting any UDP client.
sm.mu.Lock()
sm.sources["s1"] = &managedSource{id: "s1", label: "wave", addr: "127.0.0.1:44500"}
sm.sources["s2"] = &managedSource{
id: "s2", label: "mc", addr: "127.0.0.1:44501",
multicastGroup: "239.0.0.1", dataPort: 44502,
}
sm.mu.Unlock()
if !h.cal.Set(CalConfig{Source: "wave", Signal: "Adc", Scale: 0.5, Offset: -1.25, Unit: "V"}) {
t.Fatal("calibration rejected")
}
if err := sm.Save(); err != nil {
t.Fatalf("Save: %v", err)
}
data, err := os.ReadFile(path)
if err != nil {
t.Fatalf("ReadFile: %v", err)
}
srcs, cals, err := parseConfigFile(data)
if err != nil {
t.Fatalf("parseConfigFile: %v\n%s", err, data)
}
if len(srcs) != 2 {
t.Fatalf("got %d sources, want 2\n%s", len(srcs), data)
}
// Save sorts by label so the file is byte-stable across runs.
if srcs[0].Label != "mc" || srcs[1].Label != "wave" {
t.Errorf("source order = %q,%q, want mc,wave", srcs[0].Label, srcs[1].Label)
}
if len(cals) != 1 || cals[0].Signal != "Adc" || cals[0].Scale != 0.5 {
t.Fatalf("calibration round-trip failed: %+v\n%s", cals, data)
}
}
func TestSaveWithoutFilePathFails(t *testing.T) {
h := NewHub()
sm := NewSourceManager(h, "")
h.SetSourceManager(sm)
if err := sm.Save(); err == nil {
t.Error("Save() with no path = nil error, want error")
}
}
func TestLoadSeedsCalibrationTable(t *testing.T) {
h, sm, path := newTestManager(t)
// No "addr" blocks: Load must not start any UDP client during the test.
if err := os.WriteFile(path, []byte(`[
{"source":"wave","signal":"Adc","scale":0.25,"offset":2,"unit":"mV"},
{"source":"wave","signal":"Dac","scale":2}
]`), 0o644); err != nil {
t.Fatal(err)
}
if err := sm.Load(path); err != nil {
t.Fatalf("Load: %v", err)
}
got := h.cal.List()
if len(got) != 2 {
t.Fatalf("List() = %d entries, want 2", len(got))
}
if got[0].Signal != "Adc" || got[0].Unit != "mV" || got[0].Offset != 2 {
t.Errorf("Adc = %+v", got[0])
}
if sm.Path() != path {
t.Errorf("Path() = %q, want %q", sm.Path(), path)
}
}
func TestReloadReplacesCalibrationAndKeepsLiveSources(t *testing.T) {
h, sm, path := newTestManager(t)
// A live source that the file does not mention must survive the reload.
sm.mu.Lock()
sm.sources["s1"] = &managedSource{id: "s1", label: "live", addr: "127.0.0.1:44999"}
sm.mu.Unlock()
// A stale calibration that the file does not mention must be dropped.
h.cal.Set(CalConfig{Source: "stale", Signal: "Old", Scale: 9})
if err := os.WriteFile(path, []byte(`[
{"source":"wave","signal":"Adc","scale":0.5}
]`), 0o644); err != nil {
t.Fatal(err)
}
if err := sm.Reload(); err != nil {
t.Fatalf("Reload: %v", err)
}
got := h.cal.List()
if len(got) != 1 || got[0].Source != "wave" {
t.Fatalf("after Reload, calibration = %+v, want only wave/Adc", got)
}
sm.mu.RLock()
_, alive := sm.sources["s1"]
n := len(sm.sources)
sm.mu.RUnlock()
if !alive || n != 1 {
t.Errorf("live source count = %d (s1 alive=%v), want 1 / true", n, alive)
}
}
func TestReloadWithoutFilePathFails(t *testing.T) {
h := NewHub()
sm := NewSourceManager(h, "")
h.SetSourceManager(sm)
if err := sm.Reload(); err == nil {
t.Error("Reload() with no path = nil error, want error")
}
}