From dbd180880315f9fa7021c95636559893a31b6ed1 Mon Sep 17 00:00:00 2001 From: Daniel Redetzke Date: Mon, 5 Oct 2026 14:39:35 +0300 Subject: [PATCH] Live page: streamed updates and a sliding chart The server pushes each new step over server-sent events (GET /api/v1/live/stream) the moment it is sampled, so updates no longer arrive in uneven pairs. The chart slides left steadily between steps instead of jumping, and the big numbers count to their new value. Both stay still with reduced motion. --- README.md | 5 +- api.go | 43 ++++++++++++ app.css | 1 + app.js | 191 ++++++++++++++++++++++++++++++++++++--------------- main_test.go | 5 ++ speed.go | 29 +++++++- 6 files changed, 217 insertions(+), 57 deletions(-) diff --git a/README.md b/README.md index cccfca9..1fc8d17 100644 --- a/README.md +++ b/README.md @@ -267,6 +267,7 @@ signed in: POST /auth/mfa/keys/begin · /auth/mfa/keys/finish?name= · PATCH|DEL signed in: POST /auth/mfa/recovery-codes GET /status GET /stats?range=24h|7d|30d|90d GET /live?since= (speed per peer, last 2 minutes in 2-second steps) +GET /live/stream (the same as server-sent events) GET /server PATCH /server POST /server/rotate-key GET /server/detect-ip GET /peers POST /peers (returns the config and QR once) GET /peers/{id} PATCH /peers/{id} DELETE /peers/{id} @@ -298,7 +299,9 @@ downloaded, `up` is what it uploaded. `GET /live` answers `{"step": 2, "size": 60, "points": [{"t": …, "peers": {"": [down, up]}}]}` with speeds in bits per second, kept only in memory. -With `since` (unix seconds) it returns only newer steps. +With `since` (unix seconds) it returns only newer steps. `GET /live/stream` +sends the same messages as server-sent events: the history first, then one +message per new step. Latency is measured by pinging the peer's tunnel address every 30 seconds. Set it per peer with `PATCH /peers/{id}` `{"latencyCheck": "off"|"active"|"always"}` diff --git a/api.go b/api.go index 6521a8f..6acd83d 100644 --- a/api.go +++ b/api.go @@ -146,6 +146,7 @@ func (a *App) routes() http.Handler { g("GET /api/v1/status", a.status) g("GET /api/v1/stats", a.allStats) g("GET /api/v1/live", a.liveSpeeds) + g("GET /api/v1/live/stream", a.liveStream) g("GET /api/v1/server", a.getServer) g("PATCH /api/v1/server", a.patchServer) @@ -412,6 +413,48 @@ func (a *App) liveSpeeds(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusOK, map[string]any{"step": int(speedStep / time.Second), "size": speedPoints, "points": points}) } +// liveStream sends the same data as server-sent events: the current points +// first, then each new step as soon as it is sampled. The session is checked +// again with every step, so signing out ends the stream. +func (a *App) liveStream(w http.ResponseWriter, r *http.Request) { + if a.speeds == nil { + writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "live speeds are not available"}) + return + } + rc := http.NewResponseController(w) + _ = rc.SetWriteDeadline(time.Time{}) // the server's write timeout would cut the stream + w.Header().Set("Content-Type", "text/event-stream") + w.Header().Set("Cache-Control", "no-store") + w.Header().Set("X-Accel-Buffering", "no") // nginx: do not buffer + points, ch, cancel := a.speeds.Subscribe() + defer cancel() + send := func(points []SpeedPoint) bool { + b, _ := json.Marshal(map[string]any{"step": int(speedStep / time.Second), "size": speedPoints, "points": points}) + if _, err := fmt.Fprintf(w, "data: %s\n\n", b); err != nil { + return false + } + return rc.Flush() == nil + } + if !send(points) { + return + } + for { + select { + case <-r.Context().Done(): + return + case <-a.speeds.Done(): + return + case pt := <-ch: + if _, ok := a.auth.Authenticate(r); !ok { + return + } + if !send([]SpeedPoint{pt}) { + return + } + } + } +} + func (a *App) peerStats(w http.ResponseWriter, r *http.Request) { cfg := a.store.Get() if _, p := cfg.peerByID(r.PathValue("id")); p == nil { diff --git a/app.css b/app.css index 703360f..4b5f7fa 100644 --- a/app.css +++ b/app.css @@ -391,6 +391,7 @@ pre.log.tall { max-height: calc(100vh - 260px); min-height: 420px; } .chart .ldown, .chart .lup { fill: none; stroke-width: 1.75; stroke-linejoin: round; stroke-linecap: round; vector-effect: non-scaling-stroke; } .chart .ldown { stroke: var(--down); } .chart .lup { stroke: var(--up); } +.chart .plot.live svg { overflow: hidden; } .xaxis.lx span:first-child { transform: none; } .xaxis.lx span:last-child { transform: translateX(-100%); } .spark.wide { width: 96px; } diff --git a/app.js b/app.js index dc7cfdb..b563000 100644 --- a/app.js +++ b/app.js @@ -1119,42 +1119,102 @@ // Below this a peer counts as idle: keepalives and background chatter. const IDLE_BPS = 2000; - // liveChart draws download as a filled area and upload as a line over the - // last points; hover shows the values at a moment. - function liveChart(points) { - const n = points.length, W = 1000, H = 100; - const top = niceTop(Math.max(0, ...points.map((p) => Math.max(p.down, p.up))) || 1e6); - const x = (i) => (n < 2 ? W : i / (n - 1) * W).toFixed(1), y = (v) => (H - v / top * H).toFixed(2); - const line = (k) => points.map((p, i) => x(i) + ',' + y(p[k])).join(' '); + const calm = () => window.matchMedia('(prefers-reduced-motion: reduce)').matches; + + // countTo moves the number in el from its last value to v over a short + // time instead of jumping. + function countTo(el, v) { + const from = el._v ?? v; + el._v = v; + cancelAnimationFrame(el._raf); + if (calm() || from === v) { el.textContent = fmtRate(v); return; } + const t0 = performance.now(), ms = 700; + const tick = (now) => { + const k = Math.min(1, (now - t0) / ms), e = 1 - Math.pow(1 - k, 3); + el.textContent = fmtRate(from + (v - from) * e); + if (k < 1) el._raf = requestAnimationFrame(tick); + }; + el._raf = requestAnimationFrame(tick); + } + + // liveChart draws download as a filled area and upload as a line. It is + // built once and updated in place: each new step enters just beyond the + // right edge and the chart slides left by one step over the step's length, + // so it moves steadily instead of jumping. Hover shows the values at a + // moment. + function liveChart() { + const W = 1000, H = 100; + let pts = [], size = 60, step = 2, off = 0, dx = W / 59, t0 = 0, moving = false; + const area = svg('polygon', { class: 'larea' }); + const down = svg('polyline', { class: 'ldown' }); + const up = svg('polyline', { class: 'lup' }); const cursor = svg('line', { class: 'cursor', x1: 0, x2: 0, y1: 0, y2: H, visibility: 'hidden' }); - const plot = svg('svg', { viewBox: '0 0 ' + W + ' ' + H, preserveAspectRatio: 'none', 'aria-hidden': 'true' }, - n > 1 ? [ - svg('polygon', { class: 'larea', points: '0,' + H + ' ' + line('down') + ' ' + W + ',' + H }), - svg('polyline', { class: 'ldown', points: line('down') }), - svg('polyline', { class: 'lup', points: line('up') }), - ] : null, cursor); + const g = svg('g', null, area, down, up, cursor); + const plot = svg('svg', { viewBox: '0 0 ' + W + ' ' + H, preserveAspectRatio: 'none', 'aria-hidden': 'true' }, g); + const ylTop = h('div', { class: 'yl top' }), ylMid = h('div', { class: 'yl mid' }); const at = (p) => new Date(p.t * 1000).toLocaleTimeString(undefined, { hour: '2-digit', minute: '2-digit', second: '2-digit' }); - const peak = points.reduce((m, p) => p.down + p.up > m.down + m.up ? p : m, { down: 0, up: 0 }); - const idle = () => peak.t - ? ['Peak ', h('strong', null, fmtRate(peak.down + peak.up)), ' at ' + at(peak) + ' · hover the chart to see a moment.'] - : ['No traffic in the last 2 minutes.']; - const readout = h('div', { class: 'readout' }, idle()); - const wrapEl = h('div', { class: 'plot', role: 'img', 'aria-label': 'Speed of all peers over the last 2 minutes', + const idle = () => { + const peak = pts.reduce((m, p) => p.down + p.up > m.down + m.up ? p : m, { down: 0, up: 0 }); + return peak.t + ? ['Peak ', h('strong', null, fmtRate(peak.down + peak.up)), ' at ' + at(peak) + ' · hover the chart to see a moment.'] + : ['No traffic in the last 2 minutes.']; + }; + const readout = h('div', { class: 'readout' }); + let hoverT = null; + const x = (i) => off + i * dx; + const show = (i) => { + const p = pts[i]; + hoverT = p.t; + cursor.setAttribute('x1', x(i).toFixed(1)); cursor.setAttribute('x2', x(i).toFixed(1)); cursor.setAttribute('visibility', 'visible'); + readout.replaceChildren(at(p), ' · Download ', h('strong', null, fmtRate(p.down)), ' · Upload ', h('strong', null, fmtRate(p.up))); + }; + const plotEl = h('div', { class: 'plot live', role: 'img', 'aria-label': 'Speed of all peers over the last 2 minutes', onMousemove: (e) => { - const r = wrapEl.getBoundingClientRect(); - const i = Math.max(0, Math.min(n - 1, Math.round((e.clientX - r.left) / r.width * (n - 1)))); - const p = points[i]; - cursor.setAttribute('x1', x(i)); cursor.setAttribute('x2', x(i)); cursor.setAttribute('visibility', 'visible'); - readout.replaceChildren(at(p), ' · Download ', h('strong', null, fmtRate(p.down)), ' · Upload ', h('strong', null, fmtRate(p.up))); + if (!pts.length) return; + const r = plotEl.getBoundingClientRect(); + const shift = moving ? Math.min(1, (performance.now() - t0) / (step * 1000)) * dx : 0; + const xv = (e.clientX - r.left) / r.width * W + shift; + show(Math.max(0, Math.min(pts.length - 1, Math.round((xv - off) / dx)))); }, - onMouseleave: () => { cursor.setAttribute('visibility', 'hidden'); readout.replaceChildren(...idle()); } }, plot); - return h('div', null, + onMouseleave: () => { hoverT = null; cursor.setAttribute('visibility', 'hidden'); readout.replaceChildren(...idle()); } }, plot); + const el = h('div', null, readout, h('div', { class: 'chart small' }, h('div', { class: 'gl top' }), h('div', { class: 'gl mid' }), h('div', { class: 'gl base' }), - h('div', { class: 'yl top' }, fmtRate(top)), h('div', { class: 'yl mid' }, fmtRate(top / 2)), - wrapEl), + ylTop, ylMid, plotEl), h('div', { class: 'xaxis lx' }, h('span', { style: { left: '0%' } }, '2 min ago'), h('span', { style: { left: '50%' } }, '1 min ago'), h('span', { style: { left: '100%' } }, 'now'))); + + // update draws points; slide is true for a new step arriving live. + const update = (points, sz, st, slide) => { + pts = points; size = sz; step = st; dx = W / (size - 1); + moving = slide && !calm(); + const n = pts.length; + // The newest point sits at the right edge, or one step beyond it + // while sliding in. + off = W - (n - 1) * dx + (moving ? dx : 0); + const top = niceTop(Math.max(0, ...pts.map((p) => Math.max(p.down, p.up))) || 1e6); + const y = (v) => (H - v / top * H).toFixed(2); + const line = (k) => pts.map((p, i) => x(i).toFixed(1) + ',' + y(p[k])).join(' '); + if (n > 1) { + area.setAttribute('points', x(0).toFixed(1) + ',' + H + ' ' + line('down') + ' ' + x(n - 1).toFixed(1) + ',' + H); + down.setAttribute('points', line('down')); + up.setAttribute('points', line('up')); + } + ylTop.textContent = fmtRate(top); + ylMid.textContent = fmtRate(top / 2); + g.style.transition = 'none'; + g.style.transform = 'translateX(0)'; + if (moving) { + g.getBoundingClientRect(); // start the slide from 0 + t0 = performance.now(); + g.style.transition = 'transform ' + step + 's linear'; + g.style.transform = 'translateX(' + (-dx).toFixed(2) + 'px)'; + } + const i = hoverT == null ? -1 : pts.findIndex((p) => p.t === hoverT); + if (i >= 0) show(i); + else { hoverT = null; cursor.setAttribute('visibility', 'hidden'); readout.replaceChildren(...idle()); } + }; + return { el, update }; } // rateSpark draws a peer's total speed over the same window, scaled to its @@ -1169,40 +1229,43 @@ return s; } - // viewLive shows the speed of every peer right now, from the server's - // short in-memory history (the last 2 minutes in 2-second steps). + // viewLive shows the speed of every peer right now. The server streams its + // short in-memory history (the last 2 minutes in 2-second steps) and then + // each new step as it is sampled. async function viewLive(wrap) { let peers = (await api('GET', '/peers')).peers; - let live = await api('GET', '/live'); - let paused = false; - const totals = h('div', { class: 'livenow' }); - const chartBox = h('div'); + let live = { step: 2, size: 60, points: [] }; + let paused = false, es = null, retry = 0; + const down = h('div', { class: 'v' }), up = h('div', { class: 'v' }), active = h('div', { class: 'v' }); + const totals = h('div', { class: 'livenow' }, + h('div', null, h('div', { class: 'k' }, h('span', { class: 'key down' }), 'Download'), down), + h('div', null, h('div', { class: 'k' }, h('span', { class: 'key up' }), 'Upload'), up), + h('div', null, h('div', { class: 'k' }, 'Active peers'), active)); + const lc = liveChart(); const rows = h('div'); const pauseBtn = h('button', { type: 'button', class: 'btn', onClick: () => { paused = !paused; pauseBtn.textContent = paused ? 'Resume' : 'Pause'; sub.replaceChildren(...subText()); + if (paused) close(); else open(); } }, 'Pause'); const subText = () => paused ? ['Paused · the chart keeps the moment you paused'] : [h('span', { class: 'dot ok pulse' }), ' Updated every ' + live.step + ' s · speeds are averages over the step']; const sub = h('p', { class: 'sub livesub' }, subText()); - const draw = () => { + const draw = (slide) => { const pts = live.points; const sumAt = (p) => Object.values(p.peers).reduce((a, [d, u]) => ({ down: a.down + d, up: a.up + u }), { down: 0, up: 0 }); const series = pts.map((p) => ({ t: p.t, ...sumAt(p) })); const now = series[series.length - 1] || { down: 0, up: 0 }; - totals.replaceChildren( - h('div', null, h('div', { class: 'k' }, h('span', { class: 'key down' }), 'Download'), h('div', { class: 'v' }, fmtRate(now.down))), - h('div', null, h('div', { class: 'k' }, h('span', { class: 'key up' }), 'Upload'), h('div', { class: 'v' }, fmtRate(now.up))), - h('div', null, h('div', { class: 'k' }, 'Active peers'), h('div', { class: 'v' }, String(peers.filter((p) => { - const r = pts.length ? pts[pts.length - 1].peers[p.id] : null; - return r && r[0] + r[1] >= IDLE_BPS; - }).length), h('small', null, '/ ' + peers.filter((p) => p.stats.online).length + ' online')))); - chartBox.replaceChildren(liveChart(series)); - const last = pts.length ? pts[pts.length - 1].peers : {}; + countTo(down, now.down); + countTo(up, now.up); + active.replaceChildren(String(peers.filter((p) => { const r = last[p.id]; return r && r[0] + r[1] >= IDLE_BPS; }).length), + h('small', null, '/ ' + peers.filter((p) => p.stats.online).length + ' online')); + lc.update(series, live.size, live.step, slide); + const online = peers.filter((p) => p.stats.online) .map((p) => ({ p, r: last[p.id] || [0, 0], hist: pts.map((x) => { const v = x.peers[p.id]; return v ? v[0] + v[1] : 0; }) })) .sort((a, b) => (b.r[0] + b.r[1]) - (a.r[0] + a.r[1]) || a.p.name.localeCompare(b.p.name)); @@ -1222,26 +1285,44 @@ })))) : h('p', { class: 'empty' }, 'No peer is online.')); }; + // The first message of a stream is the whole history; later ones carry + // one new step each. One step more than the server keeps is held, so + // the chart's left edge stays filled while it slides. + function open() { + if (es || paused || document.hidden || !wrap.isConnected) return; + let first = true; + es = new EventSource('/api/v1/live/stream'); + es.onmessage = (e) => { + retry = 0; + const m = JSON.parse(e.data); + live = { step: m.step, size: m.size, points: first ? m.points : live.points.concat(m.points).slice(-m.size - 1) }; + draw(!first); + first = false; + }; + es.onerror = () => { + if (es.readyState !== EventSource.CLOSED) { first = true; return; } // the browser reconnects + close(); + // Refused, e.g. signed out: api() shows the sign-in page on 401. + api('GET', '/status').then(() => { if (wrap.isConnected) setTimeout(open, Math.min(30000, 2000 * 2 ** retry++)); }).catch(() => {}); + }; + } + function close() { if (es) { es.close(); es = null; } } + const onVis = () => { if (document.hidden) close(); else open(); }; + document.addEventListener('visibilitychange', onVis); + cleanups.push(() => { close(); document.removeEventListener('visibilitychange', onVis); }); + fill(wrap, h('div', { class: 'head' }, h('div', null, h('h1', null, 'Live'), sub), h('div', { class: 'actions' }, pauseBtn)), h('section', { class: 'card', 'aria-labelledby': 'lv' }, h('div', { class: 'cardhead' }, h('h2', { id: 'lv' }, 'All peers · right now')), - totals, chartBox), + totals, lc.el), h('section', { class: 'card flush', 'aria-labelledby': 'lp' }, h('div', { class: 'cardhead' }, h('h2', { id: 'lp' }, 'Peers'), h('span', { class: 'hint' }, 'Busiest first')), rows)); - draw(); - - every(2000, async () => { - if (paused) return; - try { - const more = await api('GET', '/live?since=' + (live.points.length ? live.points[live.points.length - 1].t : 0)); - live = { step: more.step, points: live.points.concat(more.points).slice(-more.size) }; - draw(); - } catch { /* keep last */ } - }); + draw(false); + open(); every(15000, async () => { try { peers = (await api('GET', '/peers')).peers; } catch { /* keep last */ } }); } diff --git a/main_test.go b/main_test.go index 7a7f488..910094f 100644 --- a/main_test.go +++ b/main_test.go @@ -1308,6 +1308,8 @@ func TestSpeeds(t *testing.T) { k.samples = []PeerSample{{PublicKey: pub, RxBytes: rx, TxBytes: tx}} sp.sample(t0.Add(time.Duration(sec) * time.Second)) } + _, ch, cancel := sp.Subscribe() + defer cancel() step(0, 1000, 1000) if n := len(sp.Since(0)); n != 0 { t.Fatalf("first sample made %d points, want 0", n) @@ -1321,6 +1323,9 @@ func TestSpeeds(t *testing.T) { if got, want := pts[0].Peers["p1"], [2]int64{8000, 1000}; got != want { t.Fatalf("speed = %v, want %v (down, up in bit/s)", got, want) } + if got := <-ch; got.T != pts[0].T { + t.Fatalf("subscriber got step %d, want %d", got.T, pts[0].T) + } if _, ok := pts[1].Peers["p1"]; ok { t.Fatal("a counter reset reported a speed") } diff --git a/speed.go b/speed.go index bba9a79..0a1a4e5 100644 --- a/speed.go +++ b/speed.go @@ -30,12 +30,32 @@ type Speeds struct { last map[string][2]int64 // raw rx, tx by public key lastAt time.Time points []SpeedPoint + subs map[chan SpeedPoint]struct{} + done chan struct{} // closed when Run returns } func newSpeeds(store *Store, kernel Kernel) *Speeds { - return &Speeds{store: store, kernel: kernel, last: map[string][2]int64{}} + return &Speeds{store: store, kernel: kernel, last: map[string][2]int64{}, subs: map[chan SpeedPoint]struct{}{}, done: make(chan struct{})} } +// Subscribe returns the current points and a channel that receives each new +// one; cancel ends the subscription. A subscriber that falls behind misses +// points rather than holding up the sampler. +func (s *Speeds) Subscribe() (points []SpeedPoint, ch <-chan SpeedPoint, cancel func()) { + c := make(chan SpeedPoint, 4) + s.mu.Lock() + defer s.mu.Unlock() + s.subs[c] = struct{}{} + return append([]SpeedPoint{}, s.points...), c, func() { + s.mu.Lock() + delete(s.subs, c) + s.mu.Unlock() + } +} + +// Done is closed when the sampler stops, so streams can end. +func (s *Speeds) Done() <-chan struct{} { return s.done } + func (s *Speeds) sample(now time.Time) { cfg := s.store.Get() samples, err := s.kernel.Sample(cfg.Server.Interface) @@ -77,6 +97,12 @@ func (s *Speeds) sample(now time.Time) { if len(s.points) > speedPoints { s.points = s.points[len(s.points)-speedPoints:] } + for c := range s.subs { + select { + case c <- pt: + default: + } + } } // Since returns the points newer than the unix time t, oldest first. @@ -93,6 +119,7 @@ func (s *Speeds) Since(t int64) []SpeedPoint { } func (s *Speeds) Run(stop <-chan struct{}) { + defer close(s.done) s.sample(time.Now()) t := time.NewTicker(speedStep) defer t.Stop()