diff --git a/client/internal/conn_mgr.go b/client/internal/conn_mgr.go index e1df336f8..c19c62721 100644 --- a/client/internal/conn_mgr.go +++ b/client/internal/conn_mgr.go @@ -176,23 +176,20 @@ func (e *ConnMgr) SetExcludeList(ctx context.Context, peerIDs map[string]bool) { } } -func (e *ConnMgr) AddPeerConn(ctx context.Context, peerKey string, conn *peer.Conn) (exists bool) { +// AddPeerConn registers the conn in the peer store and, in lazy mode, arms the wake endpoint. +// It starts no connection activity; the caller calls StartPeerConn after registering the peer +// in the status recorder. +func (e *ConnMgr) AddPeerConn(peerKey string, conn *peer.Conn) (exists bool) { if success := e.peerStore.AddPeerConn(peerKey, conn); !success { return true } if !e.isStartedWithLazyMgr() { - if err := conn.Open(ctx); err != nil { - conn.Log.Errorf("failed to open connection: %v", err) - } return } if !lazyconn.IsSupported(conn.AgentVersionString()) { - conn.Log.Warnf("peer does not support lazy connection (%s), open permanent connection", conn.AgentVersionString()) - if err := conn.Open(ctx); err != nil { - conn.Log.Errorf("failed to open connection: %v", err) - } + conn.Log.Warnf("peer does not support lazy connection (%s), will open permanent connection", conn.AgentVersionString()) return } @@ -204,18 +201,12 @@ func (e *ConnMgr) AddPeerConn(ctx context.Context, peerKey string, conn *peer.Co } excluded, err := e.lazyConnMgr.AddPeer(lazyPeerCfg) if err != nil { - conn.Log.Errorf("failed to add peer to lazyconn manager: %v", err) - if err := conn.Open(ctx); err != nil { - conn.Log.Errorf("failed to open connection: %v", err) - } + conn.Log.Errorf("failed to add peer to lazyconn manager, will open permanent connection: %v", err) return } if excluded { - conn.Log.Infof("peer is on lazy conn manager exclude list, opening connection") - if err := conn.Open(ctx); err != nil { - conn.Log.Errorf("failed to open connection: %v", err) - } + conn.Log.Infof("peer is on lazy conn manager exclude list, will open permanent connection") return } @@ -223,6 +214,23 @@ func (e *ConnMgr) AddPeerConn(ctx context.Context, peerKey string, conn *peer.Co return } +// StartPeerConn starts the activity prepared by AddPeerConn: lazy wake listening or a +// permanent connection. +func (e *ConnMgr) StartPeerConn(ctx context.Context, peerKey string) { + conn, ok := e.peerStore.PeerConn(peerKey) + if !ok { + return + } + + if e.isStartedWithLazyMgr() && e.lazyConnMgr.StartPeer(peerKey) { + return + } + + if err := conn.Open(ctx); err != nil { + conn.Log.Errorf("failed to open connection: %v", err) + } +} + func (e *ConnMgr) RemovePeerConn(peerKey string) { conn, ok := e.peerStore.Remove(peerKey) if !ok { diff --git a/client/internal/engine.go b/client/internal/engine.go index f689dd2a6..a1970a5c3 100644 --- a/client/internal/engine.go +++ b/client/internal/engine.go @@ -1771,16 +1771,20 @@ func (e *Engine) addNewPeer(peerConfig *mgmProto.RemotePeerConfig) error { return fmt.Errorf("create peer connection: %w", err) } + // Order matters: the WireGuard peer must exist before the peer becomes visible in the + // status recorder, and event sources may start only after the recorder registration. + if exists := e.connMgr.AddPeerConn(peerKey, conn); exists { + conn.Close() + return fmt.Errorf("peer already exists: %s", peerKey) + } + peerV4, peerV6 := overlayAddrsFromAllowedIPs(peerConfig.GetAllowedIps(), e.wgInterface.Address().IPv6Net) err = e.statusRecorder.AddPeer(peerKey, peerConfig.Fqdn, addrToString(peerV4), addrToString(peerV6)) if err != nil { log.Warnf("error adding peer %s to status recorder, got error: %v", peerKey, err) } - if exists := e.connMgr.AddPeerConn(e.ctx, peerKey, conn); exists { - conn.Close() - return fmt.Errorf("peer already exists: %s", peerKey) - } + e.connMgr.StartPeerConn(e.ctx, peerKey) return nil } diff --git a/client/internal/lazyconn/activity/manager.go b/client/internal/lazyconn/activity/manager.go index 4666f689f..63df5b9dd 100644 --- a/client/internal/lazyconn/activity/manager.go +++ b/client/internal/lazyconn/activity/manager.go @@ -34,12 +34,17 @@ type WgInterface interface { MTU() uint16 } +type managedListener struct { + l listener + started bool +} + type Manager struct { OnActivityChan chan Event wgIface WgInterface - peers map[peerid.ConnID]listener + peers map[peerid.ConnID]*managedListener done chan struct{} mu sync.Mutex @@ -49,13 +54,24 @@ func NewManager(wgIface WgInterface) *Manager { m := &Manager{ OnActivityChan: make(chan Event, 1), wgIface: wgIface, - peers: make(map[peerid.ConnID]listener), + peers: make(map[peerid.ConnID]*managedListener), done: make(chan struct{}), } return m } +// MonitorPeerActivity creates the peer's activity listener and starts consuming traffic on it. func (m *Manager) MonitorPeerActivity(peerCfg lazyconn.PeerConfig) error { + if err := m.CreatePeerListener(peerCfg); err != nil { + return err + } + m.StartPeerListener(peerCfg.PeerConnID) + return nil +} + +// CreatePeerListener arms the wake endpoint without starting to consume traffic; packets queue +// in the listener socket until StartPeerListener runs. +func (m *Manager) CreatePeerListener(peerCfg lazyconn.PeerConfig) error { m.mu.Lock() defer m.mu.Unlock() @@ -69,11 +85,23 @@ func (m *Manager) MonitorPeerActivity(peerCfg lazyconn.PeerConfig) error { return err } - m.peers[peerCfg.PeerConnID] = listener - go m.waitForTraffic(listener, peerCfg.PeerConnID) + m.peers[peerCfg.PeerConnID] = &managedListener{l: listener} return nil } +// StartPeerListener starts consuming traffic on the peer's armed activity listener. +func (m *Manager) StartPeerListener(peerConnID peerid.ConnID) { + m.mu.Lock() + defer m.mu.Unlock() + + ml, ok := m.peers[peerConnID] + if !ok || ml.started { + return + } + ml.started = true + go m.waitForTraffic(ml.l, peerConnID) +} + func (m *Manager) createListener(peerCfg lazyconn.PeerConfig) (listener, error) { if !m.wgIface.IsUserspaceBind() { return NewUDPListener(m.wgIface, peerCfg) @@ -91,13 +119,13 @@ func (m *Manager) RemovePeer(log *log.Entry, peerConnID peerid.ConnID) { m.mu.Lock() defer m.mu.Unlock() - listener, ok := m.peers[peerConnID] + ml, ok := m.peers[peerConnID] if !ok { return } log.Debugf("removing activity listener") delete(m.peers, peerConnID) - listener.Close() + closeListener(ml) } func (m *Manager) Close() { @@ -105,12 +133,22 @@ func (m *Manager) Close() { defer m.mu.Unlock() close(m.done) - for peerID, listener := range m.peers { + for peerID, ml := range m.peers { delete(m.peers, peerID) - listener.Close() + closeListener(ml) } } +// closeListener closes the listener. Close waits for the reader goroutine, so for a +// never-started listener the reader is launched first; outside waitForTraffic it cannot emit events. +func closeListener(ml *managedListener) { + if !ml.started { + ml.started = true + go ml.l.ReadPackets() + } + ml.l.Close() +} + func (m *Manager) waitForTraffic(l listener, peerConnID peerid.ConnID) { l.ReadPackets() diff --git a/client/internal/lazyconn/activity/manager_test.go b/client/internal/lazyconn/activity/manager_test.go index c8416a257..9bd757a9d 100644 --- a/client/internal/lazyconn/activity/manager_test.go +++ b/client/internal/lazyconn/activity/manager_test.go @@ -49,8 +49,11 @@ func (m *Manager) GetPeerListener(peerConnID peerid.ConnID) (listener, bool) { m.mu.Lock() defer m.mu.Unlock() - l, exists := m.peers[peerConnID] - return l, exists + ml, exists := m.peers[peerConnID] + if !exists { + return nil, false + } + return ml.l, true } func TestManager_MonitorPeerActivity(t *testing.T) { diff --git a/client/internal/lazyconn/manager/manager.go b/client/internal/lazyconn/manager/manager.go index 5808260cf..99585c3d6 100644 --- a/client/internal/lazyconn/manager/manager.go +++ b/client/internal/lazyconn/manager/manager.go @@ -185,6 +185,7 @@ func (m *Manager) ExcludePeer(peerConfigs []lazyconn.PeerConfig) []string { return added } +// AddPeer arms the peer's wake endpoint without starting to consume traffic on it. func (m *Manager) AddPeer(peerCfg lazyconn.PeerConfig) (bool, error) { m.managedPeersMu.Lock() defer m.managedPeersMu.Unlock() @@ -201,7 +202,7 @@ func (m *Manager) AddPeer(peerCfg lazyconn.PeerConfig) (bool, error) { return false, nil } - if err := m.activityManager.MonitorPeerActivity(peerCfg); err != nil { + if err := m.activityManager.CreatePeerListener(peerCfg); err != nil { return false, err } @@ -211,13 +212,29 @@ func (m *Manager) AddPeer(peerCfg lazyconn.PeerConfig) (bool, error) { expectedWatcher: watcherActivity, } - // Check if this peer should be activated because its HA group peers are active - if group, ok := m.shouldActivateNewPeer(peerCfg.PublicKey); ok { - peerCfg.Log.Debugf("peer belongs to active HA group %s, will activate immediately", group) - m.activateNewPeerInActiveGroup(peerCfg) + return false, nil +} + +// StartPeer starts the peer's activity listener and applies HA activation. It reports whether +// the peer is managed by the lazy manager. +func (m *Manager) StartPeer(peerKey string) bool { + m.managedPeersMu.Lock() + defer m.managedPeersMu.Unlock() + + peerCfg, ok := m.managedPeers[peerKey] + if !ok { + return false } - return false, nil + m.activityManager.StartPeerListener(peerCfg.PeerConnID) + + // Check if this peer should be activated because its HA group peers are active + if group, ok := m.shouldActivateNewPeer(peerKey); ok { + peerCfg.Log.Debugf("peer belongs to active HA group %s, will activate immediately", group) + m.activateNewPeerInActiveGroup(*peerCfg) + } + + return true } // AddActivePeers adds a list of peers to the lazy connection manager