diff --git a/README.md b/README.md index b85adf3..cccfca9 100644 --- a/README.md +++ b/README.md @@ -28,6 +28,8 @@ dependencies on the server: the binary installs, updates and removes itself. server has a global IPv6 address. - **Traffic history:** kept in `stats.json`, hourly for 48 h and daily for 400 days by default (Settings → Logs & history). +- **Live view:** the speed of every peer right now, updated every 2 seconds, + with the last 2 minutes as a chart. Kept in memory only. - **Connection history:** every online session per peer, with start, duration, address and traffic. A new session starts when a device changes networks. Country and network operator come from the free @@ -264,6 +266,7 @@ signed in: GET /auth/mfa · POST /auth/mfa/totp/setup · /auth/mfa/totp/confirm signed in: POST /auth/mfa/keys/begin · /auth/mfa/keys/finish?name= · PATCH|DELETE /auth/mfa/keys/{id} 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 /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} @@ -293,6 +296,10 @@ the newer version while there is one. Traffic is reported from the peer's point of view: `down` is what the peer 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. + 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"}` (default `off`). `active` pings only while the device sends traffic, so idle diff --git a/api.go b/api.go index 7a4ee06..6521a8f 100644 --- a/api.go +++ b/api.go @@ -21,6 +21,7 @@ type App struct { kernel Kernel recon *Reconciler stats *Stats + speeds *Speeds // nil in tests auth *Auth tls *webTLS logPath string @@ -144,6 +145,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/server", a.getServer) g("PATCH /api/v1/server", a.patchServer) @@ -398,6 +400,18 @@ func (a *App) allStats(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusOK, map[string]any{"range": rng, "points": a.stats.series(nil, rng)}) } +// liveSpeeds returns each peer's speed over the last 2 minutes; with since +// (unix seconds) only the newer steps. +func (a *App) liveSpeeds(w http.ResponseWriter, r *http.Request) { + var since int64 + fmt.Sscan(r.URL.Query().Get("since"), &since) + points := []SpeedPoint{} + if a.speeds != nil { + points = a.speeds.Since(since) + } + writeJSON(w, http.StatusOK, map[string]any{"step": int(speedStep / time.Second), "size": speedPoints, "points": points}) +} + 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 c783017..865263d 100644 --- a/app.css +++ b/app.css @@ -381,3 +381,26 @@ dialog::backdrop { background: rgba(22, 23, 26, .55); } .card h3 { margin: 0; font-size: 16px; font-weight: 600; } .saves { font-size: 12px; color: var(--ink-3); } pre.log.tall { max-height: calc(100vh - 260px); min-height: 420px; } + +/* live */ +.livenow { display: flex; flex-wrap: wrap; gap: 12px 40px; margin-top: 14px; } +.livenow .k { font-size: 13px; color: var(--ink-2); } +.livenow .v { font-size: 30px; font-weight: 600; letter-spacing: -0.01em; margin-top: 4px; font-variant-numeric: tabular-nums; } +.livenow .v small { font-size: 16px; color: var(--ink-3); font-weight: 500; margin-left: 6px; } +.chart .larea { fill: var(--down); opacity: .15; } +.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); } +.xaxis.lx span:first-child { transform: none; } +.xaxis.lx span:last-child { transform: translateX(-100%); } +.spark.wide { width: 96px; } +.spark .fill { fill: var(--down); opacity: .15; stroke: none; } +tr.idle td { color: var(--ink-3); } +tr.idle .spark polyline { stroke: var(--ink-3); } +tr.idle .spark .fill { fill: var(--ink-3); } +td .mono, td.num .mono { font-variant-numeric: tabular-nums; } +.liveoff { margin: 0; padding: 12px; font-size: 13px; color: var(--ink-3); } +.livesub { display: flex; align-items: center; gap: 6px; } +.dot.pulse { animation: pulse 2s ease-in-out infinite; } +@keyframes pulse { 50% { opacity: .35; } } +@media (prefers-reduced-motion: reduce) { .dot.pulse { animation: none; } } diff --git a/app.js b/app.js index 40eeb9d..e15e4e7 100644 --- a/app.js +++ b/app.js @@ -47,6 +47,7 @@ peers: '', server: '', settings: '', + live: '', plus: '', key: '', log: '', @@ -555,7 +556,7 @@ // ---------- shell, router ---------- - const NAV = [['#/', 'dashboard', 'Dashboard'], ['#/peers', 'peers', 'Peers'], ['#/server', 'server', 'Server'], ['#/log', 'log', 'Log'], ['#/settings', 'settings', 'Settings']]; + const NAV = [['#/', 'dashboard', 'Dashboard'], ['#/live', 'live', 'Live'], ['#/peers', 'peers', 'Peers'], ['#/server', 'server', 'Server'], ['#/log', 'log', 'Log'], ['#/settings', 'settings', 'Settings']]; let navLinks = {}; let srvBox, peerCount, verRow; @@ -617,6 +618,7 @@ const ROUTES = [ [/^#\/?$/, '#/', viewDashboard], + [/^#\/live$/, '#/live', viewLive], [/^#\/peers$/, '#/peers', viewPeers], [/^#\/peers\/new$/, '#/peers', viewPeerNew], [/^#\/peers\/([\w-]+)$/, '#/peers', viewPeer], @@ -1103,6 +1105,148 @@ every(30000, () => draw().catch(() => {})); } + // ---------- live ---------- + + // fmtRate formats a speed in bits per second. + function fmtRate(bps) { + const u = ['bit/s', 'kbit/s', 'Mbit/s', 'Gbit/s']; + let i = 0, v = bps; + while (v >= 1000 && i < u.length - 1) { v /= 1000; i++; } + const s = i === 0 ? String(Math.round(v)) : v < 10 ? v.toFixed(1) : String(Math.round(v)); + return s + ' ' + u[i]; + } + + // 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 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 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', + 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))); + }, + onMouseleave: () => { cursor.setAttribute('visibility', 'hidden'); readout.replaceChildren(...idle()); } }, plot); + return 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), + 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'))); + } + + // rateSpark draws a peer's total speed over the same window, scaled to its + // own peak. + function rateSpark(vals) { + const s = svg('svg', { class: 'spark wide', viewBox: '0 0 96 18', preserveAspectRatio: 'none', 'aria-hidden': 'true' }); + const top = Math.max(...vals, IDLE_BPS * 4); + const n = vals.length; + if (n < 2) return s; + const pts = vals.map((v, i) => (i / (n - 1) * 96).toFixed(1) + ',' + (17 - v / top * 15).toFixed(1)); + s.append(svg('polygon', { class: 'fill', points: '0,18 ' + pts.join(' ') + ' 96,18' }), svg('polyline', { points: pts.join(' ') })); + 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). + 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'); + const rows = h('div'); + const pauseBtn = h('button', { type: 'button', class: 'btn', onClick: () => { + paused = !paused; + pauseBtn.textContent = paused ? 'Resume' : 'Pause'; + sub.replaceChildren(...subText()); + } }, '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 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 : {}; + 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)); + const offline = peers.filter((p) => !p.stats.online); + const rate = (v, idle) => idle ? h('span', { class: 'muted' }, '–') : h('span', { class: 'mono' }, fmtRate(v)); + rows.replaceChildren( + online.length ? h('div', { class: 'tbl' }, h('table', { class: 'narrow' }, + h('thead', null, h('tr', null, h('th', null, 'Peer'), h('th', null, 'Last 2 minutes'), h('th', { class: 'num' }, 'Download'), h('th', { class: 'num' }, 'Upload'), h('th', null, 'Endpoint'))), + h('tbody', null, online.map(({ p, r, hist }) => { + const idle = r[0] + r[1] < IDLE_BPS; + return h('tr', { class: idle ? 'idle' : null }, + h('td', null, h('a', { href: '#/peers/' + p.id }, p.name), idle ? h('span', { class: 'tag plain' }, 'idle') : null), + h('td', null, rateSpark(hist)), + h('td', { class: 'num' }, rate(r[0], idle)), + h('td', { class: 'num' }, rate(r[1], idle)), + h('td', null, h('span', { class: 'mono' }, p.stats.endpoint ? p.stats.endpoint.replace(/:\d+$/, '') : '–'), + p.stats.location && p.stats.location.country ? h('span', { class: 'cc', title: fmtLocation(p.stats.location) }, p.stats.location.country) : null)); + })))) : h('p', { class: 'empty' }, 'No peer is online.'), + offline.length ? h('p', { class: 'liveoff' }, 'Offline: ', offline.map((p, i) => [i ? ', ' : '', h('a', { href: '#/peers/' + p.id }, p.name)])) : null); + }; + + 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), + 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 */ } + }); + every(15000, async () => { try { peers = (await api('GET', '/peers')).peers; } catch { /* keep last */ } }); + } + function describeAudit(l) { const parts = [l.msg.charAt(0).toUpperCase() + l.msg.slice(1)]; if (l.peer) parts.push(': ' + l.peer); diff --git a/kernel_other.go b/kernel_other.go index a5cb2e3..81fe495 100644 --- a/kernel_other.go +++ b/kernel_other.go @@ -5,8 +5,10 @@ package main import ( "fmt" "log/slog" + "maps" "math/rand/v2" "net/netip" + "slices" "sync" "time" ) @@ -17,6 +19,7 @@ import ( type simKernel struct { mu sync.Mutex peers map[string]*PeerSample + last time.Time // previous Sample; traffic grows with the time since } func newKernel() (Kernel, error) { @@ -53,17 +56,23 @@ func (k *simKernel) Apply(c *Config) error { func (k *simKernel) Sample(string) ([]PeerSample, error) { k.mu.Lock() defer k.mu.Unlock() + now := time.Now() + f := 1.0 + if !k.last.IsZero() { + f = now.Sub(k.last).Seconds() / 30 + } + k.last = now var out []PeerSample - i := 0 - for _, p := range k.peers { - // Every third peer stays idle; the others move some data. + for i, key := range slices.Sorted(maps.Keys(k.peers)) { + p := k.peers[key] + // Every third peer stays idle; the others move some data, scaled to + // the time since the previous sample. if i%3 != 2 { - p.TxBytes += rand.Int64N(40 << 20) - p.RxBytes += rand.Int64N(6 << 20) - p.LastHandshake = time.Now().Add(-time.Duration(rand.IntN(90)) * time.Second) + p.TxBytes += int64(float64(rand.Int64N(40<<20)) * f) + p.RxBytes += int64(float64(rand.Int64N(6<<20)) * f) + p.LastHandshake = now.Add(-time.Duration(rand.IntN(90)) * time.Second) } out = append(out, *p) - i++ } return out, nil } diff --git a/main.go b/main.go index f0384a7..ddb5b31 100644 --- a/main.go +++ b/main.go @@ -219,16 +219,18 @@ func run(configPath string) error { var stopOnce sync.Once shutdown := func() { stopOnce.Do(func() { close(stop) }) } + speeds := newSpeeds(store, kernel) auth := newAuth(store) app := &App{ - store: store, kernel: kernel, recon: recon, stats: stats, auth: auth, tls: webTLS, + store: store, kernel: kernel, recon: recon, stats: stats, speeds: speeds, auth: auth, tls: webTLS, logPath: logPath, logw: logw, geo: geo, updates: newUpdater(cfg.Updates), started: time.Now(), shutdown: shutdown, } var wg sync.WaitGroup - wg.Add(5) + wg.Add(6) go func() { defer wg.Done(); recon.Run(stop) }() go func() { defer wg.Done(); stats.Run(stop) }() + go func() { defer wg.Done(); speeds.Run(stop) }() go func() { defer wg.Done(); stats.RunPings(stop) }() go func() { defer wg.Done(); geo.Run(stop) }() go func() { defer wg.Done(); app.updates.Run(stop) }() diff --git a/main_test.go b/main_test.go index 0dd95d5..7a7f488 100644 --- a/main_test.go +++ b/main_test.go @@ -1286,3 +1286,51 @@ func TestDropSecurityKeys(t *testing.T) { t.Fatal("security key still in config.json") } } + +func TestSpeeds(t *testing.T) { + dir := t.TempDir() + store, err := openStore(filepath.Join(dir, "config.json")) + if err != nil { + t.Fatal(err) + } + key, _ := newPrivateKey() + pub := key.PublicKey().String() + if err := store.Update(func(c *Config) error { + c.Peers = append(c.Peers, Peer{ID: "p1", Name: "phone", IPv4: serverIPv4(netip.MustParsePrefix(c.Server.IPv4)).Next().String(), PublicKey: pub, Enabled: true}) + return nil + }); err != nil { + t.Fatal(err) + } + k := &fakeKernel{} + sp := newSpeeds(store, k) + t0 := time.Unix(1_800_000_000, 0) + step := func(sec int, rx, tx int64) { + k.samples = []PeerSample{{PublicKey: pub, RxBytes: rx, TxBytes: tx}} + sp.sample(t0.Add(time.Duration(sec) * time.Second)) + } + step(0, 1000, 1000) + if n := len(sp.Since(0)); n != 0 { + t.Fatalf("first sample made %d points, want 0", n) + } + step(2, 1250, 3000) // +250 up, +2000 down in 2 s + step(4, 10, 20) // counter reset: no speed for this step + pts := sp.Since(0) + if len(pts) != 2 { + t.Fatalf("got %d points, want 2", len(pts)) + } + 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 _, ok := pts[1].Peers["p1"]; ok { + t.Fatal("a counter reset reported a speed") + } + if got := sp.Since(pts[0].T); len(got) != 1 || got[0].T != pts[1].T { + t.Fatalf("Since returned %v", got) + } + for i := 0; i < speedPoints+5; i++ { + step(6+2*i, 0, 0) + } + if n := len(sp.Since(0)); n != speedPoints { + t.Fatalf("kept %d points, want %d", n, speedPoints) + } +} diff --git a/speed.go b/speed.go new file mode 100644 index 0000000..bba9a79 --- /dev/null +++ b/speed.go @@ -0,0 +1,107 @@ +package main + +import ( + "log/slog" + "sync" + "time" +) + +// Speeds keeps the last few minutes of each peer's speed in memory for the +// Live page. It reads the kernel counters every speedStep, apart from the +// traffic history in Stats, and never writes to disk. + +const ( + speedStep = 2 * time.Second + speedPoints = 60 // 2 minutes +) + +// SpeedPoint is one step: per peer ID, download and upload in bits per +// second, from the peer's point of view. +type SpeedPoint struct { + T int64 `json:"t"` + Peers map[string][2]int64 `json:"peers"` +} + +type Speeds struct { + store *Store + kernel Kernel + + mu sync.Mutex + last map[string][2]int64 // raw rx, tx by public key + lastAt time.Time + points []SpeedPoint +} + +func newSpeeds(store *Store, kernel Kernel) *Speeds { + return &Speeds{store: store, kernel: kernel, last: map[string][2]int64{}} +} + +func (s *Speeds) sample(now time.Time) { + cfg := s.store.Get() + samples, err := s.kernel.Sample(cfg.Server.Interface) + if err != nil { + slog.Debug("speed sample failed", "err", err) + return + } + idByKey := map[string]string{} + for _, p := range cfg.Peers { + if p.hasKey() { + idByKey[p.PublicKey] = p.ID + } + } + s.mu.Lock() + defer s.mu.Unlock() + secs := now.Sub(s.lastAt).Seconds() + first := s.lastAt.IsZero() + cur := map[string][2]int64{} + pt := SpeedPoint{T: now.Unix(), Peers: map[string][2]int64{}} + for _, smp := range samples { + cur[smp.PublicKey] = [2]int64{smp.RxBytes, smp.TxBytes} + id := idByKey[smp.PublicKey] + prev, ok := s.last[smp.PublicKey] + if id == "" || !ok || first { + continue + } + dRx, dTx := smp.RxBytes-prev[0], smp.TxBytes-prev[1] + if dRx < 0 || dTx < 0 { // counters were reset + continue + } + // Tx is what the server sent: the peer's download. + pt.Peers[id] = [2]int64{int64(float64(dTx*8) / secs), int64(float64(dRx*8) / secs)} + } + s.last, s.lastAt = cur, now + if first { + return + } + s.points = append(s.points, pt) + if len(s.points) > speedPoints { + s.points = s.points[len(s.points)-speedPoints:] + } +} + +// Since returns the points newer than the unix time t, oldest first. +func (s *Speeds) Since(t int64) []SpeedPoint { + s.mu.Lock() + defer s.mu.Unlock() + out := []SpeedPoint{} + for _, p := range s.points { + if p.T > t { + out = append(out, p) + } + } + return out +} + +func (s *Speeds) Run(stop <-chan struct{}) { + s.sample(time.Now()) + t := time.NewTicker(speedStep) + defer t.Stop() + for { + select { + case <-stop: + return + case now := <-t.C: + s.sample(now) + } + } +}