diff --git a/CHANGELOG.md b/CHANGELOG.md index ea3083b8..388aac8e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -332,6 +332,26 @@ 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` 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 + 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..991b3328 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,77 @@ 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. + // #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 + 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 + // 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) } - d.pollRelayedHandshakes() + 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: + } + } + extraArmed = false + case next >= 0 && !extraArmed: + extra.Reset(next) + extraArmed = true } } } @@ -6349,8 +6427,19 @@ 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. +// +// 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() { - resp, err := d.reg().PollHandshakes(d.NodeID()) + rc, nodeID := d.reg(), d.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 @@ -6391,6 +6480,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..ba1b02dd --- /dev/null +++ b/pkg/daemon/handshakepoll.go @@ -0,0 +1,438 @@ +// 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 + + // 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 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 + // 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{} + + // 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. + 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), + } +} + +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 { + // 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 + } + } + 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) + } + } +} + +// 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() + now := s.now() + s.pruneLocked(now) + var open []uint32 + for peer, w := range s.waiting { + 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) { + settled = append(settled, peer) + } else { + 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 +} + +// 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() + 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 + } + s.lastPoll = now + s.pokeDue = false + 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). +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.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 + } + if wait = handshakeOnDemandGap - s.now().Sub(s.lastPoll); wait < 0 { + wait = 0 + } + return true, wait +} + +// 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 + // re-arming its timer for a poll that cannot run). + s.mu.Lock() + s.pokeDue = false + s.mu.Unlock() + 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() + }() + } + if wait <= 0 { + return + } + 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 +// 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) +} + +// 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 +// 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..cad6c010 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.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 ea9d1961..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 @@ -200,6 +204,7 @@ 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 + // 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 +2323,45 @@ 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 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)() + } 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 + +// 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()) { + 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_handshake_poll_sched_test.go b/pkg/daemon/zz_handshake_poll_sched_test.go new file mode 100644 index 00000000..6d0b3494 --- /dev/null +++ b/pkg/daemon/zz_handshake_poll_sched_test.go @@ -0,0 +1,436 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +package daemon + +import ( + "net" + "sync/atomic" + "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() + } + // 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(burst, 0) // and nothing after that + + 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 1 to 3", 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) + } +} + +// 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/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..a83def44 --- /dev/null +++ b/tests/zz_handshake_relay_latency_test.go @@ -0,0 +1,185 @@ +// 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. 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) + } + 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) + } +} + +// 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) + } +} 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 + } + } +}