7391aac429
The server pings a peer's tunnel address every 30 s and shows the median of the last 5 minutes in the peer list (with a 1-hour sparkline) and a 24-hour chart on the peer page. Off by default; "active" pings only while the device sends traffic, "always" keeps the tunnel up.
306 lines
7.8 KiB
Go
306 lines
7.8 KiB
Go
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
|
|
}
|