Reported as samples sporadically carrying a previous packet's timestamp: holes on one side of the stream and collisions on the other, in both the Go and the MARTe2 receiver. That it appeared in both is what located it -- the shared cause is upstream of either client. Four independent defects, all of which end in a packet's values being placed at a time that is not theirs. Reassembly slot exhaustion (the "Reassembly slots full; evicting oldest" flood). Chunk size was learnt only from fragment 0, so an out-of-order burst destroyed a packet whose bytes had all arrived and left the slot occupied until the 2 s GC. Slots were keyed on the counter alone, but DATA and CONFIG number independently, so equal counters merged the two streams. The 32-byte received-mask covered 256 of the 512 fragments the client accepts, so a duplicate above 255 was counted as new and the packet was delivered with a hole of stale bytes in it. And one datagram was read per Execute(), which cannot drain a fast producer. Fixed with a pendingTail deferral, (counter, type) keying, a 64-byte mask, a 256-datagram drain, counter-age slot reclamation, and a 1 Hz aggregated warning in place of the per-eviction flood. UDPStreamer dropping whole Accumulate batches. EventSem::ResetWait is Reset-then-Wait, so a Post() landing while the sender thread was inside ServiceClients()/SendData() was destroyed by the next Reset. The batch was then skipped with dataReady false, readyFill was never cleared, and the following flush overwrote it: an entire run of RT cycles never reached the wire. The record of pending work now lives in the buffers rather than in the semaphore edge, which also removes up to UDPS_DATA_WAIT_MS of latency; genuine backpressure overwrites are counted and reported. Against the unfixed code the new test sees 2999/3000 batches never consumed. Period inflation after loss. Accumulated scalars carry no SamplingRate, so the receiver derives dt from the sender-clock gap -- but dividing it by the previous packet's sample count is only right while nothing is lost. One loss doubles the reported period, which spreads a batch a full batch past its own end and into the range the next packet claims. That is the hole and the collision, exactly. Inferring the cycle count from the estimate's own period is not a way out: it has a stable fixed point wherever gap/dt is an integer, so a real rate change locks it at the old one for good (AccumDtGTest.FollowsSustainedRateChange). The packet counter removes the ambiguity, so all three receivers now order on it: a DATA packet that does not advance the counter is dropped rather than delivered, because its values are older than data already handed over. Ordering is on the signed difference so it survives the uint32 wrap, and the sequence resets on reconnect, where the producer's counter restarts independently of ours. The loss count that falls out of the same delta feeds the period estimate as cycles = prevN * (1 + lost), which reduces exactly to gap/prevN when nothing is lost and therefore still tracks a genuine rate change. UDPSClient::AcceptDataCounter (C++), udpsprotocol.SequenceGate (Go), decode_data (C). The C client's existing gap counter was wrap-unsafe and let a stale packet rewind last_counter, which made every subsequent gap wrong; it uses the same code now. Docs/Protocol.md gains an Ordering DATA section stating the requirement for any receiver, including ones outside this repository. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1345 lines
41 KiB
Go
1345 lines
41 KiB
Go
package wshub
|
|
|
|
import (
|
|
"encoding/binary"
|
|
"encoding/json"
|
|
"log"
|
|
"math"
|
|
"net/http"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
"unsafe"
|
|
|
|
"github.com/gorilla/websocket"
|
|
"marte2/common/udpsprotocol"
|
|
)
|
|
|
|
// ─── WebSocket client ─────────────────────────────────────────────────────────
|
|
|
|
type wsMessage struct {
|
|
msgType int
|
|
data []byte
|
|
}
|
|
|
|
type wsClient struct {
|
|
hub *Hub
|
|
conn *websocket.Conn
|
|
send chan wsMessage
|
|
|
|
// window is the timespan this client is displaying, in seconds, held as
|
|
// float64 bits. The retune sweep sizes the rings from the widest window in
|
|
// use, so it must be readable from the hub goroutine while readPump writes
|
|
// it. Zero means the client has not said, and the default applies.
|
|
window atomic.Uint64
|
|
}
|
|
|
|
func (c *wsClient) setDisplayWindowSec(s float64) {
|
|
c.window.Store(math.Float64bits(s))
|
|
}
|
|
|
|
func (c *wsClient) displayWindowSec() float64 {
|
|
return math.Float64frombits(c.window.Load())
|
|
}
|
|
|
|
// sendText enqueues one JSON frame for this client, dropping it if the client
|
|
// is not draining its queue.
|
|
func (c *wsClient) sendText(msg []byte) {
|
|
select {
|
|
case c.send <- wsMessage{websocket.TextMessage, msg}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func (c *wsClient) writePump() {
|
|
pingTicker := time.NewTicker(30 * time.Second)
|
|
defer func() {
|
|
pingTicker.Stop()
|
|
c.conn.Close()
|
|
}()
|
|
for {
|
|
select {
|
|
case msg, ok := <-c.send:
|
|
if !ok {
|
|
c.conn.WriteMessage(websocket.CloseMessage, []byte{})
|
|
return
|
|
}
|
|
if err := c.conn.WriteMessage(msg.msgType, msg.data); err != nil {
|
|
return
|
|
}
|
|
case <-pingTicker.C:
|
|
if err := c.conn.WriteControl(websocket.PingMessage, []byte{},
|
|
time.Now().Add(10*time.Second)); err != nil {
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func (c *wsClient) readPump() {
|
|
defer func() {
|
|
c.hub.unregister <- c
|
|
c.conn.Close()
|
|
}()
|
|
c.conn.SetReadLimit(64 * 1024)
|
|
c.conn.SetReadDeadline(time.Now().Add(60 * time.Second))
|
|
c.conn.SetPongHandler(func(string) error {
|
|
c.conn.SetReadDeadline(time.Now().Add(60 * time.Second))
|
|
return nil
|
|
})
|
|
for {
|
|
_, msg, err := c.conn.ReadMessage()
|
|
if err != nil {
|
|
break
|
|
}
|
|
var env map[string]interface{}
|
|
if json.Unmarshal(msg, &env) == nil {
|
|
if t, ok := env["type"].(string); ok {
|
|
switch t {
|
|
case "ping":
|
|
resp, _ := json.Marshal(map[string]string{"type": "pong"})
|
|
select {
|
|
case c.send <- wsMessage{websocket.TextMessage, resp}:
|
|
default:
|
|
}
|
|
case "addSource":
|
|
label, _ := env["label"].(string)
|
|
addr, _ := env["addr"].(string)
|
|
mcastGroup, _ := env["multicastGroup"].(string)
|
|
dataPortF, _ := env["dataPort"].(float64)
|
|
if addr != "" {
|
|
select {
|
|
case c.hub.commandCh <- hubCmd{
|
|
op: "wsAddSource", label: label, addr: addr,
|
|
multicastGroup: mcastGroup, dataPort: int(dataPortF),
|
|
}:
|
|
default:
|
|
}
|
|
}
|
|
case "removeSource":
|
|
id, _ := env["id"].(string)
|
|
if id != "" {
|
|
select {
|
|
case c.hub.commandCh <- hubCmd{op: "wsRemoveSource", sourceID: id}:
|
|
default:
|
|
}
|
|
}
|
|
case "saveSources":
|
|
select {
|
|
case c.hub.commandCh <- hubCmd{op: "wsSaveSources"}:
|
|
default:
|
|
}
|
|
case "setCalibration":
|
|
source, _ := env["source"].(string)
|
|
signal, _ := env["signal"].(string)
|
|
scale, hasScale := env["scale"].(float64)
|
|
if !hasScale {
|
|
scale = 1
|
|
}
|
|
offset, _ := env["offset"].(float64)
|
|
unit, _ := env["unit"].(string)
|
|
select {
|
|
case c.hub.commandCh <- hubCmd{op: "wsSetCalibration", cal: CalConfig{
|
|
Source: source, Signal: signal,
|
|
Scale: scale, Offset: offset, Unit: unit,
|
|
}}:
|
|
default:
|
|
}
|
|
case "reloadConfig":
|
|
select {
|
|
case c.hub.commandCh <- hubCmd{op: "wsReloadConfig"}:
|
|
default:
|
|
}
|
|
case "setWindow":
|
|
// Sizes the zoom rings: the hub cannot know how far back a
|
|
// client is plotting, and a window it has not been told
|
|
// about is a window the buffers may not reach.
|
|
if sec, ok := env["seconds"].(float64); ok && sec > 0 && !math.IsInf(sec, 0) {
|
|
c.setDisplayWindowSec(sec)
|
|
}
|
|
case "setMonotonic":
|
|
enabled, _ := env["enabled"].(bool)
|
|
select {
|
|
case c.hub.commandCh <- hubCmd{op: "setMonotonic", enabled: enabled}:
|
|
default:
|
|
}
|
|
case "zoom":
|
|
c.hub.handleWSZoom(c, env)
|
|
default:
|
|
if c.hub.handleTriggerCommand(t, env) {
|
|
break
|
|
}
|
|
if c.hub.handleHistoryCommand(c, t, env) {
|
|
break
|
|
}
|
|
// Unrecognized message type — forward to DebugCh
|
|
select {
|
|
case c.hub.DebugCh <- msg:
|
|
default:
|
|
}
|
|
}
|
|
}
|
|
}
|
|
c.conn.SetReadDeadline(time.Now().Add(60 * time.Second))
|
|
}
|
|
}
|
|
|
|
// ─── Hub ─────────────────────────────────────────────────────────────────────
|
|
|
|
// allowedOrigins is the set of Origin values (scheme://host[:port]) that are
|
|
// accepted for WebSocket upgrades. If empty, same-origin is enforced by
|
|
// comparing the Origin's host to the HTTP Host header.
|
|
var allowedOrigins []string
|
|
|
|
// SetAllowedOrigins configures the WebSocket Origin allowlist. Pass an empty
|
|
// slice to enforce same-origin only (the default).
|
|
func SetAllowedOrigins(origins []string) {
|
|
allowedOrigins = origins
|
|
}
|
|
|
|
// checkOrigin validates the Origin header against the allowlist, falling back
|
|
// to a same-origin check (Origin host == Host header) when no allowlist is
|
|
// configured. Requests with no Origin header (non-browser clients) are allowed.
|
|
func checkOrigin(r *http.Request) bool {
|
|
origin := r.Header.Get("Origin")
|
|
if origin == "" {
|
|
return true // non-browser client
|
|
}
|
|
// Check explicit allowlist first.
|
|
for _, allowed := range allowedOrigins {
|
|
if origin == allowed {
|
|
return true
|
|
}
|
|
}
|
|
// Fall back to same-origin: compare the Origin's host to the Host header.
|
|
// Origin format: "scheme://host[:port]" — strip scheme.
|
|
host := origin
|
|
if idx := strings.Index(host, "://"); idx >= 0 {
|
|
host = host[idx+3:]
|
|
}
|
|
// Strip path if present.
|
|
if idx := strings.Index(host, "/"); idx >= 0 {
|
|
host = host[:idx]
|
|
}
|
|
return host == r.Host
|
|
}
|
|
|
|
var upgrader = websocket.Upgrader{
|
|
ReadBufferSize: 4096,
|
|
WriteBufferSize: 64 * 1024,
|
|
CheckOrigin: checkOrigin,
|
|
}
|
|
|
|
// sourceHubState holds all data for one active data source.
|
|
// Only accessed from the Run() goroutine.
|
|
type sourceHubState struct {
|
|
id, label, addr, connState string
|
|
signals []udpsprotocol.SignalInfo
|
|
configJS []byte
|
|
|
|
// Time-signal calibration — only accessed from Run() goroutine.
|
|
timeSigCalib map[string]float64
|
|
configSeq uint64
|
|
configSeqAtCalib uint64
|
|
|
|
// lastPktNs tracks the wall-clock time (UnixNano) of the last received packet
|
|
// per signal name. Used by the default (TimeModePacket, n>1) path to estimate
|
|
// per-element dt when only one packet arrives in a 30 Hz tick.
|
|
lastPktNs map[string]int64
|
|
|
|
// Monotonic timestamp snapping state (all accessed from Run() goroutine):
|
|
// lastFrameMeasured — uncorrected measured anchor of the previous frame.
|
|
// lastFrameEndT — corrected anchor after snapping.
|
|
// gapEMA — exponential moving average of the measured inter-frame gap.
|
|
lastFrameMeasured map[string]float64
|
|
lastFrameEndT map[string]float64
|
|
gapEMA map[string]float64
|
|
}
|
|
|
|
// taggedSample is a DataSample annotated with its source ID.
|
|
type taggedSample struct {
|
|
sourceID string
|
|
sample udpsprotocol.DataSample
|
|
}
|
|
|
|
// hubCmd carries a command to the Run() goroutine.
|
|
type hubCmd struct {
|
|
op string // "addSource","removeSource","setSourceState","updateConfig",
|
|
// "wsAddSource","wsRemoveSource","wsSaveSources",
|
|
// "wsSetCalibration","wsReloadConfig"
|
|
sourceID string
|
|
label string
|
|
addr string
|
|
state string
|
|
sigs []udpsprotocol.SignalInfo
|
|
multicastGroup string
|
|
dataPort int
|
|
enabled bool // "setMonotonic" toggle
|
|
cal CalConfig // "wsSetCalibration" payload
|
|
}
|
|
|
|
// Hub is the central broker between UDP clients and WebSocket clients.
|
|
// All map state is accessed exclusively from the Run() goroutine, except
|
|
// ringsMu/rings which are also read by HTTP handler goroutines.
|
|
type Hub struct {
|
|
clients map[*wsClient]bool
|
|
register chan *wsClient
|
|
unregister chan *wsClient
|
|
broadcastCh chan []byte
|
|
dataCh chan taggedSample
|
|
commandCh chan hubCmd
|
|
|
|
// DebugCh receives raw browser messages whose type is not handled by the hub.
|
|
DebugCh chan []byte
|
|
|
|
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
|
|
rings map[string]*sigRing // "sourceId:signalKey" → ring
|
|
|
|
// hist is the disk-backed archive behind long time windows, which hold far
|
|
// more samples than the in-memory rings can. nil when history is disabled.
|
|
// histOpenAt throttles the sweep that opens the files of signals whose
|
|
// producer declared no sampling rate; both are touched only from Run().
|
|
hist *historyWriter
|
|
histOpenAt float64
|
|
|
|
statsMu sync.RWMutex
|
|
statsMap map[string]*SourceStat
|
|
|
|
// trigger is the hub-side trigger FSM driving the oscilloscope capture mode.
|
|
// ringTuneAt throttles the sweep that keeps each ring's depth and min/max
|
|
// bucket matched to the window being displayed; both are touched only from
|
|
// Run(). ringBudgetPts is that sweep's per-signal budget; set before Run().
|
|
trigger *triggerEngine
|
|
ringTuneAt float64
|
|
// capture is the trigger double buffer's read half: the last delivered
|
|
// capture window, kept out of the rings' way so the shot being viewed
|
|
// survives the re-arm that immediately follows it.
|
|
capture captureHold
|
|
ringBudgetPts int
|
|
onClientConnectMu sync.RWMutex
|
|
onClientConnect func(send func([]byte))
|
|
|
|
// monotonicTS, when true, snaps small inter-frame timestamp deviations
|
|
// (< monotonicTolerance) to the ideal gap to eliminate jitter.
|
|
monotonicTS bool
|
|
}
|
|
|
|
// NewHub creates an initialised Hub.
|
|
func NewHub() *Hub {
|
|
return &Hub{
|
|
clients: make(map[*wsClient]bool),
|
|
register: make(chan *wsClient, 8),
|
|
unregister: make(chan *wsClient, 8),
|
|
broadcastCh: make(chan []byte, 256),
|
|
dataCh: make(chan taggedSample, 65536), // large buffer: absorbs bursts at high sample rates
|
|
commandCh: make(chan hubCmd, 64),
|
|
DebugCh: make(chan []byte, 256),
|
|
rings: make(map[string]*sigRing),
|
|
statsMap: make(map[string]*SourceStat),
|
|
trigger: newTriggerEngine(),
|
|
cal: newCalTable(),
|
|
}
|
|
}
|
|
|
|
// SetRingBudget overrides the per-signal in-memory buffer budget, in points.
|
|
// Non-positive values restore the default. It must be called before Run().
|
|
// Each point costs 16 bytes, so the budget is the memory bound per temporal
|
|
// signal. It does not limit how long a window can be held: a window too long
|
|
// to fit at full rate is stored as min/max pairs instead (see retuneRings).
|
|
func (h *Hub) SetRingBudget(n int) {
|
|
if n <= 0 {
|
|
n = defaultRingPts
|
|
}
|
|
if n < ringCapInitial {
|
|
n = ringCapInitial
|
|
}
|
|
h.ringBudgetPts = n
|
|
}
|
|
|
|
func (h *Hub) ringBudget() int {
|
|
if h.ringBudgetPts <= 0 {
|
|
return defaultRingPts
|
|
}
|
|
return h.ringBudgetPts
|
|
}
|
|
|
|
// EnableHistory turns on the disk-backed history archive. It must be called
|
|
// before Run(). A HistoryConfig with an empty Directory leaves history off.
|
|
func (h *Hub) EnableHistory(cfg HistoryConfig) error {
|
|
hw, err := newHistoryWriter(cfg)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
h.hist = hw
|
|
return nil
|
|
}
|
|
|
|
// CloseHistory flushes and closes the history files. Without it the samples
|
|
// written since the last periodic flush are on disk but unaccounted for in the
|
|
// file headers, so a restart would not see them.
|
|
func (h *Hub) CloseHistory() { h.hist.close() }
|
|
|
|
// SetOnClientConnect registers a callback invoked synchronously (from Run())
|
|
// each time a new WebSocket client connects. The callback receives a send
|
|
// function that enqueues one message to that specific client.
|
|
func (h *Hub) SetOnClientConnect(fn func(send func([]byte))) {
|
|
h.onClientConnectMu.Lock()
|
|
h.onClientConnect = fn
|
|
h.onClientConnectMu.Unlock()
|
|
}
|
|
|
|
// SetSourceManager sets the SourceManager associated with the Hub.
|
|
func (h *Hub) SetSourceManager(sm *SourceManager) {
|
|
h.sm = sm
|
|
}
|
|
|
|
// ingest routes one batch of full-resolution samples for a signal to every
|
|
// consumer that needs them at full rate: the in-memory zoom ring, the disk
|
|
// history and the trigger comparator. The live push is decimated separately by
|
|
// the caller. The ring and the archive may reduce what they store to fit their
|
|
// budget, but they are handed every sample so the reduction sees the extrema.
|
|
func (h *Hub) ingest(key string, nElem int, t, v []float64) {
|
|
if len(t) == 0 {
|
|
return
|
|
}
|
|
if rb := h.getRing(key); rb != nil {
|
|
rb.write(t, v)
|
|
}
|
|
h.hist.write(key, t, v)
|
|
h.trigger.feed(key, nElem, t, v)
|
|
}
|
|
|
|
// getRing returns the ring buffer for a fully-prefixed signal key, or nil.
|
|
func (h *Hub) getRing(key string) *sigRing {
|
|
h.ringsMu.RLock()
|
|
rb := h.rings[key]
|
|
h.ringsMu.RUnlock()
|
|
return rb
|
|
}
|
|
|
|
// zoomSlice extracts [t0, t1] for the named signals, decimating each to at most
|
|
// n points. A range inside the last trigger capture is served from the held
|
|
// copy of it, which the re-arming acquisition cannot overwrite; everything else
|
|
// comes from the live rings.
|
|
func (h *Hub) zoomSlice(t0, t1 float64, keys []string, n int) map[string]sigData {
|
|
h.ringsMu.RLock()
|
|
refs := make(map[string]*sigRing, len(keys))
|
|
for _, k := range keys {
|
|
k = strings.TrimSpace(k)
|
|
if k == "" {
|
|
continue
|
|
}
|
|
if rb, ok := h.rings[k]; ok {
|
|
refs[k] = rb
|
|
}
|
|
}
|
|
h.ringsMu.RUnlock()
|
|
|
|
result := make(map[string]sigData, len(refs))
|
|
for k, rb := range refs {
|
|
rt, rv, ok := h.capture.slice(k, t0, t1)
|
|
if !ok {
|
|
rt, rv = rb.slice(t0, t1)
|
|
}
|
|
if len(rt) == 0 {
|
|
continue
|
|
}
|
|
dt, dv := minMaxDecimate(rt, rv, n)
|
|
result[k] = sigData{T: dt, V: dv}
|
|
}
|
|
return result
|
|
}
|
|
|
|
// zoomPoints normalises the client's requested point budget: absent → 2400,
|
|
// non-positive → every sample in the range, implausibly small → 2400.
|
|
func zoomPoints(n int, present bool) int {
|
|
switch {
|
|
case !present:
|
|
return 2400
|
|
case n <= 0:
|
|
return 1 << 30 // no decimation
|
|
case n < 10:
|
|
return 2400
|
|
}
|
|
return n
|
|
}
|
|
|
|
// handleWSZoom answers a browser {"type":"zoom","reqId":..,"t0":..,"t1":..,
|
|
// "n":..,"signals":"a,b"} request, unicasting {"type":"zoom","reqId":..,
|
|
// "signals":{...}} back to the requesting client. This is the path the web SPA
|
|
// actually uses; /api/zoom is the equivalent HTTP entry point.
|
|
func (h *Hub) handleWSZoom(c *wsClient, env map[string]interface{}) {
|
|
t0, ok0 := env["t0"].(float64)
|
|
t1, ok1 := env["t1"].(float64)
|
|
if !ok0 || !ok1 || t1 <= t0 {
|
|
return
|
|
}
|
|
nF, nOK := env["n"].(float64)
|
|
n := zoomPoints(int(nF), nOK)
|
|
sigCSV, _ := env["signals"].(string)
|
|
|
|
reply, err := json.Marshal(map[string]any{
|
|
"type": "zoom",
|
|
"reqId": env["reqId"],
|
|
"signals": h.zoomSlice(t0, t1, strings.Split(sigCSV, ","), n),
|
|
})
|
|
if err != nil {
|
|
log.Printf("hub: ws zoom encode: %v", err)
|
|
return
|
|
}
|
|
c.sendText(reply)
|
|
}
|
|
|
|
// HandleZoom serves GET /api/zoom?...
|
|
func (h *Hub) HandleZoom(w http.ResponseWriter, r *http.Request) {
|
|
q := r.URL.Query()
|
|
t0, err0 := strconv.ParseFloat(q.Get("t0"), 64)
|
|
t1, err1 := strconv.ParseFloat(q.Get("t1"), 64)
|
|
if err0 != nil || err1 != nil || t1 <= t0 {
|
|
http.Error(w, "invalid t0/t1", http.StatusBadRequest)
|
|
return
|
|
}
|
|
nStr := q.Get("n")
|
|
nVal, _ := strconv.Atoi(nStr)
|
|
n := zoomPoints(nVal, nStr != "")
|
|
|
|
w.Header().Set("Content-Type", "application/json")
|
|
if err := json.NewEncoder(w).Encode(map[string]any{
|
|
"type": "zoom",
|
|
"signals": h.zoomSlice(t0, t1, strings.Split(q.Get("signals"), ","), n),
|
|
}); err != nil {
|
|
log.Printf("hub: zoom encode: %v", err)
|
|
}
|
|
}
|
|
|
|
// AddSource notifies the Hub that a new source has been registered.
|
|
func (h *Hub) AddSource(id, label, addr string) {
|
|
select {
|
|
case h.commandCh <- hubCmd{op: "addSource", sourceID: id, label: label, addr: addr}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// RemoveSource notifies the Hub that a source has been removed.
|
|
func (h *Hub) RemoveSource(id string) {
|
|
select {
|
|
case h.commandCh <- hubCmd{op: "removeSource", sourceID: id}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// SetSourceState updates the connection state of a source.
|
|
func (h *Hub) SetSourceState(id, state string) {
|
|
select {
|
|
case h.commandCh <- hubCmd{op: "setSourceState", sourceID: id, state: state}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// UpdateConfigForSource stores a new signal config for a source and broadcasts it.
|
|
func (h *Hub) UpdateConfigForSource(sourceID string, sigs []udpsprotocol.SignalInfo) {
|
|
select {
|
|
case h.commandCh <- hubCmd{op: "updateConfig", sourceID: sourceID, sigs: sigs}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// PushDataForSource enqueues a data sample from a specific source.
|
|
func (h *Hub) PushDataForSource(sourceID string, s udpsprotocol.DataSample) {
|
|
select {
|
|
case h.dataCh <- taggedSample{sourceID: sourceID, sample: s}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// broadcast enqueues a message for delivery to all WebSocket clients.
|
|
func (h *Hub) broadcast(msg []byte) {
|
|
select {
|
|
case h.broadcastCh <- msg:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// Broadcast is the exported wrapper for broadcast.
|
|
func (h *Hub) Broadcast(msg []byte) {
|
|
h.broadcast(msg)
|
|
}
|
|
|
|
// HandleWebSocket upgrades an HTTP request to a WebSocket connection.
|
|
func (h *Hub) HandleWebSocket(w http.ResponseWriter, r *http.Request) {
|
|
conn, err := upgrader.Upgrade(w, r, nil)
|
|
if err != nil {
|
|
log.Printf("ws upgrade: %v", err)
|
|
return
|
|
}
|
|
c := &wsClient{hub: h, conn: conn, send: make(chan wsMessage, 64)}
|
|
h.register <- c
|
|
go c.writePump()
|
|
go c.readPump()
|
|
}
|
|
|
|
// buildSourcesMsg serialises the current source list as a JSON "sources" message.
|
|
func buildSourcesMsg(sm map[string]*sourceHubState) []byte {
|
|
type srcInfo struct {
|
|
ID string `json:"id"`
|
|
Label string `json:"label"`
|
|
Addr string `json:"addr"`
|
|
State string `json:"state"`
|
|
}
|
|
list := make([]srcInfo, 0, len(sm))
|
|
for _, src := range sm {
|
|
list = append(list, srcInfo{ID: src.id, Label: src.label, Addr: src.addr, State: src.connState})
|
|
}
|
|
msg, _ := json.Marshal(map[string]interface{}{"type": "sources", "sources": list})
|
|
return msg
|
|
}
|
|
|
|
// buildCalibrationMsg serialises the calibration table as a "calibration"
|
|
// message. It is its own frame rather than a field on "sources" because the
|
|
// C++ BroadcastSources serialises into a fixed 4096-byte buffer that a
|
|
// calibration table would overflow.
|
|
func buildCalibrationMsg(t *calTable) []byte {
|
|
list := t.List() // never nil: the SPA replaces its table wholesale on receipt
|
|
msg, _ := json.Marshal(map[string]any{"type": "calibration", "cal": list})
|
|
return msg
|
|
}
|
|
|
|
// buildConfigAckMsg serialises a configSaved / configReloaded acknowledgement.
|
|
func buildConfigAckMsg(msgType, path string, err error) []byte {
|
|
m := map[string]any{"type": msgType, "ok": err == nil, "path": path}
|
|
if err != nil {
|
|
m["error"] = err.Error()
|
|
}
|
|
msg, _ := json.Marshal(m)
|
|
return msg
|
|
}
|
|
|
|
// Run is the hub's main goroutine. Must be started with go hub.Run().
|
|
func (h *Hub) Run() {
|
|
ticker := time.NewTicker(time.Second / 30)
|
|
defer ticker.Stop()
|
|
|
|
statsTicker := time.NewTicker(time.Second)
|
|
defer statsTicker.Stop()
|
|
|
|
// Header flushes are what make the archived samples findable again; the
|
|
// data region is written as it arrives. Ticks are ignored when history is
|
|
// off, so a disabled writer costs one no-op call per period.
|
|
flushPeriod := time.Duration(5) * time.Second
|
|
if h.hist.enabled() {
|
|
flushPeriod = time.Duration(h.hist.cfg.FlushIntervalSec) * time.Second
|
|
}
|
|
flushTicker := time.NewTicker(flushPeriod)
|
|
defer flushTicker.Stop()
|
|
|
|
sourcesMap := make(map[string]*sourceHubState)
|
|
var sourcesMsg []byte
|
|
|
|
// pending[sourceID] accumulates samples between 30 Hz ticks.
|
|
pending := make(map[string][]udpsprotocol.DataSample)
|
|
|
|
rebuildSources := func() {
|
|
sourcesMsg = buildSourcesMsg(sourcesMap)
|
|
h.broadcast(sourcesMsg)
|
|
}
|
|
|
|
for {
|
|
select {
|
|
case c := <-h.register:
|
|
h.clients[c] = true
|
|
// Send current state to the new client.
|
|
if sourcesMsg != nil {
|
|
select {
|
|
case c.send <- wsMessage{websocket.TextMessage, sourcesMsg}:
|
|
default:
|
|
}
|
|
}
|
|
for _, src := range sourcesMap {
|
|
if src.configJS != nil {
|
|
select {
|
|
case c.send <- wsMessage{websocket.TextMessage, src.configJS}:
|
|
default:
|
|
}
|
|
}
|
|
}
|
|
select {
|
|
case c.send <- wsMessage{websocket.TextMessage, h.trigger.stateMsg()}:
|
|
default:
|
|
}
|
|
monoMsg, _ := json.Marshal(map[string]any{"type": "monotonicState", "enabled": h.monotonicTS})
|
|
select {
|
|
case c.send <- wsMessage{websocket.TextMessage, monoMsg}:
|
|
default:
|
|
}
|
|
calMsg := buildCalibrationMsg(h.cal)
|
|
select {
|
|
case c.send <- wsMessage{websocket.TextMessage, calMsg}:
|
|
default:
|
|
}
|
|
if h.hist.enabled() {
|
|
if msg := h.buildHistoryInfoMsg(); msg != nil {
|
|
c.sendText(msg)
|
|
}
|
|
}
|
|
// Notify the application layer so it can replay any persistent state
|
|
// (e.g., MARTe2 connection status, forced/traced signals).
|
|
h.onClientConnectMu.RLock()
|
|
fn := h.onClientConnect
|
|
h.onClientConnectMu.RUnlock()
|
|
if fn != nil {
|
|
fn(func(msg []byte) {
|
|
select {
|
|
case c.send <- wsMessage{websocket.TextMessage, msg}:
|
|
default:
|
|
}
|
|
})
|
|
}
|
|
|
|
case c := <-h.unregister:
|
|
if _, ok := h.clients[c]; ok {
|
|
delete(h.clients, c)
|
|
close(c.send)
|
|
}
|
|
|
|
case msg := <-h.broadcastCh:
|
|
for c := range h.clients {
|
|
select {
|
|
case c.send <- wsMessage{websocket.TextMessage, msg}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
case cmd := <-h.commandCh:
|
|
switch cmd.op {
|
|
case "addSource":
|
|
sourcesMap[cmd.sourceID] = &sourceHubState{
|
|
id: cmd.sourceID,
|
|
label: cmd.label,
|
|
addr: cmd.addr,
|
|
connState: "connecting",
|
|
timeSigCalib: make(map[string]float64),
|
|
lastPktNs: make(map[string]int64),
|
|
lastFrameEndT: make(map[string]float64),
|
|
lastFrameMeasured: make(map[string]float64),
|
|
gapEMA: make(map[string]float64),
|
|
}
|
|
h.statsMu.Lock()
|
|
h.statsMap[cmd.sourceID] = &SourceStat{}
|
|
h.statsMu.Unlock()
|
|
rebuildSources()
|
|
|
|
case "removeSource":
|
|
delete(sourcesMap, cmd.sourceID)
|
|
delete(pending, cmd.sourceID)
|
|
pfxDel := cmd.sourceID + ":"
|
|
h.ringsMu.Lock()
|
|
for k := range h.rings {
|
|
if strings.HasPrefix(k, pfxDel) {
|
|
delete(h.rings, k)
|
|
}
|
|
}
|
|
h.ringsMu.Unlock()
|
|
h.statsMu.Lock()
|
|
delete(h.statsMap, cmd.sourceID)
|
|
h.statsMu.Unlock()
|
|
rebuildSources()
|
|
|
|
case "setSourceState":
|
|
if src, ok := sourcesMap[cmd.sourceID]; ok {
|
|
src.connState = cmd.state
|
|
rebuildSources()
|
|
}
|
|
|
|
case "updateConfig":
|
|
src, ok := sourcesMap[cmd.sourceID]
|
|
if !ok {
|
|
continue
|
|
}
|
|
src.signals = cmd.sigs
|
|
src.configSeq++
|
|
src.lastFrameEndT = make(map[string]float64)
|
|
cfgMsg, err := json.Marshal(map[string]any{
|
|
"type": "config",
|
|
"sourceId": cmd.sourceID,
|
|
"signals": cmd.sigs,
|
|
})
|
|
if err != nil {
|
|
log.Printf("hub: marshal config: %v", err)
|
|
continue
|
|
}
|
|
src.configJS = cfgMsg
|
|
h.broadcast(cfgMsg)
|
|
// Rebuild ring buffers for this source.
|
|
pfxUpd := cmd.sourceID + ":"
|
|
h.ringsMu.Lock()
|
|
for k := range h.rings {
|
|
if strings.HasPrefix(k, pfxUpd) {
|
|
delete(h.rings, k)
|
|
}
|
|
}
|
|
for _, sig := range cmd.sigs {
|
|
ne := sig.NumElements()
|
|
isTemporal := ne > 1 && sig.TimeMode != udpsprotocol.TimeModePacket
|
|
if isTemporal {
|
|
h.rings[pfxUpd+sig.Name] = newSigRing(ringCapInitial)
|
|
} else if ne == 1 {
|
|
h.rings[pfxUpd+sig.Name] = newSigRing(ringCapScalar)
|
|
} else {
|
|
// n>1, TimeModePacket snapshot-waveform: each packet contributes n
|
|
// elements, so this is a fast stream too and gets the same budget.
|
|
h.rings[pfxUpd+sig.Name] = newSigRing(ringCapInitial)
|
|
}
|
|
}
|
|
h.ringsMu.Unlock()
|
|
// The held capture describes rings that no longer exist. A
|
|
// restarted producer can even replay the same timestamps, so
|
|
// keeping it would answer zooms with the old run's samples.
|
|
h.capture.clear()
|
|
// Opening the archive files touches the filesystem, so keep it
|
|
// off the Run() goroutine; the write path simply drops samples
|
|
// for a key whose file is not open yet.
|
|
if h.hist.enabled() {
|
|
go func(id string, sigs []udpsprotocol.SignalInfo) {
|
|
h.hist.onSourceConfigured(id, sigs)
|
|
h.broadcast(h.buildHistoryInfoMsg())
|
|
}(cmd.sourceID, cmd.sigs)
|
|
}
|
|
|
|
case "wsAddSource":
|
|
if h.sm != nil {
|
|
go func(label, addr, mcastGroup string, dataPort int) {
|
|
h.sm.Add(label, addr, mcastGroup, dataPort)
|
|
}(cmd.label, cmd.addr, cmd.multicastGroup, cmd.dataPort)
|
|
}
|
|
|
|
case "wsRemoveSource":
|
|
if h.sm != nil {
|
|
go func(id string) { h.sm.Remove(id) }(cmd.sourceID)
|
|
}
|
|
|
|
case "wsSaveSources":
|
|
if h.sm != nil {
|
|
// Save writes to disk; run it off the Run() goroutine so a
|
|
// slow filesystem can never stall the hub loop.
|
|
go func(sm *SourceManager) {
|
|
err := sm.Save()
|
|
if err != nil {
|
|
log.Printf("hub: save config: %v", err)
|
|
}
|
|
h.broadcast(buildConfigAckMsg("configSaved", sm.Path(), err))
|
|
}(h.sm)
|
|
}
|
|
|
|
case "wsSetCalibration":
|
|
if h.cal.Set(cmd.cal) {
|
|
h.broadcast(buildCalibrationMsg(h.cal))
|
|
} else {
|
|
// No broadcast: the offending client reverts to the last
|
|
// value it was sent.
|
|
log.Printf("hub: rejected calibration %q/%q (scale=%v offset=%v)",
|
|
cmd.cal.Source, cmd.cal.Signal, cmd.cal.Scale, cmd.cal.Offset)
|
|
}
|
|
|
|
case "wsReloadConfig":
|
|
if h.sm != nil {
|
|
// Reload calls sm.Add(), which sends on commandCh; from the
|
|
// Run() goroutine that send would hit the non-blocking
|
|
// default and be dropped, so it must run elsewhere.
|
|
go func(sm *SourceManager) {
|
|
err := sm.Reload()
|
|
if err != nil {
|
|
log.Printf("hub: reload config: %v", err)
|
|
}
|
|
h.broadcast(buildConfigAckMsg("configReloaded", sm.Path(), err))
|
|
if err == nil {
|
|
h.broadcast(buildCalibrationMsg(h.cal))
|
|
}
|
|
}(h.sm)
|
|
}
|
|
|
|
case "setMonotonic":
|
|
h.monotonicTS = cmd.enabled
|
|
monoMsg, _ := json.Marshal(map[string]any{"type": "monotonicState", "enabled": h.monotonicTS})
|
|
h.broadcast(monoMsg)
|
|
}
|
|
|
|
case ts := <-h.dataCh:
|
|
pending[ts.sourceID] = append(pending[ts.sourceID], ts.sample)
|
|
|
|
case <-ticker.C:
|
|
for srcID, samples := range pending {
|
|
if len(samples) == 0 {
|
|
continue
|
|
}
|
|
src, ok := sourcesMap[srcID]
|
|
if !ok || len(src.signals) == 0 {
|
|
pending[srcID] = pending[srcID][:0]
|
|
continue
|
|
}
|
|
// Built even with no clients connected: this is also what feeds
|
|
// the rings, the disk history and the trigger, none of which may
|
|
// stop just because nobody is watching. It also keeps the push
|
|
// cursors advancing, so the first client to connect does not get
|
|
// a backlog burst. Matches the C++ StreamHub.
|
|
msg := h.buildBinaryDataMessageForSource(src, samples)
|
|
pending[srcID] = pending[srcID][:0]
|
|
if msg != nil {
|
|
for c := range h.clients {
|
|
select {
|
|
case c.send <- wsMessage{websocket.BinaryMessage, msg}:
|
|
default:
|
|
}
|
|
}
|
|
}
|
|
}
|
|
h.triggerTick()
|
|
|
|
case <-flushTicker.C:
|
|
h.hist.flushHeaders()
|
|
|
|
case <-statsTicker.C:
|
|
h.statsMu.RLock()
|
|
snap := make(map[string]StatInfo, len(h.statsMap))
|
|
for id, st := range h.statsMap {
|
|
snap[id] = st.Snapshot()
|
|
}
|
|
h.statsMu.RUnlock()
|
|
if len(snap) > 0 {
|
|
msg, _ := json.Marshal(map[string]any{"type": "stats", "sources": snap})
|
|
h.broadcast(msg)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// float64ToBytes reinterprets a []float64 as []byte without copying.
|
|
func float64ToBytes(f []float64) []byte {
|
|
if len(f) == 0 {
|
|
return nil
|
|
}
|
|
return unsafe.Slice((*byte)(unsafe.Pointer(&f[0])), len(f)*8)
|
|
}
|
|
|
|
// writeFloat64s encodes a []float64 as little-endian bytes into buf at offset
|
|
// and returns the new offset.
|
|
func writeFloat64s(buf []byte, off int, f []float64) int {
|
|
copy(buf[off:], float64ToBytes(f))
|
|
return off + len(f)*8
|
|
}
|
|
|
|
// ─── Data serialisation ───────────────────────────────────────────────────────
|
|
|
|
// maxPushPoints bounds the live push only. The zoom rings deliberately store
|
|
// every sample: decimating on the way in would cap the resolution a zoom can
|
|
// ever recover, and the browser already decimates for display.
|
|
const maxPushPoints = 50
|
|
|
|
// Ring geometry, in samples per signal (16 bytes each).
|
|
//
|
|
// defaultRingPts is the per-signal memory budget for temporal (array) signals:
|
|
// what the hub may spend keeping one signal available for zoom and for trigger
|
|
// captures. 10 M points is 160 MB. The budget buys resolution, not span —
|
|
// retuneRings buckets the input so the display window fits whatever the source
|
|
// rate is.
|
|
//
|
|
// ringCapInitial is where a ring starts, so a source that is configured but
|
|
// never sends costs nothing; the first retune sweep grows it to the budget.
|
|
//
|
|
// ringCapScalar sizes scalar signals, which arrive at the packet rate and would
|
|
// squander a budget meant for megasample streams.
|
|
const defaultRingPts = 10_000_000
|
|
const ringCapInitial = 250_000
|
|
const ringCapScalar = 100_000
|
|
|
|
// monotonicTolerance is the maximum inter-frame timestamp deviation (seconds)
|
|
// treated as jitter and snapped to the ideal gap. Larger deviations are
|
|
// preserved as genuine discontinuities (missing frames, rate changes).
|
|
const monotonicTolerance = 0.005 // 5 ms
|
|
|
|
// monotonicEMAAlpha is the smoothing factor for the inter-frame gap EMA.
|
|
// 0.01 gives a time constant of ~100 frames (~1 s at 100 Hz): fast enough to
|
|
// track real rate changes, slow enough to average out per-frame jitter.
|
|
const monotonicEMAAlpha = 0.01
|
|
|
|
// minMaxDecimate reduces (tIn, vIn) to at most threshold points the way an
|
|
// oscilloscope draws a trace it cannot show pixel-for-pixel: the range is split
|
|
// into threshold/2 equal buckets and each contributes its smallest and largest
|
|
// sample, in the order the two occurred.
|
|
//
|
|
// This is what replaced LTTB on every path here. LTTB picks the sample that
|
|
// makes the largest triangle with its neighbours, which reads as a plausible
|
|
// shape but silently drops a one-sample spike whenever a smoother neighbour
|
|
// scores higher — precisely the sample the user is looking for. The envelope
|
|
// cannot drop it: a spike is by definition its bucket's min or max. The cost is
|
|
// that a flat trace is drawn as a band rather than a line, which is how a scope
|
|
// behaves too.
|
|
//
|
|
// Both output arrays hold real samples with their real timestamps; nothing is
|
|
// interpolated or averaged.
|
|
func minMaxDecimate(tIn, vIn []float64, threshold int) ([]float64, []float64) {
|
|
n := len(tIn)
|
|
// Below four there is no room for a single min/max pair plus endpoints.
|
|
if n <= threshold || threshold < 4 {
|
|
return tIn, vIn
|
|
}
|
|
buckets := threshold / 2
|
|
outT := make([]float64, 0, threshold)
|
|
outV := make([]float64, 0, threshold)
|
|
for b := 0; b < buckets; b++ {
|
|
lo := b * n / buckets
|
|
hi := (b + 1) * n / buckets
|
|
if b == buckets-1 {
|
|
hi = n
|
|
}
|
|
if lo >= hi {
|
|
continue
|
|
}
|
|
iMin, iMax := lo, lo
|
|
for j := lo + 1; j < hi; j++ {
|
|
if vIn[j] < vIn[iMin] {
|
|
iMin = j
|
|
}
|
|
if vIn[j] > vIn[iMax] {
|
|
iMax = j
|
|
}
|
|
}
|
|
// Emit in time order so the result plots as one ascending trace.
|
|
if iMin > iMax {
|
|
iMin, iMax = iMax, iMin
|
|
}
|
|
outT = append(outT, tIn[iMin])
|
|
outV = append(outV, vIn[iMin])
|
|
// A bucket whose samples are all equal has one extreme, not two.
|
|
if iMax != iMin {
|
|
outT = append(outT, tIn[iMax])
|
|
outV = append(outV, vIn[iMax])
|
|
}
|
|
}
|
|
return outT, outV
|
|
}
|
|
|
|
type sigData struct {
|
|
T []float64 `json:"t"`
|
|
V []float64 `json:"v"`
|
|
}
|
|
|
|
type dataMsg struct {
|
|
Type string `json:"type"`
|
|
SourceID string `json:"sourceId"`
|
|
Signals map[string]sigData `json:"signals"`
|
|
}
|
|
|
|
// buildBinaryDataMessageForSource encodes a batch of samples as a compact binary frame.
|
|
func (h *Hub) buildBinaryDataMessageForSource(src *sourceHubState, batch []udpsprotocol.DataSample) []byte {
|
|
if len(batch) == 0 {
|
|
return nil
|
|
}
|
|
if src.configSeq != src.configSeqAtCalib {
|
|
src.configSeqAtCalib = src.configSeq
|
|
src.timeSigCalib = make(map[string]float64)
|
|
src.lastFrameEndT = make(map[string]float64)
|
|
src.lastFrameMeasured = make(map[string]float64)
|
|
src.gapEMA = make(map[string]float64)
|
|
}
|
|
|
|
sigs := src.signals
|
|
pfx := src.id + ":"
|
|
|
|
type pairBuf struct {
|
|
t, v []float64
|
|
}
|
|
pairs := make(map[string]pairBuf, len(sigs)*2)
|
|
|
|
for _, sig := range sigs {
|
|
n := sig.NumElements()
|
|
|
|
switch {
|
|
case n > 1 && (sig.TimeMode == udpsprotocol.TimeModeFirstSample || sig.TimeMode == udpsprotocol.TimeModeLastSample):
|
|
hasTimeSig := sig.TimeSignalIdx != udpsprotocol.NoTimeSignal && int(sig.TimeSignalIdx) < len(sigs)
|
|
var timeSigName string
|
|
timerToSec := 1e-6
|
|
if hasTimeSig {
|
|
ts := sigs[sig.TimeSignalIdx]
|
|
timeSigName = ts.Name
|
|
if ts.TypeCode == 6 {
|
|
timerToSec = 1e-9
|
|
}
|
|
}
|
|
dt := 0.0
|
|
if sig.SamplingRate > 0 {
|
|
dt = 1.0 / sig.SamplingRate
|
|
}
|
|
allT := make([]float64, 0, len(batch)*n)
|
|
allV := make([]float64, 0, len(batch)*n)
|
|
for _, s := range batch {
|
|
vals, ok := s.Values[sig.Name]
|
|
if !ok || len(vals) < n {
|
|
continue
|
|
}
|
|
var anchorTime float64
|
|
anchorIsFirstSample := sig.TimeMode == udpsprotocol.TimeModeFirstSample
|
|
if hasTimeSig {
|
|
tVals, tOk := s.Values[timeSigName]
|
|
if tOk && len(tVals) >= 1 {
|
|
timerS := tVals[0] * timerToSec
|
|
wallT := float64(s.WallTime.UnixNano()) / 1e9
|
|
if _, exists := src.timeSigCalib[timeSigName]; !exists {
|
|
src.timeSigCalib[timeSigName] = wallT - timerS
|
|
}
|
|
anchorTime = src.timeSigCalib[timeSigName] + timerS
|
|
} else {
|
|
anchorTime = float64(s.WallTime.UnixNano()) / 1e9
|
|
anchorIsFirstSample = false
|
|
}
|
|
} else {
|
|
anchorTime = float64(s.WallTime.UnixNano()) / 1e9
|
|
anchorIsFirstSample = false
|
|
}
|
|
if h.monotonicTS && dt > 0 {
|
|
nominalGap := float64(n) * dt
|
|
measuredAnchor := anchorTime
|
|
if prevMeasured, ok := src.lastFrameMeasured[sig.Name]; ok {
|
|
measuredGap := measuredAnchor - prevMeasured
|
|
prevEMA, hasEMA := src.gapEMA[sig.Name]
|
|
if !hasEMA {
|
|
prevEMA = nominalGap
|
|
}
|
|
src.gapEMA[sig.Name] = prevEMA*(1-monotonicEMAAlpha) + measuredGap*monotonicEMAAlpha
|
|
smoothedGap := src.gapEMA[sig.Name]
|
|
deviation := math.Abs(measuredGap - smoothedGap)
|
|
if deviation > 0 && deviation < monotonicTolerance {
|
|
anchorTime = src.lastFrameEndT[sig.Name] + smoothedGap
|
|
}
|
|
}
|
|
src.lastFrameMeasured[sig.Name] = measuredAnchor
|
|
src.lastFrameEndT[sig.Name] = anchorTime
|
|
}
|
|
for k := 0; k < n; k++ {
|
|
var t float64
|
|
if anchorIsFirstSample {
|
|
t = anchorTime + float64(k)*dt
|
|
} else {
|
|
t = anchorTime - float64(n-1-k)*dt
|
|
}
|
|
allT = append(allT, t)
|
|
allV = append(allV, vals[k])
|
|
}
|
|
}
|
|
h.ingest(pfx+sig.Name, n, allT, allV)
|
|
decimT, decimV := minMaxDecimate(allT, allV, maxPushPoints)
|
|
pairs[sig.Name] = pairBuf{t: decimT, v: decimV}
|
|
|
|
case sig.TimeMode == udpsprotocol.TimeModeFullArray:
|
|
hasTimeSig := sig.TimeSignalIdx != udpsprotocol.NoTimeSignal && int(sig.TimeSignalIdx) < len(sigs)
|
|
var timeSigName string
|
|
timerToSec := 1e-6
|
|
if hasTimeSig {
|
|
ts := sigs[sig.TimeSignalIdx]
|
|
timeSigName = ts.Name
|
|
if ts.TypeCode == 6 {
|
|
timerToSec = 1e-9
|
|
}
|
|
}
|
|
allT := make([]float64, 0, len(batch)*n)
|
|
allV := make([]float64, 0, len(batch)*n)
|
|
for _, s := range batch {
|
|
vals, ok := s.Values[sig.Name]
|
|
if !ok || len(vals) < n {
|
|
continue
|
|
}
|
|
if hasTimeSig {
|
|
tVals, tOk := s.Values[timeSigName]
|
|
if tOk && len(tVals) >= n {
|
|
if _, exists := src.timeSigCalib[timeSigName]; !exists {
|
|
wallT := float64(s.WallTime.UnixNano()) / 1e9
|
|
src.timeSigCalib[timeSigName] = wallT - tVals[0]*timerToSec
|
|
}
|
|
calib := src.timeSigCalib[timeSigName]
|
|
for k := 0; k < n; k++ {
|
|
allT = append(allT, calib+tVals[k]*timerToSec)
|
|
allV = append(allV, vals[k])
|
|
}
|
|
continue
|
|
}
|
|
}
|
|
wallT := float64(s.WallTime.UnixNano()) / 1e9
|
|
for k := 0; k < n; k++ {
|
|
allT = append(allT, wallT)
|
|
allV = append(allV, vals[k])
|
|
}
|
|
}
|
|
h.ingest(pfx+sig.Name, n, allT, allV)
|
|
decimT, decimV := minMaxDecimate(allT, allV, maxPushPoints)
|
|
pairs[sig.Name] = pairBuf{t: decimT, v: decimV}
|
|
|
|
case n == 1:
|
|
ts := make([]float64, 0, len(batch))
|
|
vs := make([]float64, 0, len(batch))
|
|
for _, s := range batch {
|
|
vals, ok := s.Values[sig.Name]
|
|
if !ok || len(vals) < 1 {
|
|
continue
|
|
}
|
|
ts = append(ts, float64(s.WallTime.UnixNano())/1e9)
|
|
vs = append(vs, vals[0])
|
|
}
|
|
h.ingest(pfx+sig.Name, 1, ts, vs)
|
|
pairs[sig.Name] = pairBuf{t: ts, v: vs}
|
|
|
|
default:
|
|
// n > 1, TimeModePacket: C++ sends samplingRate=0 so we interpolate
|
|
// per-element timestamps from wall-clock differences between packets.
|
|
//
|
|
// Two fixes vs the naïve approach:
|
|
// 1. Use src.lastPktNs[name] for the single-packet case so dt is
|
|
// estimated from the actual inter-packet gap, not 1/n.
|
|
// 2. Send all n elements to the browser without LTTB so sinusoidal
|
|
// waveforms are not degraded (packets arrive at ≤30 Hz, bandwidth
|
|
// is trivially acceptable).
|
|
allT := make([]float64, 0, len(batch)*n)
|
|
allV := make([]float64, 0, len(batch)*n)
|
|
for bi, s := range batch {
|
|
vals, ok := s.Values[sig.Name]
|
|
if !ok || len(vals) < n {
|
|
continue
|
|
}
|
|
wallNs := s.WallTime.UnixNano()
|
|
wallSec := float64(wallNs) / 1e9
|
|
var dtSec float64
|
|
// A gap spans the elements of every packet that went missing
|
|
// inside it as well as this packet's own, so the divisor has
|
|
// to widen with it. Without this a single loss halves the
|
|
// apparent rate and the elements overrun into the next
|
|
// packet's range. The loss count belongs to the packet the
|
|
// gap ends at.
|
|
if bi+1 < len(batch) {
|
|
// Two consecutive packets in this tick → exact dt.
|
|
span := float64(n) * float64(1+batch[bi+1].Lost)
|
|
dtSec = (float64(batch[bi+1].WallTime.UnixNano()) - float64(wallNs)) / 1e9 / span
|
|
} else if bi > 0 {
|
|
// Last of multiple packets → use diff from previous.
|
|
span := float64(n) * float64(1+s.Lost)
|
|
dtSec = (float64(wallNs) - float64(batch[bi-1].WallTime.UnixNano())) / 1e9 / span
|
|
} else if prevNs, ok2 := src.lastPktNs[sig.Name]; ok2 && prevNs > 0 && wallNs > prevNs {
|
|
// Single packet this tick → gap from the previous tick's packet.
|
|
span := float64(n) * float64(1+s.Lost)
|
|
dtSec = (float64(wallNs) - float64(prevNs)) / 1e9 / span
|
|
} else {
|
|
// Truly first packet ever — inter-packet timing unknown.
|
|
// Skip to avoid poisoning the ring with wrongly-spaced timestamps;
|
|
// lastPktNs will be recorded below so the next packet uses correct dt.
|
|
continue
|
|
}
|
|
if h.monotonicTS && dtSec > 0 {
|
|
nominalGap := float64(n) * dtSec
|
|
measuredStart := wallSec
|
|
if prevMeasured, ok := src.lastFrameMeasured[sig.Name]; ok {
|
|
measuredGap := measuredStart - prevMeasured
|
|
prevEMA, hasEMA := src.gapEMA[sig.Name]
|
|
if !hasEMA {
|
|
prevEMA = nominalGap
|
|
}
|
|
src.gapEMA[sig.Name] = prevEMA*(1-monotonicEMAAlpha) + measuredGap*monotonicEMAAlpha
|
|
smoothedGap := src.gapEMA[sig.Name]
|
|
deviation := math.Abs(measuredGap - smoothedGap)
|
|
if deviation > 0 && deviation < monotonicTolerance {
|
|
wallSec = src.lastFrameEndT[sig.Name] + smoothedGap
|
|
}
|
|
}
|
|
src.lastFrameMeasured[sig.Name] = measuredStart
|
|
src.lastFrameEndT[sig.Name] = wallSec
|
|
}
|
|
for j := 0; j < n; j++ {
|
|
allT = append(allT, wallSec+float64(j)*dtSec)
|
|
allV = append(allV, vals[j])
|
|
}
|
|
}
|
|
if len(batch) > 0 {
|
|
src.lastPktNs[sig.Name] = batch[len(batch)-1].WallTime.UnixNano()
|
|
}
|
|
if len(allT) > 0 {
|
|
h.ingest(pfx+sig.Name, n, allT, allV)
|
|
// Live push: never below one packet's worth of elements, or LTTB
|
|
// would flatten the snapshot waveform itself; never above it
|
|
// either, since anything more is just packets that piled up
|
|
// during the tick. Pushing every point unconditionally does not
|
|
// survive a fast producer: a 5 kHz x 1000-element array is 5M
|
|
// points/s on the wire and the client queue never drains.
|
|
thr := maxPushPoints
|
|
if n > thr {
|
|
thr = n
|
|
}
|
|
decimT, decimV := minMaxDecimate(allT, allV, thr)
|
|
pairs[sig.Name] = pairBuf{t: decimT, v: decimV}
|
|
}
|
|
}
|
|
}
|
|
|
|
// Compute total size and serialize
|
|
totalSize := 1 + 1 + len(src.id) + 4
|
|
for key, p := range pairs {
|
|
totalSize += 2 + len(key) + 4
|
|
totalSize += len(p.t)*8 + len(p.v)*8
|
|
}
|
|
|
|
buf := make([]byte, totalSize)
|
|
buf[0] = 1 // version
|
|
buf[1] = byte(len(src.id))
|
|
copy(buf[2:], src.id)
|
|
off := 2 + len(src.id)
|
|
binary.LittleEndian.PutUint32(buf[off:], uint32(len(pairs)))
|
|
off += 4
|
|
|
|
for key, p := range pairs {
|
|
binary.LittleEndian.PutUint16(buf[off:], uint16(len(key)))
|
|
off += 2
|
|
copy(buf[off:], key)
|
|
off += len(key)
|
|
binary.LittleEndian.PutUint32(buf[off:], uint32(len(p.t)))
|
|
off += 4
|
|
off = writeFloat64s(buf, off, p.t)
|
|
off = writeFloat64s(buf, off, p.v)
|
|
}
|
|
|
|
return buf
|
|
}
|
|
|
|
// RecordDataFragment is called by UDPClient for every incoming DATA datagram.
|
|
func (h *Hub) RecordDataFragment(sourceID string, counter uint32, nBytes int, arrivalNs int64, complete bool) {
|
|
h.statsMu.RLock()
|
|
st := h.statsMap[sourceID]
|
|
h.statsMu.RUnlock()
|
|
if st != nil {
|
|
st.RecordFragment(counter, nBytes, arrivalNs, complete)
|
|
}
|
|
}
|
|
|
|
// arrayKey returns the buffer key for element i of an array signal.
|
|
func arrayKey(name string, i int) string {
|
|
return name + "[" + itoa(i) + "]"
|
|
}
|
|
|
|
func itoa(n int) string {
|
|
if n == 0 {
|
|
return "0"
|
|
}
|
|
buf := [20]byte{}
|
|
pos := len(buf)
|
|
for n > 0 {
|
|
pos--
|
|
buf[pos] = byte('0' + n%10)
|
|
n /= 10
|
|
}
|
|
return string(buf[pos:])
|
|
}
|