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()