fixed and improved ui

This commit is contained in:
Martino Ferrari
2026-08-29 23:18:00 +02:00
parent ec0a0cdb12
commit 6b3056c612
+425
View File
@@ -0,0 +1,425 @@
# Buffer time-window & trigger logic in `Client/udpstreamer`
How the UDP Scope client acquires, buffers, times and triggers waveforms.
The pipeline has two halves that must be read together:
- the **Go hub** (`Common/Client/go/wshub/`) — owns the UDP sockets, the
full-resolution sample storage, the disk history, and the trigger FSM;
- the **browser SPA** (`static/app.js`) — owns the display buffers, the rolling
window, and the trigger capture rendering.
The same SPA is also served by `Client/webui` and talks to the C++ StreamHub,
which mirrors the Go hub's behaviour (same trigger FSM states, same binary
frames). Everything below describes the Go-hub path; the wire contracts are
identical on both.
---
## 1. End-to-end data flow
```
MARTe2 RT app ──UDP/UDPS──▶ sources.go: runSession()
│ CONFIG + DATA packets, 17-byte header, HRT timestamp per frame
udpsprotocol.ParseData() → []DataSample{HRTTimestamp, WallTime, Values}
Hub.Run() dataCh → pending[sourceID] (drained every 30 Hz tick)
buildBinaryDataMessageForSource()
├─ rebuild per-sample timestamps from TimeMode / calibration / monotonic snap
├─ h.ingest(key, n, t, v) ← FULL rate: ring.write + hist.write + trigger.feed
└─ minMaxDecimate(…, maxPushPoints=50) → WS binary v1 frame to clients
browser: onBinaryData() → pushBuffer() into per-signal circular buffers
renderDirtyPlots() (rAF loop) → buildUPlotData() → uPlot
```
The 30 Hz push is the **only** live path to the browser and it is decimated to
≤50 points/signal/tick. Everything that needs full resolution — zoom, trigger
captures, disk history — is fed independently through `ingest()` and never goes
over the wire until asked.
---
## 2. Hub-side buffers: `sigRing` (`wshub/ringbuf.go`)
One ring per `"sourceId:signalName"` key. A fixed-capacity circular buffer of
Float64 `(t, v)` pairs with a `sync.RWMutex` (writes from `Hub.Run()`, reads
from HTTP/WS handler goroutines).
### 2.1 Min/max bucketing
`bucket` is how many source samples collapse into **one min/max pair** on the
way in:
- `bucket == 1` — the stream is stored verbatim;
- `bucket > 1` — each group contributes its minimum and its maximum, emitted in
time order (`flushBucketLocked`), so the stored timestamps stay
non-decreasing (reads binary-search `rb.t`).
Bucketing is what lets an arbitrarily long window fit a fixed per-signal memory
budget at a megasample rate. Samples already stored keep the resolution they
were written at; the ring converges on a new bucket as it rolls
(`setBucket`).
### 2.2 Source-rate measurement
`sigRing` keeps its own source-sample accounting (`srcCount`, `srcT0`, `srcT1`,
reset every `srcRateWindowSec = 10 s`), because once `bucket > 1` neither `size`
nor the stored timespan measures the real incoming rate. `sourceRate()` is used
by the tuning sweep and by the history writer.
### 2.3 Ring tuning (`retuneRings`, every 1 s)
`activeWindowSec()` decides how far back the rings must reach:
1. an **armed trigger** owns the window: `cfg.windowSec + captureLagSec`
(`captureLagSec = captureMarginSec + 1/30 ≈ 0.183 s` — the capture is read
out a post-window + margin + one push tick after the trigger, so the rings
must hold that much extra or the front of the capture has already rolled);
2. otherwise the **widest window any connected client is displaying**
(`wsClient.displayWindowSec`, set by the SPA's `setWindow` command), with a
`defaultLiveWindowSec = 10 s` fallback while nobody has said;
Then per ring, with `budget = ringBudget()` (default `defaultRingPts = 10 M`,
floor `ringCapInitial = 250 k`):
- grow to the budget first (`grow()` preserves all samples, never shrinks);
- compute the needed bucket with `ringBucketFor(rate, window, capacity)`
(`ceil(2·rate·window·ringHeadroom / capacity)`, `ringHeadroom = 1.25`);
- apply it with **hysteresis**: keep the current bucket while its coverage is
between `need` and `2·need`, so a rate jittering across the boundary does not
flip the resolution every second.
The history archive is re-sized from the same window (`hist.setWindow`) so a
zoom or capture that outlives the rings can fall back to it.
### 2.4 Reading: `slice(t0, t1)`
Binary search for `t0` then `t1` over the circular layout, returning copies of
the pairs in `[t0, t1]`. Safe to use without holding the lock.
---
## 3. Disk history (`wshub/history.go`)
Optional (`EnableHistory`, `CloseHistory`), enabled by the hub configuration.
Every sample goes to disk through `ingest → hist.write` at full rate, in files
sized for the *current* window (not a retention period). It exists to back
three things the rings cannot:
- **zoom past the window**: `readRange(key, t0, t1, maxOut)`;
- **captures the rings have rolled past**: `captureRange(trigTimepre, trigTime+post)`
lifts each capture into a file of its own so nothing overwrites it before the
next trigger;
- **short captures**: `backfillCaptureHead` prepends the front of the window the
ring no longer holds (the ring only *becomes* as long as the window after a
re-tune; the archive was written straight through).
---
## 4. Live push to the browser
`Hub.Run()` drains `pending[sourceID]` on a 30 Hz ticker. Even with no client
connected the frame is built: that is what keeps feeding rings, history and the
trigger, and keeps push cursors advancing so a late client does not get a
backlog burst.
`buildBinaryDataMessageForSource` reconstructs per-sample timestamps per signal
`TimeMode`:
| Mode | Timestamp reconstruction |
|---|---|
| `FirstSample` / `LastSample` | scalar TimeSignal value × `timerToSec` (µs→s or ns→s for u64), calibrated once against `WallTime`; samples spaced by `1/SamplingRate` |
| `FullArray` | per-element TimeSignal array, calibrated once against `WallTime` |
| scalar (`n == 1`) | `WallTime` of the UDP arrival |
| `PacketTime` (default, n>1) | inter-packet wall-clock gaps divided by n (single-packet ticks use the gap from the previous tick) |
**Monotonic snapping** (optional, `setMonotonic` command / "Sync TS" checkbox):
when enabled, the inter-frame anchor gap is smoothed with an EMA
(`monotonicEMAAlpha = 0.01`, initialised from the nominal `n·dt`) and small
deviations (< `monotonicTolerance = 5 ms`) are snapped to the smoothed gap,
removing the software-dispatch jitter overlaps/gaps described in the StreamHub
docs while tracking the true hardware rate (no accumulated drift).
The live frame is a **binary v1** WS message:
```
[u8 1][u8 srcIdLen][srcId][u32 nSigs]
{[u16 keyLen][key][u32 N][f64 t×N][f64 v×N]}
```
with each signal min/max-decimated to `maxPushPoints = 50` (`minMaxDecimate`:
the range is split into `threshold/2` buckets, each contributing its min and max
in time order — a scope-style envelope that keeps glitches on screen).
---
## 5. Browser-side buffers (`static/app.js`)
### 5.1 Capacity & growth
- `MAX_CAP = 2 000 000` — hard ceiling per buffer (~32 MB/signal at Float64 t+v);
- `DEFAULT_CAP = 100 000` — starting size for scalars;
- `TEMPORAL_CAP = 500 000` — starting size for array signals (the hub pushes
≤50 pts/signal/tick, so this already covers ~5 min);
- `growBufferForWindow(buf, windowSec)`**sizes from the buffer's own span**,
not the signal's sampling rate: the incoming rate here is the hub-decimated
~1.5 kpts/s regardless of the source rate, so rate-based sizing overshot by
three orders of magnitude. Grows only when the buffer is full, to
`windowSec × 1.5` headroom, capped at `MAX_CAP`.
- `growBuffer` copies all existing samples into a larger array (preserving
circular order).
### 5.2 The window
`windowSec` (default 5 s, options 1 s … 10 min) is the rolling viewport.
Changing it:
1. updates `windowSec`;
2. `sendWindow()` → WS `setWindow` → hub `displayWindowSec` → ring re-tune;
3. grows every local buffer via `growBufferForWindow`;
4. evicts the decimation cache (a different window invalidates all cached
renderings).
The rolling "now" anchor is **data-driven, not wall-clock**:
`computePlotNow(p)` takes the newest timestamp of each contributing source and
uses the min-of-max over sources that are still active (a source lagging the
fastest by more than `windowSec` is treated as stale and excluded). This keeps
the window tracking real data regardless of clock skew between hub and browser.
### 5.3 Slicing & rendering
- `getBufferSliceRange(buf, t0, t1)` — binary search on the circular layout,
O(log n + window size);
- `getBufferSliceRangeWithBrackets` — same plus one point on each side so lines
still cross a nearly-empty zoom window;
- `supplementWithBrackets` — same bracketing for sparse server-fetched zoom data.
`buildLiveData(p)`:
1. slices every trace in `[t0, t1]`;
2. picks the **master** signal: highest `SamplingRate`, then most points;
3. decimates the master to ~2× plot width (`DECIM_MIN = 200` floor) via a
background worker (`decimateAsync`, stale-while-revalidate cache keyed per
plot/range/data-generation);
4. resamples every other trace onto the master grid with `resampleLinear`;
5. normalises Y (`applyVScaleNorm`: calibration `v·scale+offset`, then
`(y offset)/div`).
### 5.4 Zoom
A zoom pins `p.xRange` and asks the hub for hi-res data over the exact range
(WS `zoom` request or HTTP `/api/zoom`). The hub answers from the full-res
rings — or from the **held copy of the last trigger capture** (`captureHold`
double buffer) while that window is still relevant — decimated to the requested
point budget. The browser prefers the fetched data when it exists, falls back
to its own circular buffers otherwise, and always brackets with local points.
---
## 6. Hub-side trigger FSM (`wshub/trigger.go`)
### 6.1 States and configuration
```
idle ──arm──▶ armed ──edge──▶ collecting ──window elapsed──▶ triggered
▲ ▲ (pre/post latched) │
│ └───────────── rearm (normal mode, after holdoff) ◀─────────┘
└────────────── disarm / single mode stays triggered
```
Configuration (`trigConfig`, client-settable via WS `setTrigger`):
| field | meaning | clamp |
|---|---|---|
| `signalKey` | `"src:sig"` or `"src:sig[i]"` | — |
| `edge` | `rising` / `falling` / `both` | — |
| `threshold` | **raw** units (SPA converts calibrated → raw) | — |
| `windowSec` | capture window | `[1e-4, 600]` |
| `prePercent` | pre-trigger share | `[0, 100]` |
| `mode` | `normal` (auto-rearm) / `single` | — |
| `holdoffSec` | re-arm delay after a capture, double-trigger guard | `[0, 60]` |
### 6.2 Edge detection (`feed`)
Called from `ingest` with every full-resolution batch for the trigger signal.
Level tracking (`prevValue`/`prevValid`) compares consecutive samples against
the threshold; `[i]`-suffixed keys stride the flattened batch by `nElem` to
watch one column. On a qualifying edge in `armed` state: `latchWindowLocked`
freezes `trigTime` and the pre/post split (so later config edits cannot move a
capture's axis).
### 6.3 Buffer-fill gate
Before accepting an edge, the FSM checks that the trigger signal's ring reaches
back far enough that the capture will come back whole (`fillLocked`):
```
need = windowSec growth × postSec, floored at the pre-window
```
`growth` is the measured span-growth rate of the ring (`setBuffered`, refreshed
by `refreshTriggerFill` from the tick and from trigger commands). A still-filling
ring grows 1 s of span per second, so the gate reduces to the pre-window; a
full ring at a long window needs the whole window. While holding off, the level
is still tracked so the first edge after the gate opens is measured against the
right predecessor. The SPA shows the hold-off as an armed trigger with a
`bufferFill %` badge.
### 6.4 Window timing and the pending edge
`dueCapture` waits for the window on the **sample clock**, not the wall clock:
`lastT ≥ trigTime + post + captureMarginSec(0.15)`. This avoids cutting a
capture short when the stream's timestamps lag real time. Three ways it fires:
1. the samples themselves covered the window;
2. wall-clock fallback when no sample was ever seen (Force from idle);
3. `captureStallSec = 2 s` of stream silence — deliver what was collected
rather than leaving the client stuck in "collecting".
While a capture is in flight the comparator keeps running. The **first**
qualifying edge at/after `notBefore = trigTime + max(post, holdoffSec)` is
remembered (`pendingT`/`pendingValid`) and fired immediately on the automatic
`rearm()`. Without this the trigger was deaf through the whole post-window +
holdoff, which rounded sparse pulse trains up to whole periods (a 1 Hz train at
a 1 s window was caught at 0.5 Hz).
### 6.5 Holdoff and rearm
`markTriggered` moves `collecting → triggered` and, in `normal` mode (not
stopped), schedules `rearmAt = now + cfg.holdoffSec`. `dueRearm` consumes it;
`rearm()` re-arms immediately on a pending edge or returns to `armed`. The
holdoff is measured from the trigger point, overlapping the post-window rather
than adding to it.
`Force()` fires immediately at the most recent sample time (wall clock if no
sample yet) — the "Force" button.
### 6.6 Capture assembly (`buildTriggerCapture`)
On a due capture the hub builds the **binary v2** frame:
```
[u8 2][f64 trigTime][f64 preSec][f64 postSec][u32 nSig]
{[u16 keyLen][fullKey][u32 N][f64 t×N][f64 v×N]}
```
For every ring:
1. `slice(trigTimepre, trigTime+post)`;
2. `backfillCaptureHead` from disk history for the front the ring lost;
3. if it is still short by more than `shortCaptureTol = 1 %` of the window,
log it explicitly (nothing can recover data the ring never held);
4. keep the **full-resolution** slice in the `captureHold` double buffer
(so a zoom into the capture can be answered after the rings roll past);
5. min/max-decimate to `trigCapturePts = 20 000` per signal for the wire —
a 60 s window at 1 MSps is ~960 MB raw per signal and would be dropped by
the send path anyway.
The double buffer is published (`capture.publish`) only once the frame is known
good, so a shot that yielded nothing leaves the previous capture on screen.
Dropped frames (client send-queue full) are logged.
`triggerTick` (every push tick) drives the whole FSM: re-tune rings → open
pending history files → refresh the fill measurement → due capture (send + mark
triggered + `hist.captureRange`) or due rearm → broadcast state only when it
changed (`stateUnsent`).
### 6.7 WS commands
| message | effect |
|---|---|
| `setTrigger {signal, edge, threshold, windowSec, prePercent, mode, holdoffSec}` | replace config |
| `arm` / `rearm` | explicit arm (discards pending edge) |
| `disarm` | → idle |
| `trigStop {stopped}` | pause/resume auto-rearm |
| `forceTrigger` | fire now |
Every command also refreshes the buffer-fill measurement synchronously — at
1 MSps the ring crosses the fill threshold many times inside one 33 ms tick, so
waiting for the next tick would fire on a stale measurement.
---
## 7. Browser-side trigger (`static/app.js`)
### 7.1 State handling
`onTriggerState(msg)` tracks the FSM broadcast:
- **armed** — shows `bufferFill %` while the hub is holding off on the fill
gate, so a trigger that is not yet fireable does not look broken;
- **collecting** — clears the previous snapshot, latches `trigTime` (and
`preSec`/`postSec` if the hub sent them), and lets live data sweep into the
trigger axis (see 7.3);
- **triggered / idle** — bookkeeping for the Rearm/Stop buttons.
### 7.2 Capture handling
`onTriggerCapture` parses the v2 frame into `trig.snapshot[key] = {t, v}` plus
the latched `_preS`/`_postS`. It is **ignored when the client did not enable
the trigger** (`trig.enabled`), because the hub keeps an armed trigger across
client sessions and applying a foreign capture would clobber this client's zoom
and scales. On receipt: the horizontal zoom is dropped so the whole capture is
visible, but **vertical scales (V/div, offset) persist** — they are user
settings and must survive from shot to shot.
### 7.3 Rendering modes (`buildUPlotData`)
| state | renderer | source |
|---|---|---|
| collecting, no snapshot yet | `buildTrigFillData` | live buffers, drawn on the *final* trigger axis (relative seconds, `[-pre, +post]`) so the trace sweeps in from the left |
| armed, not fired | freeze last frame | — |
| snapshot present | `buildTrigData` | the capture (or a hi-res zoom reply that covers ≥98 % of the view, else the snapshot) |
| otherwise | `buildLiveData` | rolling window |
`buildTrigData` converts to trigger-relative time (`t trigT`), picks the
master by rate/count, decimates (cached per range+source-tag), resamples the
other traces, and normalises Y.
### 7.4 Threshold in calibrated units
The trigger threshold is held in **calibrated units** (what the user sees on
the Y axis). `sendTrigConfig()` inverts it through the signal's calibration
before sending: `raw = (calibrated offset)/scale`, so the hub's raw
comparator fires exactly when `GAIN·signal + OFFSET` crosses the threshold.
The threshold line (`drawTriggerMarker`) maps the same calibrated threshold
through the signal's vscale: `y_norm = (threshold offset)/div`.
---
## 8. Key constants
| constant | value | file |
|---|---|---|
| push rate | 30 Hz | `hub.go` |
| `maxPushPoints` (live) | 50 pts/signal/tick | `hub.go` |
| `trigCapturePts` (capture) | 20 000 pts/signal | `trigger.go` |
| `captureMarginSec` | 0.15 s | `trigger.go` |
| `captureStallSec` | 2.0 s | `trigger.go` |
| `autoRearmDelaySec` (default holdoff) | 0.2 s | `trigger.go` |
| `maxTriggerWindowSec` | 600 s | `trigger.go` |
| `ringBudget` default | 10 000 000 pts/signal | `hub.go` |
| `ringCapInitial` | 250 000 | `hub.go` |
| `ringCapScalar` | 100 000 | `hub.go` |
| `ringHeadroom` | 1.25 | `ringbuf.go` |
| `defaultLiveWindowSec` | 10 s | `ringbuf.go` |
| `captureLagSec` | 0.15 + 1/30 ≈ 0.183 s | `ringbuf.go` |
| `monotonicTolerance` | 5 ms | `hub.go` |
| `monotonicEMAAlpha` | 0.01 | `hub.go` |
| `MAX_CAP` (browser) | 2 000 000 pts/signal | `app.js` |
| `DEFAULT_CAP` / `TEMPORAL_CAP` | 100 000 / 500 000 | `app.js` |
| `DECIM_MIN` | 200 | `app.js` |
---
## 9. The C++ StreamHub mirror
The Go hub and the C++ StreamHub implement the same WS contracts and must stay
in sync (`AGENTS.md`): same `triggerState` FSM strings, same v2 capture frame,
same command set (`setTrigger` including `holdoffSec`, `arm`, `disarm`,
`trigStop`, `forceTrigger`), same `trigCapturePts`/`kTrigCapturePts` cap, and
the same ring/history windowing intent (`Source/Applications/StreamHub/`). A
protocol change on one side must be mirrored on the other.