mirror of
https://github.com/netbirdio/netbird.git
synced 2026-07-23 00:41:29 +02:00
[client] Create the WireGuard peer before the peer becomes routable
Split peer setup into a prepare and a start phase. AddPeerConn now only registers the conn in the peer store and, in lazy mode, arms the wake endpoint (creating the WireGuard peer) without starting any event source. The peer is then registered in the status recorder, and StartPeerConn starts the activity listener, opens the connection or applies HA activation afterwards. This closes the cold-start race where the route watcher, reacting to the freshly registered idle peer, pushed routed allowed IPs before the WireGuard peer existed: the update_only add was silently dropped while the allowed-IP refcounter recorded the prefix as installed, so nothing retried and traffic to the routed subnet stayed black-holed until the peer was woken by other means. Packets arriving on the armed wake endpoint before the listener starts queue in the socket buffer, and no status write can land before the recorder entry exists because no event source runs in between.
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user