diff --git a/CHANGELOG.md b/CHANGELOG.md index 352d9da6..f4d54d24 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -501,6 +501,63 @@ Detailed per-release notes are on the SYN-ACK shows the outbound path was working. Traffic merely received from the peer does not count as evidence, because a node whose outbound path is dead still receives its peers' keepalives. +- **A node whose IP address changes recovers in about a second instead of + about half a minute.** When the host's address changed under a running + daemon (new DHCP lease, Wi-Fi to Ethernet, a VPN taking the default route, a + container moved to another address), nothing in the daemon noticed. Peers + kept sending to the old address until their own timeouts moved them to the + relay or the node's 25 s keepalive reached them, the node's own registry + connections stayed bound to the old address until a 30 s read timeout, and + the registry kept handing out the old endpoint for a minute or more. The + daemon now checks once a second which source address the kernel would use + toward the beacon (a route lookup; nothing is sent) and, when it changes, + re-registers with the beacon, sends one authenticated probe straight to + every tunnel peer so each learns the new address from it, and sends the + registry its new endpoint and LAN addresses over a fresh connection. Only + the endpoint is sent: the registry still has the node's visibility, + hostname and trust pairs, so they are not written again (the full restore + still runs if the registry has not answered for 5 minutes or its reply + shows it lost the node). The beacon registration is repeated 31 s later, + because the beacon accepts one endpoint update per node every 30 s and + drops the rest; a move within 30 s of the last keepalive registration + otherwise reached the beacon only with the next one, up to a minute later. + Measured in a three-node Docker lab, moving one node to a new address: peer + to moved node 32 s before, under 1 s after; moved node to peer 30 s before, + no failed send after; a node with no existing tunnel timed out after 30 s + before and connected directly in 0.2 s after. Peers need no update. + - An interface that drops and returns with the same address triggers + nothing, and neither does an unrelated interface appearing (a Docker + bridge, a VPN that does not take the route), nor a switch to another + beacon by the beacon-list refresh, which may be reached over another + route: the comparison starts over with the new beacon (a change still + waiting to be announced at that moment still runs). Recoveries are at + least 10 s apart, and the gap doubles up to 2 minutes while changes keep + coming, so a flapping interface cannot flood the registry; a change seen + during the gap runs when the gap ends, unless the address has gone back. + The gaps are measured on the wall clock, so time asleep counts toward + them. + - Behind NAT the local address does not change when the public one does. + The daemon also reads the address the beacon reports seeing it at (the + reply to its beacon registration, every keepalive interval, 60 s by + default) and runs the same recovery when that IP changes. A reply counts + only if it arrives within 2 s of a registration the node sent, and a new + IP only once two replies in a row agree on it; when one reply differs, + the node registers once more at once to check it. Because a NAT that maps + a node to several public IPs looks like a stream of changes, the gap + between these recoveries starts at 1 minute and grows to 1 hour. This + path has unit tests only; it was not exercised against a real NAT. + - A host waking from sleep on a new network no longer has its two + recoveries (wake and address change) abort each other's registry calls: + registry reconnects and re-registrations, the heartbeat's included, run + one at a time, and the wake or rx-watchdog recovery is skipped when a + full re-registration succeeded in the last 5 s. + - `-no-addr-watch` (config.json `no_addr_watch`) turns the watcher off. + - Publishes `tunnel.addr_changed` (`reason`: `local_address`, + `observed_endpoint` or `registry_retry`; `previous` and `current`: the + local address, or for `observed_endpoint` the IP the beacon sees; + `peers_notified`; `registry_ok`: whether the registry accepted the + re-registration). A relay-only or compat-mode node does not probe peers + directly. - **Proxy credential hints no longer send an operator who already set `proxy_cmd` off to set it.** When the daemon re-reads its credentials with a proxy command and the proxy still rejects them (407, or Meta Muse's diff --git a/cmd/daemon/main.go b/cmd/daemon/main.go index 30600046..2e50df24 100644 --- a/cmd/daemon/main.go +++ b/cmd/daemon/main.go @@ -105,6 +105,7 @@ func main() { beaconRTTProbe := flag.Bool("beacon-rtt-probe", false, "probe beacon RTT before selection; override hash pick when >2× slower than best (ablation test, default off)") noRxWatchdog := flag.Bool("no-rx-watchdog", false, "disable the inbound-path watchdog that soft-recovers (beacon+registry re-registration) and, on a persistent wedge, exits non-zero for supervisor respawn") noPathWatch := flag.Bool("no-path-watch", false, "disable the per-peer path watchdog that probes inbound-silent peers and resets a dead peer path in place (prefer-direct sequence) without a daemon restart") + noAddrWatch := flag.Bool("no-addr-watch", false, "disable the own-address watcher that notices this host's IP address changing (or the public IP the beacon sees it at) and re-announces the node to the beacon, the registry and every tunnel peer at once") // -transport and -proxy have literal defaults: their environment // variables beat config.json (see flagSources.envOverConfig), and -help // must never print an environment value — PILOT_PROXY can hold proxy @@ -401,6 +402,7 @@ func main() { BeaconRTTProbe: *beaconRTTProbe, DisableRxWatchdog: *noRxWatchdog, DisablePathWatch: *noPathWatch, + DisableAddrWatch: *noAddrWatch, TransportMode: *transportMode, CompatBeaconURL: *compatBeacon, CompatTLSTrust: *tlsTrust, diff --git a/cmd/pilotctl/config_values.go b/cmd/pilotctl/config_values.go index 0e6c01bb..045def84 100644 --- a/cmd/pilotctl/config_values.go +++ b/cmd/pilotctl/config_values.go @@ -40,6 +40,7 @@ var configValueKinds = map[string]string{ "motd_feed_url": "String", "motd_interval": "Duration", "networks": "String", + "no_addr_watch": "Bool", "no_dataexchange": "Bool", "no_echo": "Bool", "no_eventstream": "Bool", diff --git a/pkg/daemon/addrwatch.go b/pkg/daemon/addrwatch.go new file mode 100644 index 00000000..f380048a --- /dev/null +++ b/pkg/daemon/addrwatch.go @@ -0,0 +1,630 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package daemon + +import ( + "log/slog" + "net" + "time" +) + +// Own-address watcher (L4). +// +// Failure mode this exists for: this host's IP address changes under a +// running daemon — Wi-Fi to Ethernet, a new DHCP lease, a VPN taking the +// default route, a container moved to another address. The tunnel socket is +// bound to the wildcard address and keeps working, but everyone else still +// holds the old address: peers keep sending to it, the beacon relays to it and +// the registry hands it out. Nothing in the daemon noticed, so recovery waited +// for other nodes' timeouts — a peer's blackhole flip to relay (8s silence +// plus three sends) or our own 25s NAT keepalive finally reaching the peer +// from the new address. Measured on v1.14.0-rc.1: 23–32s without traffic in +// either direction, and a registry entry still pointing at the old address +// minutes later. +// +// Detection: once a second, ask the kernel which source address it would use +// toward the beacon (a connected UDP socket; nothing is sent). That is the +// address the daemon's traffic actually leaves from, so it changes exactly +// when the daemon's path changes, and stays put when an unrelated interface +// appears (a Docker bridge, a VPN that does not take the route) or when a +// link-local address comes and goes. An IPv6 temporary address rotating is +// seen only when it is the address the route to the beacon uses, and then it +// is a real change of our source address. No route at all reads as +// "offline" and triggers nothing: there is nobody to tell until an address +// comes back, and if the same one comes back nothing has changed. +// +// A node behind NAT sees no local change when the NAT's public address +// changes. The second input covers that: the beacon's reply to a discover +// (RegisterWithBeacon, every keepalive interval) carries the address the +// beacon sees us at (TunnelManager.noteDiscoverReply). That input is noisier +// than the first. Anyone can send a datagram from the beacon's address, and +// some NATs show one node at more than one public IP (carrier-grade NAT +// pools, cloud NAT with several addresses, per-flow load balancing). So a +// reply counts only if it arrives within 2 s of a discover this node sent, +// and a new IP counts only once two replies in a row agree on it. When one +// reply names a new IP the watcher sends one more discover at once, so the +// second opinion arrives within a round trip instead of a keepalive interval. +// +// Both baselines belong to one beacon. The beacon refresh can move us to +// another one, which may be reached over another route (a private bootstrap +// beacon beside public ones) or see us at another public IP (a NAT that picks +// the public address per destination, dual WAN). A beacon switch therefore +// starts both baselines over instead of reading as a move. A local change +// still waiting to be announced at that moment is kept and still runs. +// +// Recovery (recoverFromAddrChange), in the order that matters for latency: +// 1. beacon: re-register, so relay delivery and hole-punching use the new +// endpoint (one datagram). The beacon takes at most one endpoint update +// per node every 30 s and drops the rest, so a move within 30 s of the +// last keepalive registration would not reach it until the next one; the +// registration is repeated addrBeaconReregisterDelay later to cover that; +// 2. peers: one authenticated probe straight to every tunnel peer, so each +// overwrites its entry for us from the packet's source (one datagram per +// peer, repeated only for peers that have not answered); +// 3. registry: a fresh connection — the pooled one is bound to the address +// we no longer have — and an endpoint-only re-registration +// (reRegisterEndpoint): the registry still holds our visibility, +// hostname and trust pairs, so they are not written again. +// +// Rate limit: two recoveries are always at least addrRecoverCooldown apart. +// On top of that each input has its own backoff: the gap doubles with every +// recovery in a row, up to addrRecoverCooldownMax for a local change and +// addrObservedCooldownMax for a beacon-observed one, and drops back after a +// quiet period. A flapping interface costs a handful of re-registrations and +// then one every two minutes; a NAT whose public IP keeps flapping costs one +// an hour. A change seen during the cooldown is not lost: it runs when the +// cooldown ends, unless the address has gone back to the one already +// announced. The gaps are measured on the wall clock (addrWatchNow), so time +// the host spends asleep counts. +// +// -no-addr-watch (Config.DisableAddrWatch) turns the watcher off. +const ( + // addrWatchPollInterval is how often the source address is sampled. + addrWatchPollInterval = time.Second + + // addrRecoverCooldown is the minimum gap between two recoveries, and the + // first gap of the local-address backoff. + addrRecoverCooldown = 10 * time.Second + + // addrRecoverCooldownMax caps the doubling of the local-address gap. + addrRecoverCooldownMax = 2 * time.Minute + + // addrRecoverQuiet is how long without a local-address recovery before + // that gap drops back to addrRecoverCooldown. + addrRecoverQuiet = 10 * time.Minute + + // addrObservedCooldown, addrObservedCooldownMax and addrObservedQuiet + // are the same three for the beacon-observed IP. A real change of a + // NAT's public address is rare, and an IP that keeps changing is far + // more likely a NAT that maps us to several addresses, so the gap starts + // longer and grows to an hour. The quiet period is longer than that cap, + // or a steady flap at the cap would reset the backoff every time. + addrObservedCooldown = time.Minute + addrObservedCooldownMax = time.Hour + addrObservedQuiet = 2 * addrObservedCooldownMax + + // addrObservedConfirmGap is the minimum gap between two confirming + // discovers (see wantConfirm). + addrObservedConfirmGap = 30 * time.Second + + // addrRegistryRetryMax is how many more times the registry half of a + // recovery is retried (one per cooldown) when the registry could not be + // reached right after the change. + addrRegistryRetryMax = 3 +) + +// addrAnnounceRetryDelays are the waits before the peer probe is repeated +// for peers that have not answered directly since the recovery started. A +// single datagram can be lost, and a route that has just come up may drop +// the first one. +var addrAnnounceRetryDelays = []time.Duration{time.Second, 2 * time.Second} + +// addrBeaconReregisterDelay is when, after a recovery started, the beacon +// registration is repeated: just past the beacon's 30 s per-node limit on +// endpoint updates, so the repeat is accepted even if the recovery's own +// registrations were not. +var addrBeaconReregisterDelay = 31 * time.Second + +// Reasons a recovery ran — the "reason" field of tunnel.addr_changed. +const ( + addrReasonLocal = "local_address" + addrReasonObserved = "observed_endpoint" + addrReasonRegistry = "registry_retry" +) + +// addrWatchState is the watcher's memory between ticks. All decisions are +// made by its methods on values handed in, so tests drive it without real +// interfaces or clocks. +type addrWatchState struct { + beacon string // the beacon both baselines were taken against ("" = none) + beaconAddr *net.UDPAddr // the same beacon, for one last sample on its route at a switch + + local string // source address seen on the latest sample with a route + announced string // local address the last recovery announced (or the baseline) + + observed string // beacon-observed IP the last recovery announced (or the baseline) + observedCurrent string // latest IP two replies in a row agreed on + confirmAskedFor time.Time // arrival of the reply the last confirming discover was for + lastConfirm time.Time // when that discover was sent + + registryRetries int // registry-only retries still owed + + lastRecover time.Time // the latest recovery of any kind + localGap addrBackoff // local-address and registry-retry recoveries + observedGap addrBackoff // observed-endpoint recoveries +} + +// addrBackoff is one input's run of recoveries. +type addrBackoff struct { + last time.Time + streak int // recoveries in a row without a quiet pause +} + +// noteBeacon records which beacon this tick's inputs come from. On a switch +// the baselines start over: the new beacon may be reached from another local +// address and see us at another public IP without anything having moved. +// +// A local change that is still waiting to be announced is kept, though: +// one deferred by the cooldown, or one the caller's last sample on the old +// route has just seen. Our address really did change, and the beacon switch +// tells the peers and the registry nothing. It runs when it is due and +// announces whatever address the new route uses by then. +func (st *addrWatchState) noteBeacon(beacon string) { + if beacon == st.beacon { + return + } + st.beacon = beacon + if st.local == st.announced { + st.local, st.announced = "", "" + } + st.observed, st.observedCurrent = "", "" + st.confirmAskedFor, st.lastConfirm = time.Time{}, time.Time{} +} + +// noteLocal records one sample of the local source address. ok is false when +// the host has no route (offline), which is remembered as nothing at all. +func (st *addrWatchState) noteLocal(addr string, ok bool) { + if !ok || addr == "" { + return + } + st.local = addr + if st.announced == "" { + st.announced = addr // first sample: baseline, not a change + } +} + +// noteObserved takes the two latest discover replies (newest first). An IP +// becomes where the beacon sees us only when both come from the current +// beacon and agree on it; a lone reply that disagrees changes nothing (see +// wantConfirm). Only the IP is compared: a NAT hands out a different port per +// destination and may renumber ports freely, and peers already follow a port +// change from our keepalives. +func (st *addrWatchState) noteObserved(latest, prev beaconObservation) { + if latest.endpoint == nil || prev.endpoint == nil || + latest.beacon != st.beacon || prev.beacon != st.beacon || + !latest.endpoint.IP.Equal(prev.endpoint.IP) { + return + } + st.observedCurrent = latest.endpoint.IP.String() + if st.observed == "" { + st.observed = st.observedCurrent // first agreed IP: baseline, not a change + } +} + +// wantConfirm reports whether to send a discover now to get a second opinion +// on the latest reply: it is from the current beacon, names an IP that two +// replies have not agreed on, has not had a confirming discover yet, and none +// was sent in the last addrObservedConfirmGap. The gap bounds the chain a +// beacon or NAT reporting a new IP on every reply could start. +func (st *addrWatchState) wantConfirm(latest beaconObservation, now time.Time) bool { + if latest.endpoint == nil || latest.beacon != st.beacon || + latest.endpoint.IP.String() == st.observedCurrent || latest.at.Equal(st.confirmAskedFor) { + return false + } + if recentlyAt(st.lastConfirm, now, addrObservedConfirmGap) { + return false + } + st.confirmAskedFor, st.lastConfirm = latest.at, now + return true +} + +// pending names the reason a recovery is wanted, or "" if none is. +func (st *addrWatchState) pending() string { + switch { + case st.local != st.announced: + return addrReasonLocal + case st.observed != "" && st.observedCurrent != st.observed: + return addrReasonObserved + case st.registryRetries > 0: + return addrReasonRegistry + } + return "" +} + +// addrGapParams is the backoff schedule of the input behind reason: the +// first gap, its cap, and how long without a recovery before it drops back. +func addrGapParams(reason string) (base, max, quiet time.Duration) { + if reason == addrReasonObserved { + return addrObservedCooldown, addrObservedCooldownMax, addrObservedQuiet + } + return addrRecoverCooldown, addrRecoverCooldownMax, addrRecoverQuiet +} + +// gapFor returns the backoff of the input behind reason. +func (st *addrWatchState) gapFor(reason string) *addrBackoff { + if reason == addrReasonObserved { + return &st.observedGap + } + return &st.localGap +} + +// cooldown is the gap required after the last recovery of reason's input. +func (st *addrWatchState) cooldown(reason string) time.Duration { + base, max, _ := addrGapParams(reason) + gap := base + for i := 1; i < st.gapFor(reason).streak && gap < max; i++ { + gap *= 2 + } + if gap > max { + gap = max + } + return gap +} + +// due reports whether a recovery should run now, and why. It does not change +// state; the caller follows up with began. +func (st *addrWatchState) due(now time.Time) (reason string, ok bool) { + reason = st.pending() + if reason == "" { + return "", false + } + if recentlyAt(st.lastRecover, now, addrRecoverCooldown) { + return "", false + } + if b := st.gapFor(reason); recentlyAt(b.last, now, st.cooldown(reason)) { + return "", false + } + return reason, true +} + +// recentlyAt reports whether t is set and less than d before now. The +// watcher runs on the wall clock (addrWatchNow), which can be stepped back; +// a t that lies ahead of now says nothing about how long ago it was, so it +// does not hold anything back. +func recentlyAt(t, now time.Time, d time.Duration) bool { + if t.IsZero() { + return false + } + since := now.Sub(t) + return since >= 0 && since < d +} + +// began marks a recovery for reason as started at now, and returns the +// addresses it moves us from and to: the local source address, or for +// addrReasonObserved the IP the beacon sees us at. +func (st *addrWatchState) began(now time.Time, reason string) (previous, current string) { + b := st.gapFor(reason) + if _, _, quiet := addrGapParams(reason); !b.last.IsZero() && now.Sub(b.last) >= quiet { + b.streak = 0 + } + b.streak++ + b.last = now + st.lastRecover = now + + previous, current = st.announced, st.local + if reason == addrReasonObserved { + previous, current = st.observed, st.observedCurrent + st.observed = st.observedCurrent + } + st.announced = st.local + // Our own move changes what the beacon sees too. Forget the old value + // so the replies to the recovery's own discovers become the baseline + // instead of a second change. + if reason == addrReasonLocal { + st.observed, st.observedCurrent = "", "" + } + if reason == addrReasonRegistry { + st.registryRetries-- + } else { + st.registryRetries = 0 + } + return previous, current +} + +// finished records whether the registry took the new endpoint. +func (st *addrWatchState) finished(reason string, registryOK bool) { + if registryOK { + st.registryRetries = 0 + return + } + if reason != addrReasonRegistry { + st.registryRetries = addrRegistryRetryMax + } +} + +// addrSourceFn reports the local IP the kernel would use for traffic to the +// control plane — toward beacon, or when it is nil toward the registry — and +// whether there is a route at all. +type addrSourceFn func(beacon *net.UDPAddr) (ip string, ok bool) + +// routeSourceIP is the production address source for one target: connecting +// a UDP socket performs the route lookup and binds the source address without +// sending anything. +func routeSourceIP(target *net.UDPAddr) (string, bool) { + if target == nil { + return "", false + } + c, err := net.DialUDP("udp", nil, target) + if err != nil { + return "", false // no route: offline + } + defer c.Close() + la, _ := c.LocalAddr().(*net.UDPAddr) + if la == nil || la.IP == nil || la.IP.IsUnspecified() { + return "", false + } + return la.IP.String(), true +} + +// addrWatchSource builds the address source: the route toward the beacon +// the tunnel currently uses, or, for a daemon with no beacon, toward the +// registry. +func (d *Daemon) addrWatchSource() addrSourceFn { + var registry *net.UDPAddr // resolved once; only the route to it matters + var nextResolve time.Time // a failed lookup is not repeated every second + return func(beacon *net.UDPAddr) (string, bool) { + if beacon != nil { + return routeSourceIP(beacon) + } + if registry == nil && d.config.RegistryAddr != "" && !time.Now().Before(nextResolve) { + registry, _ = net.ResolveUDPAddr("udp", d.config.RegistryAddr) + nextResolve = time.Now().Add(30 * time.Second) + } + return routeSourceIP(registry) + } +} + +func (d *Daemon) addrWatchLoop() { + if d.config.DisableAddrWatch { + return + } + // A compat-mode daemon has no UDP socket and no direct paths: all its + // traffic rides one WSS connection to the beacon, which reconnects on + // its own when the address under it changes. + if d.config.TransportMode == TransportCompat { + return + } + src := d.addrWatchSource() + st := &addrWatchState{} + ticker := time.NewTicker(addrWatchPollInterval) + defer ticker.Stop() + for { + select { + case <-d.stopCh: + return + case <-ticker.C: + d.addrWatchTick(st, src, addrWatchNow()) + } + } +} + +// addrWatchNow is the watcher's clock: the wall clock, with the monotonic +// reading stripped. Go measures the time between two readings that both +// carry a monotonic reading on the monotonic clock, and on Linux that clock +// does not advance while the host is suspended. A laptop that recovered, +// slept for an hour and woke on a new network would still be "inside" its +// 10 s cooldown, and its backoff would never see the quiet hour; a host +// waking somewhere else is exactly when the watcher matters. +func addrWatchNow() time.Time { return time.Now().Round(0) } + +// addrWatchTick takes one sample and runs a recovery if one is due. +// Extracted for testability — production drives it from addrWatchLoop. +// +// L4 panic boundary (architecture-notes/03-INVARIANTS.md §8): a panic drops +// this tick and the next one resamples from clean state. +func (d *Daemon) addrWatchTick(st *addrWatchState, src addrSourceFn, now time.Time) (fired bool) { + defer recoverLayer("L4", "addrWatchTick", d.bus, nil) + + // The beacon can change under us (beaconRefreshTick); read it once so the + // sample, the replies and the baselines all refer to the same one. + beacon := d.tunnels.BeaconUDPAddr() + beaconKey := "" + if beacon != nil { + beaconKey = beacon.String() + } + if beaconKey != st.beacon { + // One more sample on the old route before the baselines start + // over: a move in the same second as the switch shows up there, + // and noteBeacon keeps it as a change still to announce. + if st.announced != "" { + st.noteLocal(src(st.beaconAddr)) + } + st.noteBeacon(beaconKey) + st.beaconAddr = beacon + } + ip, ok := src(beacon) + st.noteLocal(ip, ok) + latest, prev := d.tunnels.beaconObservations() + st.noteObserved(latest, prev) + if ok && st.wantConfirm(latest, now) { + // One reply puts us at an IP no second reply has confirmed. Ask + // again now rather than wait a keepalive interval for the next. + d.tunnels.RegisterWithBeacon() + } + reason, due := st.due(now) + if !due { + return false + } + if !ok { + // Offline right now. Whatever is owed runs once a route is back. + return false + } + previous, current := st.began(now, reason) + if reason == addrReasonLocal { + // began dropped the observed baseline; drop the stored replies too, + // so the baseline is taken from replies that arrive after the move + // rather than from ones that describe the old address. + d.tunnels.forgetObservedEndpoints() + } + registryOK := addrWatchRecover(d, reason, previous, current, reason != addrReasonRegistry) + st.finished(reason, registryOK) + return true +} + +// addrWatchRecover is swapped by tests to observe recoveries without a live +// registry or beacon (mirrors pathWatchResetPeer). +var addrWatchRecover = func(d *Daemon, reason, previous, current string, announce bool) bool { + return d.recoverFromAddrChange(reason, previous, current, announce) +} + +// RecoverFromAddrChangeForTest runs the address-change recovery on d now, as +// the watcher does when it sees this host's address change, without the +// watcher's rate limit, and reports whether the registry accepted the +// re-registration. It exists for the end-to-end tests in ./tests, which run +// outside this package and cannot give a daemon a new address; nothing else +// should call it. +func RecoverFromAddrChangeForTest(d *Daemon) bool { + return d.recoverFromAddrChange("manual", "", "", true) +} + +// recoverFromAddrChange tells the beacon, the tunnel peers (when announce is +// set) and the registry where this node now is. See the file comment for the +// order. Returns whether the registry accepted the re-registration. +func (d *Daemon) recoverFromAddrChange(reason, previous, current string, announce bool) bool { + started := time.Now() + if announce { + slog.Warn("own address changed — re-announcing to beacon, peers and registry", + "reason", reason, "previous", previous, "current", current) + } else { + slog.Warn("own address changed — retrying the registry re-registration", "current", current) + } + + notified := 0 + if announce { + // Beacon first: it is one datagram, and a peer that cannot be + // reached directly is reached through the beacon's mapping. + d.tunnels.RegisterWithBeacon() + notified = d.announceToPeers(time.Time{}) + } + if announce && !d.stopping() { + d.bgWG.Add(1) + go func() { + defer d.bgWG.Done() + defer recoverLayer("L4", "addrAnnounceFollowUp", d.bus, nil) + d.addrAnnounceFollowUp(started) + }() + } + + // The port may be unchanged but the private addresses we advertise for + // same-LAN detection are not. + d.refreshLANAddrs() + + // Registry last: it is TCP with dial retries and may take seconds, and + // nothing above should wait for it. reestablishTransport also repeats + // the beacon registration, which costs one datagram and covers a first + // one lost while the route was still settling. It is never skipped for + // a run that just finished: that run may have registered the old address. + registryOK := d.reestablishTransport("addr-change", + reestablishOpts{freshConn: true, endpointOnly: true, always: true}) + + d.publishEvent("tunnel.addr_changed", map[string]any{ + "reason": reason, + "previous": previous, + "current": current, + "peers_notified": notified, + "registry_ok": registryOK, + }) + slog.Info("address-change recovery done", + "reason", reason, "peers_notified", notified, "registry_ok", registryOK, + "took", time.Since(started).Truncate(time.Millisecond).String()) + return registryOK +} + +// addrAnnounceFollowUp repeats the parts of a recovery that one datagram may +// not have achieved. The peer probe is repeated after each of +// addrAnnounceRetryDelays for peers that have not answered directly since the +// recovery started. The beacon registration is repeated once, +// addrBeaconReregisterDelay after the start: the beacon accepts at most one +// endpoint update per node every 30 s and drops the rest without telling us, +// so when the move came within 30 s of a keepalive registration it would +// otherwise relay to the old address until the next keepalive, up to a +// minute later. +func (d *Daemon) addrAnnounceFollowUp(started time.Time) { + sleep := func(wait time.Duration) bool { + t := time.NewTimer(wait) + defer t.Stop() + select { + case <-d.stopCh: + return false + case <-t.C: + return true + } + } + for _, wait := range addrAnnounceRetryDelays { + if !sleep(wait) { + return + } + if d.announceToPeers(started) == 0 { + break + } + } + if !sleep(time.Until(started.Add(addrBeaconReregisterDelay))) { + return + } + d.tunnels.RegisterWithBeacon() +} + +// announceToPeers sends one direct, authenticated probe to every peer with +// an established session, and returns how many were sent. With a non-zero +// since, peers that have already sent us a direct frame after that time are +// skipped: a direct frame can only have reached us at the new address, so +// that peer has learned it. +// +// A relay-only node (-relay-only / compat) never sends: a direct frame would +// show the peer the real address the node is hiding. Peers we only know +// through the beacon are skipped by SendDirectProbe, and stay reachable +// through the beacon registration. +func (d *Daemon) announceToPeers(since time.Time) int { + if d.config.RelayOnly { + return 0 + } + sent := 0 + for _, nodeID := range d.tunnels.ReadyPeerIDs() { + if !since.IsZero() && d.tunnels.LastDirectRecv(nodeID).After(since) { + continue + } + if err := d.tunnels.SendDirectPathProbe(nodeID); err != nil { + slog.Debug("address announce: probe not sent", "peer_node_id", nodeID, "err", err) + continue + } + sent++ + } + return sent +} + +// currentLANAddrs returns the private addresses advertised to the registry. +func (d *Daemon) currentLANAddrs() []string { + d.addrMu.RLock() + defer d.addrMu.RUnlock() + return d.lanAddrs +} + +// refreshLANAddrs re-collects the advertised private addresses after an +// address change. A compat daemon advertises none. +func (d *Daemon) refreshLANAddrs() { + if d.config.TransportMode == TransportCompat { + return + } + la := d.tunnels.LocalAddr() + if la == nil { + return + } + _, port, err := net.SplitHostPort(la.String()) + if err != nil || port == "" || port == "0" { + return + } + addrs := collectLANAddrs(port) + d.addrMu.Lock() + d.lanAddrs = addrs + d.addrMu.Unlock() +} diff --git a/pkg/daemon/daemon.go b/pkg/daemon/daemon.go index e7832828..4ae8fcd1 100644 --- a/pkg/daemon/daemon.go +++ b/pkg/daemon/daemon.go @@ -198,6 +198,14 @@ type Config struct { // full daemon restart. Default false (watchdog on). DisablePathWatch bool + // DisableAddrWatch turns off the own-address watcher (addrwatch.go). + // The watcher notices this host's IP address changing under the daemon + // (or the public IP the beacon sees it at) and re-announces the node to + // the beacon, the registry and every tunnel peer at once, instead of + // leaving peers to find the new address through their own timeouts. + // Default false (watcher on). + DisableAddrWatch bool + // Telemetry consent gate. When set to the telemetry endpoint URL, // the daemon initialises a telemetry client that emits signed events // (install, usage, view, review). When empty (default), the client @@ -439,6 +447,18 @@ type Daemon struct { // signature) from "whole network unreachable" (restart just loops). lastRegistryOKNano atomic.Int64 + // reestablishMu serialises registry reconnects and re-registrations: + // reestablishTransport (the resume handler, the rx watchdog and the + // address watcher) and the heartbeat's own in trustRepublishLoop. They + // can all want one at once, and one caller's forceReconnectRegistry + // closes the connection another is mid-request on. reestablishOKWall + // (wall-clock unix nanos, under the mutex) is when the last + // reestablishTransport run the registry accepted finished, and + // reestablishOKFull whether it was a full re-registration. + reestablishMu sync.Mutex + reestablishOKWall int64 + reestablishOKFull bool + // lastDialOKNano / consecutiveDialTimeouts feed the rx watchdog's // PARTIAL-wedge detector. A daemon can be "up" with rx trickling — // a couple of keepalive packets from existing peers keep PktsRecv @@ -489,7 +509,7 @@ type Daemon struct { ctx context.Context cancelCtx context.CancelFunc - lanAddrs []string // LAN addresses for same-network peer detection + lanAddrs []string // LAN addresses for same-network peer detection; guarded by addrMu once Start has spawned the loops // Endpoint cache: nodeID -> last-known endpoint (peer resilience) epCacheMu sync.RWMutex @@ -1349,6 +1369,12 @@ func (d *Daemon) Start() error { d.bgWG.Add(1) go func() { defer d.bgWG.Done(); d.pathWatchLoop() }() + // 8d. Start the address watcher (L4). Notices this host's own address + // changing and re-announces to the beacon, the registry and every + // tunnel peer at once — see addrwatch.go. + d.bgWG.Add(1) + go func() { defer d.bgWG.Done(); d.addrWatchLoop() }() + // 9. Start idle connection sweeper d.bgWG.Add(1) go func() { defer d.bgWG.Done(); d.idleSweepLoop() }() @@ -5595,7 +5621,7 @@ func (d *Daemon) ensureTunnel(nodeID uint32) error { isLoopback := realIP != nil && realIP.IsLoopback() if !isLoopback { if lanAddrs, ok := resp["lan_addrs"].([]interface{}); ok && len(lanAddrs) > 0 { - if lanAddr := matchLANSubnet(d.lanAddrs, lanAddrs); lanAddr != "" { + if lanAddr := matchLANSubnet(d.currentLANAddrs(), lanAddrs); lanAddr != "" { if !d.addrFamilyMismatch(lanAddr) { targetAddr = lanAddr slog.Info("same-LAN peer detected, using LAN address", "node_id", nodeID, "lan_addr", lanAddr) @@ -5701,7 +5727,7 @@ func (d *Daemon) trustRepublishLoop() { if errors.Is(err, errRegistryCallTimedOut) { slog.Warn("heartbeat timed out — registry connection likely half-open, forcing reconnect", "consecutive_failures", consecutiveFailures, "deadline", registryCallDeadline) - if rcErr := d.forceReconnectRegistry(); rcErr != nil { + if rcErr := d.reconnectRegistrySerialised(); rcErr != nil { slog.Warn("registry force-reconnect failed", "error", rcErr) } else { consecutiveFailures = HeartbeatReregThresh @@ -5722,7 +5748,7 @@ func (d *Daemon) trustRepublishLoop() { time.Sleep(reregBackoff + jitter) slog.Info("attempting re-registration", "backoff", reregBackoff) - d.reRegister() + _ = d.reRegisterSerialised() // logs its own failure; the next heartbeat tells consecutiveFailures = 0 // Exponential backoff: 100ms → 200ms → 400ms → ... → 30s max. @@ -5743,6 +5769,22 @@ func (d *Daemon) trustRepublishLoop() { } } +// reconnectRegistrySerialised and reRegisterSerialised are the heartbeat's +// own registry reconnect and re-registration, under reestablishMu like the +// recoveries' (reestablishTransport): run at the same time as one, either +// would replace or use the registry connection under the other's requests. +func (d *Daemon) reconnectRegistrySerialised() error { + d.reestablishMu.Lock() + defer d.reestablishMu.Unlock() + return d.forceReconnectRegistry() +} + +func (d *Daemon) reRegisterSerialised() error { + d.reestablishMu.Lock() + defer d.reestablishMu.Unlock() + return d.reRegister() +} + // tunnelKeepaliveLoop (L4) refreshes the daemon's beacon registration so the // beacon retains a fresh endpoint for hole-punching and the upstream NAT // keeps the source mapping alive. Owns L4 beacon state only; briefly RLocks @@ -5881,11 +5923,67 @@ func (d *Daemon) publishHeartbeatEvent() { }) } -// reRegister re-registers with the registry after a connection loss or registry restart. -// Checks d.stopCh between regConn calls to avoid racing with Stop(). -func (d *Daemon) reRegister() { +// reRegister re-registers with the registry after a connection loss or +// registry restart, and restores everything the registry may have lost with +// it: visibility, hostname and the trust pairs. Checks d.stopCh between +// regConn calls to avoid racing with Stop(). Returns an error when the +// registry did not accept the registration; the restore steps after it are +// best-effort and only logged. +func (d *Daemon) reRegister() error { + nodeID, _, err := d.registerEndpoint() + if err != nil { + return err + } + d.restoreRegistryState(nodeID, true) + return nil +} + +// endpointOnlyRegistryFresh bounds when reRegisterEndpoint trusts the +// registry to still hold this node: it must have answered us within this +// long. The registry reaps a node after 30 minutes without a heartbeat, and +// a reaped node comes back without its visibility or hostname. +const endpointOnlyRegistryFresh = 5 * time.Minute + +// reRegisterEndpoint is the re-registration an address change needs. The +// registry still holds this node, its visibility, hostname and trust pairs; +// only the endpoint is out of date, so only the endpoint is sent. The full +// reRegister would also write SetVisibility, SetHostname and one ReportTrust +// per trusted peer (hundreds of signed writes for a node that trusts the +// service fleet) for state the registry never lost. +// +// It falls back to the full restore when the registry may have dropped the +// node after all: no answer from it for endpointOnlyRegistryFresh, a reply +// under a different node ID, or a reply without the configured hostname +// (the registry echoes a node's hostname on every registration). +func (d *Daemon) reRegisterEndpoint() error { + if ok := d.lastRegistryOKNano.Load(); ok == 0 || time.Since(time.Unix(0, ok)) > endpointOnlyRegistryFresh { + return d.reRegister() + } + before := d.NodeID() + nodeID, resp, err := d.registerEndpoint() + if err != nil { + return err + } + if nodeID != before { + d.restoreRegistryState(nodeID, true) + return nil + } + if host, _ := resp["hostname"].(string); d.config.Hostname != "" && host != d.config.Hostname { + d.restoreRegistryState(nodeID, false) + } + return nil +} + +// registerEndpoint sends this node's current endpoint and LAN addresses to +// the registry and applies the reply. Returns the node ID the registry +// answered with, and the reply. +func (d *Daemon) registerEndpoint() (uint32, map[string]interface{}, error) { if d.stopping() { - return + return 0, nil, errDaemonStopping + } + rc := d.reg() + if rc == nil { + return 0, nil, errors.New("re-registration: no registry connection") } var registrationAddr string @@ -5911,34 +6009,34 @@ func (d *Daemon) reRegister() { d.identityMu.RLock() pubKeyB64 := crypto.EncodePublicKey(d.identity.PublicKey) d.identityMu.RUnlock() - resp, err := d.reg().RegisterWithKeyOpts(registry.RegisterOpts{ + resp, err := rc.RegisterWithKeyOpts(registry.RegisterOpts{ ListenAddr: registrationAddr, PublicKey: pubKeyB64, Owner: d.config.Owner, - LANAddrs: d.lanAddrs, + LANAddrs: d.currentLANAddrs(), Version: d.config.Version, RelayOnly: d.config.RelayOnly, // task 32 }) if err != nil { slog.Error("re-registration failed", "error", err) - return + return 0, nil, fmt.Errorf("re-registration: %w", err) } nodeIDVal, ok := resp["node_id"].(float64) if !ok { slog.Error("re-registration: missing node_id in response") - return + return 0, nil, errors.New("re-registration: missing node_id in response") } newNodeID := uint32(nodeIDVal) addrStr, ok := resp["address"].(string) if !ok { slog.Error("re-registration: missing address in response") - return + return 0, nil, errors.New("re-registration: missing address in response") } newAddr, err := protocol.ParseAddr(addrStr) if err != nil { slog.Error("re-registration: invalid address", "address", addrStr, "error", err) - return + return 0, nil, fmt.Errorf("re-registration: invalid address %q: %w", addrStr, err) } d.addrMu.Lock() @@ -5960,7 +6058,13 @@ func (d *Daemon) reRegister() { "node_id": nodeID, "reregistered": true, }) + return nodeID, resp, nil +} +// restoreRegistryState re-applies what a registry that lost this node also +// lost: visibility and hostname, and with trust set the local trust pairs +// and the beacon registration. Best-effort; failures are logged. +func (d *Daemon) restoreRegistryState(nodeID uint32, trust bool) { if d.stopping() { return } @@ -5977,7 +6081,7 @@ func (d *Daemon) reRegister() { } } - if d.stopping() { + if d.stopping() || !trust { return } diff --git a/pkg/daemon/rxwatchdog.go b/pkg/daemon/rxwatchdog.go index 95f0cc63..3bab006e 100644 --- a/pkg/daemon/rxwatchdog.go +++ b/pkg/daemon/rxwatchdog.go @@ -236,17 +236,10 @@ func (d *Daemon) rxWatchdogResume(st *rxWatchdogState, now time.Time, gap time.D "gap_seconds": int64(gap.Seconds()), }) - d.tunnels.RegisterWithBeacon() - if d.reg() != nil { - // The pooled registry conn cannot have survived the suspend, and it - // has no half-open detection of its own, so force a fresh one rather - // than waiting for a request to hang on the dead socket. - if err := d.forceReconnectRegistry(); err != nil { - slog.Warn("registry reconnect after resume failed", "error", err) - } else { - d.reRegister() - } - } + // The pooled registry conn cannot have survived the suspend, and it + // has no half-open detection of its own, so force a fresh one rather + // than waiting for a request to hang on the dead socket. + d.reestablishTransport("resume", reestablishOpts{freshConn: true}) // Fresh epoch: re-baseline so the first post-resume tick measures // recovery, not the suspend. @@ -258,6 +251,76 @@ func (d *Daemon) rxWatchdogResume(st *rxWatchdogState, now time.Time, gap time.D d.consecutiveDialTimeouts.Store(0) } +// reestablishCoalesce is how recent a run of reestablishTransport the +// registry accepted must be for a caller that may coalesce to skip its own. +const reestablishCoalesce = 5 * time.Second + +// reestablishOpts says what a reestablishTransport caller needs. +type reestablishOpts struct { + // freshConn replaces the pooled registry connection first, for callers + // that know the old one cannot have survived (a suspend, or a change of + // our own address: the TCP connection is bound to the address we no + // longer have). + freshConn bool + // endpointOnly re-registers the endpoint only (reRegisterEndpoint) + // instead of the full reRegister with its visibility, hostname and + // trust re-sync. + endpointOnly bool + // always runs even right after another run. Without it the call is + // skipped when a run the registry accepted finished less than + // reestablishCoalesce ago and did at least the same work: a full + // re-registration for a full caller. The address watcher sets it: a run + // that started before our address changed registered the old one. + always bool +} + +// reestablishTransport re-registers this node with the beacon and the +// registry. It is the one recovery routine shared by the rx watchdog's soft +// recovery, the host-resume handler and the address watcher (addrwatch.go). +// +// RegisterWithBeacon re-punches our NAT mapping and gives the beacon our +// current endpoint for relay delivery and hole-punching; the re-registration +// gives the registry our current endpoint. Runs one at a time: a second +// caller waits for the first instead of replacing the registry connection +// under its requests. Reports whether the registry accepted the +// re-registration; false when no registry is configured. +func (d *Daemon) reestablishTransport(cause string, o reestablishOpts) bool { + d.reestablishMu.Lock() + defer d.reestablishMu.Unlock() + + // Wall clock, not the monotonic one: on Linux the monotonic clock stops + // while the host is suspended, so a run from just before a suspend would + // look recent to the resume handler. + covered := o.endpointOnly || d.reestablishOKFull + if age := time.Now().UnixNano() - d.reestablishOKWall; !o.always && covered && d.reestablishOKWall != 0 && + age >= 0 && age < int64(reestablishCoalesce) { + slog.Info("transport re-established moments ago — not repeating it", "cause", cause, + "age", time.Duration(age).Truncate(time.Millisecond).String()) + return true + } + + d.tunnels.RegisterWithBeacon() + if d.reg() == nil { + return false + } + if o.freshConn { + if err := d.forceReconnectRegistry(); err != nil { + slog.Warn("registry reconnect failed", "cause", cause, "error", err) + return false + } + } + register := d.reRegister + if o.endpointOnly { + register = d.reRegisterEndpoint + } + if err := register(); err != nil { + return false // registerEndpoint logs what the registry answered + } + d.reestablishOKWall = time.Now().UnixNano() + d.reestablishOKFull = !o.endpointOnly + return true +} + // rxWatchdogTick runs one watchdog iteration. Extracted for testability — // production drives it from rxWatchdogLoop; tests drive it directly. // @@ -349,10 +412,7 @@ func (d *Daemon) rxWatchdogTick(st *rxWatchdogState, now time.Time) (action rxWa // discover reply also doubles as an active inbound probe) and the // registry. This is what a manual restart effectively did to clear // the partial wedge. - d.tunnels.RegisterWithBeacon() - if d.reg() != nil { - d.reRegister() - } + d.reestablishTransport("rx-silence", reestablishOpts{}) // Reset the dial counter so the next real dial re-tests the path: // if recovery worked the next dial succeeds (counter stays 0); if // not, failures re-accumulate to threshold and softAttempts climbs diff --git a/pkg/daemon/tunnel.go b/pkg/daemon/tunnel.go index 02926431..cc20d5d4 100644 --- a/pkg/daemon/tunnel.go +++ b/pkg/daemon/tunnel.go @@ -101,6 +101,15 @@ type TunnelManager struct { readWg sync.WaitGroup // tracks readLoop goroutine for clean shutdown closeOnce sync.Once + // Where the beacon sees this socket, from its discover replies (see + // noteDiscoverReply; read by the address watcher). lastDiscoverNano is + // when RegisterWithBeacon last sent a discover; obsLatest and obsPrev + // are the two latest replies that answered one, under obsMu. + lastDiscoverNano atomic.Int64 + obsMu sync.Mutex + obsLatest beaconObservation + obsPrev beaconObservation + // Encryption config encrypt bool // if true, attempt encrypted tunnels privKey *ecdh.PrivateKey // our X25519 private key @@ -818,6 +827,11 @@ func (tm *TunnelManager) RelayPeerIDs() []uint32 { // using the real nodeID, so the beacon knows our endpoint for punch coordination. // Thin shim over routing.Manager.RegisterWithBeacon. func (tm *TunnelManager) RegisterWithBeacon() { + // Stamped before the send: on loopback the reply can be read before + // Send returns, and it must find the window open. + if tm.routing.BeaconAddr() != nil { + tm.lastDiscoverNano.Store(time.Now().UnixNano()) + } if err := tm.routing.RegisterWithBeacon(); err != nil { slog.Warn("beacon registration failed", "error", err) return @@ -1108,7 +1122,18 @@ func (tm *TunnelManager) SendPathProbe(peerNodeID uint32) error { if pc == nil || !pc.Ready { return fmt.Errorf("path probe: no ready session for peer %d", peerNodeID) } - probe := &protocol.Packet{ + plaintext, err := tm.newPathProbePacket(peerNodeID).Marshal() + if err != nil { + return fmt.Errorf("path probe marshal: %w", err) + } + frame := tm.encryptFrame(pc, plaintext) + return tm.writeFrame(peerNodeID, addr, frame) +} + +// newPathProbePacket builds the pong-soliciting probe described on +// SendPathProbe. +func (tm *TunnelManager) newPathProbePacket(peerNodeID uint32) *protocol.Packet { + return &protocol.Packet{ Version: protocol.Version, Protocol: protocol.ProtoControl, SrcPort: protocol.PortPing, @@ -1117,12 +1142,94 @@ func (tm *TunnelManager) SendPathProbe(peerNodeID uint32) error { Dst: protocol.Addr{Node: peerNodeID}, Payload: pathProbePayload, } - plaintext, err := probe.Marshal() - if err != nil { - return fmt.Errorf("path probe marshal: %w", err) +} + +// SendDirectPathProbe sends the same probe as SendPathProbe, but straight to +// the peer's stored direct endpoint even when the peer is relay-flagged (see +// SendDirectProbe). The address watcher uses it after our own address changed: +// the peer overwrites its entry for us from the source of any authenticated +// direct frame (handleEncrypted), so one probe moves the peer to our new +// address, and its pong proves the new path in both directions. A copy sent +// through the relay would teach the peer nothing, since relayed frames arrive +// from the beacon. +func (tm *TunnelManager) SendDirectPathProbe(peerNodeID uint32) error { + return tm.SendDirectProbe(peerNodeID, tm.newPathProbePacket(peerNodeID)) +} + +// discoverReplyWindow is how soon after this node's latest discover a +// discover reply must arrive to be taken as the beacon's answer to it. +const discoverReplyWindow = 2 * time.Second + +// beaconObservation is one discover reply: where the beacon saw this socket. +type beaconObservation struct { + endpoint *net.UDPAddr + beacon string // the beacon that replied, so replies from two beacons are never compared + at time.Time // when it arrived +} + +// noteDiscoverReply records a discover reply that came from the beacon's +// address, if it arrived within discoverReplyWindow of a discover this node +// sent. The beacon replies to every discover and sends none unasked, so a +// reply outside that window is a stray or a forgery: the source address is +// all that marks it as the beacon's, and UDP does not authenticate it. +// Reports whether the reply was kept. +func (tm *TunnelManager) noteDiscoverReply(ep, beacon *net.UDPAddr, now time.Time) bool { + sent := tm.lastDiscoverNano.Load() + if ep == nil || beacon == nil || sent == 0 { + return false } - frame := tm.encryptFrame(pc, plaintext) - return tm.writeFrame(peerNodeID, addr, frame) + if age := now.Sub(time.Unix(0, sent)); age < 0 || age > discoverReplyWindow { + slog.Debug("ignoring unsolicited beacon discover reply", "observed", ep, "since_discover", age) + return false + } + tm.obsMu.Lock() + tm.obsPrev, tm.obsLatest = tm.obsLatest, beaconObservation{endpoint: ep, beacon: beacon.String(), at: now} + tm.obsMu.Unlock() + return true +} + +// beaconObservations returns the latest two kept discover replies, newest +// first; a zero value stands for one that has not arrived. +func (tm *TunnelManager) beaconObservations() (latest, prev beaconObservation) { + tm.obsMu.Lock() + defer tm.obsMu.Unlock() + return tm.obsLatest, tm.obsPrev +} + +// ObservedEndpoint returns the endpoint the latest kept discover reply +// reported seeing us at, or nil if none has arrived yet. +func (tm *TunnelManager) ObservedEndpoint() *net.UDPAddr { + latest, _ := tm.beaconObservations() + return latest.endpoint +} + +// forgetObservedEndpoints drops the stored discover replies, so the next +// ones read arrived afterwards. +func (tm *TunnelManager) forgetObservedEndpoints() { + tm.obsMu.Lock() + tm.obsLatest, tm.obsPrev = beaconObservation{}, beaconObservation{} + tm.obsMu.Unlock() +} + +// BeaconUDPAddr returns the beacon endpoint the tunnel currently uses, or +// nil when none is configured. Thin shim over routing.Manager.BeaconAddr. +func (tm *TunnelManager) BeaconUDPAddr() *net.UDPAddr { + return tm.routing.BeaconAddr() +} + +// parseDiscoverReply decodes [iplen(1)][IP(4 or 16)][port(2)], the body of a +// BeaconMsgDiscoverReply. Returns nil for a malformed body. +func parseDiscoverReply(body []byte) *net.UDPAddr { + if len(body) < 1 { + return nil + } + ipLen := int(body[0]) + if (ipLen != 4 && ipLen != 16) || len(body) < 1+ipLen+2 { + return nil + } + ip := make(net.IP, ipLen) + copy(ip, body[1:1+ipLen]) + return &net.UDPAddr{IP: ip, Port: int(binary.BigEndian.Uint16(body[1+ipLen:]))} } // getPeerPubKey returns the cached Ed25519 public key for a peer, @@ -2311,6 +2418,16 @@ func (tm *TunnelManager) handleBeaconMessage(data []byte, from *net.UDPAddr) { switch data[0] { case protocol.BeaconMsgDiscoverReply: slog.Debug("beacon discover reply on tunnel socket", "from", from) + // Remember where the beacon sees us. A node behind NAT cannot see + // its public address change locally; this reply is the only place + // it shows up (addrwatch.go compares successive values). Only a + // reply from the beacon's address to a discover we just sent + // counts, so a third party cannot feed us one at will. + if fromBeacon { + if ep := parseDiscoverReply(data[1:]); ep != nil { + tm.noteDiscoverReply(ep, tm.routing.BeaconAddr(), time.Now()) + } + } case protocol.BeaconMsgPunchCommand: if !fromBeacon { slog.Warn("dropping punch command from non-beacon source", "from", from) diff --git a/pkg/daemon/zz_addrwatch_test.go b/pkg/daemon/zz_addrwatch_test.go new file mode 100644 index 00000000..e0e3e7d2 --- /dev/null +++ b/pkg/daemon/zz_addrwatch_test.go @@ -0,0 +1,1021 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package daemon + +import ( + "encoding/binary" + "net" + "strings" + "testing" + "time" + + "github.com/pilot-protocol/common/protocol" + "github.com/pilot-protocol/pilotprotocol/pkg/daemon/udpio" +) + +// addrRecoverCall is one invocation of the stubbed recovery. +type addrRecoverCall struct { + reason, previous, current string + announce bool +} + +// swapAddrRecoverForTest stubs the recovery hook. registryOK is what each +// stubbed recovery reports. Tests using it must NOT run in parallel with each +// other (global hook), mirroring swapPathResetForTest. +func swapAddrRecoverForTest(t *testing.T, registryOK *bool) *[]addrRecoverCall { + t.Helper() + var calls []addrRecoverCall + prev := addrWatchRecover + addrWatchRecover = func(_ *Daemon, reason, previous, current string, announce bool) bool { + calls = append(calls, addrRecoverCall{reason, previous, current, announce}) + return *registryOK + } + t.Cleanup(func() { addrWatchRecover = prev }) + return &calls +} + +// fakeAddrSource is the injectable address source: tests set ip/ok between +// ticks instead of touching real interfaces. +type fakeAddrSource struct { + ip string + ok bool +} + +func (f *fakeAddrSource) fn(*net.UDPAddr) (string, bool) { return f.ip, f.ok } + +// addrWatchTestBeacon is the beacon newAddrWatchTestDaemon's tunnel uses. +var addrWatchTestBeacon = &net.UDPAddr{IP: net.IPv4(192, 0, 2, 1), Port: 9001} + +func newAddrWatchTestDaemon() *Daemon { + d := &Daemon{ + tunnels: NewTunnelManager(), + startTime: time.Now(), + stopCh: make(chan struct{}), + } + d.tunnels.routing.SetBeaconAddrUDP(addrWatchTestBeacon) + return d +} + +// observe hands d one discover reply from its current beacon, arriving at +// now in answer to a discover sent just before. +func observe(d *Daemon, ip net.IP, now time.Time) { + d.tunnels.lastDiscoverNano.Store(now.Add(-50 * time.Millisecond).UnixNano()) + d.tunnels.noteDiscoverReply(&net.UDPAddr{IP: ip, Port: 40000}, d.tunnels.BeaconUDPAddr(), now) +} + +func TestAddrWatchNoChangeNeverFires(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + for i := 0; i < 600; i++ { + if d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) { + t.Fatalf("tick %d fired with an unchanged address", i) + } + } + if len(*calls) != 0 { + t.Fatalf("recoveries = %d, want 0", len(*calls)) + } +} + +func TestAddrWatchChangeFiresOnceOnTheFirstTick(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + d.addrWatchTick(st, src.fn, now) // baseline + + src.ip = "10.0.0.77" + if !d.addrWatchTick(st, src.fn, now.Add(time.Second)) { + t.Fatal("address change did not fire on the first tick that saw it") + } + for i := 2; i < 300; i++ { + d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) + } + if len(*calls) != 1 { + t.Fatalf("recoveries = %d, want exactly 1 for one change", len(*calls)) + } + got := (*calls)[0] + want := addrRecoverCall{addrReasonLocal, "10.0.0.4", "10.0.0.77", true} + if got != want { + t.Fatalf("recovery = %+v, want %+v", got, want) + } +} + +// An interface that drops and comes back with the same address (cable +// re-plugged, Wi-Fi roaming on one network, `docker network disconnect` then +// `connect` with the same IP) changed nothing anyone else needs to hear about. +func TestAddrWatchOfflineThenSameAddressDoesNotFire(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + d.addrWatchTick(st, src.fn, now) + + for i := 1; i <= 20; i++ { + src.ip, src.ok = "", i%2 == 0 // no route on odd ticks + if src.ok { + src.ip = "10.0.0.4" + } + d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) + } + if len(*calls) != 0 { + t.Fatalf("recoveries = %d, want 0 for an interface flapping on one address", len(*calls)) + } +} + +// Going offline is not a change, and a new address after the outage fires +// only once the route is back. +func TestAddrWatchOfflineThenNewAddressFiresWhenBack(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + d.addrWatchTick(st, src.fn, now) + + src.ip, src.ok = "", false + for i := 1; i <= 5; i++ { + if d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) { + t.Fatal("fired while offline") + } + } + src.ip, src.ok = "10.0.0.77", true + if !d.addrWatchTick(st, src.fn, now.Add(6*time.Second)) { + t.Fatal("did not fire when the route came back on a new address") + } + if len(*calls) != 1 || (*calls)[0].previous != "10.0.0.4" || (*calls)[0].current != "10.0.0.77" { + t.Fatalf("calls = %+v", *calls) + } +} + +// A daemon that starts offline has nothing to compare against: the first +// address it sees is the baseline. +func TestAddrWatchFirstAddressAfterOfflineStartIsBaseline(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{} + st := &addrWatchState{} + now := time.Now() + d.addrWatchTick(st, src.fn, now) + src.ip, src.ok = "10.0.0.4", true + d.addrWatchTick(st, src.fn, now.Add(time.Second)) + if len(*calls) != 0 { + t.Fatalf("recoveries = %d, want 0", len(*calls)) + } +} + +// A second change inside the cooldown is deferred, not dropped and not run +// early. +func TestAddrWatchChangeDuringCooldownIsDeferred(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + d.addrWatchTick(st, src.fn, now) + + src.ip = "10.0.0.5" + d.addrWatchTick(st, src.fn, now.Add(time.Second)) // recovery 1 at t=1s + + src.ip = "10.0.0.6" + for i := 2; i < 11; i++ { // t=2s..10s: inside the 10s cooldown + if d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) { + t.Fatalf("fired at t=%ds, inside the cooldown", i) + } + } + if !d.addrWatchTick(st, src.fn, now.Add(11*time.Second)) { + t.Fatal("deferred change did not run when the cooldown ended") + } + if len(*calls) != 2 || (*calls)[1].previous != "10.0.0.5" || (*calls)[1].current != "10.0.0.6" { + t.Fatalf("calls = %+v", *calls) + } +} + +// A→B→A where the second hop falls inside the cooldown of an earlier +// recovery: by the time the cooldown ends the address is the one already +// announced, so there is nothing to do. +func TestAddrWatchFlapBackToAnnouncedAddressIsCancelled(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + d.addrWatchTick(st, src.fn, now) + + src.ip = "10.0.0.5" + d.addrWatchTick(st, src.fn, now.Add(time.Second)) // announces .5 + + src.ip = "10.0.0.6" + d.addrWatchTick(st, src.fn, now.Add(2*time.Second)) + src.ip = "10.0.0.5" + for i := 3; i < 120; i++ { + d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) + } + if len(*calls) != 1 { + t.Fatalf("recoveries = %d, want 1 (the excursion to .6 was never announced)", len(*calls)) + } +} + +// A flapping interface must not turn into a re-registration storm: the gap +// between recoveries doubles up to addrRecoverCooldownMax. +func TestAddrWatchFlappingIsRateLimited(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + d.addrWatchTick(st, src.fn, now) + + var fired []int + const seconds = 30 * 60 + for i := 1; i <= seconds; i++ { + // A new address every second for half an hour. + src.ip = net.IPv4(10, 0, byte(i>>8), byte(i)).String() + if d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) { + fired = append(fired, i) + } + } + // 1s, +10s, +20s, +40s, +80s, then one per 2 minutes. + wantHead := []int{1, 11, 31, 71, 151, 271, 391} + if len(fired) < len(wantHead) { + t.Fatalf("fired at %v, want it to start %v", fired, wantHead) + } + for i, w := range wantHead { + if fired[i] != w { + t.Fatalf("fired at %v, want it to start %v", fired, wantHead) + } + } + if max := 5 + seconds/int(addrRecoverCooldownMax/time.Second); len(fired) > max { + t.Fatalf("%d recoveries in %ds of flapping, want at most %d", len(fired), seconds, max) + } + if len(*calls) != len(fired) { + t.Fatalf("recoveries = %d, fired ticks = %d", len(*calls), len(fired)) + } +} + +func TestAddrWatchCooldownResetsAfterQuiet(t *testing.T) { + st := &addrWatchState{} + now := time.Now() + st.noteLocal("10.0.0.1", true) + for i, ip := range []string{"10.0.0.2", "10.0.0.3", "10.0.0.4"} { + st.noteLocal(ip, true) + st.began(now.Add(time.Duration(i)*time.Minute), addrReasonLocal) + } + if got := st.cooldown(addrReasonLocal); got != 4*addrRecoverCooldown { + t.Fatalf("cooldown after 3 recoveries = %v, want %v", got, 4*addrRecoverCooldown) + } + st.noteLocal("10.0.0.5", true) + st.began(now.Add(2*time.Minute+addrRecoverQuiet), addrReasonLocal) + if got := st.cooldown(addrReasonLocal); got != addrRecoverCooldown { + t.Fatalf("cooldown after a quiet period = %v, want %v", got, addrRecoverCooldown) + } +} + +// The NAT case: the local address is unchanged but the beacon reports us at +// a different public IP. +func TestAddrWatchObservedEndpointChangeFires(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "192.168.1.20", ok: true} + st := &addrWatchState{} + now := time.Now() + tick := func(sec int) bool { return d.addrWatchTick(st, src.fn, now.Add(time.Duration(sec)*time.Second)) } + a, b := net.IPv4(203, 0, 113, 7), net.IPv4(198, 51, 100, 9) + + observe(d, a, now) + tick(0) + observe(d, a, now.Add(time.Second)) + tick(1) // two agreeing replies: the baseline + + // Same IP, different port: a NAT renumbering ports is not a move. + d.tunnels.lastDiscoverNano.Store(now.Add(2 * time.Second).UnixNano()) + d.tunnels.noteDiscoverReply(&net.UDPAddr{IP: a, Port: 40001}, addrWatchTestBeacon, now.Add(2*time.Second)) + if tick(2) { + t.Fatal("fired on a port-only change of the observed endpoint") + } + + observe(d, b, now.Add(3*time.Second)) + observe(d, b, now.Add(3*time.Second+100*time.Millisecond)) + if !tick(3) { + t.Fatal("did not fire when the beacon reported a new public IP twice") + } + for i := 4; i < 200; i++ { + tick(i) + } + if len(*calls) != 1 || (*calls)[0].reason != addrReasonObserved || !(*calls)[0].announce { + t.Fatalf("calls = %+v, want one %q recovery", *calls, addrReasonObserved) + } + // The event reports the IPs the beacon saw, not the unchanged local one. + if c := (*calls)[0]; c.previous != a.String() || c.current != b.String() { + t.Fatalf("recovery moved %q -> %q, want %v -> %v", c.previous, c.current, a, b) + } +} + +// One reply is not enough: a datagram from the beacon's address may be +// forged, and some NATs show us at several public IPs. Only two replies in a +// row that agree move the observed IP. +func TestAddrWatchObservedNeedsTwoAgreeingReplies(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "192.168.1.20", ok: true} + st := &addrWatchState{} + now := time.Now() + sec := 0 + step := func(ip net.IP) bool { + sec++ + at := now.Add(time.Duration(sec) * time.Second) + observe(d, ip, at) + return d.addrWatchTick(st, src.fn, at) + } + a, b := net.IPv4(203, 0, 113, 7), net.IPv4(198, 51, 100, 9) + step(a) + step(a) // baseline + + // A, B, A, B, A, ...: no two replies in a row agree on B. + for i := 0; i < 20; i++ { + ip := b + if i%2 == 1 { + ip = a + } + if step(ip) { + t.Fatalf("fired on reply %d of an alternating sequence", i) + } + } + if len(*calls) != 0 { + t.Fatalf("recoveries = %d, want 0", len(*calls)) + } + step(b) + if !step(b) { + t.Fatal("two agreeing replies on a new IP did not fire") + } +} + +// A lone reply naming a new IP is checked at once with one more discover, +// not left until the next keepalive registration a minute later. That +// discover is sent once per reply and at most once per addrObservedConfirmGap. +func TestAddrWatchLoneReplyAsksForConfirmation(t *testing.T) { + t.Parallel() + r := newAddrAnnounceRig(t) + d := r.d + beacon := d.tunnels.BeaconUDPAddr() + src := &fakeAddrSource{ip: "192.168.1.20", ok: true} + st := &addrWatchState{} + now := time.Now() + reply := func(ip net.IP, at time.Time) { + d.tunnels.lastDiscoverNano.Store(at.Add(-50 * time.Millisecond).UnixNano()) + d.tunnels.noteDiscoverReply(&net.UDPAddr{IP: ip, Port: 40000}, beacon, at) + } + confirmSent := func() bool { + msg := readOne(r.beacon, 300*time.Millisecond) + return len(msg) == 5 && msg[0] == protocol.BeaconMsgDiscover + } + a, b, c := net.IPv4(203, 0, 113, 7), net.IPv4(198, 51, 100, 9), net.IPv4(198, 51, 100, 10) + + reply(a, now) + d.addrWatchTick(st, src.fn, now) + if !confirmSent() { + t.Fatal("first reply was not confirmed with a second discover") + } + reply(a, now.Add(100*time.Millisecond)) + d.addrWatchTick(st, src.fn, now.Add(time.Second)) + if confirmSent() { + t.Fatal("sent a confirming discover for an IP two replies already agree on") + } + + reply(b, now.Add(60*time.Second)) + d.addrWatchTick(st, src.fn, now.Add(61*time.Second)) + if !confirmSent() { + t.Fatal("a reply naming a new IP was not confirmed") + } + d.addrWatchTick(st, src.fn, now.Add(62*time.Second)) + if confirmSent() { + t.Fatal("asked twice about the same reply") + } + reply(c, now.Add(63*time.Second)) + d.addrWatchTick(st, src.fn, now.Add(64*time.Second)) + if confirmSent() { + t.Fatal("confirming discovers less than addrObservedConfirmGap apart") + } + d.addrWatchTick(st, src.fn, now.Add(61*time.Second+addrObservedConfirmGap)) + if !confirmSent() { + t.Fatal("no confirming discover once the gap had passed") + } +} + +// When our own address moves, the beacon's view moves with it. That must not +// count as a second change and run a second recovery after the cooldown. +func TestAddrWatchLocalChangeDoesNotDoubleFireOnObserved(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.78.0.3", ok: true} + st := &addrWatchState{} + now := time.Now() + observe(d, net.IPv4(10, 78, 0, 3), now) + observe(d, net.IPv4(10, 78, 0, 3), now) + d.addrWatchTick(st, src.fn, now) + + src.ip = "10.78.0.77" + d.addrWatchTick(st, src.fn, now.Add(time.Second)) + if d.tunnels.ObservedEndpoint() != nil { + t.Fatal("the pre-move beacon replies were kept; they would be re-read as the new baseline") + } + d.addrWatchTick(st, src.fn, now.Add(2*time.Second)) + // The beacon's replies to the recovery's registrations land. + observe(d, net.IPv4(10, 78, 0, 77), now.Add(2*time.Second)) + observe(d, net.IPv4(10, 78, 0, 77), now.Add(2*time.Second)) + for i := 3; i < 200; i++ { + d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) + } + if len(*calls) != 1 { + t.Fatalf("recoveries = %d, want 1: %+v", len(*calls), *calls) + } +} + +// A→B→A on the observed IP inside the cooldown is cancelled, as it is for +// the local address: by the time the cooldown ends the beacon sees us where +// the last recovery already said we are. +func TestAddrWatchObservedFlapBackIsCancelled(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "192.168.1.20", ok: true} + st := &addrWatchState{} + now := time.Now() + at := func(sec int) time.Time { return now.Add(time.Duration(sec) * time.Second) } + twice := func(ip net.IP, sec int) { + observe(d, ip, at(sec)) + observe(d, ip, at(sec).Add(100*time.Millisecond)) + } + a, b := net.IPv4(203, 0, 113, 7), net.IPv4(198, 51, 100, 9) + twice(a, 0) + d.addrWatchTick(st, src.fn, at(0)) + twice(b, 1) + if !d.addrWatchTick(st, src.fn, at(1)) { // announces B + t.Fatal("first change did not fire") + } + twice(a, 10) // back to A inside the cooldown... + d.addrWatchTick(st, src.fn, at(10)) + twice(b, 20) // ...and to B again before it ends + for i := 20; i < 600; i++ { + d.addrWatchTick(st, src.fn, at(i)) + } + if len(*calls) != 1 { + t.Fatalf("recoveries = %d, want 1 (the excursion back to A was never announced): %+v", len(*calls), *calls) + } +} + +// A NAT that keeps moving us between public IPs (an address pool, per-flow +// balancing) must not cost a recovery every two minutes forever: the +// observed IP's gap grows to addrObservedCooldownMax. +func TestAddrWatchObservedFlappingBacksOffToTheLongCap(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "192.168.1.20", ok: true} + st := &addrWatchState{} + now := time.Now() + ips := []net.IP{net.IPv4(203, 0, 113, 7), net.IPv4(198, 51, 100, 9)} + + var fired []time.Duration + const hours = 12 + for sec := 0; sec <= hours*3600; sec++ { + at := now.Add(time.Duration(sec) * time.Second) + if sec%60 == 0 { // a keepalive registration; the IP flips every 2 minutes + ip := ips[(sec/120)%2] + observe(d, ip, at) + observe(d, ip, at.Add(100*time.Millisecond)) + } + if d.addrWatchTick(st, src.fn, at) { + fired = append(fired, time.Duration(sec)*time.Second) + } + } + if len(*calls) != len(fired) { + t.Fatalf("recoveries = %d, fired ticks = %d", len(*calls), len(fired)) + } + // 1m, 2m, 4m, ... 32m, then the 1h cap: about 7 in the first two hours + // and one an hour after that. The local-address schedule (2 minute cap) + // would have run over 300. + if max := 8 + hours; len(fired) > max { + t.Fatalf("%d recoveries in %dh of a flapping observed IP, want at most %d: %v", len(fired), hours, max, fired) + } + for i := 1; i < len(fired); i++ { + if gap := fired[i] - fired[i-1]; gap < addrObservedCooldown { + t.Fatalf("recoveries %v apart, under the observed-IP minimum %v: %v", gap, addrObservedCooldown, fired) + } + } + var late []time.Duration + for i := 1; i < len(fired); i++ { + if fired[i] > 4*time.Hour { + late = append(late, fired[i]-fired[i-1]) + } + } + for _, gap := range late { + if gap < addrObservedCooldownMax { + t.Fatalf("gap %v after four hours of flapping, want the %v cap: %v", gap, addrObservedCooldownMax, fired) + } + } +} + +// A beacon switch is not a move. The new beacon can be reached over another +// route (a private bootstrap beacon beside public ones) and see us at another +// public IP (NAT that picks the address per destination, dual WAN); both +// baselines start over, and the old beacon's stored replies are never +// compared with the new one's. +func TestAddrWatchBeaconSwitchIsNotAMove(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + private := &net.UDPAddr{IP: net.IPv4(10, 128, 0, 5), Port: 9001} + src := func(b *net.UDPAddr) (string, bool) { + if b.String() == private.String() { + return "10.128.0.20", true + } + return "192.168.1.20", true + } + st := &addrWatchState{} + now := time.Now() + at := func(sec int) time.Time { return now.Add(time.Duration(sec) * time.Second) } + + observe(d, net.IPv4(203, 0, 113, 7), at(0)) + observe(d, net.IPv4(203, 0, 113, 7), at(0)) + for i := 0; i < 5; i++ { + d.addrWatchTick(st, src, at(i)) + } + + // beaconRefreshTick moves us to the private beacon. Its stored replies + // are still the public beacon's for a while. + d.tunnels.routing.SetBeaconAddrUDP(private) + for i := 5; i < 10; i++ { + d.addrWatchTick(st, src, at(i)) + } + observe(d, net.IPv4(10, 128, 0, 20), at(10)) + observe(d, net.IPv4(10, 128, 0, 20), at(10)) + for i := 10; i < 600; i++ { + d.addrWatchTick(st, src, at(i)) + } + if len(*calls) != 0 { + t.Fatalf("a beacon switch ran %d recoveries, want 0: %+v", len(*calls), *calls) + } + + // Moves seen through the new beacon are still caught. + observe(d, net.IPv4(10, 128, 0, 21), at(600)) + observe(d, net.IPv4(10, 128, 0, 21), at(600)) + if !d.addrWatchTick(st, src, at(600)) { + t.Fatal("an observed change through the new beacon did not fire") + } +} + +// If the registry could not be reached right after the change, only the +// registry half is retried, a bounded number of times, one per cooldown. +func TestAddrWatchRegistryRetryIsBounded(t *testing.T) { + ok := false + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + d.addrWatchTick(st, src.fn, now) + src.ip = "10.0.0.77" + for i := 1; i < 3600; i++ { + d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) + } + if len(*calls) != 1+addrRegistryRetryMax { + t.Fatalf("recoveries = %d, want 1 + %d registry retries", len(*calls), addrRegistryRetryMax) + } + for i, c := range (*calls)[1:] { + if c.reason != addrReasonRegistry || c.announce { + t.Fatalf("retry %d = %+v, want a registry-only retry", i, c) + } + } + + // A retry that succeeds stops the retries. + ok = false + *calls = nil + st = &addrWatchState{} + d.addrWatchTick(st, src.fn, now) + src.ip = "10.0.0.78" + d.addrWatchTick(st, src.fn, now.Add(time.Second)) + ok = true + for i := 2; i < 3600; i++ { + d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) + } + if len(*calls) != 2 { + t.Fatalf("recoveries = %d, want 2 (the change, then one successful retry)", len(*calls)) + } +} + +// The production source: the kernel's choice of source address for a target. +// Toward loopback that is loopback, every time; an unrelated interface +// appearing, a link-local address or a rotating IPv6 temporary address cannot +// change the answer unless the route to the target itself moves. +func TestRouteSourceIP(t *testing.T) { + t.Parallel() + target := &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 9001} + for i := 0; i < 3; i++ { + ip, ok := routeSourceIP(target) + if !ok || ip != "127.0.0.1" { + t.Fatalf("routeSourceIP(loopback) = %q, %v; want 127.0.0.1, true", ip, ok) + } + } + if ip, ok := routeSourceIP(nil); ok || ip != "" { + t.Fatalf("routeSourceIP(nil) = %q, %v; want no address", ip, ok) + } +} + +func TestParseDiscoverReply(t *testing.T) { + t.Parallel() + v4 := []byte{4, 203, 0, 113, 7, 0x9c, 0x40} + if got := parseDiscoverReply(v4); got == nil || got.String() != "203.0.113.7:40000" { + t.Fatalf("v4 = %v", got) + } + v6 := append([]byte{16}, net.ParseIP("2001:db8::1")...) + v6 = binary.BigEndian.AppendUint16(v6, 4000) + if got := parseDiscoverReply(v6); got == nil || got.String() != "[2001:db8::1]:4000" { + t.Fatalf("v6 = %v", got) + } + for _, bad := range [][]byte{nil, {}, {4}, {4, 1, 2, 3, 4}, {4, 1, 2, 3, 4, 0}, {5, 1, 2, 3, 4, 5, 0, 1}, {16, 1, 2}} { + if got := parseDiscoverReply(bad); got != nil { + t.Fatalf("parseDiscoverReply(%v) = %v, want nil", bad, got) + } + } +} + +// Only the beacon's own reply may set the observed endpoint: anyone can send +// a datagram that looks like a discover reply to the tunnel port. +func TestObservedEndpointOnlyFromBeacon(t *testing.T) { + t.Parallel() + tm := NewTunnelManager() + if err := tm.SetBeaconAddr("192.0.2.1:9001"); err != nil { + t.Fatalf("SetBeaconAddr: %v", err) + } + reply := []byte{protocol.BeaconMsgDiscoverReply, 4, 203, 0, 113, 7, 0x9c, 0x40} + tm.RegisterWithBeacon() // opens the reply window; no socket, so nothing is sent + + tm.handleBeaconMessage(reply, &net.UDPAddr{IP: net.IPv4(198, 51, 100, 66), Port: 9001}) + if got := tm.ObservedEndpoint(); got != nil { + t.Fatalf("observed endpoint set from a non-beacon source: %v", got) + } + tm.handleBeaconMessage(reply, &net.UDPAddr{IP: net.IPv4(192, 0, 2, 1), Port: 9001}) + if got := tm.ObservedEndpoint(); got == nil || got.String() != "203.0.113.7:40000" { + t.Fatalf("observed endpoint = %v, want 203.0.113.7:40000", got) + } +} + +// A reply from the beacon's address counts only as the answer to a discover +// this node sent less than discoverReplyWindow earlier. The beacon never +// replies unasked, and the source address of a UDP datagram is easy to forge. +func TestDiscoverReplyMustAnswerADiscover(t *testing.T) { + t.Parallel() + tm := NewTunnelManager() + beacon := &net.UDPAddr{IP: net.IPv4(192, 0, 2, 1), Port: 9001} + tm.routing.SetBeaconAddrUDP(beacon) + ep := &net.UDPAddr{IP: net.IPv4(203, 0, 113, 7), Port: 40000} + now := time.Now() + + if tm.noteDiscoverReply(ep, beacon, now) { + t.Fatal("kept a reply although no discover was ever sent") + } + tm.lastDiscoverNano.Store(now.UnixNano()) + if tm.noteDiscoverReply(ep, beacon, now.Add(discoverReplyWindow+time.Millisecond)) { + t.Fatal("kept a reply that arrived after the window") + } + if tm.noteDiscoverReply(ep, beacon, now.Add(-time.Millisecond)) { + t.Fatal("kept a reply that arrived before the discover was sent") + } + if !tm.noteDiscoverReply(ep, beacon, now.Add(300*time.Millisecond)) { + t.Fatal("dropped the answer to a discover") + } + latest, prev := tm.beaconObservations() + if latest.endpoint != ep || latest.beacon != beacon.String() || prev.endpoint != nil { + t.Fatalf("stored latest=%+v prev=%+v", latest, prev) + } +} + +// addrAnnounceRig is a daemon with a real loopback tunnel socket and one +// peer that is a plain UDP listener, so the test sees exactly which datagrams +// the announce puts on the wire and where. +type addrAnnounceRig struct { + d *Daemon + peer *net.UDPConn + beacon *net.UDPConn +} + +const addrAnnouncePeer = 77 + +func newAddrAnnounceRig(t *testing.T) *addrAnnounceRig { + t.Helper() + listen := func() *net.UDPConn { + c, err := net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)}) + if err != nil { + t.Fatalf("listen: %v", err) + } + t.Cleanup(func() { c.Close() }) + return c + } + conn, peer, beacon := listen(), listen(), listen() + + d := newAddrWatchTestDaemon() + d.bus = newInProcessBus(func() uint32 { return 1 }) + d.tunnels.sock = udpio.WrapConn(conn) + d.tunnels.routing.SetSocket(d.tunnels.sock) + if err := d.tunnels.EnableEncryption(); err != nil { + t.Fatalf("EnableEncryption: %v", err) + } + if err := d.tunnels.SetBeaconAddr(beacon.LocalAddr().String()); err != nil { + t.Fatalf("SetBeaconAddr: %v", err) + } + // Install the session before the address so AddPeer's key exchange + // (which would also write to the peer) is the only extra datagram. + pc := fakePC(t) + pc.Authenticated = true + d.tunnels.envelope.Install(addrAnnouncePeer, pc) + d.tunnels.mu.Lock() + d.tunnels.peers[addrAnnouncePeer] = peer.LocalAddr().(*net.UDPAddr) + d.tunnels.mu.Unlock() + return &addrAnnounceRig{d: d, peer: peer, beacon: beacon} +} + +// readOne returns the next datagram on c, or nil if none arrives in time. +func readOne(c *net.UDPConn, wait time.Duration) []byte { + buf := make([]byte, 2048) + _ = c.SetReadDeadline(time.Now().Add(wait)) + n, _, err := c.ReadFromUDP(buf) + if err != nil { + return nil + } + return buf[:n] +} + +func TestAnnounceToPeersSendsDirectEncryptedProbe(t *testing.T) { + t.Parallel() + r := newAddrAnnounceRig(t) + if got := r.d.announceToPeers(time.Time{}); got != 1 { + t.Fatalf("announced to %d peers, want 1", got) + } + frame := readOne(r.peer, 2*time.Second) + if frame == nil { + t.Fatal("peer received nothing") + } + if len(frame) < 4 || string(frame[:4]) != string(protocol.TunnelMagicSecure[:]) { + t.Fatalf("peer received a frame that is not an encrypted tunnel frame: % x", frame[:min(len(frame), 8)]) + } +} + +// A peer we flipped to relay while our own address was gone (send errors +// during the outage) must still be told directly: a relayed frame reaches it +// from the beacon and teaches it nothing. +func TestAnnounceToPeersBypassesRelayFlag(t *testing.T) { + t.Parallel() + r := newAddrAnnounceRig(t) + r.d.tunnels.SetRelayPeerPinned(addrAnnouncePeer, true) + if got := r.d.announceToPeers(time.Time{}); got != 1 { + t.Fatalf("announced to %d peers, want 1", got) + } + if readOne(r.peer, 2*time.Second) == nil { + t.Fatal("relay-flagged peer was not probed directly") + } + if got := readOne(r.beacon, 200*time.Millisecond); got != nil { + t.Fatalf("announce went through the beacon (%d bytes)", len(got)) + } +} + +// A relay-only node hides its real address; a direct probe would reveal it. +func TestAnnounceToPeersRelayOnlyNodeStaysSilent(t *testing.T) { + t.Parallel() + r := newAddrAnnounceRig(t) + r.d.config.RelayOnly = true + if got := r.d.announceToPeers(time.Time{}); got != 0 { + t.Fatalf("relay-only node announced to %d peers, want 0", got) + } + if got := readOne(r.peer, 200*time.Millisecond); got != nil { + t.Fatal("relay-only node sent a direct frame to a peer") + } +} + +// A peer whose only known address is the beacon placeholder has no direct +// endpoint to probe. +func TestAnnounceToPeersSkipsBeaconPlaceholder(t *testing.T) { + t.Parallel() + r := newAddrAnnounceRig(t) + r.d.tunnels.mu.Lock() + r.d.tunnels.peers[addrAnnouncePeer] = r.beacon.LocalAddr().(*net.UDPAddr) + r.d.tunnels.mu.Unlock() + if got := r.d.announceToPeers(time.Time{}); got != 0 { + t.Fatalf("announced to %d peers, want 0", got) + } +} + +// The retry rounds leave alone a peer that has already answered directly +// since the recovery started. +func TestAnnounceToPeersSkipsPeersThatAnswered(t *testing.T) { + t.Parallel() + r := newAddrAnnounceRig(t) + started := time.Now() + if got := r.d.announceToPeers(started); got != 1 { + t.Fatalf("unanswered peer: announced to %d, want 1", got) + } + r.d.tunnels.routing.RecordDirectRecv(addrAnnouncePeer, started.Add(time.Millisecond)) + if got := r.d.announceToPeers(started); got != 0 { + t.Fatalf("answered peer: announced to %d, want 0", got) + } +} + +// The whole recovery without a registry: beacon registration, peer probe and +// the event on the bus. +func TestRecoverFromAddrChangePublishesEventAndNotifies(t *testing.T) { + t.Parallel() + r := newAddrAnnounceRig(t) + events, cancel := r.d.bus.Subscribe("tunnel.addr_changed") + defer cancel() + + if r.d.recoverFromAddrChange(addrReasonLocal, "10.0.0.4", "10.0.0.77", true) { + t.Fatal("reported registry success with no registry configured") + } + close(r.d.stopCh) // end the retry goroutine + r.d.bgWG.Wait() + + if readOne(r.peer, 2*time.Second) == nil { + t.Fatal("peer was not probed") + } + msg := readOne(r.beacon, 2*time.Second) + if len(msg) != 5 || msg[0] != protocol.BeaconMsgDiscover { + t.Fatalf("beacon did not receive a discover: % x", msg) + } + select { + case ev := <-events: + p := ev.Payload + if p["reason"] != addrReasonLocal || p["previous"] != "10.0.0.4" || p["current"] != "10.0.0.77" || + p["peers_notified"] != 1 || p["registry_ok"] != false { + t.Fatalf("event payload = %+v", p) + } + case <-time.After(2 * time.Second): + t.Fatal("no tunnel.addr_changed event") + } +} + +// The beacon takes at most one endpoint update per node every 30 s and drops +// the rest. A move within 30 s of a keepalive registration has the +// recovery's own registrations dropped, so the recovery registers once more +// after addrBeaconReregisterDelay. +func TestRecoverFromAddrChangeRepeatsBeaconRegistrationAfterTheBeaconLimit(t *testing.T) { + // Not parallel: swaps package-level delays. + prevRetries, prevDelay := addrAnnounceRetryDelays, addrBeaconReregisterDelay + addrAnnounceRetryDelays = []time.Duration{10 * time.Millisecond} + addrBeaconReregisterDelay = 600 * time.Millisecond + t.Cleanup(func() { addrAnnounceRetryDelays, addrBeaconReregisterDelay = prevRetries, prevDelay }) + + r := newAddrAnnounceRig(t) + defer func() { + close(r.d.stopCh) + r.d.bgWG.Wait() + }() + start := time.Now() + r.d.recoverFromAddrChange(addrReasonLocal, "10.0.0.4", "10.0.0.77", true) + + // The recovery's own registrations are already queued at the beacon. + immediate := 0 + for readOne(r.beacon, 150*time.Millisecond) != nil { + immediate++ + } + if immediate == 0 { + t.Fatal("the recovery did not register with the beacon") + } + msg := readOne(r.beacon, 3*time.Second) + if len(msg) != 5 || msg[0] != protocol.BeaconMsgDiscover { + t.Fatalf("no beacon registration after the recovery's own (got % x)", msg) + } + if since := time.Since(start); since < addrBeaconReregisterDelay { + t.Fatalf("repeat registration %v after the start, want it after %v (the beacon would drop it)", since, addrBeaconReregisterDelay) + } +} + +// -no-addr-watch (Config.DisableAddrWatch) keeps the watcher from running at +// all; without it the loop runs until the daemon stops. +func TestAddrWatchLoopHonoursDisableAddrWatch(t *testing.T) { + t.Parallel() + d := newAddrWatchTestDaemon() + d.config.DisableAddrWatch = true + done := make(chan struct{}) + go func() { d.addrWatchLoop(); close(done) }() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("addrWatchLoop ran with DisableAddrWatch set") + } + + d = newAddrWatchTestDaemon() + done = make(chan struct{}) + go func() { d.addrWatchLoop(); close(done) }() + select { + case <-done: + t.Fatal("addrWatchLoop returned at once without DisableAddrWatch") + case <-time.After(1500 * time.Millisecond): + } + close(d.stopCh) + <-done +} + +// A move to .5, then to .6 inside the cooldown, then a beacon switch: the +// switch starts the baselines over, but .6 was never announced and must +// still be once the cooldown ends. +func TestAddrWatchChangeDeferredByCooldownSurvivesBeaconSwitch(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + at := func(sec int) time.Time { return now.Add(time.Duration(sec) * time.Second) } + d.addrWatchTick(st, src.fn, at(0)) + + src.ip = "10.0.0.5" + d.addrWatchTick(st, src.fn, at(1)) // announces .5 + src.ip = "10.0.0.6" + d.addrWatchTick(st, src.fn, at(2)) // deferred by the cooldown + d.tunnels.routing.SetBeaconAddrUDP(&net.UDPAddr{IP: net.IPv4(192, 0, 2, 2), Port: 9001}) + for i := 3; i < 120; i++ { + d.addrWatchTick(st, src.fn, at(i)) + } + if len(*calls) != 2 { + t.Fatalf("recoveries = %d, want 2 (.5, then .6 once the cooldown ended): %+v", len(*calls), *calls) + } + if c := (*calls)[1]; c.reason != addrReasonLocal || c.previous != "10.0.0.5" || c.current != "10.0.0.6" { + t.Fatalf("second recovery = %+v, want .5 -> .6", c) + } +} + +// A move and a beacon switch in the same tick: the old route shows the move, +// so the switch does not swallow it as a new baseline. +func TestAddrWatchMoveInTheSameTickAsABeaconSwitchRuns(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now() + for i := 0; i < 5; i++ { + d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) + } + + src.ip = "10.0.0.77" + d.tunnels.routing.SetBeaconAddrUDP(&net.UDPAddr{IP: net.IPv4(192, 0, 2, 2), Port: 9001}) + if !d.addrWatchTick(st, src.fn, now.Add(5*time.Second)) { + t.Fatal("a move in the same tick as a beacon switch did not run a recovery") + } + for i := 6; i < 120; i++ { + d.addrWatchTick(st, src.fn, now.Add(time.Duration(i)*time.Second)) + } + if len(*calls) != 1 || (*calls)[0].previous != "10.0.0.4" || (*calls)[0].current != "10.0.0.77" { + t.Fatalf("calls = %+v, want one .4 -> .77 recovery", *calls) + } +} + +// The watcher's cooldowns must count the time the host spent asleep. Go +// takes the difference of two readings that both carry a monotonic reading +// on the monotonic clock, which on Linux stops during suspend. A test cannot +// suspend the host, so this checks the clock the loop hands addrWatchTick: +// it must carry no monotonic reading, which makes every gap the watcher +// measures a wall-clock one. +func TestAddrWatchClockIsTheWallClock(t *testing.T) { + t.Parallel() + if s := addrWatchNow().String(); strings.Contains(s, " m=") { + t.Fatalf("addrWatchNow() = %s carries a monotonic reading; a cooldown would not count a suspend", s) + } +} + +// The wall clock can be stepped back (NTP, a VM restored from a snapshot). +// A last recovery that now lies in the future says nothing about how long +// ago it was, and must not hold a real change back until the clock catches +// up. +func TestAddrWatchClockSteppedBackDoesNotBlockRecovery(t *testing.T) { + ok := true + calls := swapAddrRecoverForTest(t, &ok) + d := newAddrWatchTestDaemon() + src := &fakeAddrSource{ip: "10.0.0.4", ok: true} + st := &addrWatchState{} + now := time.Now().Round(0) + d.addrWatchTick(st, src.fn, now) + src.ip = "10.0.0.5" + d.addrWatchTick(st, src.fn, now.Add(time.Second)) + + stepped := now.Add(-time.Hour) + src.ip = "10.0.0.6" + if !d.addrWatchTick(st, src.fn, stepped) { + t.Fatalf("a change after the clock stepped back an hour waited: %+v", *calls) + } +} diff --git a/pkg/daemon/zz_reestablish_transport_test.go b/pkg/daemon/zz_reestablish_transport_test.go new file mode 100644 index 00000000..105b9b6c --- /dev/null +++ b/pkg/daemon/zz_reestablish_transport_test.go @@ -0,0 +1,369 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package daemon + +import ( + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/pilot-protocol/common/crypto" + "github.com/pilot-protocol/common/protocol" + registry "github.com/pilot-protocol/common/registry/client" +) + +// registryLog counts the requests a fake registry received, by type. +type registryLog struct { + mu sync.Mutex + count map[string]int +} + +func (l *registryLog) add(typ string) { + l.mu.Lock() + defer l.mu.Unlock() + if l.count == nil { + l.count = map[string]int{} + } + l.count[typ]++ +} + +func (l *registryLog) get(typ string) int { + l.mu.Lock() + defer l.mu.Unlock() + return l.count[typ] +} + +func (l *registryLog) reset() { + l.mu.Lock() + defer l.mu.Unlock() + l.count = nil +} + +// registerOK is a fake registry's reply to a register for node 7. +func registerOK(hostname string) map[string]interface{} { + resp := map[string]interface{}{ + "type": "register_ok", + "node_id": float64(7), + "address": protocol.Addr{Node: 7}.String(), + } + if hostname != "" { + resp["hostname"] = hostname + } + return resp +} + +// newRegistryTestDaemon is a daemon registered as node 7 with the fake +// registry at addr, which it can also dial again (forceReconnectRegistry). +func newRegistryTestDaemon(t *testing.T, addr string, cfg Config) *Daemon { + t.Helper() + id, err := crypto.GenerateIdentity() + if err != nil { + t.Fatalf("gen identity: %v", err) + } + cfg.RegistryAddr = addr + if cfg.ListenAddr == "" { + cfg.ListenAddr = "127.0.0.1:4000" + } + d := New(cfg) + d.identity = id + rc, err := registry.Dial(addr) + if err != nil { + t.Fatalf("dial fake registry: %v", err) + } + rc.SetSigner(func(challenge string) string { return "test-sig:" + challenge }) + d.regConn.Store(rc) + t.Cleanup(func() { + if c := d.reg(); c != nil { + c.Close() + } + }) + d.setNodeID_testhelper(7) + d.lastRegistryOKNano.Store(time.Now().UnixNano()) + return d +} + +// An address change leaves the registry holding everything about this node +// except where it is. The re-registration must say only that: a node that +// trusts the service fleet would otherwise send a signed ReportTrust per +// trusted peer, plus SetVisibility and SetHostname, on every address change. +func TestReRegisterEndpointWritesOnlyTheEndpoint(t *testing.T) { + t.Parallel() + var log registryLog + hostname := "edge-1" + echo := hostname + var echoMu sync.Mutex + nodeID := float64(7) + addr, stop := serveFakeRegistry(t, func(req map[string]interface{}) map[string]interface{} { + typ, _ := req["type"].(string) + log.add(typ) + if typ != "register" { + return map[string]interface{}{"type": typ + "_ok"} + } + echoMu.Lock() + defer echoMu.Unlock() + resp := registerOK(echo) + resp["node_id"] = nodeID + resp["address"] = protocol.Addr{Node: uint32(nodeID)}.String() + return resp + }) + defer stop() + d := newRegistryTestDaemon(t, addr, Config{Public: true, Hostname: hostname}) + fs := installFakeHandshake(d) + fs.trustedRecs = []HandshakeTrustRecord{{NodeID: 41}, {NodeID: 42}, {NodeID: 43}} + + if err := d.reRegisterEndpoint(); err != nil { + t.Fatalf("reRegisterEndpoint: %v", err) + } + if log.get("register") != 1 || log.get("report_trust") != 0 || log.get("set_visibility") != 0 || log.get("set_hostname") != 0 { + t.Fatalf("endpoint-only re-registration sent %v, want one register and nothing else", log.count) + } + + // The registry may have reaped the node while it could not reach us: + // no answer for endpointOnlyRegistryFresh means the full restore. + log.reset() + d.lastRegistryOKNano.Store(time.Now().Add(-endpointOnlyRegistryFresh - time.Minute).UnixNano()) + if err := d.reRegisterEndpoint(); err != nil { + t.Fatalf("reRegisterEndpoint after silence: %v", err) + } + if log.get("report_trust") != 3 || log.get("set_visibility") != 1 || log.get("set_hostname") != 1 { + t.Fatalf("after a long silence sent %v, want the full restore", log.count) + } + + // A reply without our hostname: the node came back without it, and + // without its visibility. Trust pairs are keyed by node ID and survive. + log.reset() + echoMu.Lock() + echo = "" + echoMu.Unlock() + if err := d.reRegisterEndpoint(); err != nil { + t.Fatalf("reRegisterEndpoint: %v", err) + } + if log.get("set_visibility") != 1 || log.get("set_hostname") != 1 || log.get("report_trust") != 0 { + t.Fatalf("hostname missing from the reply: sent %v, want visibility and hostname restored", log.count) + } + + // A different node ID: the registry did not know this key any more. + log.reset() + echoMu.Lock() + echo, nodeID = hostname, 8 + echoMu.Unlock() + if err := d.reRegisterEndpoint(); err != nil { + t.Fatalf("reRegisterEndpoint: %v", err) + } + if log.get("report_trust") != 3 { + t.Fatalf("new node ID: sent %v, want the trust pairs re-synced", log.count) + } +} + +// The address-change recovery uses the endpoint-only re-registration; the +// resume handler keeps the full one, since a long suspend can outlast the +// registry's reap threshold. +func TestAddrChangeRecoveryRegistersEndpointOnly(t *testing.T) { + t.Parallel() + var log registryLog + addr, stop := serveFakeRegistry(t, func(req map[string]interface{}) map[string]interface{} { + typ, _ := req["type"].(string) + log.add(typ) + if typ == "register" { + return registerOK("") + } + return map[string]interface{}{"type": typ + "_ok"} + }) + defer stop() + d := newRegistryTestDaemon(t, addr, Config{Public: true}) + fs := installFakeHandshake(d) + fs.trustedRecs = []HandshakeTrustRecord{{NodeID: 41}, {NodeID: 42}} + + if !d.recoverFromAddrChange(addrReasonLocal, "10.0.0.4", "10.0.0.77", false) { + t.Fatal("recovery reported a registry failure") + } + if log.get("register") != 1 || log.get("report_trust") != 0 || log.get("set_visibility") != 0 { + t.Fatalf("address-change recovery sent %v, want one register and nothing else", log.count) + } + + log.reset() + d.reestablishOKWall = 0 // not coalesced with the run above + if !d.reestablishTransport("resume", reestablishOpts{freshConn: true}) { + t.Fatal("resume re-establishment failed") + } + if log.get("report_trust") != 2 || log.get("set_visibility") != 1 { + t.Fatalf("resume sent %v, want the full re-registration", log.count) + } +} + +// registry_ok must say what the registry answered. A heartbeat that lands +// while the re-registration fails moves lastRegistryOKNano too, and must not +// make the failure read as success. +func TestReestablishTransportReportsAFailedRegistration(t *testing.T) { + t.Parallel() + var d *Daemon + var ready sync.WaitGroup + ready.Add(1) + addr, stop := serveFakeRegistry(t, func(req map[string]interface{}) map[string]interface{} { + if req["type"] == "register" { + ready.Wait() + d.lastRegistryOKNano.Store(time.Now().UnixNano()) // a heartbeat, concurrently + return map[string]interface{}{"type": "error", "error": "registry overloaded"} + } + return map[string]interface{}{"type": "ok"} + }) + defer stop() + d = newRegistryTestDaemon(t, addr, Config{}) + ready.Done() + + for _, o := range []reestablishOpts{{}, {freshConn: true, endpointOnly: true, always: true}} { + if d.reestablishTransport("test", o) { + t.Fatalf("reestablishTransport(%+v) reported success for a rejected re-registration", o) + } + } +} + +// Two recoveries at once — the resume handler and the address watcher both +// fire when a laptop wakes on a new network — must not abort each other: +// each forceReconnectRegistry closes the connection the other may be in the +// middle of a request on. They run one after the other. +func TestReestablishTransportRunsOneAtATime(t *testing.T) { + t.Parallel() + var inFlight, overlapped atomic.Int32 + registering := make(chan struct{}, 4) + addr, stop := serveFakeRegistry(t, func(req map[string]interface{}) map[string]interface{} { + if req["type"] != "register" { + return map[string]interface{}{"type": "ok"} + } + if inFlight.Add(1) > 1 { + overlapped.Store(1) + } + registering <- struct{}{} + time.Sleep(300 * time.Millisecond) + inFlight.Add(-1) + return registerOK("") + }) + defer stop() + d := newRegistryTestDaemon(t, addr, Config{}) + + results := make(chan bool, 2) + go func() { results <- d.reestablishTransport("resume", reestablishOpts{freshConn: true}) }() + <-registering // the first run is waiting on its register + go func() { + results <- d.reestablishTransport("addr-change", reestablishOpts{freshConn: true, endpointOnly: true, always: true}) + }() + for i := 0; i < 2; i++ { + select { + case ok := <-results: + if !ok { + t.Fatal("a re-registration failed: the other run replaced the registry connection under it") + } + case <-time.After(30 * time.Second): + t.Fatal("re-establishment did not finish") + } + } + if overlapped.Load() != 0 { + t.Fatal("two re-registrations were in flight at once") + } +} + +// A run the registry accepted moments ago did the work a resume or rx-silence +// recovery wants, so they skip theirs. The address watcher never skips: the +// run before may have registered the address we just left. +func TestReestablishTransportCoalescesRecentRuns(t *testing.T) { + t.Parallel() + var log registryLog + addr, stop := serveFakeRegistry(t, func(req map[string]interface{}) map[string]interface{} { + typ, _ := req["type"].(string) + log.add(typ) + if typ == "register" { + return registerOK("") + } + return map[string]interface{}{"type": "ok"} + }) + defer stop() + d := newRegistryTestDaemon(t, addr, Config{}) + + run := func(cause string, o reestablishOpts) { + t.Helper() + if !d.reestablishTransport(cause, o) { + t.Fatalf("%s: reported failure", cause) + } + } + run("resume", reestablishOpts{freshConn: true}) + run("resume", reestablishOpts{freshConn: true}) + run("rx-silence", reestablishOpts{}) + if got := log.get("register"); got != 1 { + t.Fatalf("registers = %d after three runs within %v, want 1", got, reestablishCoalesce) + } + run("addr-change", reestablishOpts{freshConn: true, endpointOnly: true, always: true}) + if got := log.get("register"); got != 2 { + t.Fatalf("registers = %d, want the address change to run regardless", got) + } + // The address change re-registered the endpoint only. A resume wants the + // full re-registration, so it does not count as done. + run("resume", reestablishOpts{freshConn: true}) + if got := log.get("register"); got != 3 { + t.Fatalf("registers = %d, want a full caller to run after an endpoint-only run", got) + } + run("rx-silence", reestablishOpts{}) + if got := log.get("register"); got != 3 { + t.Fatalf("registers = %d, want rx-silence skipped right after a full run", got) + } + + d.reestablishMu.Lock() + d.reestablishOKWall = time.Now().Add(-2 * reestablishCoalesce).UnixNano() + d.reestablishMu.Unlock() + run("resume", reestablishOpts{freshConn: true}) + if got := log.get("register"); got != 4 { + t.Fatalf("registers = %d, want a run once the last one is older than %v", got, reestablishCoalesce) + } +} + +// The heartbeat's own reconnect and re-registration (trustRepublishLoop) +// take reestablishMu too, so they cannot replace or use the registry +// connection under a recovery's requests either: while a recovery holds the +// mutex they wait for it. +func TestHeartbeatRegistryCallsWaitForARecovery(t *testing.T) { + t.Parallel() + var log registryLog + addr, stop := serveFakeRegistry(t, func(req map[string]interface{}) map[string]interface{} { + typ, _ := req["type"].(string) + log.add(typ) + if typ == "register" { + return registerOK("") + } + return map[string]interface{}{"type": "ok"} + }) + defer stop() + d := newRegistryTestDaemon(t, addr, Config{}) + + for _, call := range []struct { + name string + fn func() error + }{ + {"reconnect", d.reconnectRegistrySerialised}, + {"re-register", d.reRegisterSerialised}, + } { + before := d.reg() + registers := log.get("register") + d.reestablishMu.Lock() // a recovery in progress + done := make(chan error, 1) + go func() { done <- call.fn() }() + select { + case err := <-done: + d.reestablishMu.Unlock() + t.Fatalf("heartbeat %s ran while a recovery held the lock (err %v)", call.name, err) + case <-time.After(200 * time.Millisecond): + } + if d.reg() != before || log.get("register") != registers { + d.reestablishMu.Unlock() + t.Fatalf("heartbeat %s touched the registry while a recovery held the lock", call.name) + } + d.reestablishMu.Unlock() + select { + case err := <-done: + if err != nil { + t.Fatalf("heartbeat %s: %v", call.name, err) + } + case <-time.After(10 * time.Second): + t.Fatalf("heartbeat %s never ran once the recovery let go", call.name) + } + } +} diff --git a/pkg/daemon/zz_test_helpers_registry_test.go b/pkg/daemon/zz_test_helpers_registry_test.go index b6b01066..a0caa604 100644 --- a/pkg/daemon/zz_test_helpers_registry_test.go +++ b/pkg/daemon/zz_test_helpers_registry_test.go @@ -24,6 +24,38 @@ import ( // is still consumed by handshake_dispatch_test.go and a few daemon-side // info/wire tests, so we keep it here in a dedicated _test.go. func startFakeRegistry(t *testing.T, respond func(req map[string]interface{}) map[string]interface{}) (*registry.Client, func()) { + t.Helper() + addr, stop := serveFakeRegistry(t, respond) + client, err := registry.Dial(addr) + if err != nil { + stop() + t.Fatalf("dial fake registry: %v", err) + } + + // Configure a no-op test signer. After PILOT-128 the registry client + // refuses to send any signed request when no signer is set. The fake + // registry doesn't verify signatures, so any non-empty string suffices. + // Tests that need the no-signer error path should call client.SetSigner(nil). + client.SetSigner(func(challenge string) string { + return "test-sig:" + challenge + }) + + var once sync.Once + cleanup := func() { + once.Do(func() { + client.Close() + stop() + }) + } + return client, cleanup +} + +// serveFakeRegistry is the server half of startFakeRegistry, for tests that +// need the address itself (a daemon that dials its own registry connections, +// such as forceReconnectRegistry). It accepts any number of connections. +// The returned stop closes the listener and waits for every connection +// handler to return. +func serveFakeRegistry(t *testing.T, respond func(req map[string]interface{}) map[string]interface{}) (string, func()) { t.Helper() lis, err := net.Listen("tcp", "127.0.0.1:0") if err != nil { @@ -32,9 +64,9 @@ func startFakeRegistry(t *testing.T, respond func(req map[string]interface{}) ma var ( wg sync.WaitGroup - quit = make(chan struct{}) - quitMu sync.Mutex + mu sync.Mutex closed bool + conns []net.Conn ) wg.Add(1) @@ -45,6 +77,16 @@ func startFakeRegistry(t *testing.T, respond func(req map[string]interface{}) ma if err != nil { return } + // Server-side conns are closed by stop, so a client the test + // never closes cannot keep a handler (and stop) waiting. + mu.Lock() + if closed { + mu.Unlock() + conn.Close() + return + } + conns = append(conns, conn) + mu.Unlock() wg.Add(1) go func(c net.Conn) { defer wg.Done() @@ -84,32 +126,19 @@ func startFakeRegistry(t *testing.T, respond func(req map[string]interface{}) ma } }() - client, err := registry.Dial(lis.Addr().String()) - if err != nil { - lis.Close() - t.Fatalf("dial fake registry: %v", err) - } - - // Configure a no-op test signer. After PILOT-128 the registry client - // refuses to send any signed request when no signer is set. The fake - // registry doesn't verify signatures, so any non-empty string suffices. - // Tests that need the no-signer error path should call client.SetSigner(nil). - client.SetSigner(func(challenge string) string { - return "test-sig:" + challenge - }) - - cleanup := func() { - quitMu.Lock() + stop := func() { + mu.Lock() if closed { - quitMu.Unlock() + mu.Unlock() return } closed = true - close(quit) - quitMu.Unlock() - client.Close() lis.Close() + for _, c := range conns { + c.Close() + } + mu.Unlock() wg.Wait() } - return client, cleanup + return lis.Addr().String(), stop } diff --git a/tests/zz_addr_change_recovery_test.go b/tests/zz_addr_change_recovery_test.go new file mode 100644 index 00000000..877146db --- /dev/null +++ b/tests/zz_addr_change_recovery_test.go @@ -0,0 +1,149 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package tests + +import ( + "net" + "testing" + "time" + + "github.com/pilot-protocol/pilotprotocol/pkg/daemon" + "github.com/pilot-protocol/pilotprotocol/pkg/daemon/keyexchange" +) + +// peerEndpoint returns the endpoint d currently holds for nodeID, or "". +func peerEndpoint(d *daemon.Daemon, nodeID uint32) string { + for _, p := range d.Tunnels().PeerList() { + if p.NodeID == nodeID { + return p.Endpoint + } + } + return "" +} + +// TestAddrChangeRecoveryTeachesPeerNewSource checks the peer half of the +// address-change recovery: after a node's address changed, the peer must +// learn where the node now is from the recovery itself, not 25 s later from +// the node's next NAT keepalive. +// +// The harness runs every daemon on loopback and cannot give one a new +// address, so the address change itself is not reproduced here (that is +// covered by the Docker lab run described in the changelog entry). What is +// reproduced is the state the peer is left in: B's entry for A is pointed at +// an address that never answers, which is exactly what B holds after A has +// moved. The recovery on A then has to put A's real source address back into +// B's table, before anything else in the test could have. +func TestAddrChangeRecoveryTeachesPeerNewSource(t *testing.T) { + env := NewTestEnv(t) + encrypted := func(cfg *daemon.Config) { cfg.Encrypt = true } + a := env.AddDaemon(encrypted) + b := env.AddDaemon(encrypted) + + // Establish the encrypted tunnel in both directions with one echo. + ln, err := b.Driver.Listen(8140) + if err != nil { + t.Fatalf("listen: %v", err) + } + go func() { + c, err := ln.Accept() + if err != nil { + return + } + defer c.Close() + buf := make([]byte, 16) + n, _ := c.Read(buf) + _, _ = c.Write(buf[:n]) + }() + conn, err := a.Driver.DialAddr(b.Daemon.Addr(), 8140) + if err != nil { + t.Fatalf("dial: %v", err) + } + if _, err := conn.Write([]byte("ping")); err != nil { + t.Fatalf("write: %v", err) + } + buf := make([]byte, 16) + if _, err := conn.Read(buf); err != nil { + t.Fatalf("read echo: %v", err) + } + conn.Close() + + aID, bID := a.Daemon.NodeID(), b.Daemon.NodeID() + _, aPort, err := net.SplitHostPort(a.Daemon.TunnelAddr().String()) + if err != nil { + t.Fatalf("tunnel addr: %v", err) + } + real := peerEndpoint(b.Daemon, aID) + if _, port, _ := net.SplitHostPort(real); port != aPort { + t.Fatalf("B's endpoint for A before the test = %q, want A's tunnel port %s", real, aPort) + } + + // "A moved": B still holds an address A is no longer at. A loopback + // socket that never answers stands in for the old address; it stays open + // so B's sends to it vanish silently, as they would at a real old + // address, instead of drawing ICMP errors. + dead, err := net.ListenUDP("udp", &net.UDPAddr{IP: net.IPv4(127, 0, 0, 1)}) + if err != nil { + t.Fatalf("reserve stale addr: %v", err) + } + defer dead.Close() + stale := dead.LocalAddr().(*net.UDPAddr) + // Let the echo connection's teardown finish first: any frame A still + // sends for it would correct B's entry and hide what the recovery does. + // Not much longer, though; see naturalHeal below. + time.Sleep(2 * time.Second) + poisoned := time.Now() + b.Daemon.Tunnels().AddPeer(aID, stale) + + // The earliest anything but the recovery could correct B: AddPeer + // starts a key exchange with A (also sent through the beacon, so it + // reaches A), retransmitted every RekeyRetransmitInterval, and A may + // answer one from its real address once it has heard nothing from B for + // KeyExchangeReplyStaleThreshold. With the 2 s pause above, the first + // send A could answer goes out about 4 s after the poisoning. (In + // practice A's next keepalive is what corrects B, over 20 s later; this + // is the lower bound.) The recovery has to have healed B before then, or + // the test could not tell which one did it. A slow runner only eats into + // the recovery's share of those 4 s, which needs a few milliseconds. + lastFromB, ok := a.Daemon.Tunnels().LastInboundDecrypt(bID) + if !ok { + t.Fatal("A has never decrypted a frame from B; the echo did not run over the tunnel") + } + naturalHeal := poisoned + for naturalHeal.Sub(lastFromB) < keyexchange.KeyExchangeReplyStaleThreshold { + naturalHeal = naturalHeal.Add(keyexchange.RekeyRetransmitInterval) + } + deadline := naturalHeal.Add(-250 * time.Millisecond) + + time.Sleep(100 * time.Millisecond) + if got := peerEndpoint(b.Daemon, aID); got != stale.String() { + t.Fatalf("B's endpoint for A healed before the recovery ran (%q): echo traffic was still in flight", got) + } + + done := make(chan struct{}) + start := time.Now() + go func() { + defer close(done) + daemon.RecoverFromAddrChangeForTest(a.Daemon) + }() + defer func() { + // The registry half of the recovery may still be running; let it + // finish before the environment is torn down. + select { + case <-done: + case <-time.After(30 * time.Second): + t.Error("recovery did not return") + } + }() + + for peerEndpoint(b.Daemon, aID) != real { + if time.Now().After(deadline) { + t.Fatalf("B still holds %q for A %v after A's recovery started, want %q "+ + "(B's own key exchange could correct it from %v after the poisoning, so a later heal proves nothing)", + peerEndpoint(b.Daemon, aID), time.Since(start).Truncate(time.Millisecond), real, + naturalHeal.Sub(poisoned).Truncate(time.Millisecond)) + } + time.Sleep(5 * time.Millisecond) + } + t.Logf("B learned A's source address %v after the recovery started (B alone could not have before %v after the poisoning)", + time.Since(start).Truncate(time.Millisecond), naturalHeal.Sub(poisoned).Truncate(time.Millisecond)) +}