Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
112 changes: 101 additions & 11 deletions pkg/daemon/daemon.go
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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)
}
}()
}
}
}
Expand Down Expand Up @@ -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))))
Comment thread
github-advanced-security[bot] marked this conversation as resolved.
Fixed
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
}
}
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading