Compare commits

..

1 Commits

Author SHA1 Message Date
pascal
ec3d627a69 reduce hashing on relevant groups calc 2026-07-24 15:06:24 +02:00
8 changed files with 59 additions and 77 deletions

View File

@@ -464,8 +464,6 @@ func Test_RemovePeer(t *testing.T) {
}
func Test_ConnectPeers(t *testing.T) {
t.Setenv("NB_DISABLE_EBPF_WG_PROXY", "true")
peer1ifaceName := fmt.Sprintf("utun%d", WgIntNumber+400)
peer1wgIP := netip.MustParsePrefix("10.99.99.17/30")
peer1Key, _ := wgtypes.GeneratePrivateKey()
@@ -507,8 +505,12 @@ func Test_ConnectPeers(t *testing.T) {
t.Fatal(err)
}
localIP1 := "127.0.0.1"
peer1endpoint, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", localIP1, peer1wgPort))
localIP, err := getLocalIP()
if err != nil {
t.Fatal(err)
}
peer1endpoint, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", localIP, peer1wgPort))
if err != nil {
t.Fatal(err)
}
@@ -544,8 +546,7 @@ func Test_ConnectPeers(t *testing.T) {
t.Fatal(err)
}
localIP2 := "127.0.0.1"
peer2endpoint, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", localIP2, peer2wgPort))
peer2endpoint, err := net.ResolveUDPAddr("udp", fmt.Sprintf("%s:%d", localIP, peer2wgPort))
if err != nil {
t.Fatal(err)
}
@@ -620,3 +621,28 @@ func getPeer(ifaceName, peerPubKey string) (wgtypes.Peer, error) {
}
return wgtypes.Peer{}, fmt.Errorf("peer not found")
}
func getLocalIP() (string, error) {
// Get all interfaces
addrs, err := net.InterfaceAddrs()
if err != nil {
return "", err
}
for _, addr := range addrs {
ipNet, ok := addr.(*net.IPNet)
if !ok {
continue
}
if ipNet.IP.IsLoopback() {
continue
}
if ipNet.IP.To4() == nil {
continue
}
return ipNet.IP.String(), nil
}
return "", fmt.Errorf("no local IP found")
}

View File

