diff --git a/Test/E2E/suite/client/.gitignore b/Test/E2E/suite/client/.gitignore new file mode 100644 index 0000000..aa071d2 --- /dev/null +++ b/Test/E2E/suite/client/.gitignore @@ -0,0 +1,2 @@ +chain-client +configcheck/configcheck diff --git a/Test/E2E/suite/client/configcheck/main.go b/Test/E2E/suite/client/configcheck/main.go new file mode 100644 index 0000000..0e2d413 --- /dev/null +++ b/Test/E2E/suite/client/configcheck/main.go @@ -0,0 +1,232 @@ +// configcheck exercises the calibration and config-persistence WebSocket frames +// against a StreamHub (either the Go hub or the C++ StreamHub) and exits +// non-zero if the hub's replies do not match the protocol. +package main + +import ( + "encoding/json" + "flag" + "fmt" + "os" + "time" + + "github.com/gorilla/websocket" +) + +type calEntry struct { + Source string `json:"source"` + Signal string `json:"signal"` + Scale float64 `json:"scale"` + Offset float64 `json:"offset"` + Unit string `json:"unit"` +} + +type frame struct { + Type string `json:"type"` + Cal []calEntry `json:"cal"` + OK bool `json:"ok"` + Path string `json:"path"` + Error string `json:"error"` +} + +type conn struct { + ws *websocket.Conn + timeout time.Duration + // readCh is a long-lived goroutine that forwards frames; nil until startReader. + readCh chan readResult +} + +type readResult struct { + f frame + err error +} + +// startReader launches a background goroutine that reads all text frames from +// the WebSocket and forwards them on readCh. This avoids setting a read +// deadline on the underlying connection, which permanently poisons gorilla +// websocket after a timeout fires. +func (c *conn) startReader() { + c.readCh = make(chan readResult, 32) + go func() { + for { + mt, data, err := c.ws.ReadMessage() + if err != nil { + c.readCh <- readResult{err: fmt.Errorf("ws read: %w", err)} + return + } + if mt != websocket.TextMessage { + continue + } + var f frame + if jsonErr := json.Unmarshal(data, &f); jsonErr != nil { + continue + } + c.readCh <- readResult{f: f} + } + }() +} + +// next reads from the background reader until a frame with the wanted type +// arrives, or the per-frame deadline passes. +func (c *conn) next(want string) (frame, error) { + deadline := time.NewTimer(c.timeout) + defer deadline.Stop() + for { + select { + case <-deadline.C: + return frame{}, fmt.Errorf("timeout waiting for %q", want) + case r, ok := <-c.readCh: + if !ok { + return frame{}, fmt.Errorf("reader closed while waiting for %q", want) + } + if r.err != nil { + return frame{}, fmt.Errorf("read while waiting for %q: %w", want, r.err) + } + if r.f.Type == want { + return r.f, nil + } + } + } +} + +// nextWithin reads from the background reader until a frame with the wanted +// type arrives within d, returning (frame, true) or (frame{}, false). Unlike +// next() it does NOT return an error on timeout, making it suitable for the +// "must NOT arrive" assertion. +func (c *conn) nextWithin(want string, d time.Duration) (frame, bool) { + deadline := time.NewTimer(d) + defer deadline.Stop() + for { + select { + case <-deadline.C: + return frame{}, false + case r, ok := <-c.readCh: + if !ok || r.err != nil { + return frame{}, false + } + if r.f.Type == want { + return r.f, true + } + } + } +} + +func (c *conn) send(v interface{}) error { + data, err := json.Marshal(v) + if err != nil { + return err + } + return c.ws.WriteMessage(websocket.TextMessage, data) +} + +func findCal(list []calEntry, source, signal string) (calEntry, bool) { + for _, e := range list { + if e.Source == source && e.Signal == signal { + return e, true + } + } + return calEntry{}, false +} + +func run(url, source, signal string, timeout time.Duration) error { + ws, _, err := websocket.DefaultDialer.Dial(url, nil) + if err != nil { + return fmt.Errorf("dial %s: %w", url, err) + } + defer ws.Close() + c := &conn{ws: ws, timeout: timeout} + c.startReader() + + // 1. The hub sends a calibration frame on connect, even when empty. + if _, err := c.next("calibration"); err != nil { + return fmt.Errorf("on connect: %w", err) + } + + // 2. setCalibration is accepted and echoed back to every client. + want := calEntry{Source: source, Signal: signal, Scale: 0.5, Offset: -1.25, Unit: "V"} + if err := c.send(map[string]interface{}{ + "type": "setCalibration", + "source": want.Source, + "signal": want.Signal, + "scale": want.Scale, + "offset": want.Offset, + "unit": want.Unit, + }); err != nil { + return err + } + f, err := c.next("calibration") + if err != nil { + return fmt.Errorf("after setCalibration: %w", err) + } + got, ok := findCal(f.Cal, source, signal) + if !ok { + return fmt.Errorf("setCalibration: entry %s/%s missing from broadcast", source, signal) + } + if got != want { + return fmt.Errorf("setCalibration: got %+v, want %+v", got, want) + } + + // 3. An invalid entry (scale = 0) must be rejected: no broadcast follows. + if err := c.send(map[string]interface{}{ + "type": "setCalibration", "source": source, "signal": signal, + "scale": 0.0, "offset": 0.0, "unit": "", + }); err != nil { + return err + } + if _, accepted := c.nextWithin("calibration", 500*time.Millisecond); accepted { + return fmt.Errorf("setCalibration with scale=0 was accepted, must be rejected") + } + + // 4. saveSources acknowledges with configSaved. + if err := c.send(map[string]string{"type": "saveSources"}); err != nil { + return err + } + f, err = c.next("configSaved") + if err != nil { + return err + } + if !f.OK { + return fmt.Errorf("configSaved: ok=false, error=%q", f.Error) + } + if f.Path == "" { + return fmt.Errorf("configSaved: ok=true but path is empty") + } + + // 5. reloadConfig acknowledges and re-broadcasts the saved calibration. + if err := c.send(map[string]string{"type": "reloadConfig"}); err != nil { + return err + } + f, err = c.next("configReloaded") + if err != nil { + return err + } + if !f.OK { + return fmt.Errorf("configReloaded: ok=false, error=%q", f.Error) + } + f, err = c.next("calibration") + if err != nil { + return fmt.Errorf("after reloadConfig: %w", err) + } + got, ok = findCal(f.Cal, source, signal) + if !ok { + return fmt.Errorf("reloadConfig: entry %s/%s did not survive the round-trip", source, signal) + } + if got != want { + return fmt.Errorf("reloadConfig: got %+v, want %+v", got, want) + } + return nil +} + +func main() { + url := flag.String("url", "ws://127.0.0.1:8090/ws", "hub WebSocket URL") + source := flag.String("source", "cfgcheck", "calibration source label to use") + signal := flag.String("signal", "Probe", "calibration signal name to use") + timeout := flag.Duration("timeout", 5*time.Second, "per-frame timeout") + flag.Parse() + + if err := run(*url, *source, *signal, *timeout); err != nil { + fmt.Fprintf(os.Stderr, "configcheck FAIL: %v\n", err) + os.Exit(1) + } + fmt.Println("configcheck OK") +}