diff --git a/Common/Client/go/wshub/hub.go b/Common/Client/go/wshub/hub.go index 833ffa3..a67f609 100644 --- a/Common/Client/go/wshub/hub.go +++ b/Common/Client/go/wshub/hub.go @@ -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(), } } diff --git a/Common/Client/go/wshub/sources.go b/Common/Client/go/wshub/sources.go index c5f04b8..23a6e77 100644 --- a/Common/Client/go/wshub/sources.go +++ b/Common/Client/go/wshub/sources.go @@ -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 diff --git a/Common/Client/go/wshub/sources_test.go b/Common/Client/go/wshub/sources_test.go new file mode 100644 index 0000000..ee3c6ae --- /dev/null +++ b/Common/Client/go/wshub/sources_test.go @@ -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") + } +}