@@ -841,8 +841,6 @@ func (conn *Conn) enableWgWatcherIfNeeded(enabledTime time.Time) {
return
}
conn.endpointUpdater.EnableKeepAlive()
watcher := NewWGWatcher(conn.Log, conn.config.WgConfig.WgInterface, conn.config.Key, conn.dumpState)
watcher.PrepareInitialHandshake()
@@ -890,7 +888,6 @@ func (conn *Conn) resetEndpoint() {
return
}
conn.Log.Infof("reset wg endpoint")
conn.endpointUpdater.EnableKeepAlive()
if conn.wgWatcher != nil {
conn.wgWatcher.Reset()
}
@@ -948,12 +945,7 @@ func (conn *Conn) onWGHandshakeSuccess(when time.Time) {
func (conn *Conn) onWGCheckSuccess() {
conn.mu.Lock()
conn.wgTimeouts = 0
presharedKey := conn.presharedKey(conn.rosenpassRemoteKey)
conn.mu.Unlock()
if err := conn.endpointUpdater.DisableKeepAlive(presharedKey); err != nil {
conn.Log.Warnf("failed to disable WireGuard keepalive: %v", err)
}
}
// recordConnectionMetrics records connection stage timestamps as metrics

View File

@@ -316,16 +316,11 @@ func newWGTimeoutTestConn(rosenpassEnabled bool, disconnected *[]string) *Conn {
cfg.RosenpassConfig = RosenpassConfig{PubKey: []byte("dummykey")}
}
connLog := log.WithField("peer", cfg.Key)
endpointUpdater := NewEndpointUpdater(connLog, cfg.WgConfig, false)
endpointUpdater.keepAlive = 0
conn := &Conn{
ctx: context.Background(),
config: cfg,
Log: connLog,
metricsStages: &MetricsStages{},
endpointUpdater: endpointUpdater,
ctx: context.Background(),
config: cfg,
Log: log.WithField("peer", cfg.Key),
metricsStages: &MetricsStages{},
}
conn.SetOnDisconnected(func(remotePeer string) {
*disconnected = append(*disconnected, remotePeer)

View File

@@ -20,11 +20,10 @@ type EndpointUpdater struct {
wgConfig WgConfig
initiator bool
// mu protects cancelFunc and keepAlive
// mu protects cancelFunc
mu sync.Mutex
cancelFunc func()
updateWg sync.WaitGroup
keepAlive time.Duration
}
func NewEndpointUpdater(log *logrus.Entry, wgConfig WgConfig, initiator bool) *EndpointUpdater {
@@ -32,7 +31,6 @@ func NewEndpointUpdater(log *logrus.Entry, wgConfig WgConfig, initiator bool) *E
log: log,
wgConfig: wgConfig,
initiator: initiator,
keepAlive: defaultWgKeepAlive,
}
}
@@ -75,31 +73,6 @@ func (e *EndpointUpdater) RemoveEndpointAddress() error {
return e.wgConfig.WgInterface.RemoveEndpointAddress(e.wgConfig.RemoteKey)
}
func (e *EndpointUpdater) DisableKeepAlive(presharedKey *wgtypes.Key) error {
if !isWgKeepAliveDisabled() {
return nil
}
e.mu.Lock()
defer e.mu.Unlock()
if e.keepAlive == 0 {
return nil
}
e.waitForCloseTheDelayedUpdate()
e.keepAlive = 0
e.log.Debugf("disable WireGuard persistent keepalive")
return e.updateWireGuardPeer(nil, presharedKey)
}
func (e *EndpointUpdater) EnableKeepAlive() {
e.mu.Lock()
defer e.mu.Unlock()
e.keepAlive = defaultWgKeepAlive
}
func (e *EndpointUpdater) configureAsInitiator(addr *net.UDPAddr, presharedKey *wgtypes.Key) error {
if err := e.updateWireGuardPeer(addr, presharedKey); err != nil {
return err
@@ -154,7 +127,7 @@ func (e *EndpointUpdater) updateWireGuardPeer(endpoint *net.UDPAddr, presharedKe
return e.wgConfig.WgInterface.UpdatePeer(
e.wgConfig.RemoteKey,
e.wgConfig.AllowedIps,
e.keepAlive,
defaultWgKeepAlive,
endpoint,
presharedKey,
)

View File

@@ -7,9 +7,8 @@ import (
)
const (
EnvKeyNBForceRelay = "NB_FORCE_RELAY"
EnvKeyNBHomeRelayServers = "NB_HOME_RELAY_SERVERS"
EnvKeyNBDisableWgKeepAlive = "NB_DISABLE_WG_KEEP_ALIVE"
EnvKeyNBForceRelay = "NB_FORCE_RELAY"
EnvKeyNBHomeRelayServers = "NB_HOME_RELAY_SERVERS"
)
func IsForceRelayed() bool {
@@ -19,10 +18,6 @@ func IsForceRelayed() bool {
return strings.EqualFold(os.Getenv(EnvKeyNBForceRelay), "true")
}
func isWgKeepAliveDisabled() bool {
return strings.EqualFold(os.Getenv(EnvKeyNBDisableWgKeepAlive), "true")
}
// OverrideRelayURLs returns the relay server URL list set in
// NB_HOME_RELAY_SERVERS (comma-separated) and a boolean indicating whether
// the override is active. When the env var is unset, the boolean is false

View File

@@ -109,15 +109,10 @@ func (w *WGWatcher) periodicHandshakeCheck(ctx context.Context, onDisconnectedFn
}
lastHandshake = *handshake
w.stateDump.WGcheckSuccess()
if isWgKeepAliveDisabled() {
w.log.Debugf("WireGuard watcher waiting for peer reset")
continue
}
resetTime := time.Until(handshake.Add(checkPeriod))
timer.Reset(resetTime)
w.stateDump.WGcheckSuccess()
w.log.Debugf("WireGuard watcher reset timer: %v", resetTime)
case <-w.resetCh:

View File

@@ -321,10 +321,19 @@ func (a *Account) getPeersGroupsPoliciesRoutes(
relevantPeerIDs[peerID] = a.GetPeer(peerID).ToComponent()
addRelevantGroup := func(groupID string) *Group {
if g, ok := relevantGroupIDs[groupID]; ok {
return g
}
g := a.GetGroup(groupID)
relevantGroupIDs[groupID] = g
return g
}
peerGroupSet := make(map[string]struct{}, 8)
for groupID, group := range a.Groups {
if slices.Contains(group.Peers, peerID) {
relevantGroupIDs[groupID] = a.GetGroup(groupID)
relevantGroupIDs[groupID] = group
peerGroupSet[groupID] = struct{}{}
}
}
@@ -355,15 +364,12 @@ func (a *Account) getPeersGroupsPoliciesRoutes(
continue
}
for _, groupID := range r.PeerGroups {
relevantGroupIDs[groupID] = a.GetGroup(groupID)
}
for _, groupID := range r.Groups {
relevantGroupIDs[groupID] = a.GetGroup(groupID)
addRelevantGroup(groupID)
}
if r.Enabled {
for _, groupID := range r.AccessControlGroups {
relevantGroupIDs[groupID] = a.GetGroup(groupID)
addRelevantGroup(groupID)
routeAccessControlGroups[groupID] = struct{}{}
}
}
@@ -389,7 +395,7 @@ func (a *Account) getPeersGroupsPoliciesRoutes(
}
}
for _, groupID := range r.PeerGroups {
g := a.GetGroup(groupID)
g := addRelevantGroup(groupID)
if g == nil {
continue
}
@@ -424,10 +430,10 @@ func (a *Account) getPeersGroupsPoliciesRoutes(
if _, needed := routeAccessControlGroups[destGroupID]; needed {
policyRelevant = true
for _, srcGroupID := range rule.Sources {
relevantGroupIDs[srcGroupID] = a.GetGroup(srcGroupID)
addRelevantGroup(srcGroupID)
}
for _, dstGroupID := range rule.Destinations {
relevantGroupIDs[dstGroupID] = a.GetGroup(dstGroupID)
addRelevantGroup(dstGroupID)
}
break
}
@@ -463,7 +469,7 @@ func (a *Account) getPeersGroupsPoliciesRoutes(
}
}
for _, dstGroupID := range rule.Destinations {
relevantGroupIDs[dstGroupID] = a.GetGroup(dstGroupID)
addRelevantGroup(dstGroupID)
}
}
@@ -475,7 +481,7 @@ func (a *Account) getPeersGroupsPoliciesRoutes(
}
}
for _, srcGroupID := range rule.Sources {
relevantGroupIDs[srcGroupID] = a.GetGroup(srcGroupID)
addRelevantGroup(srcGroupID)
}
if rule.Protocol == PolicyRuleProtocolNetbirdSSH {

View File

@@ -24,7 +24,7 @@ import (
)
// ErrSharedSockStopped indicates that shared socket has been stopped
var ErrSharedSockStopped = fmt.Errorf("shared socket stopped")
var ErrSharedSockStopped = fmt.Errorf("shared socked stopped")
// SharedSocket is a net.PacketConn that initiates two raw sockets (ipv4 and ipv6) and listens to UDP packets filtered
// by BPF instructions (e.g., IncomingSTUNFilter that checks and sends only STUN packets to the listeners (ReadFrom)).