Peers: optional latency check with per-peer setting
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.
This commit is contained in:
+305
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user