From bcfdfdbec9059492ad52609f96086692aaeed1fd Mon Sep 17 00:00:00 2001 From: Teo Calin Date: Thu, 1 Oct 2026 18:39:50 +0300 Subject: [PATCH 1/4] daemon: relayed trust handshakes complete in seconds instead of two minutes Between two private nodes the handshake request and its answer are parked at the registry and each side learns of them only by polling, once per keepalive interval (60s): ~59s for the request to show up on the target, ~57s more for the approval to reach the requester. Keep the 60s poll as the baseline and add three bounded triggers (handshakepoll.go): - After sending a handshake request the node polls every 2s until the peer answers or becomes trusted, for at most 2 minutes. Automatic (dial-driven) requests cannot restart the window for the same peer within 10 minutes. - The IPC handlers behind pending / trust / approve / reject / wait-for-trust poll first, at most once per 2s and bounded to 3s when the registry is slow. - A beacon notify ([0x0A][kind], accepted only from the beacon's address) triggers one poll from a token bucket (burst 3, one per 5s), never sooner than 2s after the previous poll. Released beacons do not send it yet. One extra timer in the existing poll loop, armed only while a request is outstanding or a notify is owed a poll; no goroutine per request. The loop's startup jitter no longer delays these triggers. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 16 + pkg/daemon/daemon.go | 108 +++++- pkg/daemon/handshakepoll.go | 310 ++++++++++++++++++ pkg/daemon/ipc.go | 16 +- pkg/daemon/tunnel.go | 30 ++ .../zz_coverage_pkg_daemon_round2_test.go | 4 +- pkg/daemon/zz_handshake_poll_sched_test.go | 263 +++++++++++++++ tests/testenv.go | 14 + tests/zz_handshake_relay_latency_test.go | 130 ++++++++ 9 files changed, 876 insertions(+), 15 deletions(-) create mode 100644 pkg/daemon/handshakepoll.go create mode 100644 pkg/daemon/zz_handshake_poll_sched_test.go create mode 100644 tests/zz_handshake_relay_latency_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index ea3083b8..1a53610b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -332,6 +332,22 @@ Detailed per-release notes are on the over its own limit no longer uses up the tokens other peers need (previously every SYN it had rejected still took a shared token). `-syn-whitelist` is unchanged. +- **A trust handshake between two private nodes takes seconds, not two + minutes.** Neither side can reach the other before trust exists, so the + request and the answer are parked at the registry until each node polls — + once per keepalive interval (60s). Measured before: request visible on the + target after ~59s, approval back at the requester ~57s later. Now: + - a node that has sent a handshake request polls every 2s until it is + answered, for at most 2 minutes; + - `pilotctl pending`, `trust`, `approve`, `reject` and waiting for trust + poll first (at most one such poll every 2s), so a relayed request is + there as soon as someone looks; + - the beacon can tell a node that something is waiting for it (a two-byte + notify that carries nothing else; needs the matching registry/beacon + release), and the node polls at once — at most 3 polls in a burst and + one per 5s after that, whatever arrives. + An idle node still makes one poll per minute, as before. Nodes and + servers that are not updated keep working at the old pace for their part. - **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/pkg/daemon/daemon.go b/pkg/daemon/daemon.go index 6b0d6a86..ae82735e 100644 --- a/pkg/daemon/daemon.go +++ b/pkg/daemon/daemon.go @@ -417,6 +417,10 @@ type Daemon struct { // IPC calls are NOT throttled. autoHandshakeLastAttempt sync.Map + // hsPoll schedules relayed-handshake polls: the 60s baseline plus the + // bounded fast, on-demand and poke-triggered polls (handshakepoll.go). + hsPoll *handshakePollSched + // outbound records the peers this node recently dialed or sent a // trust handshake to, so the private-node SYN gate can admit their // dial-back replies (see replywindow.go). @@ -635,6 +639,7 @@ func New(cfg Config) *Daemon { tunnels: NewTunnelManager(), ports: NewPortManager(), stopCh: make(chan struct{}), + hsPoll: newHandshakePollSched(), synTokens: cfg.synRateLimit(), synLastFill: time.Now(), perSrcSYN: make(map[uint32]*srcSYNBucket), @@ -1040,6 +1045,7 @@ func (d *Daemon) Start() error { // the race under §4.8 stress with -race). d.bus is constructed in // New() so it's safe to publish here. d.tunnels.SetEventBus(d.bus) + d.tunnels.SetBeaconNotifyHandler(d.handshakePoke) // 3. Start UDP listener for tunnel traffic. Compat-mode daemons // skip this — the WSS transport is dialed after register, once we @@ -2300,6 +2306,18 @@ func (d *Daemon) HandshakeSendRequest(nodeID uint32, reason string) error { return d.handshakes.SendRequest(nodeID, reason) } +// handshakeRequestSent starts fast polling for the answer to a handshake +// request that just left. The answer of a private peer comes back through +// the registry whichever way the request went, so this does not need to +// know whether the request was relayed. A peer that is already trusted +// (auto-approved over a direct connection) starts nothing. +func (d *Daemon) handshakeRequestSent(nodeID uint32, explicit bool) { + if d.handshakes == nil || d.handshakes.IsTrusted(nodeID) { + return + } + d.hsPoll.requestSent(nodeID, explicit) +} + // RegisterHandshakeService installs the daemon-wide HandshakeService // implementation provided by the handshake plugin (T3.3). Called from // the composition root after constructing the plugin's Service. @@ -3940,7 +3958,11 @@ func (d *Daemon) dialConnectionLocked(ctx context.Context, dstAddr protocol.Addr // effort, just like before. ErrHandshakeInFlight short-circuit // is the dedup hit, which is the success case. if d.shouldAutoHandshake(dstAddr.Node) { - go func() { _ = d.HandshakeSendRequest(dstAddr.Node, "") }() + go func() { + if d.HandshakeSendRequest(dstAddr.Node, "") == nil { + d.handshakeRequestSent(dstAddr.Node, false) + } + }() } } } @@ -5508,21 +5530,70 @@ func (d *Daemon) tunnelKeepaliveLoop() { // Owns no transport state. Independent of trustRepublishLoop's failure // tracking — handshake polling is best-effort and survives transient // registry hiccups on its own. +// +// The keepalive-interval tick is the baseline and the only thing an idle +// node runs. A second timer is armed only while a request this node sent is +// unanswered, or a beacon poke is owed a poll (see handshakepoll.go), and is +// stopped again as soon as neither holds. func (d *Daemon) handshakePollLoop() { - // Independent jitter so this loop does not align with the others. - time.Sleep(time.Duration(rand.Int63n(int64(5 * time.Second)))) - - ticker := time.NewTicker(d.config.keepaliveInterval()) - defer ticker.Stop() + // Independent jitter so this loop does not align with the others. The + // baseline ticker starts once it has passed; requests, pokes and Stop + // are served during it. + jitter := time.NewTimer(time.Duration(rand.Int63n(int64(5 * time.Second)))) + defer jitter.Stop() + var ticker *time.Ticker + var tick <-chan time.Time + defer func() { + if ticker != nil { + ticker.Stop() + } + }() + extra := time.NewTimer(time.Hour) + extra.Stop() + defer extra.Stop() + extraArmed := false + trusted := func(nodeID uint32) bool { + return d.handshakes != nil && d.handshakes.IsTrusted(nodeID) + } for { select { case <-d.stopCh: return - case <-ticker.C: - if d.reg() == nil { - continue + case <-jitter.C: + ticker = time.NewTicker(d.config.keepaliveInterval()) + tick = ticker.C + case <-tick: + d.pollHandshakes(0, 0) + case <-extra.C: + extraArmed = false + d.pollHandshakes(handshakeFastPollInterval/2, 0) + case <-d.hsPoll.wake: + if due, wait := d.hsPoll.pokeWait(); due && wait == 0 { + d.pollHandshakes(handshakeOnDemandGap, 0) + } + } + + // Arm the extra timer for whichever comes first: a poll owed to a + // poke, or the next fast poll while a request is outstanding. + next := time.Duration(-1) + if due, wait := d.hsPoll.pokeWait(); due { + next = wait + } + if d.hsPoll.fastActive(trusted) && (next < 0 || handshakeFastPollInterval < next) { + next = handshakeFastPollInterval + } + switch { + case next < 0 && extraArmed: + if !extra.Stop() { + select { + case <-extra.C: + default: + } } - d.pollRelayedHandshakes() + extraArmed = false + case next >= 0 && !extraArmed: + extra.Reset(next) + extraArmed = true } } } @@ -6349,8 +6420,20 @@ func (d *Daemon) lookupPeerPubKey(nodeID uint32) (ed25519.PublicKey, error) { // pollRelayedHandshakes checks the registry for handshake requests and // responses relayed to this node and processes them. -func (d *Daemon) pollRelayedHandshakes() { - resp, err := d.reg().PollHandshakes(d.NodeID()) +// +// timeout > 0 bounds the registry call (a local client is waiting on it); +// 0 leaves it unbounded, as the background loop always ran it. +func (d *Daemon) pollRelayedHandshakes(timeout time.Duration) { + rc, nodeID := d.reg(), d.NodeID() + var resp map[string]interface{} + var err error + if timeout > 0 { + resp, err = withRegistryDeadline(timeout, func() (map[string]interface{}, error) { + return rc.PollHandshakes(nodeID) + }) + } else { + resp, err = rc.PollHandshakes(nodeID) + } if err != nil { slog.Debug("poll handshakes failed", "error", err) return @@ -6391,6 +6474,7 @@ func (d *Daemon) pollRelayedHandshakes() { } fromNodeID := uint32(fromIDVal) accept, _ := respMsg["accept"].(bool) + d.hsPoll.answered(fromNodeID) if accept { slog.Info("relayed handshake approval received", "from_node_id", fromNodeID) diff --git a/pkg/daemon/handshakepoll.go b/pkg/daemon/handshakepoll.go new file mode 100644 index 00000000..007fcfc2 --- /dev/null +++ b/pkg/daemon/handshakepoll.go @@ -0,0 +1,310 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package daemon + +import ( + "sync" + "sync/atomic" + "time" +) + +// Relayed-handshake poll scheduling (L11). +// +// A trust handshake between two private nodes cannot go direct: neither +// side can resolve the other before trust exists, so the request and its +// answer are parked at the registry and each node learns of them only when +// it polls (pollRelayedHandshakes). With one poll per keepalive interval +// (60s) a manual handshake took about two minutes end to end. +// +// The 60s poll stays as the baseline and is all an idle node ever does. +// Three things add polls, each bounded: +// +// - Waiting on an answer: after this node sends a handshake request it +// polls every handshakeFastPollInterval until the peer answers, becomes +// trusted, or handshakeFastPollWindow has passed. +// - On demand: a local client asking for pending requests or the trust +// list, approving, rejecting or waiting for trust triggers one poll +// first, at most one per handshakeOnDemandGap. +// - Poke: the beacon can tell this node that something is waiting for it +// at the registry (beaconMsgNotify). A poke triggers one poll, from a +// small token bucket, and never sooner than handshakeOnDemandGap after +// the previous poll. +const ( + // handshakeFastPollInterval is the poll period while a request this + // node sent is unanswered. + handshakeFastPollInterval = 2 * time.Second + + // handshakeFastPollWindow bounds fast polling per request. + handshakeFastPollWindow = 2 * time.Minute + + // handshakeAutoRearm is how long a peer that was the target of an + // automatic (dial-driven) handshake cannot start another fast window. + // An application redialing a peer that never answers re-sends the + // request every autoHandshakeCooldown; without this it would hold the + // fast window open forever. Explicit requests are not held back. + handshakeAutoRearm = 10 * time.Minute + + // handshakeMaxWaiting caps the peers tracked at once. Fast polling is + // one timer however many peers are waiting, so this bounds memory, not + // request rate. + handshakeMaxWaiting = 64 + + // handshakeOnDemandGap is the minimum spacing of polls triggered by a + // local client or a beacon poke. + handshakeOnDemandGap = 2 * time.Second + + // handshakeOnDemandTimeout bounds how long a local client's request is + // held up by its poll when the registry is slow or unreachable. + handshakeOnDemandTimeout = 3 * time.Second + + // handshakePokeBurst / handshakePokeRefill: token bucket for polls + // triggered by beacon pokes. A flood of pokes costs at most the burst + // plus one poll per refill period. + handshakePokeBurst = 3 + handshakePokeRefill = 5 * time.Second +) + +// handshakePollSched decides when relayed handshakes are polled. It holds +// no timers and makes no calls; handshakePollLoop and the IPC handlers +// drive it, and tests drive it with their own clock. +type handshakePollSched struct { + mu sync.Mutex + now func() time.Time + + // lastPoll is when the most recent poll started, whatever triggered it. + lastPoll time.Time + + // waiting maps a peer we sent a request to → the end of its fast-poll + // window. Entries stay (expired) until rearm passes, so the automatic + // path cannot restart a window early. + waiting map[uint32]handshakeWait + + // pokeDue is set when a poke was accepted but a poll had just run; the + // loop polls once handshakeOnDemandGap has passed since lastPoll. + pokeDue bool + + pokeTokens int + pokeFill time.Time + + // wake nudges handshakePollLoop to re-evaluate its timers (capacity 1, + // never blocks the sender). + wake chan struct{} + + // sem serializes polls so two triggers never overlap (capacity 1). + sem chan struct{} + + // polls counts registry polls, for tests and the info reply. + polls atomic.Uint64 +} + +type handshakeWait struct { + until time.Time // fast polling ends + rearm time.Time // automatic handshakes may not restart the window before this +} + +func newHandshakePollSched() *handshakePollSched { + return &handshakePollSched{ + now: time.Now, + waiting: make(map[uint32]handshakeWait), + pokeTokens: handshakePokeBurst, + wake: make(chan struct{}, 1), + sem: make(chan struct{}, 1), + } +} + +func (s *handshakePollSched) nudge() { + select { + case s.wake <- struct{}{}: + default: + } +} + +// requestSent starts the fast-poll window for a peer we just sent a +// handshake request to. explicit is true for a request a local client asked +// for; an automatic one cannot restart a window for the same peer within +// handshakeAutoRearm. +func (s *handshakePollSched) requestSent(peer uint32, explicit bool) { + s.mu.Lock() + now := s.now() + w, tracked := s.waiting[peer] + switch { + case tracked && now.Before(w.until): + // Already waiting on this peer; a repeat does not extend it. + s.mu.Unlock() + return + case tracked && !explicit && now.Before(w.rearm): + s.mu.Unlock() + return + case !tracked && len(s.waiting) >= handshakeMaxWaiting: + s.pruneLocked(now) + if len(s.waiting) >= handshakeMaxWaiting { + s.mu.Unlock() + return + } + } + s.waiting[peer] = handshakeWait{until: now.Add(handshakeFastPollWindow), rearm: now.Add(handshakeAutoRearm)} + s.mu.Unlock() + s.nudge() +} + +// answered ends the wait for a peer whose answer arrived. +func (s *handshakePollSched) answered(peer uint32) { + s.mu.Lock() + if w, ok := s.waiting[peer]; ok { + w.until = time.Time{} + s.waiting[peer] = w + } + s.mu.Unlock() +} + +func (s *handshakePollSched) pruneLocked(now time.Time) { + for peer, w := range s.waiting { + if !now.Before(w.until) && !now.Before(w.rearm) { + delete(s.waiting, peer) + } + } +} + +// fastActive reports whether any request is still inside its fast-poll +// window. Peers that trusted reports as trusted are settled first. +func (s *handshakePollSched) fastActive(trusted func(uint32) bool) bool { + s.mu.Lock() + defer s.mu.Unlock() + now := s.now() + s.pruneLocked(now) + active := false + for peer, w := range s.waiting { + if !now.Before(w.until) { + continue + } + if trusted != nil && trusted(peer) { + w.until = time.Time{} + s.waiting[peer] = w + continue + } + active = true + } + return active +} + +// claim reports whether a poll may start now: no poll started within +// minGap. It records the start when it says yes. +func (s *handshakePollSched) claim(minGap time.Duration) bool { + s.mu.Lock() + defer s.mu.Unlock() + now := s.now() + if !s.lastPoll.IsZero() && now.Sub(s.lastPoll) < minGap { + return false + } + s.lastPoll = now + s.pokeDue = false + return true +} + +// poke handles a beacon notification. It returns true when the loop should +// be woken: the poke took a token and a poll is now due (immediately, or +// once handshakeOnDemandGap has passed since the last one). +func (s *handshakePollSched) poke() bool { + s.mu.Lock() + defer s.mu.Unlock() + now := s.now() + if s.pokeDue { + return false // one poll is already owed; it will cover this poke too + } + if s.pokeFill.IsZero() { + s.pokeFill = now + } + if n := int(now.Sub(s.pokeFill) / handshakePokeRefill); n > 0 { + s.pokeTokens += n + if s.pokeTokens >= handshakePokeBurst { + s.pokeTokens = handshakePokeBurst + s.pokeFill = now + } else { + s.pokeFill = s.pokeFill.Add(time.Duration(n) * handshakePokeRefill) + } + } + if s.pokeTokens <= 0 { + return false + } + s.pokeTokens-- + s.pokeDue = true + return true +} + +// pokeWait reports whether a poke-triggered poll is owed and how long the +// loop must still wait before running it. +func (s *handshakePollSched) pokeWait() (due bool, wait time.Duration) { + s.mu.Lock() + defer s.mu.Unlock() + if !s.pokeDue { + return false, 0 + } + if s.lastPoll.IsZero() { + return true, 0 + } + if wait = handshakeOnDemandGap - s.now().Sub(s.lastPoll); wait < 0 { + wait = 0 + } + return true, wait +} + +// pollHandshakes runs one relayed-handshake poll unless one started within +// minGap. Polls never overlap; a caller arriving while one is running waits +// for it (up to timeout when timeout > 0) and then finds it fresh enough. +// timeout also bounds the registry call itself; 0 means no bound, which is +// what the background loop has always used. +func (d *Daemon) pollHandshakes(minGap, timeout time.Duration) { + s := d.hsPoll + if d.reg() == nil { + // Nothing to poll; do not leave a poke owed (the loop would keep + // re-arming its timer for a poll that cannot run). + s.mu.Lock() + s.pokeDue = false + s.mu.Unlock() + return + } + if timeout > 0 { + t := time.NewTimer(timeout) + defer t.Stop() + select { + case s.sem <- struct{}{}: + case <-t.C: + return + case <-d.stopCh: + return + } + } else { + select { + case s.sem <- struct{}{}: + case <-d.stopCh: + return + } + } + defer func() { <-s.sem }() + if !s.claim(minGap) { + return + } + s.polls.Add(1) + d.pollRelayedHandshakes(timeout) +} + +// pollHandshakesOnDemand is called before a local client is shown, or acts +// on, handshake state, so a relayed request or answer parked at the +// registry is visible at once instead of after the next background poll. +func (d *Daemon) pollHandshakesOnDemand() { + d.pollHandshakes(handshakeOnDemandGap, handshakeOnDemandTimeout) +} + +// handshakePoke is the tunnel layer's callback for a beacon notification +// that something is waiting for this node at the registry. Runs on the +// tunnel read loop, so it only records the poke; the poll loop does the +// work. +func (d *Daemon) handshakePoke() { + if d.hsPoll.poke() { + d.hsPoll.nudge() + } +} + +// RelayedHandshakePolls returns how many relayed-handshake polls this +// daemon has sent to the registry since it started. +func (d *Daemon) RelayedHandshakePolls() uint64 { return d.hsPoll.polls.Load() } diff --git a/pkg/daemon/ipc.go b/pkg/daemon/ipc.go index ef088d61..6fc18e39 100644 --- a/pkg/daemon/ipc.go +++ b/pkg/daemon/ipc.go @@ -1955,7 +1955,13 @@ func (s *IPCServer) handleHandshake(conn *ipcConn, reqID uint64, payload []byte) if len(rest) > 4 { justification = string(rest[4:]) } - if err := s.daemon.HandshakeSendRequest(nodeID, justification); err != nil { + err := s.daemon.HandshakeSendRequest(nodeID, justification) + if err == nil || errors.Is(err, ErrHandshakeInFlight) { + // Poll for the answer every couple of seconds for a while + // instead of once a minute. + s.daemon.handshakeRequestSent(nodeID, true) + } + if err != nil { if errors.Is(err, ErrHandshakeInFlight) { // Soft success: another caller is already running the // same handshake. Reply with a distinct status so @@ -1985,6 +1991,10 @@ func (s *IPCServer) handleHandshake(conn *ipcConn, reqID uint64, payload []byte) return } nodeID := binary.BigEndian.Uint32(rest[0:4]) + // A request relayed through the registry is only known here once + // polled; fetch it now so it can be answered without waiting for + // the next background poll. Same for reject, pending and trusted. + s.daemon.pollHandshakesOnDemand() if err := s.daemon.handshakes.ApproveHandshake(nodeID); err != nil { s.sendError(conn, reqID, fmt.Sprintf("handshake approve: %v", err)) return @@ -2005,6 +2015,7 @@ func (s *IPCServer) handleHandshake(conn *ipcConn, reqID uint64, payload []byte) if len(rest) > 4 { reason = string(rest[4:]) } + s.daemon.pollHandshakesOnDemand() if err := s.daemon.handshakes.RejectHandshake(nodeID, reason); err != nil { s.sendError(conn, reqID, fmt.Sprintf("handshake reject: %v", err)) return @@ -2016,6 +2027,7 @@ func (s *IPCServer) handleHandshake(conn *ipcConn, reqID uint64, payload []byte) s.ipcWriteHandshakeOK(conn, reqID, data) case SubHandshakePending: + s.daemon.pollHandshakesOnDemand() pending := s.daemon.handshakes.PendingRequests() list := make([]map[string]interface{}, len(pending)) for i, p := range pending { @@ -2032,6 +2044,7 @@ func (s *IPCServer) handleHandshake(conn *ipcConn, reqID uint64, payload []byte) s.ipcWriteHandshakeOK(conn, reqID, data) case SubHandshakeTrusted: + s.daemon.pollHandshakesOnDemand() trusted := s.daemon.handshakes.TrustedPeers() list := make([]map[string]interface{}, len(trusted)) for i, t := range trusted { @@ -2076,6 +2089,7 @@ func (s *IPCServer) handleHandshake(conn *ipcConn, reqID uint64, payload []byte) if timeoutMs > 30000 { timeoutMs = 30000 } + s.daemon.pollHandshakesOnDemand() ok := s.daemon.handshakes.WaitForTrust(nodeID, time.Duration(timeoutMs)*time.Millisecond) data, _ := json.Marshal(map[string]interface{}{ "node_id": nodeID, diff --git a/pkg/daemon/tunnel.go b/pkg/daemon/tunnel.go index ea9d1961..8bfade16 100644 --- a/pkg/daemon/tunnel.go +++ b/pkg/daemon/tunnel.go @@ -200,6 +200,10 @@ type TunnelManager struct { // or only the session layer (decrypt/key desync). Stamped at Listen / // ConnectCompat so the age is measured from transport start, not epoch. LastRecvNano int64 + + // onBeaconNotify is run when the beacon says the registry holds + // something for this node (beaconMsgNotify). + onBeaconNotify atomic.Pointer[func()] // P1-008: packets dropped from the per-peer pending queue while waiting // for key exchange. Exposed so operators can tell a congested overlay // apart from a silent crypto stall. @@ -2318,11 +2322,37 @@ func (tm *TunnelManager) handleBeaconMessage(data []byte, from *net.UDPAddr) { return } tm.handleRelayDeliver(data[1:]) + case beaconMsgNotify: + // Carries nothing but "the registry holds something for you". + // Only the beacon may say so, and what the daemon does with it is + // rate-limited (handshakepoll.go), so a forged or repeated notify + // costs at most a few registry polls. + if !fromBeacon { + slog.Debug("dropping notify from non-beacon source", "from", from) + return + } + if fn := tm.onBeaconNotify.Load(); fn != nil { + (*fn)() + } default: slog.Debug("unknown beacon message on tunnel socket", "type", data[0], "from", from) } } +// beaconMsgNotify is a beacon → node message, [0x0A][kind(1)], telling the +// node that the registry is holding something for it. kind 0x01 is a relayed +// trust-handshake request or answer; the node polls for it. It carries no +// node ID, address or payload. Daemons that predate it log it as an unknown +// beacon message and drop it. (Must match beacon.MsgNotify; the other beacon +// message types live in common/protocol and this one should move there.) +const beaconMsgNotify byte = 0x0A + +// SetBeaconNotifyHandler installs the callback run for each beacon notify. +// It is called on the tunnel read loop and must not block. +func (tm *TunnelManager) SetBeaconNotifyHandler(fn func()) { + tm.onBeaconNotify.Store(&fn) +} + // handlePunchCommand processes a beacon punch command, sending a punch packet // to the specified target to create a NAT mapping. Thin shim over // routing.Manager.HandlePunchCommand. diff --git a/pkg/daemon/zz_coverage_pkg_daemon_round2_test.go b/pkg/daemon/zz_coverage_pkg_daemon_round2_test.go index eff4addb..4dea5bf0 100644 --- a/pkg/daemon/zz_coverage_pkg_daemon_round2_test.go +++ b/pkg/daemon/zz_coverage_pkg_daemon_round2_test.go @@ -165,7 +165,7 @@ func TestPollRelayedHandshakesEmptyNoOp(t *testing.T) { registerSelfOnRegistry(t, d) // No panic, no side effect. The empty-list branches are still walked. - d.pollRelayedHandshakes() + d.pollRelayedHandshakes(0) } // pollRelayedHandshakesHasService verifies that when a HandshakeService is @@ -208,7 +208,7 @@ func TestPollRelayedHandshakesWithServiceNoOp(t *testing.T) { // Empty mailbox — the code walks the empty requests/responses lists // and the service-non-nil guards but does not invoke any Process* method. - d.pollRelayedHandshakes() + d.pollRelayedHandshakes(0) if svc.requests.Load() != 0 || svc.approvals.Load() != 0 || svc.rejections.Load() != 0 { t.Fatalf("unexpected service calls: req=%d approve=%d reject=%d", diff --git a/pkg/daemon/zz_handshake_poll_sched_test.go b/pkg/daemon/zz_handshake_poll_sched_test.go new file mode 100644 index 00000000..56904460 --- /dev/null +++ b/pkg/daemon/zz_handshake_poll_sched_test.go @@ -0,0 +1,263 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package daemon + +import ( + "net" + "testing" + "time" +) + +// schedAt returns a scheduler on a clock the test advances by hand. +func schedAt(start time.Time) (*handshakePollSched, func(time.Duration)) { + now := start + s := newHandshakePollSched() + s.now = func() time.Time { return now } + return s, func(d time.Duration) { now = now.Add(d) } +} + +func TestHandshakeFastPollWindowStartsAndEnds(t *testing.T) { + t.Parallel() + s, advance := schedAt(time.Unix(1_700_000_000, 0)) + never := func(uint32) bool { return false } + + if s.fastActive(never) { + t.Fatal("idle scheduler reports a fast window") + } + s.requestSent(7, true) + select { + case <-s.wake: + default: + t.Fatal("starting a fast window did not wake the loop") + } + if !s.fastActive(never) { + t.Fatal("no fast window after a request was sent") + } + + // A repeat while waiting does not push the end out. + advance(handshakeFastPollWindow - time.Second) + s.requestSent(7, true) + if !s.fastActive(never) { + t.Fatal("window ended early") + } + advance(time.Second) + if s.fastActive(never) { + t.Fatalf("window still open %s after the request", handshakeFastPollWindow) + } + + // The peer's answer ends it at once; so does the peer becoming trusted. + s.requestSent(7, true) + s.answered(7) + if s.fastActive(never) { + t.Fatal("window still open after the answer arrived") + } + s.requestSent(8, true) + if s.fastActive(func(peer uint32) bool { return peer == 8 }) { + t.Fatal("window still open for a peer that is now trusted") + } +} + +// An application redialing a peer that never answers re-sends the automatic +// handshake every autoHandshakeCooldown. That must not keep fast polling on. +func TestHandshakeFastPollAutomaticRequestCannotHoldWindowOpen(t *testing.T) { + t.Parallel() + s, advance := schedAt(time.Unix(1_700_000_000, 0)) + never := func(uint32) bool { return false } + + s.requestSent(7, false) + fast := 0 + for elapsed := time.Duration(0); elapsed < handshakeAutoRearm; elapsed += autoHandshakeCooldown { + if s.fastActive(never) { + fast++ + } + advance(autoHandshakeCooldown) + s.requestSent(7, false) + } + if want := int(handshakeFastPollWindow / autoHandshakeCooldown); fast != want { + t.Fatalf("fast window open at %d of the %s checks, want %d (one window)", fast, autoHandshakeCooldown, want) + } + if !s.fastActive(never) { + t.Fatalf("automatic request could not start a window %s after the last one", handshakeAutoRearm) + } + + // An explicit request is never held back. + advance(handshakeFastPollWindow) + s.requestSent(7, true) + if !s.fastActive(never) { + t.Fatal("explicit request did not start a window") + } +} + +func TestHandshakeFastPollTracksBoundedPeers(t *testing.T) { + t.Parallel() + s, _ := schedAt(time.Unix(1_700_000_000, 0)) + for peer := uint32(1); peer <= 4*handshakeMaxWaiting; peer++ { + s.requestSent(peer, true) + } + if n := len(s.waiting); n != handshakeMaxWaiting { + t.Fatalf("tracking %d peers, want the cap %d", n, handshakeMaxWaiting) + } +} + +func TestHandshakePollClaimRateLimit(t *testing.T) { + t.Parallel() + s, advance := schedAt(time.Unix(1_700_000_000, 0)) + + granted := 0 + for i := 0; i < 600; i++ { // a client asking every 100ms for a minute + if s.claim(handshakeOnDemandGap) { + granted++ + } + advance(100 * time.Millisecond) + } + if want := int(time.Minute / handshakeOnDemandGap); granted != want { + t.Fatalf("%d on-demand polls in a minute of constant asking, want %d", granted, want) + } + // The baseline tick is never refused. + if !s.claim(0) { + t.Fatal("baseline poll refused") + } +} + +func TestHandshakePokeRateLimit(t *testing.T) { + t.Parallel() + s, advance := schedAt(time.Unix(1_700_000_000, 0)) + + // A flood: 100 pokes a second for a minute. The loop polls whenever a + // poke is owed a poll and the gap since the last poll has passed. + polls := 0 + for i := 0; i < 6000; i++ { + s.poke() + if due, wait := s.pokeWait(); due && wait == 0 && s.claim(handshakeOnDemandGap) { + polls++ + } + advance(10 * time.Millisecond) + } + if max := handshakePokeBurst + int(time.Minute/handshakePokeRefill); polls > max || polls < max-1 { + t.Fatalf("%d polls from a one-minute poke flood, want about %d (burst + one per %s)", polls, max, handshakePokeRefill) + } + + // Quiet for a while: the bucket refills to the burst, no further. + advance(time.Hour) + polls = 0 + for i := 0; i < 50; i++ { + s.poke() + if due, _ := s.pokeWait(); due { + advance(handshakeOnDemandGap) + if s.claim(handshakeOnDemandGap) { + polls++ + } + } + } + if polls < handshakePokeBurst || polls > handshakePokeBurst+int(50*handshakeOnDemandGap/handshakePokeRefill)+1 { + t.Fatalf("%d polls after an idle hour, want the burst %d plus refill", polls, handshakePokeBurst) + } +} + +func TestHandshakePokeWaitsOutARecentPoll(t *testing.T) { + t.Parallel() + s, advance := schedAt(time.Unix(1_700_000_000, 0)) + + if !s.claim(0) { + t.Fatal("first poll refused") + } + advance(500 * time.Millisecond) + if !s.poke() { + t.Fatal("poke refused with a full bucket") + } + due, wait := s.pokeWait() + if !due || wait != handshakeOnDemandGap-500*time.Millisecond { + t.Fatalf("poke 500ms after a poll: due=%v wait=%s, want a poll owed in %s", due, wait, handshakeOnDemandGap-500*time.Millisecond) + } + // Further pokes while one poll is owed take no more tokens. + tokens := s.pokeTokens + for i := 0; i < 10; i++ { + if s.poke() { + t.Fatal("a second poll was scheduled while one is owed") + } + } + if s.pokeTokens != tokens { + t.Fatalf("pokes while a poll is owed spent tokens: %d -> %d", tokens, s.pokeTokens) + } + // Any poll settles the debt. + advance(wait) + if !s.claim(handshakeOnDemandGap) { + t.Fatal("owed poll refused after the gap") + } + if due, _ := s.pokeWait(); due { + t.Fatal("poll still owed after one ran") + } +} + +// A notify is honoured only when it comes from the beacon's address. +func TestBeaconNotifyOnlyFromBeacon(t *testing.T) { + t.Parallel() + tm := NewTunnelManager() + if err := tm.SetBeaconAddr("192.0.2.10:9001"); err != nil { + t.Fatal(err) + } + calls := 0 + tm.SetBeaconNotifyHandler(func() { calls++ }) + + frame := []byte{beaconMsgNotify, 0x01} + tm.handleBeaconMessage(frame, &net.UDPAddr{IP: net.ParseIP("192.0.2.66"), Port: 9001}) + tm.handleBeaconMessage(frame, &net.UDPAddr{IP: net.ParseIP("192.0.2.10"), Port: 4444}) + if calls != 0 { + t.Fatalf("notify from a non-beacon source ran the handler %d times", calls) + } + tm.handleBeaconMessage(frame, &net.UDPAddr{IP: net.ParseIP("192.0.2.10"), Port: 9001}) + if calls != 1 { + t.Fatalf("notify from the beacon ran the handler %d times, want 1", calls) + } +} + +// TestHandshakePollLoopExtraPollsAreBounded drives the real loop against a +// registry with the baseline tick out of the way (1h): an idle loop polls +// nothing, a burst of pokes costs one poll, a request in flight is polled +// at the fast period, and the loop goes quiet again once it is answered. +func TestHandshakePollLoopExtraPollsAreBounded(t *testing.T) { + t.Parallel() + reg, rc := startTestRegistry(t) + t.Cleanup(func() { reg.Close() }) + t.Cleanup(func() { rc.Close() }) + + d := New(Config{KeepaliveInterval: time.Hour}) + d.regConn.Store(rc) + done := make(chan struct{}) + go func() { d.handshakePollLoop(); close(done) }() + t.Cleanup(func() { close(d.stopCh); <-done }) + + waitPolls := func(want uint64, within time.Duration) { + t.Helper() + deadline := time.Now().Add(within) + for d.RelayedHandshakePolls() < want && time.Now().Before(deadline) { + time.Sleep(20 * time.Millisecond) + } + if got := d.RelayedHandshakePolls(); got != want { + t.Fatalf("registry polls = %d, want %d", got, want) + } + } + + time.Sleep(time.Second) + waitPolls(0, 0) // idle: nothing beyond the baseline + + for i := 0; i < 200; i++ { + d.handshakePoke() + } + waitPolls(1, 2*time.Second) // served even during the startup jitter + time.Sleep(handshakeOnDemandGap + time.Second) + waitPolls(1, 0) // the rest of the burst bought nothing + + d.hsPoll.requestSent(42, true) + time.Sleep(2*handshakeFastPollInterval + handshakeFastPollInterval/2) + if got := d.RelayedHandshakePolls(); got < 2 || got > 4 { + t.Fatalf("registry polls = %d after 2.5 fast periods with a request in flight, want 3 (1 + 2)", got) + } + d.hsPoll.answered(42) + time.Sleep(handshakeFastPollInterval + 200*time.Millisecond) // at most one already-armed poll + settled := d.RelayedHandshakePolls() + time.Sleep(2 * handshakeFastPollInterval) + if got := d.RelayedHandshakePolls(); got != settled { + t.Fatalf("loop kept polling after the answer: %d -> %d", settled, got) + } +} diff --git a/tests/testenv.go b/tests/testenv.go index 5af3b543..20d80780 100644 --- a/tests/testenv.go +++ b/tests/testenv.go @@ -84,6 +84,10 @@ type TestEnv struct { RegistryAddr string AdminToken string + // HandshakeNotify is true when the registry and beacon in use can + // prompt a node to poll for a relayed handshake. + HandshakeNotify bool + daemons []*daemon.Daemon drivers []*driver.Driver runtimes []*pluginsruntime.Runtime @@ -136,6 +140,16 @@ func NewTestEnv(t *testing.T) *TestEnv { } env.RegistryAddr = resolveLocalAddr(env.Registry.Addr()) + // Same wiring as cmd/rendezvous: the beacon prompts the recipient of a + // relayed handshake to poll. Asserted so the harness builds against + // beacon / rendezvous releases that predate the hook. + if b, ok := interface{}(env.Beacon).(interface{ NotifyNode(uint32) error }); ok { + if r, ok := interface{}(env.Registry).(interface{ SetHandshakeNotifier(func(uint32)) }); ok { + r.SetHandshakeNotifier(func(nodeID uint32) { _ = b.NotifyNode(nodeID) }) + env.HandshakeNotify = true + } + } + t.Cleanup(func() { env.Close() }) diff --git a/tests/zz_handshake_relay_latency_test.go b/tests/zz_handshake_relay_latency_test.go new file mode 100644 index 00000000..706a71c3 --- /dev/null +++ b/tests/zz_handshake_relay_latency_test.go @@ -0,0 +1,130 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package tests + +import ( + "testing" + "time" + + "github.com/pilot-protocol/pilotprotocol/pkg/daemon" +) + +// TestRelayedHandshakeCompletesInSeconds runs a manual trust handshake +// between two private nodes. Neither can resolve the other before trust +// exists, so the request and the approval both travel through the registry +// and each side only learns of them by polling. With one poll per keepalive +// interval (60s) the request reached the target after up to a minute and +// the approval took up to another minute to come back. +// +// The request must be visible as soon as the target's operator asks for +// pending requests, and the approval must reach the requester within a few +// fast-poll periods, without anyone prodding the requester. The thresholds +// are far below one poll interval and far above the 2s fast-poll period. +func TestRelayedHandshakeCompletesInSeconds(t *testing.T) { + requireRealNetwork(t) + t.Parallel() + env := NewTestEnv(t) + + private := func(c *daemon.Config) { c.Public = false } + a := env.AddDaemon(private) + b := env.AddDaemon(private) + const limit = 20 * time.Second + + start := time.Now() + if _, err := a.Driver.Handshake(b.Daemon.NodeID(), "relay latency test"); err != nil { + t.Fatalf("handshake: %v", err) + } + + // Target side: the operator asks what is pending. + seen := false + for time.Since(start) < limit && !seen { + pending, err := b.Driver.PendingHandshakes() + if err != nil { + t.Fatalf("pending: %v", err) + } + for _, p := range asList(pending["pending"]) { + if rec, ok := p.(map[string]interface{}); ok && uint32(rec["node_id"].(float64)) == a.Daemon.NodeID() { + seen = true + } + } + if !seen { + time.Sleep(200 * time.Millisecond) + } + } + requestLatency := time.Since(start) + if !seen { + t.Fatalf("relayed request not in the target's pending list after %s", limit) + } + + approved := time.Now() + if _, err := b.Driver.ApproveHandshake(a.Daemon.NodeID()); err != nil { + t.Fatalf("approve: %v", err) + } + + // Requester side: nobody asks it anything; its own polling must pick + // the approval up. + for !a.Daemon.HandshakeService().IsTrusted(b.Daemon.NodeID()) { + if time.Since(approved) > limit { + t.Fatalf("relayed approval did not reach the requester within %s", limit) + } + time.Sleep(100 * time.Millisecond) + } + approvalLatency := time.Since(approved) + t.Logf("request visible after %s, approval received after %s; registry polls: requester %d, target %d", + requestLatency.Truncate(time.Millisecond), approvalLatency.Truncate(time.Millisecond), + a.Daemon.RelayedHandshakePolls(), b.Daemon.RelayedHandshakePolls()) + + // Bounded cost: the target polled only when asked (at most one per + // 2s), the requester only at the fast-poll period while it waited. + elapsed := time.Since(start) + if max := uint64(elapsed/time.Second) + 2; a.Daemon.RelayedHandshakePolls() > max || b.Daemon.RelayedHandshakePolls() > max { + t.Errorf("too many registry polls in %s: requester %d, target %d (max %d each)", + elapsed.Truncate(time.Millisecond), a.Daemon.RelayedHandshakePolls(), b.Daemon.RelayedHandshakePolls(), max) + } + + // Once trusted, the requester stops polling fast. + settled := a.Daemon.RelayedHandshakePolls() + time.Sleep(3 * time.Second) + if got := a.Daemon.RelayedHandshakePolls(); got > settled+1 { + t.Errorf("requester kept fast-polling after trust: %d polls in 3s", got-settled) + } +} + +func asList(v interface{}) []interface{} { + l, _ := v.([]interface{}) + return l +} + +// TestRelayedHandshakeReachesUnattendedTarget: nobody asks the target +// anything. The registry has the beacon prompt it, and it polls at once +// instead of at its next once-a-minute tick. The relayed approval reaches +// the requester the same way. +func TestRelayedHandshakeReachesUnattendedTarget(t *testing.T) { + requireRealNetwork(t) + t.Parallel() + env := NewTestEnv(t) + if !env.HandshakeNotify { + t.Skip("needs a registry and beacon with the handshake notify hook (rendezvous SetHandshakeNotifier, beacon NotifyNode)") + } + + private := func(c *daemon.Config) { c.Public = false } + a := env.AddDaemon(private) + b := env.AddDaemon(private) + const limit = 20 * time.Second + + start := time.Now() + if _, err := a.Driver.Handshake(b.Daemon.NodeID(), "unattended target"); err != nil { + t.Fatalf("handshake: %v", err) + } + for b.Daemon.HandshakeService().PendingCount() == 0 { + if time.Since(start) > limit { + t.Fatalf("relayed request did not reach the unattended target within %s", limit) + } + time.Sleep(50 * time.Millisecond) + } + t.Logf("request reached the target after %s; target polls: %d", + time.Since(start).Truncate(time.Millisecond), b.Daemon.RelayedHandshakePolls()) + if polls := b.Daemon.RelayedHandshakePolls(); polls > 2 { + t.Errorf("target polled the registry %d times for one request", polls) + } +} From 05041c2518af40d2bdf9e95e222870957e608326 Mon Sep 17 00:00:00 2001 From: Teo Calin Date: Thu, 1 Oct 2026 19:23:25 +0300 Subject: [PATCH 2/4] tests: wait for the trust store write before reading it TestHandshakeTrustPersistence and TestHandshakeTrustLoadVerify stat/read trust.json the moment the peer shows up in the trusted list. The handshake plugin marks the peer trusted in memory and writes the store just after, so the tests depended on losing that race. The trust-list request now does a registry poll first, which lined the read up with the gap every time. Co-Authored-By: Claude Opus 5.5 --- tests/zz_handshake_test.go | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/tests/zz_handshake_test.go b/tests/zz_handshake_test.go index 6e2e278e..df2f1bfd 100644 --- a/tests/zz_handshake_test.go +++ b/tests/zz_handshake_test.go @@ -485,7 +485,9 @@ func TestHandshakeTrustPersistence(t *testing.T) { } } - // Verify trust store file exists + // Verify trust store file exists. The store is written just after the + // peer is marked trusted in memory, so give the write a moment. + waitForFile(trustPathA, 2*time.Second) if _, err := os.Stat(trustPathA); err != nil { t.Fatalf("trust store not created: %v", err) } @@ -615,6 +617,7 @@ func TestHandshakeTrustLoadVerify(t *testing.T) { } // Verify trust file was created and contains correct data + waitForFile(trustPath, 2*time.Second) data, err := os.ReadFile(trustPath) if err != nil { t.Fatalf("read trust file: %v", err) @@ -752,3 +755,13 @@ func TestHandshakeTrustLoadFromDisk(t *testing.T) { var _ = driver.Connect // keep driver import var _ = os.Remove // keep os import + +// waitForFile returns once path exists or the timeout has passed; the +// caller's own check reports a missing file. +func waitForFile(path string, timeout time.Duration) { + for deadline := time.Now().Add(timeout); time.Now().Before(deadline); time.Sleep(10 * time.Millisecond) { + if _, err := os.Stat(path); err == nil { + return + } + } +} From e3d54a077b29ca5cff53f2900b22392cba0355de Mon Sep 17 00:00:00 2001 From: Teo Calin Date: Fri, 2 Oct 2026 12:43:36 +0300 Subject: [PATCH 3/4] daemon: justify the poll-loop jitter for gosec; the poke-burst test allows the one deferred poll A burst of pokes is served by one poll at once and, if any poke lands after that poll began, one more after handshakeOnDemandGap. The test required exactly one and failed in CI when the scheduler let part of the burst land late (registry polls = 2, want 1). It now checks the bound that holds: one or two polls for the burst, and none after. Co-Authored-By: Claude Opus 5.5 --- pkg/daemon/daemon.go | 1 + pkg/daemon/zz_handshake_poll_sched_test.go | 20 ++++++++++++++++---- 2 files changed, 17 insertions(+), 4 deletions(-) diff --git a/pkg/daemon/daemon.go b/pkg/daemon/daemon.go index ae82735e..e2f1a883 100644 --- a/pkg/daemon/daemon.go +++ b/pkg/daemon/daemon.go @@ -5539,6 +5539,7 @@ func (d *Daemon) handshakePollLoop() { // Independent jitter so this loop does not align with the others. The // baseline ticker starts once it has passed; requests, pokes and Stop // are served during it. + // #nosec G404 -- startup-jitter scheduling only, not a security decision jitter := time.NewTimer(time.Duration(rand.Int63n(int64(5 * time.Second)))) defer jitter.Stop() var ticker *time.Ticker diff --git a/pkg/daemon/zz_handshake_poll_sched_test.go b/pkg/daemon/zz_handshake_poll_sched_test.go index 56904460..90453381 100644 --- a/pkg/daemon/zz_handshake_poll_sched_test.go +++ b/pkg/daemon/zz_handshake_poll_sched_test.go @@ -244,14 +244,26 @@ func TestHandshakePollLoopExtraPollsAreBounded(t *testing.T) { for i := 0; i < 200; i++ { d.handshakePoke() } - waitPolls(1, 2*time.Second) // served even during the startup jitter + // The first poke is served at once, even during the startup jitter. A + // poke that lands after that poll began is owed one more, deferred by + // handshakeOnDemandGap; whether any of the burst does is a matter of + // scheduling. So the whole burst buys one poll or two, never more. + deadline := time.Now().Add(2 * time.Second) + for d.RelayedHandshakePolls() < 1 && time.Now().Before(deadline) { + time.Sleep(20 * time.Millisecond) + } + time.Sleep(handshakeOnDemandGap + time.Second) + burst := d.RelayedHandshakePolls() + if burst < 1 || burst > 2 { + t.Fatalf("a burst of 200 pokes caused %d registry polls, want 1 or 2", burst) + } time.Sleep(handshakeOnDemandGap + time.Second) - waitPolls(1, 0) // the rest of the burst bought nothing + waitPolls(burst, 0) // and nothing after that d.hsPoll.requestSent(42, true) time.Sleep(2*handshakeFastPollInterval + handshakeFastPollInterval/2) - if got := d.RelayedHandshakePolls(); got < 2 || got > 4 { - t.Fatalf("registry polls = %d after 2.5 fast periods with a request in flight, want 3 (1 + 2)", got) + if got := d.RelayedHandshakePolls() - burst; got < 1 || got > 3 { + t.Fatalf("%d registry polls in 2.5 fast periods with a request in flight, want 2", got) } d.hsPoll.answered(42) time.Sleep(handshakeFastPollInterval + 200*time.Millisecond) // at most one already-armed poll From 22a49dc4da56b201c45e50732ced6cb66fbe106e Mon Sep 17 00:00:00 2001 From: Teo Calin Date: Fri, 2 Oct 2026 13:08:31 +0300 Subject: [PATCH 4/4] daemon: handshake polls are never abandoned; a trusted peer costs no poll Fixes from an independent review of this branch. - An on-demand poll gave up after 3s but had already asked the registry, which empties the node's handshake inbox as it answers: a reply that came later was thrown away, and with it the requests and approvals it carried. A poll now runs on its own goroutine and always completes; a caller only stops waiting for it. - At most one poll is in flight. A caller that arrived behind a stuck poll used to wait 3s and then start another with its own 3s; now it waits for the one in flight, counted from when that poll started, and starts nothing. One slow registry call occupies one pooled connection, not one per caller. - Wait-for-trust polled before checking trust. pilotctl calls it with a zero timeout before every send, connect and ping, so a busy node made up to 30 polls a minute with no handshake in flight and stalled 3s per command when the registry hung. It polls only while a request this node sent that peer is unanswered. - fastActive asked the handshake plugin about trust with the scheduler's lock held; the plugin can hold its own lock across a registry lookup and the tunnel read loop takes the scheduler's lock for every notify. The question is now asked with the lock released. - A beacon notify must be exactly [0x0A][0x01]; a bare type byte or an unknown kind no longer triggers a poll. - The extra timer polls only if a poke or an unanswered request is still owed one, and timer-driven polls keep to the 2s gap. - Peers kept only for their rearm time no longer fill the 64-entry table and deny a new request its fast polling. Co-Authored-By: Claude Opus 5.5 --- CHANGELOG.md | 10 +- pkg/daemon/daemon.go | 25 ++- pkg/daemon/handshakepoll.go | 204 ++++++++++++++---- pkg/daemon/ipc.go | 2 +- pkg/daemon/tunnel.go | 15 +- .../zz_coverage_pkg_daemon_round2_test.go | 4 +- pkg/daemon/zz_handshake_poll_sched_test.go | 163 +++++++++++++- tests/zz_handshake_relay_latency_test.go | 59 ++++- 8 files changed, 422 insertions(+), 60 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1a53610b..388aac8e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -339,9 +339,13 @@ Detailed per-release notes are on the target after ~59s, approval back at the requester ~57s later. Now: - a node that has sent a handshake request polls every 2s until it is answered, for at most 2 minutes; - - `pilotctl pending`, `trust`, `approve`, `reject` and waiting for trust - poll first (at most one such poll every 2s), so a relayed request is - there as soon as someone looks; + - `pilotctl pending`, `trust`, `approve` and `reject` poll first (at most + one such poll every 2s), so a relayed request is there as soon as + someone looks. Waiting for trust polls only while a request this node + sent that peer is unanswered, so checking a peer that is already trusted + — which pilotctl does before every send — costs nothing. The caller is + held at most 3s by a slow registry; the poll itself still completes and + delivers what it fetched; - the beacon can tell a node that something is waiting for it (a two-byte notify that carries nothing else; needs the matching registry/beacon release), and the node polls at once — at most 3 polls in a burst and diff --git a/pkg/daemon/daemon.go b/pkg/daemon/daemon.go index e2f1a883..991b3328 100644 --- a/pkg/daemon/daemon.go +++ b/pkg/daemon/daemon.go @@ -5567,7 +5567,13 @@ func (d *Daemon) handshakePollLoop() { d.pollHandshakes(0, 0) case <-extra.C: extraArmed = false - d.pollHandshakes(handshakeFastPollInterval/2, 0) + // The timer was armed for a poke or a request in flight; poll + // only if one of them is still owed. A request answered while + // the timer ran needs nothing more. + due, wait := d.hsPoll.pokeWait() + if (due && wait == 0) || d.hsPoll.fastActive(trusted) { + d.pollHandshakes(handshakeOnDemandGap-handshakePollSlack, 0) + } case <-d.hsPoll.wake: if due, wait := d.hsPoll.pokeWait(); due && wait == 0 { d.pollHandshakes(handshakeOnDemandGap, 0) @@ -6424,17 +6430,16 @@ func (d *Daemon) lookupPeerPubKey(nodeID uint32) (ed25519.PublicKey, error) { // // timeout > 0 bounds the registry call (a local client is waiting on it); // 0 leaves it unbounded, as the background loop always ran it. -func (d *Daemon) pollRelayedHandshakes(timeout time.Duration) { +func (d *Daemon) pollRelayedHandshakes() { rc, nodeID := d.reg(), d.NodeID() - var resp map[string]interface{} - var err error - if timeout > 0 { - resp, err = withRegistryDeadline(timeout, func() (map[string]interface{}, error) { - return rc.PollHandshakes(nodeID) - }) - } else { - resp, err = rc.PollHandshakes(nodeID) + if rc == nil { + return } + // No deadline of our own. The registry empties this node's handshake + // inbox as it answers; walking away from a slow reply would drop the + // requests and approvals in it. The client's own read deadline bounds + // the call, and callers that cannot wait stop waiting (pollHandshakes). + resp, err := rc.PollHandshakes(nodeID) if err != nil { slog.Debug("poll handshakes failed", "error", err) return diff --git a/pkg/daemon/handshakepoll.go b/pkg/daemon/handshakepoll.go index 007fcfc2..ba1b02dd 100644 --- a/pkg/daemon/handshakepoll.go +++ b/pkg/daemon/handshakepoll.go @@ -53,8 +53,16 @@ const ( // local client or a beacon poke. handshakeOnDemandGap = 2 * time.Second + // handshakePollSlack is how much sooner than handshakeOnDemandGap a + // timer-driven poll may start, so a timer armed for exactly the gap is + // not refused for firing a hair early by the scheduler's clock. + handshakePollSlack = 100 * time.Millisecond + // handshakeOnDemandTimeout bounds how long a local client's request is - // held up by its poll when the registry is slow or unreachable. + // held up by a poll when the registry is slow or unreachable, counted + // from when that poll started. The poll itself is not cut short: the + // registry empties a node's handshake inbox as it answers, so a reply + // nobody waits for would be a request or an approval lost for good. handshakeOnDemandTimeout = 3 * time.Second // handshakePokeBurst / handshakePokeRefill: token bucket for polls @@ -90,10 +98,15 @@ type handshakePollSched struct { // never blocks the sender). wake chan struct{} - // sem serializes polls so two triggers never overlap (capacity 1). - sem chan struct{} + // running is non-nil while a poll is in flight and is closed when it + // finishes. At most one poll is in flight; a trigger that arrives + // meanwhile waits for that one instead of starting another. + running chan struct{} + + // run performs one poll. It is d.pollRelayedHandshakes; tests replace it. + run func() - // polls counts registry polls, for tests and the info reply. + // polls counts registry polls, for tests. polls atomic.Uint64 } @@ -108,7 +121,6 @@ func newHandshakePollSched() *handshakePollSched { waiting: make(map[uint32]handshakeWait), pokeTokens: handshakePokeBurst, wake: make(chan struct{}, 1), - sem: make(chan struct{}, 1), } } @@ -137,6 +149,16 @@ func (s *handshakePollSched) requestSent(peer uint32, explicit bool) { return case !tracked && len(s.waiting) >= handshakeMaxWaiting: s.pruneLocked(now) + if len(s.waiting) >= handshakeMaxWaiting { + // Still full of peers whose window has closed and that are + // only kept for their rearm time. Let those go rather than + // refuse a new request its fast polling. + for p, w := range s.waiting { + if !now.Before(w.until) { + delete(s.waiting, p) + } + } + } if len(s.waiting) >= handshakeMaxWaiting { s.mu.Unlock() return @@ -165,24 +187,52 @@ func (s *handshakePollSched) pruneLocked(now time.Time) { } } +// waitingOn reports whether a request this node sent to peer is still inside +// its fast-poll window. +func (s *handshakePollSched) waitingOn(peer uint32) bool { + s.mu.Lock() + defer s.mu.Unlock() + w, ok := s.waiting[peer] + return ok && s.now().Before(w.until) +} + // fastActive reports whether any request is still inside its fast-poll // window. Peers that trusted reports as trusted are settled first. +// +// trusted is called without s.mu held. It takes the handshake plugin's +// lock, which the plugin can hold across a registry lookup, and s.mu is also +// taken by poke on the tunnel read loop: calling it under s.mu would stall +// inbound packets for as long as that lookup took. func (s *handshakePollSched) fastActive(trusted func(uint32) bool) bool { s.mu.Lock() - defer s.mu.Unlock() now := s.now() s.pruneLocked(now) - active := false + var open []uint32 for peer, w := range s.waiting { - if !now.Before(w.until) { - continue + if now.Before(w.until) { + open = append(open, peer) } + } + s.mu.Unlock() + + active := false + var settled []uint32 + for _, peer := range open { if trusted != nil && trusted(peer) { - w.until = time.Time{} - s.waiting[peer] = w - continue + settled = append(settled, peer) + } else { + active = true } - active = true + } + if len(settled) > 0 { + s.mu.Lock() + for _, peer := range settled { + if w, ok := s.waiting[peer]; ok { + w.until = time.Time{} + s.waiting[peer] = w + } + } + s.mu.Unlock() } return active } @@ -192,6 +242,10 @@ func (s *handshakePollSched) fastActive(trusted func(uint32) bool) bool { func (s *handshakePollSched) claim(minGap time.Duration) bool { s.mu.Lock() defer s.mu.Unlock() + return s.claimLocked(minGap) +} + +func (s *handshakePollSched) claimLocked(minGap time.Duration) bool { now := s.now() if !s.lastPoll.IsZero() && now.Sub(s.lastPoll) < minGap { return false @@ -201,6 +255,43 @@ func (s *handshakePollSched) claim(minGap time.Duration) bool { return true } +// begin is what a trigger calls when it wants a poll. It returns the channel +// that is closed when the poll covering this trigger finishes — nil if none +// is needed because one started within minGap — and whether the caller is +// the one that must run it (and call end when it is done). While a poll is +// in flight every trigger gets that poll's channel and starts nothing. +func (s *handshakePollSched) begin(minGap time.Duration) (done chan struct{}, start bool) { + s.mu.Lock() + defer s.mu.Unlock() + if s.running != nil { + return s.running, false + } + if !s.claimLocked(minGap) { + return nil, false + } + s.running = make(chan struct{}) + return s.running, true +} + +// end marks the poll in flight as finished and wakes the loop, which may owe +// a poll to a poke that arrived while this one was running. +func (s *handshakePollSched) end() { + s.mu.Lock() + if s.running != nil { + close(s.running) + s.running = nil + } + s.mu.Unlock() + s.nudge() +} + +// sinceLastPoll is how long ago the most recent poll started. +func (s *handshakePollSched) sinceLastPoll() time.Duration { + s.mu.Lock() + defer s.mu.Unlock() + return s.now().Sub(s.lastPoll) +} + // poke handles a beacon notification. It returns true when the loop should // be woken: the poke took a token and a poll is now due (immediately, or // once handshakeOnDemandGap has passed since the last one). @@ -239,6 +330,11 @@ func (s *handshakePollSched) pokeWait() (due bool, wait time.Duration) { if !s.pokeDue { return false, 0 } + if s.running != nil { + // A poll is in flight; end() wakes the loop when it finishes. + // Until then there is nothing to do but look again later. + return true, handshakeOnDemandGap + } if s.lastPoll.IsZero() { return true, 0 } @@ -248,12 +344,18 @@ func (s *handshakePollSched) pokeWait() (due bool, wait time.Duration) { return true, wait } -// pollHandshakes runs one relayed-handshake poll unless one started within -// minGap. Polls never overlap; a caller arriving while one is running waits -// for it (up to timeout when timeout > 0) and then finds it fresh enough. -// timeout also bounds the registry call itself; 0 means no bound, which is -// what the background loop has always used. -func (d *Daemon) pollHandshakes(minGap, timeout time.Duration) { +// pollHandshakes makes sure a relayed-handshake poll has run recently: it +// starts one unless one is in flight or started within minGap. Polls never +// overlap. +// +// The poll runs on its own goroutine and always runs to completion, +// processing whatever the registry returns however late. wait is how long +// the caller is prepared to be held up, counted from when the poll it is +// waiting for started; 0 means not at all, which is what the background loop +// uses. Giving up on waiting does not cancel the poll: the registry empties +// the node's handshake inbox as it answers, so a reply that was abandoned +// would be handshakes lost. +func (d *Daemon) pollHandshakes(minGap, wait time.Duration) { s := d.hsPoll if d.reg() == nil { // Nothing to poll; do not leave a poke owed (the loop would keep @@ -263,29 +365,38 @@ func (d *Daemon) pollHandshakes(minGap, timeout time.Duration) { s.mu.Unlock() return } - if timeout > 0 { - t := time.NewTimer(timeout) - defer t.Stop() - select { - case s.sem <- struct{}{}: - case <-t.C: - return - case <-d.stopCh: - return - } - } else { - select { - case s.sem <- struct{}{}: - case <-d.stopCh: - return + done, start := s.begin(minGap) + if done == nil { + return + } + if start { + s.polls.Add(1) + run := s.run + if run == nil { + run = d.pollRelayedHandshakes } + go func() { + defer s.end() + defer recoverLayer("L11", "pollRelayedHandshakes", d.bus, nil) + run() + }() } - defer func() { <-s.sem }() - if !s.claim(minGap) { + if wait <= 0 { return } - s.polls.Add(1) - d.pollRelayedHandshakes(timeout) + remaining := wait - s.sinceLastPoll() + if remaining <= 0 { + // The poll in flight has already been running longer than a + // client should be held; do not add to it. + return + } + t := time.NewTimer(remaining) + defer t.Stop() + select { + case <-done: + case <-t.C: + case <-d.stopCh: + } } // pollHandshakesOnDemand is called before a local client is shown, or acts @@ -295,6 +406,23 @@ func (d *Daemon) pollHandshakesOnDemand() { d.pollHandshakes(handshakeOnDemandGap, handshakeOnDemandTimeout) } +// pollHandshakesForTrustWait is the on-demand poll for wait-for-trust. It +// polls only when there is an answer to wait for: the peer is not trusted +// yet and a request this node sent it is still outstanding. pilotctl asks +// wait-for-trust with a zero timeout before every send, connect and ping, so +// polling unconditionally here put a registry round trip — and, with the +// registry hung, a three-second stall — in front of every command to a peer +// that was trusted all along. +func (d *Daemon) pollHandshakesForTrustWait(peer uint32) { + if d.handshakes != nil && d.handshakes.IsTrusted(peer) { + return + } + if !d.hsPoll.waitingOn(peer) { + return + } + d.pollHandshakesOnDemand() +} + // handshakePoke is the tunnel layer's callback for a beacon notification // that something is waiting for this node at the registry. Runs on the // tunnel read loop, so it only records the poke; the poll loop does the diff --git a/pkg/daemon/ipc.go b/pkg/daemon/ipc.go index 6fc18e39..cad6c010 100644 --- a/pkg/daemon/ipc.go +++ b/pkg/daemon/ipc.go @@ -2089,7 +2089,7 @@ func (s *IPCServer) handleHandshake(conn *ipcConn, reqID uint64, payload []byte) if timeoutMs > 30000 { timeoutMs = 30000 } - s.daemon.pollHandshakesOnDemand() + s.daemon.pollHandshakesForTrustWait(nodeID) ok := s.daemon.handshakes.WaitForTrust(nodeID, time.Duration(timeoutMs)*time.Millisecond) data, _ := json.Marshal(map[string]interface{}{ "node_id": nodeID, diff --git a/pkg/daemon/tunnel.go b/pkg/daemon/tunnel.go index 8bfade16..34aeb1e8 100644 --- a/pkg/daemon/tunnel.go +++ b/pkg/daemon/tunnel.go @@ -170,6 +170,10 @@ type TunnelManager struct { trustGate func(nodeID uint32) bool peerTrustFn func(nodeID uint32) bool + // onBeaconNotify is run when the beacon says the registry holds + // something for this node (beaconMsgNotify). + onBeaconNotify atomic.Pointer[func()] + // Event bus — replaces inline tm.webhook.Emit calls. Set via // SetEventBus from daemon during construction. May be nil in // tests; tm.publishEvent handles that. Webhook delivery is a @@ -201,9 +205,6 @@ type TunnelManager struct { // ConnectCompat so the age is measured from transport start, not epoch. LastRecvNano int64 - // onBeaconNotify is run when the beacon says the registry holds - // something for this node (beaconMsgNotify). - onBeaconNotify atomic.Pointer[func()] // P1-008: packets dropped from the per-peer pending queue while waiting // for key exchange. Exposed so operators can tell a congested overlay // apart from a silent crypto stall. @@ -2331,6 +2332,11 @@ func (tm *TunnelManager) handleBeaconMessage(data []byte, from *net.UDPAddr) { slog.Debug("dropping notify from non-beacon source", "from", from) return } + if len(data) != 2 || data[1] != beaconNotifyHandshake { + // A kind this daemon does not know is not a reason to poll. + slog.Debug("dropping beacon notify of unknown kind", "len", len(data)) + return + } if fn := tm.onBeaconNotify.Load(); fn != nil { (*fn)() } @@ -2347,6 +2353,9 @@ func (tm *TunnelManager) handleBeaconMessage(data []byte, from *net.UDPAddr) { // message types live in common/protocol and this one should move there.) const beaconMsgNotify byte = 0x0A +// beaconNotifyHandshake is the notify kind for a relayed trust handshake. +const beaconNotifyHandshake byte = 0x01 + // SetBeaconNotifyHandler installs the callback run for each beacon notify. // It is called on the tunnel read loop and must not block. func (tm *TunnelManager) SetBeaconNotifyHandler(fn func()) { diff --git a/pkg/daemon/zz_coverage_pkg_daemon_round2_test.go b/pkg/daemon/zz_coverage_pkg_daemon_round2_test.go index 4dea5bf0..eff4addb 100644 --- a/pkg/daemon/zz_coverage_pkg_daemon_round2_test.go +++ b/pkg/daemon/zz_coverage_pkg_daemon_round2_test.go @@ -165,7 +165,7 @@ func TestPollRelayedHandshakesEmptyNoOp(t *testing.T) { registerSelfOnRegistry(t, d) // No panic, no side effect. The empty-list branches are still walked. - d.pollRelayedHandshakes(0) + d.pollRelayedHandshakes() } // pollRelayedHandshakesHasService verifies that when a HandshakeService is @@ -208,7 +208,7 @@ func TestPollRelayedHandshakesWithServiceNoOp(t *testing.T) { // Empty mailbox — the code walks the empty requests/responses lists // and the service-non-nil guards but does not invoke any Process* method. - d.pollRelayedHandshakes(0) + d.pollRelayedHandshakes() if svc.requests.Load() != 0 || svc.approvals.Load() != 0 || svc.rejections.Load() != 0 { t.Fatalf("unexpected service calls: req=%d approve=%d reject=%d", diff --git a/pkg/daemon/zz_handshake_poll_sched_test.go b/pkg/daemon/zz_handshake_poll_sched_test.go index 90453381..6d0b3494 100644 --- a/pkg/daemon/zz_handshake_poll_sched_test.go +++ b/pkg/daemon/zz_handshake_poll_sched_test.go @@ -4,6 +4,7 @@ package daemon import ( "net" + "sync/atomic" "testing" "time" ) @@ -263,7 +264,7 @@ func TestHandshakePollLoopExtraPollsAreBounded(t *testing.T) { d.hsPoll.requestSent(42, true) time.Sleep(2*handshakeFastPollInterval + handshakeFastPollInterval/2) if got := d.RelayedHandshakePolls() - burst; got < 1 || got > 3 { - t.Fatalf("%d registry polls in 2.5 fast periods with a request in flight, want 2", got) + t.Fatalf("%d registry polls in 2.5 fast periods with a request in flight, want 1 to 3", got) } d.hsPoll.answered(42) time.Sleep(handshakeFastPollInterval + 200*time.Millisecond) // at most one already-armed poll @@ -273,3 +274,163 @@ func TestHandshakePollLoopExtraPollsAreBounded(t *testing.T) { t.Fatalf("loop kept polling after the answer: %d -> %d", settled, got) } } + +// A poll the caller stops waiting for must still run to completion. The +// registry empties a node's handshake inbox as it answers, so a reply that +// arrives after its caller has gone and is then thrown away is a request or +// an approval lost for good — `pilotctl pending` against a registry that took +// four seconds lost whatever was in the inbox. +func TestOnDemandPollOutlivesTheCallerThatGaveUp(t *testing.T) { + t.Parallel() + reg, rc := startTestRegistry(t) + t.Cleanup(func() { reg.Close() }) + t.Cleanup(func() { rc.Close() }) + d := New(Config{KeepaliveInterval: time.Hour}) + d.regConn.Store(rc) + + release := make(chan struct{}) + var started, finished atomic.Int32 + d.hsPoll.run = func() { + started.Add(1) + <-release // a registry that is slow to answer + finished.Add(1) + } + + const wait = 80 * time.Millisecond + begin := time.Now() + d.pollHandshakes(handshakeOnDemandGap, wait) + if held := time.Since(begin); held > wait+500*time.Millisecond { + t.Fatalf("caller was held %v, want about %v", held, wait) + } + if started.Load() != 1 || finished.Load() != 0 { + t.Fatalf("after the caller gave up: started=%d finished=%d, want the poll still in flight", started.Load(), finished.Load()) + } + + // Callers arriving while it is in flight start no second poll, and one + // arriving after the wait has already been used up is not held at all. + time.Sleep(wait) + begin = time.Now() + for i := 0; i < 5; i++ { + d.pollHandshakes(handshakeOnDemandGap, wait) + } + if held := time.Since(begin); held > wait { + t.Fatalf("five callers behind a poll already running past its wait were held %v in total", held) + } + if n := started.Load(); n != 1 { + t.Fatalf("%d polls started while one was in flight, want 1", n) + } + + // The reply lands: the poll finishes its work. + close(release) + deadline := time.Now().Add(2 * time.Second) + for finished.Load() != 1 && time.Now().Before(deadline) { + time.Sleep(5 * time.Millisecond) + } + if finished.Load() != 1 { + t.Fatal("the poll its caller stopped waiting for never completed") + } + if got := d.RelayedHandshakePolls(); got != 1 { + t.Fatalf("registry polls = %d, want 1", got) + } +} + +// fastActive asks the handshake plugin whether a peer is trusted, and the +// plugin can hold its lock across a registry lookup. The tunnel read loop +// takes the scheduler's lock for every beacon notify, so that question must +// not be asked with the scheduler's lock held. +func TestFastActiveAsksAboutTrustWithoutHoldingTheLock(t *testing.T) { + t.Parallel() + s := newHandshakePollSched() + s.requestSent(7, true) + s.requestSent(8, true) + + result := make(chan bool, 1) + go func() { + result <- s.fastActive(func(peer uint32) bool { + s.poke() // what the read loop does; deadlocks if s.mu is held + return peer == 7 + }) + }() + select { + case active := <-result: + if !active { + t.Fatal("peer 8 is still waiting: fastActive should report true") + } + case <-time.After(2 * time.Second): + t.Fatal("fastActive called the trust check with the scheduler's lock held") + } + if s.waitingOn(7) { + t.Fatal("peer 7 was reported trusted and should be settled") + } + if !s.waitingOn(8) { + t.Fatal("peer 8 should still be waited on") + } +} + +// Peers stay tracked for handshakeAutoRearm after their window closes, so +// that an automatic handshake cannot reopen it. Sixty-four of those must not +// stop a new request from getting its fast polling. +func TestSettledPeersDoNotCrowdOutANewRequest(t *testing.T) { + t.Parallel() + s := newHandshakePollSched() + for peer := uint32(1); peer <= handshakeMaxWaiting; peer++ { + s.requestSent(peer, true) + s.answered(peer) + } + s.requestSent(1000, true) + if !s.waitingOn(1000) { + t.Fatalf("with %d answered peers tracked, a new request got no fast polling", handshakeMaxWaiting) + } + + // Still a cap on requests that are genuinely outstanding. + s = newHandshakePollSched() + for peer := uint32(1); peer <= handshakeMaxWaiting; peer++ { + s.requestSent(peer, true) + } + s.requestSent(1000, true) + if s.waitingOn(1000) { + t.Fatalf("a request beyond %d outstanding ones was tracked", handshakeMaxWaiting) + } +} + +// Wait-for-trust polls only when there is an answer to wait for. pilotctl +// calls it with a zero timeout before every send, connect and ping. +func TestWaitingOnTracksOnlyOutstandingRequests(t *testing.T) { + t.Parallel() + s := newHandshakePollSched() + if s.waitingOn(5) { + t.Fatal("no request was sent to peer 5") + } + s.requestSent(5, true) + if !s.waitingOn(5) { + t.Fatal("a request to peer 5 is outstanding") + } + s.answered(5) + if s.waitingOn(5) { + t.Fatal("peer 5 answered") + } +} + +// Only the handshake kind triggers a poll. A bare type byte, or a kind this +// daemon does not know, is dropped. +func TestBeaconNotifyOfUnknownKindIsDropped(t *testing.T) { + t.Parallel() + tm := NewTunnelManager() + if err := tm.SetBeaconAddr("192.0.2.10:9001"); err != nil { + t.Fatal(err) + } + calls := 0 + tm.SetBeaconNotifyHandler(func() { calls++ }) + beacon := &net.UDPAddr{IP: net.ParseIP("192.0.2.10"), Port: 9001} + + tm.handleBeaconMessage([]byte{beaconMsgNotify}, beacon) + tm.handleBeaconMessage([]byte{beaconMsgNotify, 0x02}, beacon) + tm.handleBeaconMessage([]byte{beaconMsgNotify, beaconNotifyHandshake, 0x00}, beacon) + if calls != 0 { + t.Fatalf("a malformed or unknown notify ran the handler %d times", calls) + } + tm.handleBeaconMessage([]byte{beaconMsgNotify, beaconNotifyHandshake}, beacon) + if calls != 1 { + t.Fatalf("a handshake notify ran the handler %d times, want 1", calls) + } +} diff --git a/tests/zz_handshake_relay_latency_test.go b/tests/zz_handshake_relay_latency_test.go index 706a71c3..a83def44 100644 --- a/tests/zz_handshake_relay_latency_test.go +++ b/tests/zz_handshake_relay_latency_test.go @@ -61,8 +61,12 @@ func TestRelayedHandshakeCompletesInSeconds(t *testing.T) { t.Fatalf("approve: %v", err) } - // Requester side: nobody asks it anything; its own polling must pick - // the approval up. + // Requester side: nobody asks it anything. In this harness both nodes + // sit on loopback, so once the target has approved it can usually reach + // the requester directly and the approval arrives that way, with the + // requester's fast poll as the fallback. The fast poll on its own is + // exercised by the scheduler tests in pkg/daemon and was measured + // between containers that cannot reach each other directly. for !a.Daemon.HandshakeService().IsTrusted(b.Daemon.NodeID()) { if time.Since(approved) > limit { t.Fatalf("relayed approval did not reach the requester within %s", limit) @@ -128,3 +132,54 @@ func TestRelayedHandshakeReachesUnattendedTarget(t *testing.T) { t.Errorf("target polled the registry %d times for one request", polls) } } + +// TestWaitForTrustOnTrustedPeerDoesNotPollRegistry pins that asking whether +// a peer is trusted costs no registry round trip once it is. pilotctl asks +// with a zero timeout before every send-message, send-file, connect, ping +// and publish; polling there put up to thirty registry requests a minute on +// a busy node with no handshake in flight, and a three-second stall in front +// of each command whenever the registry was slow. +func TestWaitForTrustOnTrustedPeerDoesNotPollRegistry(t *testing.T) { + requireRealNetwork(t) + t.Parallel() + env := NewTestEnv(t) + a := env.AddDaemon() + b := env.AddDaemon() + + if _, err := a.Driver.Handshake(b.Daemon.NodeID(), "trusted-wait test"); err != nil { + t.Fatalf("A handshake: %v", err) + } + if _, err := b.Driver.Handshake(a.Daemon.NodeID(), "trusted-wait test"); err != nil { + t.Fatalf("B handshake: %v", err) + } + deadline := time.Now().Add(15 * time.Second) + for { + resp, err := a.Driver.WaitForTrust(b.Daemon.NodeID(), 1000) + if trusted, _ := resp["trusted"].(bool); err == nil && trusted { + break + } + if time.Now().After(deadline) { + t.Fatal("A and B did not become trusted") + } + } + + before := a.Daemon.RelayedHandshakePolls() + for i := 0; i < 5; i++ { + start := time.Now() + resp, err := a.Driver.WaitForTrust(b.Daemon.NodeID(), 0) + if err != nil { + t.Fatalf("WaitForTrust: %v", err) + } + if trusted, _ := resp["trusted"].(bool); !trusted { + t.Fatal("B should be trusted") + } + if took := time.Since(start); took > time.Second { + t.Fatalf("WaitForTrust on a trusted peer took %v", took) + } + // Past the on-demand gap, so a poll would be allowed if one were asked for. + time.Sleep(2100 * time.Millisecond) + } + if got := a.Daemon.RelayedHandshakePolls(); got != before { + t.Fatalf("five trust checks on a trusted peer caused %d registry polls, want 0", got-before) + } +}