Add connection history per peer with country and network lookup

- The stats sampler records sessions per peer: start, end, address and
  traffic. A session ends when the peer goes quiet or is disabled; a new
  one starts when the device changes networks. Stored in stats.json and
  kept as long as the daily traffic history (max 1000 per peer).
- Country and network operator come from the free DB-IP Lite databases
  (CC BY 4.0), downloaded monthly and looked up locally, so peer
  addresses never leave the server. Settings → Data retention can switch
  this off, which deletes the databases.
- API: GET /peers/{id}/sessions; peer stats include the current location;
  settings include the database status.
- Web UI and iOS app: connection history card, location line, country
  code in the peer list (web), switch in data retention.

Claude-Session: https://claude.ai/code/session_01RAnLbyQZ5ZTA7KqwXP98nw
This commit is contained in:
Daniel Redetzke
2026-10-03 19:21:10 +03:00
parent f31bb360c9
commit 37ab26b415
15 changed files with 707 additions and 21 deletions
+4 -2
View File
@@ -13,6 +13,7 @@ manages peers (add, change, disable, remove) and records traffic per peer.
as `wg syncconf`, so connected peers stay connected. as `wg syncconf`, so connected peers stay connected.
- **Logs:** written to `GHOSTWIRE.jsonl`, rotated at 10 MB with 5 old files kept by default (Settings → Data retention). - **Logs:** written to `GHOSTWIRE.jsonl`, rotated at 10 MB with 5 old files kept by default (Settings → Data retention).
- **Traffic history:** kept in `stats.json`: hourly for 48 h and daily for 400 days by default (Settings → Data retention). - **Traffic history:** kept in `stats.json`: hourly for 48 h and daily for 400 days by default (Settings → Data retention).
- **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 [DB-IP Lite](https://db-ip.com) databases (CC BY 4.0). GHOSTWIRE downloads them monthly (about 20 MB) and looks addresses up locally, so peer addresses never leave the server. You can switch this off under Settings → Data retention.
- **Client private keys are never stored.** A config is shown once, as a - **Client private keys are never stored.** A config is shown once, as a
download or QR code. "Issue new config" makes new keys. download or QR code. "Issue new config" makes new keys.
@@ -109,7 +110,8 @@ After editing `config.json` by hand, run `sudo systemctl reload ghostwire`.
|---|---| |---|---|
| `GHOSTWIRE` | the program | | `GHOSTWIRE` | the program |
| `config.json` | all settings, server key, peers, token hashes (0600) | | `config.json` | all settings, server key, peers, token hashes (0600) |
| `stats.json` | traffic history per peer | | `stats.json` | traffic and connection history per peer |
| `geo-country.mmdb`, `geo-asn.mmdb` | DB-IP Lite databases for country and network lookups |
| `GHOSTWIRE.jsonl` | log, one JSON object per line. Changes carry `"audit":true` | | `GHOSTWIRE.jsonl` | log, one JSON object per line. Changes carry `"audit":true` |
| `acme/`, `tls/` | certificates | | `acme/`, `tls/` | certificates |
@@ -128,7 +130,7 @@ GET /server PATCH /server POST /server/rotate-key GET /server
GET /peers POST /peers (returns the config and QR once) GET /peers POST /peers (returns the config and QR once)
GET /peers/{id} PATCH /peers/{id} DELETE /peers/{id} GET /peers/{id} PATCH /peers/{id} DELETE /peers/{id}
POST /peers/{id}/enable | /disable | /issue-config POST /peers/{id}/enable | /disable | /issue-config
GET /peers/{id}/stats?range=… GET /peers/{id}/stats?range=… GET /peers/{id}/sessions?limit=100
GET /settings PATCH /settings POST /restart GET /settings PATCH /settings POST /restart
GET /logs?level=&limit=&audit=1 GET /logs/download GET /logs?level=&limit=&audit=1 GET /logs/download
admin: GET|POST /tokens · DELETE /tokens/{id} · GET /backup · POST /restore admin: GET|POST /tokens · DELETE /tokens/{id} · GET /backup · POST /restore
+27
View File
@@ -26,6 +26,7 @@ type App struct {
tls *webTLS tls *webTLS
logPath string logPath string
logw *rotatingWriter // nil in tests logw *rotatingWriter // nil in tests
geo *Geo // nil in tests
started time.Time started time.Time
shutdown func() // graceful stop; systemd restarts the service shutdown func() // graceful stop; systemd restarts the service
} }
@@ -124,6 +125,7 @@ func (a *App) routes() http.Handler {
g("POST /api/v1/peers/{id}/disable", a.setEnabled(false)) g("POST /api/v1/peers/{id}/disable", a.setEnabled(false))
g("POST /api/v1/peers/{id}/issue-config", a.issueConfig) 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}/stats", a.peerStats)
g("GET /api/v1/peers/{id}/sessions", a.peerSessions)
// Full-access tokens (the iOS app) may change app settings and read logs. // Full-access tokens (the iOS app) may change app settings and read logs.
// Password, tokens and backups stay with the admin account. // Password, tokens and backups stay with the admin account.
@@ -806,6 +808,7 @@ func (a *App) getSettings(w http.ResponseWriter, r *http.Request) {
"web": cfg.Web, "web": cfg.Web,
"log": cfg.Log, "log": cfg.Log,
"stats": cfg.Stats, "stats": cfg.Stats,
"geo": a.geoStatus(),
"adminUsername": cfg.Admin.Username, "adminUsername": cfg.Admin.Username,
"fingerprint": a.tls.Fingerprint(), "fingerprint": a.tls.Fingerprint(),
"logPath": a.logPath, "logPath": a.logPath,
@@ -980,4 +983,28 @@ func (a *App) applyRuntime(c *Config) {
if a.logw != nil { if a.logw != nil {
a.logw.SetLimits(c.Log.MaxSizeMB, c.Log.MaxFiles) a.logw.SetLimits(c.Log.MaxSizeMB, c.Log.MaxFiles)
} }
a.geo.SetEnabled(c.Stats.geoEnabled())
}
func (a *App) geoStatus() GeoStatus {
if a.geo == nil {
return GeoStatus{}
}
return a.geo.Status()
}
// peerSessions returns the connection history, newest first.
func (a *App) peerSessions(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
}
limit := 100
if _, err := fmt.Sscan(r.URL.Query().Get("limit"), &limit); err != nil || limit < 1 || limit > 1000 {
limit = 100
}
writeJSON(w, http.StatusOK, map[string]any{
"sessions": a.stats.Sessions(r.PathValue("id"), limit),
"attribution": "IP geolocation by DB-IP (https://db-ip.com), CC BY 4.0",
})
} }
+1
View File
@@ -94,6 +94,7 @@ h1 { margin: 0; font-size: 26px; font-weight: 600; letter-spacing: -0.01em; over
.pills { display: flex; flex-wrap: wrap; gap: 6px; } .pills { display: flex; flex-wrap: wrap; gap: 6px; }
/* status */ /* status */
.cc { display: inline-block; margin-left: 6px; padding: 0 5px; border: 1px solid var(--line); border-radius: 4px; font-size: 11px; color: var(--ink-2); vertical-align: 1px; }
.badge { display: inline-flex; align-items: center; gap: 6px; font-size: 12px; font-weight: 500; padding: 3px 10px; border-radius: 999px; background: #efefeb; color: #3d3e42; white-space: nowrap; } .badge { display: inline-flex; align-items: center; gap: 6px; font-size: 12px; font-weight: 500; padding: 3px 10px; border-radius: 999px; background: #efefeb; color: #3d3e42; white-space: nowrap; }
.dot { width: 8px; height: 8px; border-radius: 50%; background: #9a9b97; display: inline-block; flex: none; } .dot { width: 8px; height: 8px; border-radius: 50%; background: #9a9b97; display: inline-block; flex: none; }
.dot.ok { background: var(--good); } .dot.ok { background: var(--good); }
+46 -4
View File
@@ -124,6 +124,23 @@
return ts + ' ' + String(l.level).padEnd(5) + ' ' + l.msg + (rest ? ' ' + rest : ''); return ts + ' ' + String(l.level).padEnd(5) + ' ' + l.msg + (rest ? ' ' + rest : '');
} }
// "Germany · Deutsche Telekom AG", "Local network" or "".
function fmtLocation(g) {
if (!g) return '';
return [g.countryName || g.country, g.network].filter(Boolean).join(' · ');
}
function fmtDuration(sec) {
if (sec < 60) return 'under 1 min';
const m = Math.round(sec / 60);
if (m < 60) return m + ' min';
const h = Math.floor(m / 60);
if (h < 48) return h + ' h ' + (m % 60) + ' min';
return Math.round(h / 24) + ' days';
}
const fmtStamp = (iso) => new Date(iso).toLocaleString(undefined, { day: 'numeric', month: 'short', hour: '2-digit', minute: '2-digit' });
const badge = (st) => h('span', { class: 'badge' }, h('span', { class: st.dot }), st.label); const badge = (st) => h('span', { class: 'badge' }, h('span', { class: st.dot }), st.label);
// ---------- API ---------- // ---------- API ----------
@@ -503,7 +520,8 @@
h('td', null, h('a', { href: '#/peers/' + p.id }, h('strong', null, p.name)), p.note ? h('div', { class: 'note' }, p.note) : null), h('td', null, h('a', { href: '#/peers/' + p.id }, h('strong', null, p.name)), p.note ? h('div', { class: 'note' }, p.note) : null),
h('td', { class: 'mono' }, p.ipv4), h('td', { class: 'mono' }, p.ipv4),
h('td', null, badge(peerState(p))), h('td', null, badge(peerState(p))),
h('td', { class: 'mono muted' }, p.stats.endpoint || '–'), 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),
h('td', { class: 'num' }, fmtBytes(p.stats.down30d)), h('td', { class: 'num' }, fmtBytes(p.stats.down30d)),
h('td', { class: 'num' }, fmtBytes(p.stats.up30d)), 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) })), 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) })),
@@ -663,7 +681,8 @@
// ---------- peer detail ---------- // ---------- peer detail ----------
async function viewPeer(wrap, id) { async function viewPeer(wrap, id) {
const [p, srv] = await Promise.all([api('GET', '/peers/' + id), api('GET', '/server')]); const [p, srv, sess] = await Promise.all([api('GET', '/peers/' + id), api('GET', '/server'), api('GET', '/peers/' + id + '/sessions?limit=50')]);
const sessions = sess.sessions;
let range = '7d'; let range = '7d';
const st = peerState(p); const st = peerState(p);
const traffic = h('div'); const traffic = h('div');
@@ -757,6 +776,7 @@
h('dl', { class: 'kv' }, h('dl', { class: 'kv' },
h('dt', null, 'Tunnel address'), h('dd', { class: 'mono' }, p.ipv4 + '/32', p.ipv6 ? [h('br'), p.ipv6 + '/128'] : null), h('dt', null, 'Tunnel address'), h('dd', { class: 'mono' }, p.ipv4 + '/32', p.ipv6 ? [h('br'), p.ipv6 + '/128'] : null),
h('dt', null, 'Endpoint'), h('dd', { class: 'mono' }, p.stats.endpoint || '–'), 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, 'Latest handshake'), h('dd', null, ago(p.stats.lastHandshake)),
h('dt', null, 'Public key'), h('dd', { class: 'mono' }, p.publicKey), 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, 'Preshared key'), h('dd', null, p.hasPresharedKey ? 'Set' : 'None'),
@@ -769,6 +789,24 @@
h('button', { type: 'button', class: 'btn', onClick: pasteKey }, 'Use a key from the device…')), h('button', { type: 'button', class: 'btn', onClick: pasteKey }, 'Use a key from the device…')),
h('p', { class: 'hint', style: { margin: '12px 0 0' } }, p.configIssued ? 'Last issued ' + fmtDate(p.configIssued) + '.' : 'Created with a key from the device.'))), h('p', { class: 'hint', style: { margin: '12px 0 0' } }, p.configIssued ? 'Last issued ' + fmtDate(p.configIssued) + '.' : 'Created with a key from the device.'))),
h('section', { class: 'card flush', 'aria-labelledby': 'hist' },
h('div', { class: 'cardhead' }, h('h2', { id: 'hist' }, 'Connection history'),
h('span', { class: 'hint' }, 'Newest first · a new row starts when the device changes networks')),
sessions.length ? h('div', { class: 'tbl' }, h('table', null,
h('thead', null, h('tr', null, h('th', null, 'Started'), h('th', null, 'Duration'), h('th', null, 'From'), h('th', null, 'Address'),
h('th', { class: 'num' }, 'Download'), h('th', { class: 'num' }, 'Upload'))),
h('tbody', null, sessions.map((se) => h('tr', null,
h('td', null, fmtStamp(se.start)),
h('td', null, se.open ? [h('span', { class: 'badge' }, h('span', { class: 'dot ok' }), 'Online now'), ' ', fmtDuration(se.seconds)] : fmtDuration(se.seconds)),
h('td', null, fmtLocation(se.geo) || h('span', { class: 'muted' }, 'Unknown')),
h('td', { class: 'mono muted' }, se.ip),
h('td', { class: 'num' }, fmtBytes(se.down)),
h('td', { class: 'num' }, fmtBytes(se.up)))))))
: h('p', { class: 'empty' }, 'No connections recorded yet.'),
h('p', { class: 'hint', style: { margin: '4px 12px 12px' } }, 'Country and network: ',
h('a', { href: 'https://db-ip.com', target: '_blank', rel: 'noopener' }, 'IP Geolocation by DB-IP'),
'. Kept as long as the daily traffic history.')),
h('form', { class: 'card', onSubmit: save }, h('form', { class: 'card', onSubmit: save },
h('h2', null, 'Settings'), 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 and address changes apply immediately. DNS, AllowedIPs and keepalive are part of the client config: they take effect after the config is issued again.'),
@@ -1051,6 +1089,8 @@
const logFiles = h('input', { id: 'rf', type: 'number', min: '1', max: '100', value: s.log.maxFiles, inputMode: 'numeric' }); const logFiles = h('input', { id: 'rf', type: 'number', min: '1', max: '100', value: s.log.maxFiles, inputMode: 'numeric' });
const hourly = presetSelect('rh', s.stats.hourlyHours, [[24, '1 day'], [48, '2 days'], [168, '7 days'], [336, '14 days'], [744, '31 days']], 'hours'); const hourly = presetSelect('rh', s.stats.hourlyHours, [[24, '1 day'], [48, '2 days'], [168, '7 days'], [336, '14 days'], [744, '31 days']], 'hours');
const daily = presetSelect('rd', s.stats.dailyDays, [[30, '30 days'], [90, '90 days'], [180, '6 months'], [400, '13 months'], [730, '2 years'], [1825, '5 years'], [3660, '10 years']], 'days'); const daily = presetSelect('rd', s.stats.dailyDays, [[30, '30 days'], [90, '90 days'], [180, '6 months'], [400, '13 months'], [730, '2 years'], [1825, '5 years'], [3660, '10 years']], 'days');
const geo = h('input', { type: 'checkbox', checked: s.stats.geoip !== false });
const geoStatus = s.geo && s.geo.updated ? 'Database from ' + fmtDate(s.geo.updated) + '.' : 'Not downloaded yet.';
const diskHint = h('span', { class: 'hint' }); const diskHint = h('span', { class: 'hint' });
const drawDiskHint = () => { const drawDiskHint = () => {
const mb = Number(logSize.value) * (Number(logFiles.value) + 1); const mb = Number(logSize.value) * (Number(logFiles.value) + 1);
@@ -1069,7 +1109,7 @@
try { try {
await api('PATCH', '/settings', { await api('PATCH', '/settings', {
log: { ...s.log, maxSizeMB: next.maxSizeMB, maxFiles: next.maxFiles }, log: { ...s.log, maxSizeMB: next.maxSizeMB, maxFiles: next.maxFiles },
stats: { hourlyHours: next.hourlyHours, dailyDays: next.dailyDays }, stats: { hourlyHours: next.hourlyHours, dailyDays: next.dailyDays, geoip: geo.checked },
}); });
toast('Retention saved'); toast('Retention saved');
render(); render();
@@ -1144,7 +1184,9 @@
h('div', { class: 'field' }, h('label', { htmlFor: 'rs' }, 'Log file size (MB)'), logSize, h('span', { class: 'hint' }, 'The log starts a new file at this size. 1–1000')), h('div', { class: 'field' }, h('label', { htmlFor: 'rs' }, 'Log file size (MB)'), logSize, h('span', { class: 'hint' }, 'The log starts a new file at this size. 1–1000')),
h('div', { class: 'field' }, h('label', { htmlFor: 'rf' }, 'Old log files kept'), logFiles, diskHint), h('div', { class: 'field' }, h('label', { htmlFor: 'rf' }, 'Old log files kept'), logFiles, diskHint),
h('div', { class: 'field' }, h('label', { htmlFor: 'rh' }, 'Hourly traffic history'), hourly, h('span', { class: 'hint' }, 'Used by the 24-hour charts')), h('div', { class: 'field' }, h('label', { htmlFor: 'rh' }, 'Hourly traffic history'), hourly, h('span', { class: 'hint' }, 'Used by the 24-hour charts')),
h('div', { class: 'field' }, h('label', { htmlFor: 'rd' }, 'Daily traffic history'), daily, h('span', { class: 'hint' }, 'Used by the 7- and 30-day charts. All-time totals are always kept'))), h('div', { class: 'field' }, h('label', { htmlFor: 'rd' }, 'Daily traffic history'), daily, h('span', { class: 'hint' }, 'Used by the 7- and 30-day charts and the connection history. All-time totals are always kept'))),
h('label', { class: 'check section' }, geo, h('span', null, 'Show country and network of peer addresses', h('br'),
h('span', { class: 'hint' }, 'Downloads the free DB-IP Lite databases (about 20 MB) once a month and looks addresses up on this server only. ' + geoStatus))),
retErr, retErr,
h('div', { class: 'formfoot' }, h('button', { type: 'submit', class: 'btn primary' }, 'Save retention'))), h('div', { class: 'formfoot' }, h('button', { type: 'submit', class: 'btn primary' }, 'Save retention'))),
+6 -1
View File
@@ -32,9 +32,14 @@ type Config struct {
// StatsConfig sets how long traffic history is kept in stats.json. // StatsConfig sets how long traffic history is kept in stats.json.
type StatsConfig struct { type StatsConfig struct {
HourlyHours int `json:"hourlyHours"` // hourly buckets, for the 24 h charts HourlyHours int `json:"hourlyHours"` // hourly buckets, for the 24 h charts
DailyDays int `json:"dailyDays"` // daily buckets, for the 7/30/90 day charts DailyDays int `json:"dailyDays"` // daily buckets and connection history
// GeoIP looks up country and network of peer addresses in the DB-IP Lite
// databases, downloaded monthly. Default on.
GeoIP *bool `json:"geoip,omitempty"`
} }
func (c StatsConfig) geoEnabled() bool { return c.GeoIP == nil || *c.GeoIP }
// Limits for the retention settings. // Limits for the retention settings.
const ( const (
minLogSizeMB, maxLogSizeMB = 1, 1000 minLogSizeMB, maxLogSizeMB = 1, 1000
+270
View File
@@ -0,0 +1,270 @@
package main
import (
"compress/gzip"
"context"
"errors"
"fmt"
"io"
"log/slog"
"net"
"net/http"
"net/netip"
"os"
"path/filepath"
"sync"
"sync/atomic"
"time"
"github.com/oschwald/maxminddb-golang"
)
// Country and network (autonomous system) of peer endpoints come from the
// free DB-IP Lite databases (CC BY 4.0, https://db-ip.com). They are
// downloaded once a month and searched locally, so endpoint addresses never
// leave the server.
// GeoInfo describes where an address is.
type GeoInfo struct {
Country string `json:"country,omitempty"` // ISO code, e.g. "DE"
CountryName string `json:"countryName,omitempty"` // e.g. "Germany"
ASN uint `json:"asn,omitempty"`
Network string `json:"network,omitempty"` // operator, e.g. "Deutsche Telekom AG"
}
const (
geoMaxAge = 32 * 24 * time.Hour // DB-IP publishes monthly
geoCheckFreq = 24 * time.Hour
geoBaseURL = "https://download.db-ip.com/free/dbip-%s-lite-%s.mmdb.gz"
)
type Geo struct {
dir string
enabled atomic.Bool
kick chan struct{}
mu sync.RWMutex
country *maxminddb.Reader
asn *maxminddb.Reader
}
func newGeo(dir string, enabled bool) *Geo {
g := &Geo{dir: dir, kick: make(chan struct{}, 1)}
g.enabled.Store(enabled)
if enabled {
g.open()
}
return g
}
func (g *Geo) path(kind string) string { return filepath.Join(g.dir, "geo-"+kind+".mmdb") }
// open (re)loads whichever database files exist.
func (g *Geo) open() {
load := func(kind string) *maxminddb.Reader {
r, err := maxminddb.Open(g.path(kind))
if err != nil {
if !errors.Is(err, os.ErrNotExist) {
slog.Warn("geo database unreadable", "file", g.path(kind), "err", err)
}
return nil
}
return r
}
c, a := load("country"), load("asn")
g.mu.Lock()
old := []*maxminddb.Reader{g.country, g.asn}
g.country, g.asn = c, a
g.mu.Unlock()
for _, r := range old {
if r != nil {
r.Close()
}
}
}
// SetEnabled turns lookups and monthly downloads on or off.
func (g *Geo) SetEnabled(on bool) {
if g == nil || g.enabled.Swap(on) == on {
return
}
select {
case g.kick <- struct{}{}:
default:
}
}
// Lookup returns where ipPort ("203.0.113.7:51820" or a bare address) is.
// Private addresses are reported as the local network.
func (g *Geo) Lookup(ipPort string) *GeoInfo {
host := ipPort
if h, _, err := net.SplitHostPort(ipPort); err == nil {
host = h
}
ip, err := netip.ParseAddr(host)
if err != nil {
return nil
}
ip = ip.Unmap()
if ip.IsPrivate() || ip.IsLoopback() || ip.IsLinkLocalUnicast() || ip.IsUnspecified() || isCGNAT(ip) {
return &GeoInfo{Network: "Local network"}
}
if g == nil || !g.enabled.Load() {
return nil
}
g.mu.RLock()
defer g.mu.RUnlock()
info := GeoInfo{}
if g.country != nil {
var rec struct {
Country struct {
ISOCode string `maxminddb:"iso_code"`
Names map[string]string `maxminddb:"names"`
} `maxminddb:"country"`
}
if g.country.Lookup(net.IP(ip.AsSlice()), &rec) == nil {
info.Country, info.CountryName = rec.Country.ISOCode, rec.Country.Names["en"]
}
}
if g.asn != nil {
var rec struct {
Number uint `maxminddb:"autonomous_system_number"`
Org string `maxminddb:"autonomous_system_organization"`
}
if g.asn.Lookup(net.IP(ip.AsSlice()), &rec) == nil {
info.ASN, info.Network = rec.Number, rec.Org
}
}
if info == (GeoInfo{}) {
return nil
}
return &info
}
var cgnat = netip.MustParsePrefix("100.64.0.0/10")
func isCGNAT(ip netip.Addr) bool { return cgnat.Contains(ip) }
// GeoStatus is shown in the settings.
type GeoStatus struct {
Enabled bool `json:"enabled"`
Updated *time.Time `json:"updated"` // date of the country database file
}
func (g *Geo) Status() GeoStatus {
st := GeoStatus{Enabled: g.enabled.Load()}
if fi, err := os.Stat(g.path("country")); err == nil {
t := fi.ModTime()
st.Updated = &t
}
return st
}
// Run downloads missing or outdated databases daily while enabled, and
// deletes them when the feature is switched off.
func (g *Geo) Run(stop <-chan struct{}) {
t := time.NewTicker(geoCheckFreq)
defer t.Stop()
for {
g.maintain()
select {
case <-stop:
return
case <-t.C:
case <-g.kick:
}
}
}
func (g *Geo) maintain() {
if !g.enabled.Load() {
g.mu.Lock()
for _, r := range []*maxminddb.Reader{g.country, g.asn} {
if r != nil {
r.Close()
}
}
g.country, g.asn = nil, nil
g.mu.Unlock()
for _, kind := range []string{"country", "asn"} {
_ = os.Remove(g.path(kind))
}
return
}
changed := false
for _, kind := range []string{"country", "asn"} {
if fi, err := os.Stat(g.path(kind)); err == nil && time.Since(fi.ModTime()) < geoMaxAge {
continue
}
if err := g.download(kind); err != nil {
slog.Warn("geo database download failed", "db", kind, "err", err)
continue
}
changed = true
slog.Info("geo database updated", "db", kind)
}
if changed {
g.open()
}
}
// download fetches this month's database, or last month's if this month's
// is not published yet, and replaces the local file atomically.
func (g *Geo) download(kind string) error {
now := time.Now().UTC()
var lastErr error
for _, month := range []time.Time{now, now.AddDate(0, -1, 0)} {
url := fmt.Sprintf(geoBaseURL, kind, month.Format("2006-01"))
if lastErr = g.fetch(url, g.path(kind)); lastErr == nil {
return nil
}
}
return lastErr
}
func (g *Geo) fetch(url, dest string) error {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
resp, err := http.DefaultClient.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("%s: HTTP %d", url, resp.StatusCode)
}
zr, err := gzip.NewReader(resp.Body)
if err != nil {
return err
}
tmp, err := os.CreateTemp(filepath.Dir(dest), ".geo-*")
if err != nil {
return err
}
defer os.Remove(tmp.Name())
if _, err := io.Copy(tmp, io.LimitReader(zr, 512<<20)); err != nil {
tmp.Close()
return err
}
if err := tmp.Close(); err != nil {
return err
}
// Refuse a file that is not a readable database.
r, err := maxminddb.Open(tmp.Name())
if err != nil {
return fmt.Errorf("downloaded file is not a valid database: %w", err)
}
r.Close()
return os.Rename(tmp.Name(), dest)
}
func (g *Geo) Close() {
g.mu.Lock()
defer g.mu.Unlock()
for _, r := range []*maxminddb.Reader{g.country, g.asn} {
if r != nil {
r.Close()
}
}
}
+1
View File
@@ -17,6 +17,7 @@ require (
github.com/mdlayher/genetlink v1.3.2 // indirect github.com/mdlayher/genetlink v1.3.2 // indirect
github.com/mdlayher/netlink v1.7.3-0.20250113171957-fbb4dce95f42 // indirect github.com/mdlayher/netlink v1.7.3-0.20250113171957-fbb4dce95f42 // indirect
github.com/mdlayher/socket v0.5.1 // 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 github.com/vishvananda/netns v0.0.5 // indirect
golang.org/x/net v0.58.0 // indirect golang.org/x/net v0.58.0 // indirect
golang.org/x/sync v0.23.0 // indirect golang.org/x/sync v0.23.0 // indirect
+2
View File
@@ -10,6 +10,8 @@ github.com/mdlayher/socket v0.5.1 h1:VZaqt6RkGkt2OE9l3GcC6nZkqD3xKeQLyfleW/uBcos
github.com/mdlayher/socket v0.5.1/go.mod h1:TjPLHI1UgwEv5J1B5q0zTZq12A/6H7nKmtTanQE37IQ= github.com/mdlayher/socket v0.5.1/go.mod h1:TjPLHI1UgwEv5J1B5q0zTZq12A/6H7nKmtTanQE37IQ=
github.com/mikioh/ipaddr v0.0.0-20190404000644-d465c8ab6721 h1:RlZweED6sbSArvlE924+mUcZuXKLBHA35U7LN621Bws= github.com/mikioh/ipaddr v0.0.0-20190404000644-d465c8ab6721 h1:RlZweED6sbSArvlE924+mUcZuXKLBHA35U7LN621Bws=
github.com/mikioh/ipaddr v0.0.0-20190404000644-d465c8ab6721/go.mod h1:Ickgr2WtCLZ2MDGd4Gr0geeCH5HybhRJbonOgQpvSxc= 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/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e h1:MRM5ITcdelLK2j1vwZ3Je0FKVCfqOLp5zO6trqMLYs0= 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/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e/go.mod h1:XV66xRDqSt+GTGFMVlhk3ULuV0y9ZmzeVGR4mloJI3M=
github.com/vishvananda/netlink v1.3.1 h1:3AEMt62VKqz90r0tmNhog0r/PpWKmrEShJU0wJW6bV0= github.com/vishvananda/netlink v1.3.1 h1:3AEMt62VKqz90r0tmNhog0r/PpWKmrEShJU0wJW6bV0=
+10
View File
@@ -33,6 +33,16 @@ func fmtDate(_ date: Date?) -> String {
return date.formatted(date: .abbreviated, time: .omitted) return date.formatted(date: .abbreviated, time: .omitted)
} }
/// "35 min", "2 h 5 min", "3 days".
func fmtDuration(_ seconds: Int64) -> String {
if seconds < 60 { return "under 1 min" }
let m = Int((Double(seconds) / 60).rounded())
if m < 60 { return "\(m) min" }
let h = m / 60
if h < 48 { return "\(h) h \(m % 60) min" }
return "\(Int((Double(h) / 24).rounded())) days"
}
/// Label of a chart point: "3 h ago" for hours, "Sat 3 Oct" for days. /// Label of a chart point: "3 h ago" for hours, "Sat 3 Oct" for days.
func pointLabel(_ p: StatPoint, range: String) -> String { func pointLabel(_ p: StatPoint, range: String) -> String {
if range == "24h" { if range == "24h" {
+39
View File
@@ -61,6 +61,38 @@ nonisolated struct PeerStats: Decodable, Hashable {
let lastHandshake: Date? let lastHandshake: Date?
let endpoint: String let endpoint: String
let down24h, up24h, down30d, up30d, downTotal, upTotal: Int64 let down24h, up24h, down30d, up30d, downTotal, upTotal: Int64
let location: GeoInfo?
}
/// Country and network of an address, from the server's DB-IP lookup.
nonisolated struct GeoInfo: Decodable, Hashable {
let country: String?
let countryName: String?
let asn: Int?
let network: String?
/// "Germany · Deutsche Telekom AG" or "Local network".
var label: String {
[countryName ?? country, network].compactMap { $0 }.filter { !$0.isEmpty }.joined(separator: " · ")
}
}
/// One row of a peer's connection history.
nonisolated struct ConnSession: Decodable, Identifiable, Hashable {
let start: Date
let end: Date
let open: Bool
let seconds: Int64
let endpoint: String
let ip: String
let geo: GeoInfo?
let down: Int64
let up: Int64
var id: String { "\(start.timeIntervalSince1970)-\(ip)" }
}
nonisolated struct SessionsResponse: Decodable {
let sessions: [ConnSession]
} }
nonisolated struct Peer: Decodable, Identifiable, Hashable { nonisolated struct Peer: Decodable, Identifiable, Hashable {
@@ -166,6 +198,12 @@ nonisolated struct LogSettings: Codable, Equatable {
nonisolated struct StatsSettings: Codable, Equatable { nonisolated struct StatsSettings: Codable, Equatable {
var hourlyHours: Int var hourlyHours: Int
var dailyDays: Int var dailyDays: Int
var geoip: Bool?
}
nonisolated struct GeoStatus: Decodable {
let enabled: Bool
let updated: Date?
} }
nonisolated struct AppSettings: Decodable { nonisolated struct AppSettings: Decodable {
@@ -175,6 +213,7 @@ nonisolated struct AppSettings: Decodable {
var adminUsername: String var adminUsername: String
var fingerprint: String var fingerprint: String
var logPath: String var logPath: String
var geo: GeoStatus?
} }
nonisolated struct SettingsResult: Decodable { nonisolated struct SettingsResult: Decodable {
+73
View File
@@ -15,8 +15,11 @@ struct PeerDetailView: View {
@State private var askKey = false @State private var askKey = false
@State private var deviceKey = "" @State private var deviceKey = ""
@State private var editing = false @State private var editing = false
@State private var sessions: [ConnSession] = []
@State private var allSessions = false
var body: some View { var body: some View {
ScrollViewReader { proxy in
ScrollView { ScrollView {
VStack(spacing: 16) { VStack(spacing: 16) {
if let error { Notice(text: error, isError: true) } if let error { Notice(text: error, isError: true) }
@@ -24,6 +27,7 @@ struct PeerDetailView: View {
header(p) header(p)
traffic traffic
connection(p) connection(p)
history.id("history")
clientConfig(p) clientConfig(p)
settings(p) settings(p)
} else if error == nil { } else if error == nil {
@@ -33,6 +37,13 @@ struct PeerDetailView: View {
.padding(16) .padding(16)
} }
.background(Color.gwGround) .background(Color.gwGround)
#if DEBUG
// Development: `-scrollToHistory YES` for screenshots of the history.
.task(id: sessions.count) {
if UserDefaults.standard.bool(forKey: "scrollToHistory"), !sessions.isEmpty { proxy.scrollTo("history", anchor: .top) }
}
#endif
}
.navigationTitle(peer?.name ?? "Peer") .navigationTitle(peer?.name ?? "Peer")
.navigationBarTitleDisplayMode(.inline) .navigationBarTitleDisplayMode(.inline)
.toolbar { .toolbar {
@@ -107,6 +118,7 @@ struct PeerDetailView: View {
SectionTitle(text: "Connection") SectionTitle(text: "Connection")
KV(key: "Tunnel address", value: p.ipv4 + "/32" + (p.ipv6.map { "\n" + $0 + "/128" } ?? ""), mono: true) KV(key: "Tunnel address", value: p.ipv4 + "/32" + (p.ipv6.map { "\n" + $0 + "/128" } ?? ""), mono: true)
KV(key: "Endpoint", value: p.stats.endpoint.isEmpty ? "–" : p.stats.endpoint, mono: true) KV(key: "Endpoint", value: p.stats.endpoint.isEmpty ? "–" : p.stats.endpoint, mono: true)
KV(key: "Location", value: p.stats.location?.label.isEmpty == false ? p.stats.location!.label : "–")
KV(key: "Latest handshake", value: ago(p.stats.lastHandshake)) KV(key: "Latest handshake", value: ago(p.stats.lastHandshake))
KV(key: "Public key", value: p.publicKey, mono: true) KV(key: "Public key", value: p.publicKey, mono: true)
KV(key: "Preshared key", value: p.hasPresharedKey ? "Set" : "None") KV(key: "Preshared key", value: p.hasPresharedKey ? "Set" : "None")
@@ -115,6 +127,66 @@ struct PeerDetailView: View {
.card() .card()
} }
private var history: some View {
VStack(alignment: .leading, spacing: 0) {
SectionTitle(text: "Connection history")
Text("Newest first. A new row starts when the device changes networks.")
.font(.caption)
.foregroundStyle(Color.gwText2)
.padding(.bottom, 8)
if sessions.isEmpty {
Text("No connections recorded yet.").font(.footnote).foregroundStyle(Color.gwText2).padding(.vertical, 6)
}
let shown = allSessions ? sessions : Array(sessions.prefix(8))
ForEach(shown) { se in
HStack(alignment: .top, spacing: 12) {
VStack(alignment: .leading, spacing: 4) {
HStack(spacing: 6) {
Text(se.start.formatted(.dateTime.day().month(.abbreviated).hour().minute()))
.font(.subheadline.weight(.medium))
if se.open {
HStack(spacing: 4) {
Circle().fill(Color.gwGood).frame(width: 6, height: 6)
Text("online")
}
.font(.caption.weight(.medium))
.padding(.horizontal, 6).padding(.vertical, 2)
.background(Color.gwBadge, in: Capsule())
}
}
Text(se.geo?.label.isEmpty == false ? se.geo!.label : "Unknown location")
.font(.footnote)
.foregroundStyle(se.geo == nil ? Color.gwText2 : Color.gwText)
Text(se.ip + " · " + fmtDuration(se.seconds))
.font(.mono(.caption))
.foregroundStyle(Color.gwText2)
}
Spacer(minLength: 8)
VStack(alignment: .trailing, spacing: 4) {
Text("↓ " + fmtBytes(se.down)).font(.footnote.monospacedDigit())
Text("↑ " + fmtBytes(se.up)).font(.footnote.monospacedDigit()).foregroundStyle(Color.gwText2)
}
}
.padding(.vertical, 8)
.accessibilityElement(children: .combine)
if se.id != shown.last?.id { Divider() }
}
if sessions.count > 8 {
Button(allSessions ? "Show fewer" : "Show all \(sessions.count)") { allSessions.toggle() }
.font(.footnote.weight(.medium))
.padding(.top, 8)
}
HStack(spacing: 4) {
Text("Country and network:")
Link("IP Geolocation by DB-IP", destination: URL(string: "https://db-ip.com")!).underline()
}
.font(.caption2)
.foregroundStyle(Color.gwText2)
.padding(.top, 10)
}
.card()
}
private func clientConfig(_ p: Peer) -> some View { private func clientConfig(_ p: Peer) -> some View {
VStack(alignment: .leading, spacing: 12) { VStack(alignment: .leading, spacing: 12) {
SectionTitle(text: "Client configuration") SectionTitle(text: "Client configuration")
@@ -154,6 +226,7 @@ struct PeerDetailView: View {
(peer, server) = try await (p, s) (peer, server) = try await (p, s)
error = nil error = nil
await loadStats() await loadStats()
if let r: SessionsResponse = try? await api.get("/peers/\(peerID)/sessions?limit=100") { sessions = r.sessions }
} catch { } catch {
self.error = session.message(for: error) self.error = session.message(for: error)
} }
+12 -1
View File
@@ -139,6 +139,12 @@ struct SettingsView: View {
Picker("Daily traffic history", selection: s.dailyDays) { Picker("Daily traffic history", selection: s.dailyDays) {
ForEach(options(Self.daily, current: stats!.dailyDays, unit: "days"), id: \.0) { Text($0.1).tag($0.0) } ForEach(options(Self.daily, current: stats!.dailyDays, unit: "days"), id: \.0) { Text($0.1).tag($0.0) }
} }
Toggle(isOn: Binding(get: { stats!.geoip ?? true }, set: { stats!.geoip = $0 })) {
VStack(alignment: .leading, spacing: 2) {
Text("Show country and network")
Text(geoStatusText).font(.caption).foregroundStyle(Color.gwText2)
}
}
Button("Save retention") { Button("Save retention") {
guard let old = settings, let log, let stats else { return } guard let old = settings, let log, let stats else { return }
if log.maxFiles < old.log.maxFiles || stats.hourlyHours < old.stats.hourlyHours || stats.dailyDays < old.stats.dailyDays { if log.maxFiles < old.log.maxFiles || stats.hourlyHours < old.stats.hourlyHours || stats.dailyDays < old.stats.dailyDays {
@@ -151,7 +157,7 @@ struct SettingsView: View {
} header: { } header: {
Text("Data retention") Text("Data retention")
} footer: { } footer: {
Text("The log uses up to \(log!.maxSizeMB * (log!.maxFiles + 1)) MB on disk. All-time traffic totals are always kept. Applies immediately.") Text("The log uses up to \(log!.maxSizeMB * (log!.maxFiles + 1)) MB on disk. Connection history is kept as long as the daily traffic history; all-time totals are always kept. Country and network come from the free DB-IP Lite databases, downloaded monthly and looked up on the server only. Applies immediately.")
} }
} }
@@ -171,6 +177,11 @@ struct SettingsView: View {
} }
} }
private var geoStatusText: String {
guard let updated = settings?.geo?.updated else { return "Database not downloaded yet" }
return "Database from \(fmtDate(updated))"
}
private func options(_ presets: [(Int, String)], current: Int, unit: String) -> [(Int, String)] { private func options(_ presets: [(Int, String)], current: Int, unit: String) -> [(Int, String)] {
presets.contains { $0.0 == current } ? presets : (presets + [(current, "\(current) \(unit)")]).sorted { $0.0 < $1.0 } presets.contains { $0.0 == current } ? presets : (presets + [(current, "\(current) \(unit)")]).sorted { $0.0 < $1.0 }
} }
+7 -2
View File
@@ -187,6 +187,10 @@ func run(configPath string) error {
return fmt.Errorf("open stats: %w", err) return fmt.Errorf("open stats: %w", err)
} }
geo := newGeo(dataDir, cfg.Stats.geoEnabled())
defer geo.Close()
stats.geo = geo
webTLS, err := setupTLS(cfg, dataDir) webTLS, err := setupTLS(cfg, dataDir)
if err != nil { if err != nil {
slog.Error("tls setup failed", "err", err) slog.Error("tls setup failed", "err", err)
@@ -200,13 +204,14 @@ func run(configPath string) error {
auth := newAuth(store) auth := newAuth(store)
app := &App{ app := &App{
store: store, kernel: kernel, recon: recon, stats: stats, auth: auth, tls: webTLS, store: store, kernel: kernel, recon: recon, stats: stats, auth: auth, tls: webTLS,
logPath: logPath, logw: logw, started: time.Now(), shutdown: shutdown, logPath: logPath, logw: logw, geo: geo, started: time.Now(), shutdown: shutdown,
} }
var wg sync.WaitGroup var wg sync.WaitGroup
wg.Add(2) wg.Add(3)
go func() { defer wg.Done(); recon.Run(stop) }() go func() { defer wg.Done(); recon.Run(stop) }()
go func() { defer wg.Done(); stats.Run(stop) }() go func() { defer wg.Done(); stats.Run(stop) }()
go func() { defer wg.Done(); geo.Run(stop) }()
go func() { go func() {
t := time.NewTicker(10 * time.Minute) t := time.NewTicker(10 * time.Minute)
defer t.Stop() defer t.Stop()
+72
View File
@@ -430,3 +430,75 @@ func TestRetentionValidation(t *testing.T) {
} }
} }
} }
func TestSessions(t *testing.T) {
s := &Stats{data: statsFile{Peers: map[string]*peerStats{}}}
ps := &peerStats{}
t0 := time.Date(2026, 10, 3, 12, 0, 0, 0, time.UTC)
at := func(min int) time.Time { return t0.Add(time.Duration(min) * time.Minute) }
smp := func(min int, ep string) PeerSample { return PeerSample{Endpoint: ep, LastHandshake: at(min)} }
s.track(ps, smp(0, "192.168.1.20:5000"), 10, 100, at(0)) // online from home
s.track(ps, smp(2, "192.168.1.20:5001"), 10, 100, at(2)) // same network, new port
if len(ps.Sessions) != 1 || ps.Sessions[0].Rx != 20 || ps.Sessions[0].Endpoint != "192.168.1.20:5001" {
t.Fatalf("same network should continue the session: %+v", ps.Sessions)
}
if g := ps.Sessions[0].Geo; g == nil || g.Network != "Local network" {
t.Fatalf("private address should be the local network: %+v", g)
}
s.track(ps, smp(4, "198.51.100.7:6000"), 5, 50, at(4)) // roamed to mobile
if len(ps.Sessions) != 2 || ps.Sessions[0].Open || !ps.Sessions[1].Open {
t.Fatalf("roaming should start a new session: %+v", ps.Sessions)
}
s.track(ps, smp(4, "198.51.100.7:6000"), 0, 0, at(10)) // no handshake for 6 min
if ps.openSession() != nil {
t.Fatal("session should end when the peer goes quiet")
}
if ps.Sessions[1].End != at(4) {
t.Fatalf("end should be the last time seen online, got %v", ps.Sessions[1].End)
}
// A peer missing from the kernel (disabled) gets its session closed, and
// sessions older than the daily retention are pruned.
ps.Sessions = append(ps.Sessions, connSession{Start: at(20), End: at(20), Open: true, Endpoint: "198.51.100.7:6000"})
s.data.Peers["p"] = ps
s.prune(StatsConfig{HourlyHours: 24, DailyDays: 7}, t0.AddDate(0, 0, 30))
if len(ps.Sessions) != 1 || !ps.Sessions[0].Open {
t.Fatalf("prune should keep only the open session: %+v", ps.Sessions)
}
if v := s.Sessions("p", 10); len(v) != 1 || v[0].IP != "198.51.100.7" {
t.Fatalf("sessions view: %+v", v)
}
}
func TestGeoLookupWithoutDatabase(t *testing.T) {
var g *Geo // no databases: private addresses still resolve
if info := g.Lookup("10.1.2.3:51820"); info == nil || info.Network != "Local network" {
t.Fatalf("private: %+v", info)
}
if info := g.Lookup("[2001:db8::1]:51820"); info != nil {
t.Fatalf("public without database should be unknown: %+v", info)
}
g2 := newGeo(t.TempDir(), false)
if info := g2.Lookup("203.0.113.9:1"); info != nil {
t.Fatalf("disabled: %+v", info)
}
}
// TestGeoDatabase runs only with GHOSTWIRE_GEO_DIR pointing at a folder with
// downloaded geo-country.mmdb and geo-asn.mmdb.
func TestGeoDatabase(t *testing.T) {
dir := os.Getenv("GHOSTWIRE_GEO_DIR")
if dir == "" {
t.Skip("GHOSTWIRE_GEO_DIR not set")
}
g := newGeo(dir, true)
defer g.Close()
for _, ep := range []string{"9.9.9.9:53", "1.1.1.1:53", "[2620:fe::fe]:53"} {
info := g.Lookup(ep)
if info == nil || info.Country == "" || info.Network == "" {
t.Errorf("%s: %+v", ep, info)
}
t.Logf("%s → %+v", ep, info)
}
}
+137 -11
View File
@@ -4,6 +4,7 @@ import (
"encoding/json" "encoding/json"
"errors" "errors"
"log/slog" "log/slog"
"net"
"os" "os"
"sync" "sync"
"time" "time"
@@ -27,14 +28,77 @@ type bucket struct {
} }
type peerStats struct { type peerStats struct {
LastRx int64 `json:"lastRx"` // last raw counter values LastRx int64 `json:"lastRx"` // last raw counter values
LastTx int64 `json:"lastTx"` LastTx int64 `json:"lastTx"`
TotalRx int64 `json:"totalRx"` TotalRx int64 `json:"totalRx"`
TotalTx int64 `json:"totalTx"` TotalTx int64 `json:"totalTx"`
LastHandshake time.Time `json:"lastHandshake"` LastHandshake time.Time `json:"lastHandshake"`
Endpoint string `json:"endpoint"` Endpoint string `json:"endpoint"`
Hourly []bucket `json:"hourly"` Hourly []bucket `json:"hourly"`
Daily []bucket `json:"daily"` Daily []bucket `json:"daily"`
Sessions []connSession `json:"sessions,omitempty"`
}
// session is one stretch of a peer being online from one address. When the
// device changes networks (Wi-Fi to mobile), a new session starts.
type connSession struct {
Start time.Time `json:"start"`
End time.Time `json:"end"` // last time the peer was seen online
Open bool `json:"open,omitempty"`
Endpoint string `json:"endpoint"`
Rx int64 `json:"rx"`
Tx int64 `json:"tx"`
Geo *GeoInfo `json:"geo,omitempty"` // looked up when the session started
}
const maxSessions = 1000 // per peer, besides the retention window
func (ps *peerStats) openSession() *connSession {
if n := len(ps.Sessions); n > 0 && ps.Sessions[n-1].Open {
return &ps.Sessions[n-1]
}
return nil
}
func hostOf(endpoint string) string {
if h, _, err := net.SplitHostPort(endpoint); err == nil {
return h
}
return endpoint
}
// track updates the session list after a sample. Online means a handshake
// within the online window, as in the peer lists.
func (s *Stats) track(ps *peerStats, smp PeerSample, rx, tx int64, now time.Time) {
cur := ps.openSession()
online := !smp.LastHandshake.IsZero() && now.Sub(smp.LastHandshake) < onlineWindow
if !online {
if cur != nil {
cur.Open = false
}
return
}
if cur != nil && smp.Endpoint != "" && hostOf(cur.Endpoint) != hostOf(smp.Endpoint) {
cur.Open = false // roamed to another network
cur = nil
}
if cur == nil {
start := smp.LastHandshake
if start.After(now) {
start = now
}
ps.Sessions = append(ps.Sessions, connSession{Start: start, Open: true, Endpoint: smp.Endpoint, Geo: s.geo.Lookup(smp.Endpoint)})
if len(ps.Sessions) > maxSessions {
ps.Sessions = append([]connSession(nil), ps.Sessions[len(ps.Sessions)-maxSessions:]...)
}
cur = &ps.Sessions[len(ps.Sessions)-1]
}
if smp.Endpoint != "" {
cur.Endpoint = smp.Endpoint // same network, the port may change
}
cur.End = now
cur.Rx += rx
cur.Tx += tx
} }
type statsFile struct { type statsFile struct {
@@ -49,6 +113,7 @@ type Stats struct {
dirty bool dirty bool
store *Store store *Store
kernel Kernel kernel Kernel
geo *Geo // nil: no country and network lookups
} }
func openStats(path string, store *Store, k Kernel) (*Stats, error) { func openStats(path string, store *Store, k Kernel) (*Stats, error) {
@@ -106,10 +171,18 @@ func (s *Stats) prune(c StatsConfig, now time.Time) {
y, m, d := now.Date() y, m, d := now.Date()
dayCut := time.Date(y, m, d-(c.DailyDays-1), 0, 0, 0, 0, now.Location()).Unix() dayCut := time.Date(y, m, d-(c.DailyDays-1), 0, 0, 0, 0, now.Location()).Unix()
for _, ps := range s.data.Peers { for _, ps := range s.data.Peers {
h, dl := len(ps.Hourly), len(ps.Daily) h, dl, sl := len(ps.Hourly), len(ps.Daily), len(ps.Sessions)
ps.Hourly = dropBefore(ps.Hourly, hourCut) ps.Hourly = dropBefore(ps.Hourly, hourCut)
ps.Daily = dropBefore(ps.Daily, dayCut) ps.Daily = dropBefore(ps.Daily, dayCut)
if len(ps.Hourly) != h || len(ps.Daily) != dl { // Connection history is kept as long as the daily traffic history.
i := 0
for i < len(ps.Sessions) && !ps.Sessions[i].Open && ps.Sessions[i].End.Unix() < dayCut {
i++
}
if i > 0 {
ps.Sessions = append([]connSession(nil), ps.Sessions[i:]...)
}
if len(ps.Hourly) != h || len(ps.Daily) != dl || len(ps.Sessions) != sl {
s.dirty = true s.dirty = true
} }
} }
@@ -131,6 +204,7 @@ func (s *Stats) sample() {
now := time.Now() now := time.Now()
s.mu.Lock() s.mu.Lock()
defer s.mu.Unlock() defer s.mu.Unlock()
seen := map[string]bool{}
for _, smp := range samples { for _, smp := range samples {
id := idByKey[smp.PublicKey] id := idByKey[smp.PublicKey]
if id == "" { if id == "" {
@@ -141,6 +215,7 @@ func (s *Stats) sample() {
ps = &peerStats{} ps = &peerStats{}
s.data.Peers[id] = ps s.data.Peers[id] = ps
} }
seen[id] = true
dRx, dTx := smp.RxBytes-ps.LastRx, smp.TxBytes-ps.LastTx dRx, dTx := smp.RxBytes-ps.LastRx, smp.TxBytes-ps.LastTx
if dRx < 0 || dTx < 0 { // counters were reset if dRx < 0 || dTx < 0 { // counters were reset
dRx, dTx = smp.RxBytes, smp.TxBytes dRx, dTx = smp.RxBytes, smp.TxBytes
@@ -158,12 +233,19 @@ func (s *Stats) sample() {
if smp.Endpoint != "" { if smp.Endpoint != "" {
ps.Endpoint = smp.Endpoint ps.Endpoint = smp.Endpoint
} }
s.track(ps, smp, dRx, dTx, now)
s.dirty = true s.dirty = true
} }
for id := range s.data.Peers { for id, ps := range s.data.Peers {
if !exists[id] { if !exists[id] {
delete(s.data.Peers, id) delete(s.data.Peers, id)
s.dirty = true s.dirty = true
continue
}
// Disabled peers are not in the kernel any more: end their session.
if cur := ps.openSession(); cur != nil && !seen[id] {
cur.Open = false
s.dirty = true
} }
} }
s.prune(cfg.Stats, now) s.prune(cfg.Stats, now)
@@ -274,6 +356,7 @@ type PeerSummary struct {
Up30d int64 `json:"up30d"` Up30d int64 `json:"up30d"`
DownTotal int64 `json:"downTotal"` DownTotal int64 `json:"downTotal"`
UpTotal int64 `json:"upTotal"` UpTotal int64 `json:"upTotal"`
Location *GeoInfo `json:"location"` // of the current or last endpoint
} }
func sumPoints(pts []Point) (down, up int64) { func sumPoints(pts []Point) (down, up int64) {
@@ -298,6 +381,49 @@ func (s *Stats) Summary(id string) PeerSummary {
out.LastHandshake = &t out.LastHandshake = &t
out.Online = time.Since(t) < onlineWindow out.Online = time.Since(t) < onlineWindow
} }
if cur := ps.openSession(); cur != nil && cur.Geo != nil && hostOf(cur.Endpoint) == hostOf(ps.Endpoint) {
out.Location = cur.Geo
}
}
if out.Location == nil && out.Endpoint != "" {
out.Location = s.geo.Lookup(out.Endpoint)
}
return out
}
// SessionView is one row of a peer's connection history, from the peer's
// point of view (down = downloaded by the peer).
type SessionView struct {
Start time.Time `json:"start"`
End time.Time `json:"end"`
Open bool `json:"open"`
Seconds int64 `json:"seconds"`
Endpoint string `json:"endpoint"`
IP string `json:"ip"`
Geo *GeoInfo `json:"geo"`
Down int64 `json:"down"`
Up int64 `json:"up"`
}
// Sessions returns up to limit sessions, newest first.
func (s *Stats) Sessions(id string, limit int) []SessionView {
s.mu.Lock()
var list []connSession
if ps := s.data.Peers[id]; ps != nil {
list = append(list, ps.Sessions...)
}
s.mu.Unlock()
out := []SessionView{}
for i := len(list) - 1; i >= 0 && len(out) < limit; i-- {
se := list[i]
geo := se.Geo
if geo == nil {
geo = s.geo.Lookup(se.Endpoint) // database was missing when it started
}
out = append(out, SessionView{
Start: se.Start, End: se.End, Open: se.Open, Seconds: int64(se.End.Sub(se.Start).Seconds()),
Endpoint: se.Endpoint, IP: hostOf(se.Endpoint), Geo: geo, Down: se.Tx, Up: se.Rx,
})
} }
return out return out
} }