diff --git a/README.md b/README.md index 46a6d07..e4d09aa 100644 --- a/README.md +++ b/README.md @@ -221,6 +221,7 @@ GET /peers POST /peers (returns the config and QR once) GET /peers/{id} PATCH /peers/{id} DELETE /peers/{id} POST /peers/{id}/enable | /disable | /issue-config GET /peers/{id}/stats?range=… GET /peers/{id}/sessions?limit=100 +GET /peers/{id}/latency (24 h, one point per 5 minutes) GET /peers/{id}/setup (not read-only) DELETE /peers/{id}/setup GET /settings PATCH /settings POST /restart GET /logs?level=&limit=&audit=1 GET /logs/download @@ -236,6 +237,15 @@ link, the peer's current keys keep working until the link is opened. Traffic is reported from the peer's point of view: `down` is what the peer downloaded, `up` is what it uploaded. +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 +phones are not woken up; `always` keeps the tunnel up, so the peer always shows +as online. The service opens an unprivileged ICMP socket, which needs its group +in the sysctl `net.ipv4.ping_group_range` (systemd allows all groups by +default). Devices that block ping, such as Windows with its default firewall, +show no reply. + ## Firewall note GHOSTWIRE's rules sit in their own nftables table. An accept there cannot diff --git a/api.go b/api.go index 0d3b786..24b4f1f 100644 --- a/api.go +++ b/api.go @@ -1,6 +1,7 @@ package main import ( + "cmp" "context" "encoding/json" "errors" @@ -123,6 +124,7 @@ func (a *App) routes() http.Handler { g("POST /api/v1/peers/{id}/disable", a.setEnabled(false)) g("POST /api/v1/peers/{id}/issue-config", a.issueConfig) g("GET /api/v1/peers/{id}/stats", a.peerStats) + g("GET /api/v1/peers/{id}/latency", a.peerLatency) g("GET /api/v1/peers/{id}/sessions", a.peerSessions) g("GET /api/v1/peers/{id}/setup", a.getSetup) g("DELETE /api/v1/peers/{id}/setup", a.revokeSetup) @@ -269,6 +271,9 @@ func (a *App) status(w http.ResponseWriter, r *http.Request) { ac.Detail = applyErr.Error() } checks = append(checks, ac) + if pc, ok := a.stats.PingCheck(cfg); ok { + checks = append(checks, pc) + } healthy := true for _, c := range checks { healthy = healthy && c.OK @@ -319,6 +324,26 @@ func (a *App) peerStats(w http.ResponseWriter, r *http.Request) { writeJSON(w, http.StatusOK, map[string]any{"range": rng, "points": a.stats.series([]string{r.PathValue("id")}, rng)}) } +func (a *App) peerLatency(w http.ResponseWriter, r *http.Request) { + if _, p := a.store.Get().peerByID(r.PathValue("id")); p == nil { + writeJSON(w, http.StatusNotFound, map[string]string{"error": "no such peer"}) + return + } + writeJSON(w, http.StatusOK, map[string]any{"stepSeconds": int(latencyStep.Seconds()), "points": a.stats.LatencyHistory(r.PathValue("id"))}) +} + +// latencyMode turns the API value into the stored one: "off" (or nothing) +// is stored empty. +func latencyMode(v string) (string, error) { + if v == "off" { + return latencyOff, nil + } + if !validLatencyCheck(v) { + return "", badRequest("latencyCheck must be off, active or always") + } + return v, nil +} + // --- server --- type serverView struct { @@ -517,6 +542,7 @@ type peerView struct { EffDNS []string `json:"effectiveDNS"` EffAllowed []string `json:"effectiveAllowedIPs"` EffKeepalive int `json:"effectiveKeepalive"` + LatencyCheck string `json:"latencyCheck"` // off | active | always Created time.Time `json:"created"` ConfigIssued *time.Time `json:"configIssued"` Setup *setupView `json:"setup"` // null = no pending setup link @@ -528,7 +554,8 @@ func (a *App) peerView(c *Config, p *Peer) peerView { ID: p.ID, Name: p.Name, Note: p.Note, Enabled: p.Enabled, PublicKey: p.PublicKey, HasPSK: p.PresharedKey != "", IPv4: p.IPv4, DNS: p.DNS, AllowedIPs: p.AllowedIPs, Keepalive: p.Keepalive, EffDNS: peerDNS(c, p), EffAllowed: peerAllowedIPs(c, p), EffKeepalive: peerKeepalive(c, p), - Created: p.Created, ConfigIssued: p.ConfigIssued, Setup: viewSetup(p.Setup), Stats: a.stats.Summary(p.ID), + LatencyCheck: cmp.Or(p.LatencyCheck, "off"), + Created: p.Created, ConfigIssued: p.ConfigIssued, Setup: viewSetup(p.Setup), Stats: a.stats.Summary(p.ID), } if c.Server.IPv6Enabled { v.IPv6 = mapIPv6(netip.MustParsePrefix(c.Server.IPv6), netip.MustParseAddr(p.IPv4)).String() @@ -610,6 +637,7 @@ func (a *App) createPeer(w http.ResponseWriter, r *http.Request) { AllowedIPs []string `json:"allowedIPs"` Keepalive *int `json:"keepalive"` PresharedKey *bool `json:"presharedKey"` + LatencyCheck string `json:"latencyCheck"` setupRequest } if err := readJSON(r, &in); err != nil { @@ -617,6 +645,11 @@ func (a *App) createPeer(w http.ResponseWriter, r *http.Request) { return } in.Name = strings.TrimSpace(in.Name) + lc, err := latencyMode(in.LatencyCheck) + if err != nil { + writeErr(w, err) + return + } link, err := in.newLink() if err != nil { writeErr(w, err) @@ -624,7 +657,7 @@ func (a *App) createPeer(w http.ResponseWriter, r *http.Request) { } p := Peer{ ID: newID(), Name: in.Name, Note: strings.TrimSpace(in.Note), Enabled: true, - DNS: in.DNS, AllowedIPs: in.AllowedIPs, Keepalive: in.Keepalive, Created: time.Now().UTC(), + DNS: in.DNS, AllowedIPs: in.AllowedIPs, Keepalive: in.Keepalive, LatencyCheck: lc, Created: time.Now().UTC(), } // With a link, the keys are made when the link is opened. var priv string @@ -735,6 +768,18 @@ func (a *App) patchPeer(w http.ResponseWriter, r *http.Request) { } changed = append(changed, "keepalive") } + if raw, ok := m["latencyCheck"]; ok { + var v string + if err := json.Unmarshal(raw, &v); err != nil { + return badRequest("latencyCheck: %v", err) + } + lc, err := latencyMode(v) + if err != nil { + return err + } + p.LatencyCheck = lc + changed = append(changed, "latencyCheck") + } p.Name = strings.TrimSpace(p.Name) p.Note = strings.TrimSpace(p.Note) name = p.Name diff --git a/app.css b/app.css index 45ee616..82399e5 100644 --- a/app.css +++ b/app.css @@ -181,6 +181,20 @@ fieldset { border: 0; margin: 0; padding: 0; min-width: 0; display: flex; flex-d .chart .grp.on span.total { background: var(--down-strong); } .xaxis { display: flex; justify-content: space-between; margin: 8px 0 0 56px; font-size: 11px; color: var(--axis-ink); } .chart.loading { opacity: .5; } +.chart .plot { position: absolute; left: 56px; right: 0; top: 0; bottom: 1px; } +.chart .plot svg { width: 100%; height: 100%; display: block; overflow: visible; } +.chart .band, .key.band { fill: var(--down); background: var(--down); opacity: .2; } +.chart .med { fill: none; stroke: var(--down); stroke-width: 1.75; stroke-linejoin: round; stroke-linecap: round; vector-effect: non-scaling-stroke; } +.chart .cursor { stroke: var(--axis); stroke-width: 1; vector-effect: non-scaling-stroke; } + +/* latency in the peer list */ +.lat { display: inline-flex; vertical-align: middle; align-items: center; justify-content: flex-end; gap: 8px; } +.lat .mono { min-width: 6ch; text-align: right; } +.latnote { font-size: 12px; color: var(--ink-3); } +.lat.stale, .latnote.stale { opacity: .5; } +.spark { width: 56px; height: 18px; display: block; flex: none; } +.spark polyline { fill: none; stroke: var(--down); stroke-width: 1.5; stroke-linejoin: round; stroke-linecap: round; } +.spark circle { fill: var(--down); } /* config / log blocks */ pre.code, pre.log { margin: 0; padding: 14px; background: var(--ink); color: #e6e6e1; border-radius: 10px; font-family: var(--mono); font-size: 12px; line-height: 1.6; overflow-x: auto; white-space: pre; } diff --git a/app.js b/app.js index 03c8464..13e59bd 100644 --- a/app.js +++ b/app.js @@ -148,6 +148,54 @@ const badge = (st) => h('span', { class: 'badge' }, h('span', { class: st.dot }), st.label); + // svg builds an SVG element; attrs are set as attributes. + function svg(tag, attrs, ...kids) { + const el = document.createElementNS('http://www.w3.org/2000/svg', tag); + for (const [k, v] of Object.entries(attrs || {})) el.setAttribute(k, v); + add(el, kids); + return el; + } + + const fmtMs = (ms) => (ms < 10 ? ms.toFixed(1) : String(Math.round(ms))) + ' ms'; + + // A latency value is stale when the server stopped pinging, e.g. because + // the device went idle; it is then shown greyed out. + const latStale = (l) => Date.now() - Date.parse(l.at) > 3 * 30000; + + // latState describes a peer's latency for the lists: null when there is + // nothing to show. + function latState(p) { + if (!p.enabled || !p.publicKey) return null; + if (p.latencyCheck === 'off') return { note: 'Check off', title: 'The latency check is off for this peer' }; + const l = p.stats.latency; + const idle = p.latencyCheck === 'active' ? 'Pinged only while the device sends traffic' : 'Not measured yet'; + if (!l) return { note: '–', title: idle }; + const stale = latStale(l); + const when = stale ? ' · measured ' + ago(l.at) : ''; + if (l.ms == null) return { note: 'No ping reply', stale, title: 'The device does not answer ping. Windows blocks it in its firewall by default.' + when }; + return { ms: l.ms, spark: l.spark, stale, title: 'Median of the last 5 minutes · ' + fmtMs(l.min) + '–' + fmtMs(l.max) + ' · ' + l.loss + ' % loss' + when }; + } + + // sparkline draws the medians of the last hour; gaps are skipped. + function sparkline(vals) { + const pts = (vals || []).map((v, i) => [i, v]).filter(([, v]) => v != null); + const s = svg('svg', { class: 'spark', viewBox: '0 0 56 18', 'aria-hidden': 'true' }); + if (pts.length < 2) return s; + const lo = Math.min(...pts.map(([, v]) => v)), hi = Math.max(...pts.map(([, v]) => v)); + const span = Math.max(hi - lo, hi * 0.2, 1); + const xy = ([i, v]) => [(i / (vals.length - 1) * 54 + 1).toFixed(1), (16 - (v - lo) / span * 14).toFixed(1)]; + const [ex, ey] = xy(pts[pts.length - 1]); + s.append(svg('polyline', { points: pts.map((p) => xy(p).join(',')).join(' ') }), svg('circle', { cx: ex, cy: ey, r: 2 })); + return s; + } + + function latCell(p) { + const st = latState(p); + if (!st) return h('td', { class: 'num muted' }, '–'); + if (st.ms == null) return h('td', { class: 'num' }, h('span', { class: st.stale ? 'latnote stale' : 'latnote', title: st.title }, st.note)); + return h('td', { class: 'num' }, h('span', { class: st.stale ? 'lat stale' : 'lat', title: st.title }, sparkline(st.spark), h('span', { class: 'mono' }, fmtMs(st.ms)))); + } + // ---------- API ---------- async function api(method, path, body) { @@ -346,6 +394,63 @@ h('span', null, range === '24h' ? 'now' : pointLabel(points[points.length - 1].t, range)))); } + // latencyChart draws the median as a line over a min–max band, one point + // per 5 minutes; steps without replies leave a gap. Hover shows the values. + function latencyChart(points) { + const ok = (p) => p.sent > p.lost; + const peak = Math.max(0, ...points.filter(ok).map((p) => p.max)); + const top = peak > 0 ? niceTop(peak) : 100; + const n = points.length, W = 1000, H = 100; + const x = (i) => (i + 0.5) / n * W, y = (v) => (H - v / top * H).toFixed(2); + const segs = []; + points.forEach((p, i) => { + if (!ok(p)) return; + const last = segs[segs.length - 1]; + if (last && last[last.length - 1] === i - 1) last.push(i); else segs.push([i]); + }); + const half = W / n * 0.4; + const shapes = segs.flatMap((seg) => { + // A lone point gets a short flat stretch so it stays visible. + const xs = seg.length > 1 ? seg.map(x) : [x(seg[0]) - half, x(seg[0]) + half]; + const at = (k) => points[seg[Math.min(k, seg.length - 1)]]; + const upper = xs.map((xv, k) => xv.toFixed(1) + ',' + y(at(k).max)); + const lower = xs.map((xv, k) => xv.toFixed(1) + ',' + y(at(k).min)).reverse(); + return [ + svg('polygon', { class: 'band', points: upper.concat(lower).join(' ') }), + svg('polyline', { class: 'med', points: xs.map((xv, k) => xv.toFixed(1) + ',' + y(at(k).med)).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' }, shapes, cursor); + const label = (p) => { + const d = new Date(p.t * 1000); + return d.toLocaleTimeString(undefined, { hour: '2-digit', minute: '2-digit' }); + }; + const withData = points.filter(ok); + const idle = () => withData.length + ? ['Median over 24 hours ', h('strong', null, fmtMs(withData.map((p) => p.med).sort((a, b) => a - b)[withData.length >> 1])), ' · hover the chart to see a time.'] + : ['No measurements in the last 24 hours.']; + const readout = h('div', { class: 'readout' }, idle()); + const wrapEl = h('div', { class: 'plot', role: 'img', 'aria-label': 'Latency over the last 24 hours', + onMousemove: (e) => { + const r = wrapEl.getBoundingClientRect(); + const i = Math.max(0, Math.min(n - 1, Math.floor((e.clientX - r.left) / r.width * n))); + const p = points[i]; + cursor.setAttribute('x1', x(i)); cursor.setAttribute('x2', x(i)); cursor.setAttribute('visibility', 'visible'); + readout.replaceChildren(...(ok(p) + ? [label(p), ' · median ', h('strong', null, fmtMs(p.med)), ' · ' + fmtMs(p.min) + '–' + fmtMs(p.max) + ' · ' + Math.round(p.lost * 100 / p.sent) + ' % loss'] + : [label(p), p.sent ? ' · no reply' : ' · not measured'])); + }, + 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' }, fmtMs(top)), h('div', { class: 'yl mid' }, fmtMs(top / 2)), + wrapEl), + h('div', { class: 'xaxis' }, h('span', null, '24 h ago'), h('span', null, '12 h ago'), h('span', null, 'now'))); + } + // ---------- shell, router ---------- const NAV = [['#/', 'dashboard', 'Dashboard'], ['#/peers', 'peers', 'Peers'], ['#/server', 'server', 'Server'], ['#/settings', 'settings', 'Settings']]; @@ -574,6 +679,7 @@ h('td', null, badge(peerState(p))), h('td', { class: 'mono muted' }, p.stats.endpoint || '–', p.stats.location && p.stats.location.country ? h('span', { class: 'cc', title: fmtLocation(p.stats.location) }, p.stats.location.country) : null), + latCell(p), h('td', { class: 'num' }, fmtBytes(p.stats.down30d)), h('td', { class: 'num' }, fmtBytes(p.stats.up30d)), h('td', null, h('input', { type: 'checkbox', class: 'sw', checked: p.enabled, 'aria-label': (p.enabled ? 'Disable ' : 'Enable ') + p.name, onChange: (e) => toggle(p, e.target.checked) })), @@ -593,9 +699,9 @@ pills), h('section', { class: 'card flush' }, h('div', { class: 'tbl' }, h('table', null, h('thead', null, h('tr', null, ['Name', 'Address', 'Status', 'Endpoint'].map((t) => h('th', null, t)), - h('th', { class: 'num' }, 'Download, 30 d'), h('th', { class: 'num' }, 'Upload, 30 d'), h('th', null, 'Enabled'), h('th', null, h('span', { class: 'sr' }, 'Actions')))), + h('th', { class: 'num' }, 'Latency'), h('th', { class: 'num' }, 'Download, 30 d'), h('th', { class: 'num' }, 'Upload, 30 d'), h('th', null, 'Enabled'), h('th', null, h('span', { class: 'sr' }, 'Actions')))), tbody), empty)), - h('p', { class: 'muted', style: { margin: '0', fontSize: '13px' } }, 'Online means a handshake in the last 3 minutes. Download and Upload are measured from the peer\'s side. Changes apply live without disconnecting other peers.')); + h('p', { class: 'muted', style: { margin: '0', fontSize: '13px' } }, 'Online means a handshake in the last 3 minutes. Latency is the round trip from the server through the tunnel to the device and back, median of the last 5 minutes; turn it on in a peer\'s settings. Download and Upload are measured from the peer\'s side. Changes apply live without disconnecting other peers.')); every(15000, async () => { try { data = await api('GET', '/peers'); drawRows(); } catch { /* keep last */ } }); } @@ -745,6 +851,21 @@ const totals = h('div', { class: 'legend-row' }); const pills = h('div', { class: 'pills', role: 'group', 'aria-label': 'Time range' }); + const showLat = p.latencyCheck !== 'off' || p.stats.latency; + const latBox = h('div'); + async function drawLatency() { + if (!showLat) return; + try { latBox.replaceChildren(latencyChart((await api('GET', '/peers/' + id + '/latency')).points)); } catch (e) { latBox.replaceChildren(h('p', { class: 'err-text' }, e.message)); } + } + const latText = () => { + const l = p.stats.latency; + if (p.latencyCheck === 'off' && !l) return h('span', { class: 'muted' }, 'Check off'); + if (!l) return h('span', { class: 'muted' }, p.latencyCheck === 'active' ? 'Not measured yet · pinged only while the device sends traffic' : 'Not measured yet'); + const when = latStale(l) || p.latencyCheck === 'off' ? h('span', { class: 'muted' }, ' · measured ' + ago(l.at)) : null; + if (l.ms == null) return [h('span', null, 'No ping reply'), when]; + return [h('span', { class: 'mono' }, fmtMs(l.ms)), h('span', { class: 'muted' }, ' median · ' + fmtMs(l.min) + '–' + fmtMs(l.max) + ' · ' + l.loss + ' % loss'), when]; + }; + async function drawTraffic() { pills.replaceChildren(...[['24h', '24 h'], ['7d', '7 days'], ['30d', '30 days']].map(([k, t]) => h('button', { type: 'button', class: range === k ? 'pill on' : 'pill', 'aria-pressed': String(range === k), onClick: () => { range = k; drawTraffic(); } }, t))); @@ -829,11 +950,19 @@ const dns = choice({ id: 'dns', label: 'DNS', ...ch.dns, placeholder: '9.9.9.9' }); const allowed = choice({ id: 'ai', label: 'Client AllowedIPs', ...ch.allowed, hint: 'Used in the client config', placeholder: '0.0.0.0/0, ::/0' }); const ka = choice({ id: 'ka', label: 'Persistent keepalive', ...ch.ka, placeholder: 'Seconds' }); + const LAT_HINTS = { + off: 'The server never pings this peer.', + active: 'Pings every 30 s while the device sends traffic. An idle device is left alone.', + always: 'Pings every 30 s, even when idle. This keeps the tunnel up, so the peer always shows as Online. Best for servers and routers.', + }; + const latHint = h('span', { class: 'hint' }, LAT_HINTS[p.latencyCheck]); + const lc = h('select', { id: 'lc', onChange: () => { latHint.textContent = LAT_HINTS[lc.value]; } }, + [['off', 'Off'], ['active', 'While the device is active'], ['always', 'Always']].map(([v, t]) => h('option', { value: v, selected: v === p.latencyCheck }, t))); const save = async (e) => { e.preventDefault(); err.textContent = ''; try { - const body = { name: name.value, note: note.value, ipv4: ip.value.trim(), ...overrides(dns, allowed, ka, srv) }; + const body = { name: name.value, note: note.value, ipv4: ip.value.trim(), latencyCheck: lc.value, ...overrides(dns, allowed, ka, srv) }; const res = await api('PATCH', '/peers/' + id, body); applied(res, 'Saved'); render(); @@ -856,6 +985,13 @@ h('div', { class: 'cardhead' }, h('div', null, h('h2', { id: 'traffic' }, 'Traffic'), totals), pills), traffic), + showLat ? h('section', { class: 'card', 'aria-labelledby': 'lat' }, + h('div', { class: 'cardhead' }, h('h2', { id: 'lat' }, 'Latency · last 24 hours'), + h('div', { class: 'legend-row', style: { marginTop: '0' } }, + h('span', null, h('span', { class: 'key down' }), 'Median'), + h('span', null, h('span', { class: 'key band' }), 'Min–max'))), + latBox) : null, + h('div', { class: 'cols' }, h('section', { class: 'card' }, h('h2', null, 'Connection'), @@ -864,6 +1000,7 @@ h('dt', null, 'Endpoint'), h('dd', { class: 'mono' }, p.stats.endpoint || '–'), h('dt', null, 'Location'), h('dd', null, fmtLocation(p.stats.location) || '–'), h('dt', null, 'Latest handshake'), h('dd', null, ago(p.stats.lastHandshake)), + h('dt', null, 'Latency'), h('dd', null, latText()), h('dt', null, 'Public key'), h('dd', { class: 'mono' }, p.publicKey || '–'), h('dt', null, 'Preshared key'), h('dd', null, p.hasPresharedKey ? 'Set' : 'None'), h('dt', null, 'All-time traffic'), h('dd', null, 'Download ' + fmtBytes(p.stats.downTotal) + ' · Upload ' + fmtBytes(p.stats.upTotal)))), @@ -894,17 +1031,18 @@ h('form', { class: 'card', onSubmit: save }, h('h2', null, 'Settings'), - h('p', { class: 'lead' }, 'Name and address changes apply immediately. DNS, AllowedIPs and keepalive are part of the client config: they take effect after the config is issued again.'), + h('p', { class: 'lead' }, 'Name, address and latency check changes apply immediately. DNS, AllowedIPs and keepalive are part of the client config: they take effect after the config is issued again.'), h('div', { class: 'grid' }, h('div', { class: 'field' }, h('label', { htmlFor: 'n' }, 'Name'), name, h('span', { class: 'hint' }, 'Letters, numbers, . _ @ - · max 32')), h('div', { class: 'field' }, h('label', { htmlFor: 'no' }, 'Note'), note), h('div', { class: 'field' }, h('label', { htmlFor: 'ip' }, 'IPv4 address'), ip, h('span', { class: 'hint' }, 'Changing it requires a new client config')), - ka.el, allowed.el, dns.el), + ka.el, allowed.el, dns.el, + h('div', { class: 'field' }, h('label', { htmlFor: 'lc' }, 'Latency check'), lc, latHint)), err, h('div', { class: 'formfoot' }, h('button', { type: 'button', class: 'btn', onClick: render }, 'Cancel'), h('button', { type: 'submit', class: 'btn primary' }, 'Save changes')))); - await drawTraffic(); + await Promise.all([drawTraffic(), drawLatency()]); } // ---------- server ---------- diff --git a/config.go b/config.go index 99ef24e..2344731 100644 --- a/config.go +++ b/config.go @@ -106,16 +106,19 @@ type ClientDefaults struct { // Peer is one client. Its private key is never stored: it is shown once when // the config is issued. type Peer struct { - ID string `json:"id"` - Name string `json:"name"` - Note string `json:"note"` - Enabled bool `json:"enabled"` - PublicKey string `json:"publicKey"` - PresharedKey string `json:"presharedKey,omitempty"` - IPv4 string `json:"ipv4"` - DNS []string `json:"dns,omitempty"` // nil = server default - AllowedIPs []string `json:"allowedIPs,omitempty"` // nil = server default - Keepalive *int `json:"keepalive,omitempty"` // nil = server default + ID string `json:"id"` + Name string `json:"name"` + Note string `json:"note"` + Enabled bool `json:"enabled"` + PublicKey string `json:"publicKey"` + PresharedKey string `json:"presharedKey,omitempty"` + IPv4 string `json:"ipv4"` + DNS []string `json:"dns,omitempty"` // nil = server default + AllowedIPs []string `json:"allowedIPs,omitempty"` // nil = server default + Keepalive *int `json:"keepalive,omitempty"` // nil = server default + // LatencyCheck says when the server pings the peer through the tunnel: + // "" (off), "active" (while the device sends traffic) or "always". + LatencyCheck string `json:"latencyCheck,omitempty"` Created time.Time `json:"created"` ConfigIssued *time.Time `json:"configIssued,omitempty"` // Setup is a pending one-time setup link. A peer created with a link has @@ -366,6 +369,9 @@ func (c *Config) validate() error { if p.Keepalive != nil && (*p.Keepalive < 0 || *p.Keepalive > 3600) { return fmt.Errorf("peer %q: keepalive must be 0–3600 seconds", p.Name) } + if !validLatencyCheck(p.LatencyCheck) { + return fmt.Errorf("peer %q: latency check must be off, active or always", p.Name) + } } return nil } diff --git a/go.mod b/go.mod index fd1d0d0..0b5b4e2 100644 --- a/go.mod +++ b/go.mod @@ -4,9 +4,11 @@ go 1.27.1 require ( github.com/google/nftables v0.3.0 + github.com/oschwald/maxminddb-golang v1.13.1 github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e github.com/vishvananda/netlink v1.3.1 golang.org/x/crypto v0.57.0 + golang.org/x/net v0.58.0 golang.org/x/sys v0.48.0 golang.org/x/term v0.46.0 golang.zx2c4.com/wireguard/wgctrl v0.0.0-20241231184526-a9ab2273dd10 @@ -17,9 +19,7 @@ require ( github.com/mdlayher/genetlink v1.3.2 // indirect github.com/mdlayher/netlink v1.7.3-0.20250113171957-fbb4dce95f42 // indirect github.com/mdlayher/socket v0.5.1 // indirect - github.com/oschwald/maxminddb-golang v1.13.1 // indirect github.com/vishvananda/netns v0.0.5 // indirect - golang.org/x/net v0.58.0 // indirect golang.org/x/sync v0.23.0 // indirect golang.org/x/text v0.42.0 // indirect golang.zx2c4.com/wireguard v0.0.0-20231211153847-12269c276173 // indirect diff --git a/go.sum b/go.sum index 87cd4f6..89b21ca 100644 --- a/go.sum +++ b/go.sum @@ -1,3 +1,5 @@ +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/nftables v0.3.0 h1:bkyZ0cbpVeMHXOrtlFc8ISmfVqq5gPJukoYieyVmITg= @@ -12,8 +14,12 @@ github.com/mikioh/ipaddr v0.0.0-20190404000644-d465c8ab6721 h1:RlZweED6sbSArvlE9 github.com/mikioh/ipaddr v0.0.0-20190404000644-d465c8ab6721/go.mod h1:Ickgr2WtCLZ2MDGd4Gr0geeCH5HybhRJbonOgQpvSxc= github.com/oschwald/maxminddb-golang v1.13.1 h1:G3wwjdN9JmIK2o/ermkHM+98oX5fS+k5MbwsmL4MRQE= github.com/oschwald/maxminddb-golang v1.13.1/go.mod h1:K4pgV9N/GcK694KSTmVSDTODk4IsCNThNdTmnaBZ/F8= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e h1:MRM5ITcdelLK2j1vwZ3Je0FKVCfqOLp5zO6trqMLYs0= github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e/go.mod h1:XV66xRDqSt+GTGFMVlhk3ULuV0y9ZmzeVGR4mloJI3M= +github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsTg= +github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= github.com/vishvananda/netlink v1.3.1 h1:3AEMt62VKqz90r0tmNhog0r/PpWKmrEShJU0wJW6bV0= github.com/vishvananda/netlink v1.3.1/go.mod h1:ARtKouGSTGchR8aMwmkzC0qiNPrrWO5JS/XMVl45+b4= github.com/vishvananda/netns v0.0.5 h1:DfiHV+j8bA32MFM7bfEunvT8IAqQ/NzSJHtcmW5zdEY= @@ -36,3 +42,5 @@ golang.zx2c4.com/wireguard v0.0.0-20231211153847-12269c276173 h1:/jFs0duh4rdb8uI golang.zx2c4.com/wireguard v0.0.0-20231211153847-12269c276173/go.mod h1:tkCQ4FQXmpAgYVh++1cq16/dH4QJtmvpRv19DWGAHSA= golang.zx2c4.com/wireguard/wgctrl v0.0.0-20241231184526-a9ab2273dd10 h1:3GDAcqdIg1ozBNLgPy4SLT84nfcBjr6rhGtXYtrkWLU= golang.zx2c4.com/wireguard/wgctrl v0.0.0-20241231184526-a9ab2273dd10/go.mod h1:T97yPqesLiNrOYxkwmhMI0ZIlJDm+p0PMR8eRVeR5tQ= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/kernel.go b/kernel.go index ed62cc8..c3dbd32 100644 --- a/kernel.go +++ b/kernel.go @@ -2,6 +2,7 @@ package main import ( "log/slog" + "net/netip" "sync" "time" ) @@ -30,6 +31,9 @@ type Kernel interface { Sample(iface string) ([]PeerSample, error) Checks(c *Config) []Check Uplink(c *Config, v6 bool) string + // Ping sends one echo request to each address and returns the round-trip + // times of the replies that came within timeout. + Ping(dsts []netip.Addr, timeout time.Duration) (map[netip.Addr]time.Duration, error) Down(c *Config) error Close() error } diff --git a/kernel_other.go b/kernel_other.go index 01946f1..a5cb2e3 100644 --- a/kernel_other.go +++ b/kernel_other.go @@ -6,6 +6,7 @@ import ( "fmt" "log/slog" "math/rand/v2" + "net/netip" "sync" "time" ) @@ -89,3 +90,18 @@ func (k *simKernel) Uplink(c *Config, v6 bool) string { } func (k *simKernel) Down(*Config) error { return nil } + +// Ping invents round trips from the address, so each peer keeps its own +// typical latency. Every seventh address never answers, like a Windows PC. +func (k *simKernel) Ping(dsts []netip.Addr, _ time.Duration) (map[netip.Addr]time.Duration, error) { + out := map[netip.Addr]time.Duration{} + for _, d := range dsts { + last := int(d.As4()[3]) + if last%7 == 4 { + continue + } + base := 8 + last*37%180 + out[d] = time.Duration(base*1000+rand.IntN(base*300+1)) * time.Microsecond + } + return out, nil +} diff --git a/latency.go b/latency.go new file mode 100644 index 0000000..601335b --- /dev/null +++ b/latency.go @@ -0,0 +1,305 @@ +package main + +import ( + "log/slog" + "net/netip" + "slices" + "time" +) + +// Latency is measured by pinging a peer's tunnel address from the server: +// WireGuard itself reports no round-trip time. A ping to an idle device +// wakes it and starts new handshakes, which would keep it "online" for good, +// so the "active" mode pings only while the device sends traffic of its own. + +const ( + latencyOff = "" + latencyActive = "active" + latencyAlways = "always" + + pingInterval = 30 * time.Second + pingTimeout = 2 * time.Second + latencyStep = 5 * time.Minute // one point of the latency history + latencyKeep = 24 * time.Hour + latencyWindow = 5 * time.Minute // the current value is the median of this window + activeWindow = 2 * time.Minute + // activeRxBytes is what a peer must send in one sample interval to count + // as active. Ping replies (~128 bytes) and keepalives stay well below it. + activeRxBytes = 1024 +) + +func validLatencyCheck(m string) bool { + return m == latencyOff || m == latencyActive || m == latencyAlways +} + +// latBucket is one step of a peer's latency history, in milliseconds. The +// single round trips are kept only while the step is open; a closed step +// keeps its summary. +type latBucket struct { + T int64 `json:"t"` + Sent int `json:"sent"` + Lost int `json:"lost"` + Min float64 `json:"min,omitempty"` + Med float64 `json:"med,omitempty"` + Max float64 `json:"max,omitempty"` + RTTs []float64 `json:"rtts,omitempty"` +} + +// summary returns the bucket with Min, Med and Max filled in and no RTTs. +func (b latBucket) summary() latBucket { + if len(b.RTTs) > 0 { + b.Min, b.Med, b.Max = spread(b.RTTs) + } + b.RTTs = nil + return b +} + +// spread returns the minimum, median and maximum of v (not empty). +func spread(v []float64) (lo, med, hi float64) { + s := slices.Clone(v) + slices.Sort(s) + n := len(s) + med = s[n/2] + if n%2 == 0 { + med = (s[n/2-1] + s[n/2]) / 2 + } + return s[0], med, s[n-1] +} + +type latSample struct { + t time.Time + ms float64 + ok bool +} + +// latLive is the part of the latency state that is not saved. +type latLive struct { + recent []latSample // pings of the last latencyWindow + lastActive time.Time // last sample interval with traffic from the peer +} + +func (s *Stats) liveFor(id string) *latLive { + l := s.live[id] + if l == nil { + l = &latLive{} + s.live[id] = l + } + return l +} + +// record adds one ping result. s.mu must be held. +func (s *Stats) record(id string, now time.Time, rtt time.Duration, ok bool) { + ps := s.data.Peers[id] + if ps == nil { + ps = &peerStats{} + s.data.Peers[id] = ps + } + ms := float64(rtt.Microseconds()) / 1000 + t := now.Truncate(latencyStep).Unix() + if n := len(ps.Latency); n == 0 || ps.Latency[n-1].T != t { + if n > 0 { + ps.Latency[n-1] = ps.Latency[n-1].summary() + } + ps.Latency = append(ps.Latency, latBucket{T: t}) + } + b := &ps.Latency[len(ps.Latency)-1] + b.Sent++ + if ok { + b.RTTs = append(b.RTTs, ms) + } else { + b.Lost++ + } + l := s.liveFor(id) + i := 0 + for i < len(l.recent) && now.Sub(l.recent[i].t) >= latencyWindow { + i++ + } + l.recent = append(l.recent[i:], latSample{t: now, ms: ms, ok: ok}) + s.dirty = true +} + +// pingTargets returns the tunnel addresses to ping now, by peer ID. +func (s *Stats) pingTargets(c *Config, now time.Time) map[netip.Addr]string { + s.mu.Lock() + defer s.mu.Unlock() + out := map[netip.Addr]string{} + for _, p := range c.Peers { + if !p.Enabled || !p.hasKey() || p.LatencyCheck == latencyOff { + continue + } + ps := s.data.Peers[p.ID] + if ps == nil || ps.LastHandshake.IsZero() { + continue // never connected: the kernel has no endpoint to send to + } + if p.LatencyCheck == latencyActive && now.Sub(s.liveFor(p.ID).lastActive) > activeWindow { + continue + } + if a, err := netip.ParseAddr(p.IPv4); err == nil { + out[a] = p.ID + } + } + return out +} + +func (s *Stats) pingRound() { + targets := s.pingTargets(s.store.Get(), time.Now()) + if len(targets) == 0 { + return + } + dsts := make([]netip.Addr, 0, len(targets)) + for a := range targets { + dsts = append(dsts, a) + } + start := time.Now() + rtts, err := s.kernel.Ping(dsts, pingTimeout) + s.mu.Lock() + defer s.mu.Unlock() + if (err == nil) != (s.pingErr == nil) { + if err != nil { + slog.Warn("latency check failed", "err", err) + } else { + slog.Info("latency check works again") + } + } + s.pingErr = err + if err != nil { + return + } + for a, id := range targets { + rtt, ok := rtts[a] + s.record(id, start, rtt, ok) + } +} + +// RunPings runs the latency checks until stop is closed. +func (s *Stats) RunPings(stop <-chan struct{}) { + t := time.NewTicker(pingInterval) + defer t.Stop() + for { + select { + case <-stop: + return + case <-t.C: + s.pingRound() + } + } +} + +// PingCheck is the health line for the latency check; ok is false when no +// peer has the check turned on. +func (s *Stats) PingCheck(c *Config) (Check, bool) { + on := slices.ContainsFunc(c.Peers, func(p Peer) bool { return p.Enabled && p.LatencyCheck != latencyOff }) + if !on { + return Check{}, false + } + s.mu.Lock() + defer s.mu.Unlock() + if s.pingErr != nil { + return Check{Name: "Latency check", OK: false, Detail: s.pingErr.Error()}, true + } + return Check{Name: "Latency check", OK: true, Detail: "ping through the tunnel works"}, true +} + +// LatencyView is a peer's current latency, in milliseconds. +type LatencyView struct { + MS *float64 `json:"ms"` // median; null = no reply + Min float64 `json:"min"` + Max float64 `json:"max"` + Loss int `json:"loss"` // percent of pings without a reply + At time.Time `json:"at"` // newest ping + Spark []*float64 `json:"spark"` // medians of the last hour, oldest first; null = no reply +} + +// latencyView returns the current latency, or nil when the peer was never +// pinged. s.mu must be held. +func (s *Stats) latencyView(id string, ps *peerStats, now time.Time) *LatencyView { + if len(ps.Latency) == 0 { + return nil + } + v := &LatencyView{} + var rtts []float64 + var sent, lost int + if l := s.live[id]; l != nil && len(l.recent) > 0 { + v.At = l.recent[len(l.recent)-1].t + for _, x := range l.recent { + sent++ + if x.ok { + rtts = append(rtts, x.ms) + } else { + lost++ + } + } + } else { + // Nothing pinged since the start: show the newest saved step. + b := ps.Latency[len(ps.Latency)-1] + v.At = time.Unix(b.T, 0).Add(latencyStep) + if v.At.After(now) { + v.At = now + } + rtts, sent, lost = b.RTTs, b.Sent, b.Lost + if len(rtts) == 0 && b.Sent > b.Lost { + rtts = []float64{b.Min, b.Med, b.Max} + } + } + if len(rtts) > 0 { + var med float64 + v.Min, med, v.Max = spread(rtts) + v.MS = &med + } + if sent > 0 { + v.Loss = lost * 100 / sent + } + for _, b := range s.latencySeries(ps, now, 12) { + if b.Sent > b.Lost { + m := b.Med + v.Spark = append(v.Spark, &m) + } else { + v.Spark = append(v.Spark, nil) + } + } + return v +} + +// latencySeries returns n steps ending with the current one; steps without +// pings have Sent 0. s.mu must be held. +func (s *Stats) latencySeries(ps *peerStats, now time.Time, n int) []latBucket { + last := now.Truncate(latencyStep) + out := make([]latBucket, n) + idx := map[int64]int{} + for i := range out { + t := last.Add(-time.Duration(n-1-i) * latencyStep).Unix() + out[i].T = t + idx[t] = i + } + for _, b := range ps.Latency { + if i, ok := idx[b.T]; ok { + out[i] = b.summary() + } + } + return out +} + +// LatencyHistory returns the latency of the last 24 hours, one point per step. +func (s *Stats) LatencyHistory(id string) []latBucket { + s.mu.Lock() + defer s.mu.Unlock() + ps := s.data.Peers[id] + if ps == nil { + ps = &peerStats{} + } + return s.latencySeries(ps, time.Now(), int(latencyKeep/latencyStep)) +} + +// pruneLatency drops history older than latencyKeep. s.mu must be held. +func pruneLatency(ps *peerStats, now time.Time) bool { + cut := now.Add(-latencyKeep).Truncate(latencyStep).Unix() + i := 0 + for i < len(ps.Latency) && ps.Latency[i].T < cut { + i++ + } + if i == 0 { + return false + } + ps.Latency = append([]latBucket(nil), ps.Latency[i:]...) + return true +} diff --git a/main.go b/main.go index 8e97462..de1afd6 100644 --- a/main.go +++ b/main.go @@ -208,9 +208,10 @@ func run(configPath string) error { } var wg sync.WaitGroup - wg.Add(3) + wg.Add(4) go func() { defer wg.Done(); recon.Run(stop) }() go func() { defer wg.Done(); stats.Run(stop) }() + go func() { defer wg.Done(); stats.RunPings(stop) }() go func() { defer wg.Done(); geo.Run(stop) }() go func() { t := time.NewTicker(10 * time.Minute) diff --git a/main_test.go b/main_test.go index ba61172..c34339c 100644 --- a/main_test.go +++ b/main_test.go @@ -162,6 +162,13 @@ func (k *fakeKernel) Checks(*Config) []Check { return []Check{{"fak func (k *fakeKernel) Uplink(*Config, bool) string { return "eth0" } func (k *fakeKernel) Down(*Config) error { return nil } func (k *fakeKernel) Close() error { return nil } +func (k *fakeKernel) Ping(dsts []netip.Addr, _ time.Duration) (map[netip.Addr]time.Duration, error) { + out := map[netip.Addr]time.Duration{} + for _, d := range dsts { + out[d] = 20 * time.Millisecond + } + return out, nil +} func TestStatsDeltas(t *testing.T) { dir := t.TempDir() @@ -681,3 +688,53 @@ func TestInstallQuestions(t *testing.T) { } } } + +func TestLatency(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() + ip := serverIPv4(netip.MustParsePrefix(store.Get().Server.IPv4)).Next() + if err := store.Update(func(c *Config) error { + c.Peers = append(c.Peers, Peer{ID: "p1", Name: "phone", IPv4: ip.String(), PublicKey: pub, Enabled: true, LatencyCheck: latencyActive}) + return nil + }); err != nil { + t.Fatal(err) + } + if err := store.Update(func(c *Config) error { c.Peers[0].LatencyCheck = "sometimes"; return nil }); err == nil { + t.Fatal("invalid latency check accepted") + } + k := &fakeKernel{} + st, _ := openStats(filepath.Join(dir, "stats.json"), store, k) + now := time.Now() + k.samples = []PeerSample{{PublicKey: pub, RxBytes: 100, LastHandshake: now}} + st.sample() + k.samples[0].RxBytes += 200 // below activeRxBytes: idle, e.g. keepalives + st.sample() + if got := st.pingTargets(store.Get(), now); len(got) != 0 { + t.Fatalf("idle peer pinged: %v", got) + } + k.samples[0].RxBytes += 50_000 + st.sample() + if got := st.pingTargets(store.Get(), time.Now()); got[ip] != "p1" { + t.Fatalf("active peer not pinged: %v", got) + } + + st.mu.Lock() + for _, ms := range []int{30, 10, 20} { + st.record("p1", now, time.Duration(ms)*time.Millisecond, true) + } + st.record("p1", now, 0, false) + st.mu.Unlock() + l := st.Summary("p1").Latency + if l == nil || l.MS == nil || *l.MS != 20 || l.Min != 10 || l.Max != 30 || l.Loss != 25 { + t.Fatalf("latency %+v", l) + } + h := st.LatencyHistory("p1") + if last := h[len(h)-1]; last.Sent != 4 || last.Lost != 1 || last.Med != 20 || last.RTTs != nil { + t.Fatalf("history %+v", last) + } +} diff --git a/ping_linux.go b/ping_linux.go new file mode 100644 index 0000000..980364b --- /dev/null +++ b/ping_linux.go @@ -0,0 +1,90 @@ +//go:build linux + +package main + +import ( + "errors" + "math/rand/v2" + "net" + "net/netip" + "time" + + "golang.org/x/net/icmp" + "golang.org/x/net/ipv4" +) + +// listenICMP opens an unprivileged ICMP socket, which the kernel allows for +// the groups in net.ipv4.ping_group_range (systemd opens it to all groups). +// Running as root, a raw socket works too. +func listenICMP() (c *icmp.PacketConn, raw bool, err error) { + if c, err = icmp.ListenPacket("udp4", "0.0.0.0"); err == nil { + return c, false, nil + } + if c, err2 := icmp.ListenPacket("ip4:icmp", "0.0.0.0"); err2 == nil { + return c, true, nil + } + return nil, false, errors.New("cannot open an ICMP socket (" + err.Error() + "): add the service's group to the sysctl net.ipv4.ping_group_range") +} + +func (k *linuxKernel) Ping(dsts []netip.Addr, timeout time.Duration) (map[netip.Addr]time.Duration, error) { + conn, raw, err := listenICMP() + if err != nil { + return nil, err + } + defer conn.Close() + // On an unprivileged socket the kernel sets the ID and delivers only the + // socket's own replies; on a raw socket the ID tells them apart. + id := rand.IntN(0xffff) + 1 + bySeq := map[int]netip.Addr{} + sent := map[netip.Addr]time.Time{} + for i, d := range dsts { + seq := i + 1 + b, err := (&icmp.Message{Type: ipv4.ICMPTypeEcho, Body: &icmp.Echo{ID: id, Seq: seq, Data: []byte(appName)}}).Marshal(nil) + if err != nil { + return nil, err + } + var to net.Addr = &net.UDPAddr{IP: d.AsSlice()} + if raw { + to = &net.IPAddr{IP: d.AsSlice()} + } + sent[d] = time.Now() + if _, err := conn.WriteTo(b, to); err != nil { + delete(sent, d) // e.g. the peer has no endpoint: counts as lost + continue + } + bySeq[seq] = d + } + out := map[netip.Addr]time.Duration{} + _ = conn.SetReadDeadline(time.Now().Add(timeout)) + buf := make([]byte, 1500) + for len(out) < len(sent) { + n, from, err := conn.ReadFrom(buf) + if err != nil { + break // deadline reached + } + now := time.Now() + m, err := icmp.ParseMessage(1, buf[:n]) + if err != nil || m.Type != ipv4.ICMPTypeEchoReply { + continue + } + e, ok := m.Body.(*icmp.Echo) + if !ok || (raw && e.ID != id) { + continue + } + var src netip.Addr + switch a := from.(type) { + case *net.UDPAddr: + src, _ = netip.AddrFromSlice(a.IP) + case *net.IPAddr: + src, _ = netip.AddrFromSlice(a.IP) + } + src = src.Unmap() + if d, ok := bySeq[e.Seq]; !ok || d != src { + continue + } + if _, dup := out[src]; !dup { + out[src] = now.Sub(sent[src]) + } + } + return out, nil +} diff --git a/stats.go b/stats.go index d8a2217..c22a40c 100644 --- a/stats.go +++ b/stats.go @@ -37,6 +37,7 @@ type peerStats struct { Hourly []bucket `json:"hourly"` Daily []bucket `json:"daily"` Sessions []connSession `json:"sessions,omitempty"` + Latency []latBucket `json:"latency,omitempty"` } // session is one stretch of a peer being online from one address. When the @@ -114,10 +115,13 @@ type Stats struct { store *Store kernel Kernel geo *Geo // nil: no country and network lookups + + live map[string]*latLive // by peer ID + pingErr error // of the last latency check } func openStats(path string, store *Store, k Kernel) (*Stats, error) { - s := &Stats{path: path, store: store, kernel: k, data: statsFile{Version: 1, Peers: map[string]*peerStats{}}} + s := &Stats{path: path, store: store, kernel: k, live: map[string]*latLive{}, data: statsFile{Version: 1, Peers: map[string]*peerStats{}}} b, err := os.ReadFile(path) if errors.Is(err, os.ErrNotExist) { return s, nil @@ -182,7 +186,7 @@ func (s *Stats) prune(c StatsConfig, now time.Time) { if i > 0 { ps.Sessions = append([]connSession(nil), ps.Sessions[i:]...) } - if len(ps.Hourly) != h || len(ps.Daily) != dl || len(ps.Sessions) != sl { + if pruneLatency(ps, now) || len(ps.Hourly) != h || len(ps.Daily) != dl || len(ps.Sessions) != sl { s.dirty = true } } @@ -223,6 +227,9 @@ func (s *Stats) sample() { dRx, dTx = smp.RxBytes, smp.TxBytes } ps.LastRx, ps.LastTx = smp.RxBytes, smp.TxBytes + if dRx >= activeRxBytes { + s.liveFor(id).lastActive = now + } if dRx > 0 || dTx > 0 { ps.TotalRx += dRx ps.TotalTx += dTx @@ -241,6 +248,7 @@ func (s *Stats) sample() { for id, ps := range s.data.Peers { if !exists[id] { delete(s.data.Peers, id) + delete(s.live, id) s.dirty = true continue } @@ -349,16 +357,17 @@ func (s *Stats) series(ids []string, rng string) []Point { // PeerSummary is the live state shown in peer lists. type PeerSummary struct { - Online bool `json:"online"` - LastHandshake *time.Time `json:"lastHandshake"` - Endpoint string `json:"endpoint"` - Down24h int64 `json:"down24h"` - Up24h int64 `json:"up24h"` - Down30d int64 `json:"down30d"` - Up30d int64 `json:"up30d"` - DownTotal int64 `json:"downTotal"` - UpTotal int64 `json:"upTotal"` - Location *GeoInfo `json:"location"` // of the current or last endpoint + Online bool `json:"online"` + LastHandshake *time.Time `json:"lastHandshake"` + Endpoint string `json:"endpoint"` + Down24h int64 `json:"down24h"` + Up24h int64 `json:"up24h"` + Down30d int64 `json:"down30d"` + Up30d int64 `json:"up30d"` + DownTotal int64 `json:"downTotal"` + UpTotal int64 `json:"upTotal"` + Location *GeoInfo `json:"location"` // of the current or last endpoint + Latency *LatencyView `json:"latency"` // null = never pinged } func sumPoints(pts []Point) (down, up int64) { @@ -386,6 +395,7 @@ func (s *Stats) Summary(id string) PeerSummary { if cur := ps.openSession(); cur != nil && cur.Geo != nil && hostOf(cur.Endpoint) == hostOf(ps.Endpoint) { out.Location = cur.Geo } + out.Latency = s.latencyView(id, ps, time.Now()) } if out.Location == nil && out.Endpoint != "" { out.Location = s.geo.Lookup(out.Endpoint)