mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-06 05:29:07 +02:00
[management] Refresh only affected peers on IPv6 settings changes (#8051)
Account settings changes refreshed every peer, even an IPv6 group toggle that re-addresses a few, and the refresh goroutine got the request context, so it could be cancelled when the handler returned. Group paths that reconcile IPv6 addresses only walked the changed group, so peers reaching a re-addressed peer through its other groups missed the new address. The IPv6 reconcile now returns the peers whose address changed and callers pass them as changed peers. An IPv6-only settings change dispatches affected peers, adding every IPv6 holder on a range change since the interface prefix comes from the range. IPv4 range and account-wide changes keep the full refresh with a detached context.
This commit is contained in:
@@ -334,6 +334,9 @@ func (am *DefaultAccountManager) UpdateAccountSettings(ctx context.Context, acco
|
||||
var groupChangesAffectPeers bool
|
||||
var reloadReverseProxy bool
|
||||
var effectiveOldNetworkRange netip.Prefix
|
||||
var ipv6Changed bool
|
||||
var ipv6Snap *affectedpeers.Snapshot
|
||||
var ipv6Change affectedpeers.Change
|
||||
|
||||
err = am.Store.ExecuteInTransaction(ctx, func(transaction store.Store) error {
|
||||
var groupsUpdated bool
|
||||
@@ -379,10 +382,10 @@ func (am *DefaultAccountManager) UpdateAccountSettings(ctx context.Context, acco
|
||||
}
|
||||
|
||||
if ipv6SettingsChanged(oldSettings, newSettings) {
|
||||
if err = am.updatePeerIPv6Addresses(ctx, transaction, accountID, newSettings); err != nil {
|
||||
if ipv6Change, err = am.applyIPv6SettingsChange(ctx, transaction, accountID, oldSettings, newSettings); err != nil {
|
||||
return err
|
||||
}
|
||||
updateAccountPeers = true
|
||||
ipv6Changed = true
|
||||
}
|
||||
|
||||
if oldSettings.RoutingPeerDNSResolutionEnabled != newSettings.RoutingPeerDNSResolutionEnabled ||
|
||||
@@ -419,12 +422,20 @@ func (am *DefaultAccountManager) UpdateAccountSettings(ctx context.Context, acco
|
||||
return err
|
||||
}
|
||||
|
||||
if updateAccountPeers || groupsUpdated {
|
||||
if updateAccountPeers || groupsUpdated || ipv6Changed {
|
||||
if err = transaction.IncrementNetworkSerial(ctx, accountID); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
// A full account refresh already covers the IPv6 change, so the affected-peers
|
||||
// snapshot is only needed when nothing account-wide changed.
|
||||
if ipv6Changed && !updateAccountPeers && !groupChangesAffectPeers {
|
||||
if ipv6Snap, err = affectedpeers.Load(ctx, transaction, accountID, ipv6Change); err != nil {
|
||||
return fmt.Errorf("load affected peers: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
@@ -486,13 +497,34 @@ func (am *DefaultAccountManager) UpdateAccountSettings(ctx context.Context, acco
|
||||
}
|
||||
}
|
||||
|
||||
if updateAccountPeers || extraSettingsChanged || groupChangesAffectPeers {
|
||||
go am.UpdateAccountPeers(ctx, accountID, types.UpdateReason{Resource: types.UpdateResourceAccountSettings, Operation: types.UpdateOperationUpdate})
|
||||
switch {
|
||||
case updateAccountPeers || extraSettingsChanged || groupChangesAffectPeers:
|
||||
go am.UpdateAccountPeers(context.WithoutCancel(ctx), accountID, types.UpdateReason{Resource: types.UpdateResourceAccountSettings, Operation: types.UpdateOperationUpdate})
|
||||
case ipv6Snap != nil:
|
||||
am.ExpandAndUpdateAffected(ctx, accountID, ipv6Snap, ipv6Change)
|
||||
}
|
||||
|
||||
return newSettings, nil
|
||||
}
|
||||
|
||||
// applyIPv6SettingsChange reconciles peer IPv6 addresses for new IPv6 settings and
|
||||
// returns the affected-peers change: peers whose address changed refresh together
|
||||
// with every peer that reaches them. On a range change every peer holding an address
|
||||
// also refreshes itself, since its interface prefix comes from the account range even
|
||||
// when its address stays inside the new one.
|
||||
func (am *DefaultAccountManager) applyIPv6SettingsChange(ctx context.Context, transaction store.Store, accountID string, oldSettings, newSettings *types.Settings) (affectedpeers.Change, error) {
|
||||
result, err := am.updatePeerIPv6Addresses(ctx, transaction, accountID, newSettings)
|
||||
if err != nil {
|
||||
return affectedpeers.Change{}, err
|
||||
}
|
||||
|
||||
change := affectedpeers.Change{ChangedPeerIDs: result.changed}
|
||||
if oldSettings.NetworkRangeV6 != newSettings.NetworkRangeV6 {
|
||||
change.OutputPeerIDs = result.withIPv6
|
||||
}
|
||||
return change, nil
|
||||
}
|
||||
|
||||
func ipv6SettingsChanged(old, updated *types.Settings) bool {
|
||||
if old.NetworkRangeV6 != updated.NetworkRangeV6 {
|
||||
return true
|
||||
@@ -1742,9 +1774,11 @@ func (am *DefaultAccountManager) SyncUserJWTGroups(ctx context.Context, userAuth
|
||||
|
||||
change.LinkGroups = allGroupChanges
|
||||
|
||||
if err = am.reconcileIPv6ForGroupChanges(ctx, transaction, userAuth.AccountId, allGroupChanges); err != nil {
|
||||
ipv6Changed, err := am.reconcileIPv6ForGroupChanges(ctx, transaction, userAuth.AccountId, allGroupChanges)
|
||||
if err != nil {
|
||||
return fmt.Errorf("reconcile IPv6 for group changes: %w", err)
|
||||
}
|
||||
change.ChangedPeerIDs = append(change.ChangedPeerIDs, ipv6Changed...)
|
||||
|
||||
if err = transaction.IncrementNetworkSerial(ctx, userAuth.AccountId); err != nil {
|
||||
return fmt.Errorf("error incrementing network serial: %w", err)
|
||||
@@ -2334,7 +2368,8 @@ func (am *DefaultAccountManager) propagateUserGroupMemberships(ctx context.Conte
|
||||
return false, false, err
|
||||
}
|
||||
|
||||
if err = am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, updatedGroups); err != nil {
|
||||
ipv6Changed, err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, updatedGroups)
|
||||
if err != nil {
|
||||
return false, false, fmt.Errorf("reconcile IPv6 for group changes: %w", err)
|
||||
}
|
||||
|
||||
@@ -2343,7 +2378,7 @@ func (am *DefaultAccountManager) propagateUserGroupMemberships(ctx context.Conte
|
||||
return false, false, fmt.Errorf("error checking if group changes affect peers: %w", err)
|
||||
}
|
||||
|
||||
return len(updatedGroups) > 0, peersAffected, nil
|
||||
return len(updatedGroups) > 0, peersAffected || len(ipv6Changed) > 0, nil
|
||||
}
|
||||
|
||||
// propagateAutoGroupsForUsers adds each user's peers to their AutoGroups where not already present.
|
||||
@@ -2440,56 +2475,78 @@ func (am *DefaultAccountManager) checkIPv6Collision(ctx context.Context, transac
|
||||
return nil
|
||||
}
|
||||
|
||||
func (am *DefaultAccountManager) updatePeerIPv6Addresses(ctx context.Context, transaction store.Store, accountID string, settings *types.Settings) error {
|
||||
// ipv6Reassignment reports the outcome of an IPv6 address reconciliation.
|
||||
type ipv6Reassignment struct {
|
||||
// changed are the peers whose IPv6 address was assigned, removed or reallocated.
|
||||
changed []string
|
||||
// withIPv6 are all peers holding an IPv6 address after the reconciliation.
|
||||
withIPv6 []string
|
||||
}
|
||||
|
||||
func (am *DefaultAccountManager) updatePeerIPv6Addresses(ctx context.Context, transaction store.Store, accountID string, settings *types.Settings) (ipv6Reassignment, error) {
|
||||
peers, err := transaction.GetAccountPeers(ctx, store.LockingStrengthUpdate, accountID, "", "", "")
|
||||
if err != nil {
|
||||
return fmt.Errorf("get peers: %w", err)
|
||||
return ipv6Reassignment{}, fmt.Errorf("get peers: %w", err)
|
||||
}
|
||||
|
||||
network, err := transaction.GetAccountNetwork(ctx, store.LockingStrengthUpdate, accountID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("get network: %w", err)
|
||||
return ipv6Reassignment{}, fmt.Errorf("get network: %w", err)
|
||||
}
|
||||
|
||||
if err := am.ensureIPv6Subnet(ctx, transaction, accountID, settings, network); err != nil {
|
||||
return err
|
||||
return ipv6Reassignment{}, err
|
||||
}
|
||||
|
||||
allowedPeers, err := am.buildIPv6AllowedPeers(ctx, transaction, accountID, settings)
|
||||
if err != nil {
|
||||
return err
|
||||
return ipv6Reassignment{}, err
|
||||
}
|
||||
|
||||
v6Prefix, err := netip.ParsePrefix(network.NetV6.String())
|
||||
if err != nil {
|
||||
return fmt.Errorf("parse IPv6 prefix: %w", err)
|
||||
return ipv6Reassignment{}, fmt.Errorf("parse IPv6 prefix: %w", err)
|
||||
}
|
||||
|
||||
if err := am.assignPeerIPv6Addresses(ctx, transaction, accountID, peers, network, allowedPeers, v6Prefix); err != nil {
|
||||
return err
|
||||
changed, err := am.assignPeerIPv6Addresses(ctx, transaction, accountID, peers, network, allowedPeers, v6Prefix)
|
||||
if err != nil {
|
||||
return ipv6Reassignment{}, err
|
||||
}
|
||||
|
||||
log.WithContext(ctx).Infof("updated IPv6 addresses for %d peers in account %s (groups=%d)",
|
||||
len(peers), accountID, len(settings.IPv6EnabledGroups))
|
||||
result := ipv6Reassignment{changed: changed}
|
||||
for _, peer := range peers {
|
||||
if peer.IPv6.IsValid() {
|
||||
result.withIPv6 = append(result.withIPv6, peer.ID)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
log.WithContext(ctx).Infof("updated IPv6 addresses for %d of %d peers in account %s (groups=%d)",
|
||||
len(changed), len(peers), accountID, len(settings.IPv6EnabledGroups))
|
||||
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// reconcileIPv6ForGroupChanges checks whether the given group IDs overlap with
|
||||
// the account's IPv6EnabledGroups. If they do, it runs a full IPv6 address
|
||||
// reconciliation so that peers gaining or losing membership in an IPv6-enabled
|
||||
// group get their addresses assigned or removed.
|
||||
func (am *DefaultAccountManager) reconcileIPv6ForGroupChanges(ctx context.Context, transaction store.Store, accountID string, groupIDs []string) error {
|
||||
// group get their addresses assigned or removed. It returns the peers whose IPv6
|
||||
// address changed, which callers pass as changed peers so every peer that can
|
||||
// reach them refreshes.
|
||||
func (am *DefaultAccountManager) reconcileIPv6ForGroupChanges(ctx context.Context, transaction store.Store, accountID string, groupIDs []string) ([]string, error) {
|
||||
settings, err := transaction.GetAccountSettings(ctx, store.LockingStrengthNone, accountID)
|
||||
if err != nil {
|
||||
return fmt.Errorf("get account settings: %w", err)
|
||||
return nil, fmt.Errorf("get account settings: %w", err)
|
||||
}
|
||||
|
||||
if !ipv6ReconcileNeeded(settings, groupIDs) {
|
||||
return nil
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
return am.updatePeerIPv6Addresses(ctx, transaction, accountID, settings)
|
||||
result, err := am.updatePeerIPv6Addresses(ctx, transaction, accountID, settings)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return result.changed, nil
|
||||
}
|
||||
|
||||
// ipv6ReconcileNeeded reports whether changes to the given groups trigger an IPv6
|
||||
@@ -2528,7 +2585,7 @@ func (am *DefaultAccountManager) assignPeerIPv6Addresses(
|
||||
ctx context.Context, transaction store.Store, accountID string,
|
||||
peers []*nbpeer.Peer, network *types.Network,
|
||||
allowedPeers map[string]struct{}, v6Prefix netip.Prefix,
|
||||
) error {
|
||||
) ([]string, error) {
|
||||
takenV6 := make(map[netip.Addr]struct{})
|
||||
for _, peer := range peers {
|
||||
if _, ok := allowedPeers[peer.ID]; ok && peer.IPv6.IsValid() && network.NetV6.Contains(peer.IPv6.AsSlice()) {
|
||||
@@ -2536,6 +2593,7 @@ func (am *DefaultAccountManager) assignPeerIPv6Addresses(
|
||||
}
|
||||
}
|
||||
|
||||
var changed []string
|
||||
for _, peer := range peers {
|
||||
_, allowed := allowedPeers[peer.ID]
|
||||
oldIPv6 := peer.IPv6
|
||||
@@ -2545,7 +2603,7 @@ func (am *DefaultAccountManager) assignPeerIPv6Addresses(
|
||||
} else if !peer.IPv6.IsValid() || !network.NetV6.Contains(peer.IPv6.AsSlice()) {
|
||||
newIP, err := allocateIPv6WithRetry(v6Prefix, takenV6, peer.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
peer.IPv6 = newIP
|
||||
}
|
||||
@@ -2555,10 +2613,11 @@ func (am *DefaultAccountManager) assignPeerIPv6Addresses(
|
||||
}
|
||||
|
||||
if err := transaction.SavePeer(ctx, accountID, peer); err != nil {
|
||||
return fmt.Errorf("save peer %s: %w", peer.ID, err)
|
||||
return nil, fmt.Errorf("save peer %s: %w", peer.ID, err)
|
||||
}
|
||||
changed = append(changed, peer.ID)
|
||||
}
|
||||
return nil
|
||||
return changed, nil
|
||||
}
|
||||
|
||||
func allocateIPv6WithRetry(prefix netip.Prefix, taken map[netip.Addr]struct{}, peerID string) (netip.Addr, error) {
|
||||
|
||||
@@ -0,0 +1,243 @@
|
||||
package server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/netip"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/netbirdio/netbird/management/internals/controllers/network_map"
|
||||
nbpeer "github.com/netbirdio/netbird/management/server/peer"
|
||||
"github.com/netbirdio/netbird/management/server/store"
|
||||
"github.com/netbirdio/netbird/management/server/types"
|
||||
)
|
||||
|
||||
const (
|
||||
ipv6GroupA = "ipv6-grp-a"
|
||||
ipv6GroupB = "ipv6-grp-b"
|
||||
ipv6GroupC = "ipv6-grp-c"
|
||||
ipv6GroupD = "ipv6-grp-d"
|
||||
)
|
||||
|
||||
// ipv6AffectedTest holds three peers: peer1 in group A, peer2 in group B, peer3 in
|
||||
// group C, with a single A<->B policy. peer3 is unrelated to peer1 and peer2. Group D
|
||||
// is empty and referenced by nothing.
|
||||
type ipv6AffectedTest struct {
|
||||
manager *DefaultAccountManager
|
||||
accountID string
|
||||
peer1, peer2, peer3 *nbpeer.Peer
|
||||
updMsg1, updMsg2, updMsg3 <-chan *network_map.UpdateMessage
|
||||
}
|
||||
|
||||
func setupIPv6AffectedTest(t *testing.T, ipv6Groups []string) *ipv6AffectedTest {
|
||||
t.Helper()
|
||||
|
||||
manager, updateManager, account, peer1, peer2, peer3 := setupNetworkMapTest(t)
|
||||
ctx := context.Background()
|
||||
accountID := account.Id
|
||||
|
||||
policies, err := manager.Store.GetAccountPolicies(ctx, store.LockingStrengthNone, accountID)
|
||||
require.NoError(t, err)
|
||||
for _, p := range policies {
|
||||
require.NoError(t, manager.Store.DeletePolicy(ctx, accountID, p.ID))
|
||||
}
|
||||
|
||||
for _, g := range []*types.Group{
|
||||
{ID: ipv6GroupA, Name: "IPv6-A", Peers: []string{peer1.ID}},
|
||||
{ID: ipv6GroupB, Name: "IPv6-B", Peers: []string{peer2.ID}},
|
||||
{ID: ipv6GroupC, Name: "IPv6-C", Peers: []string{peer3.ID}},
|
||||
{ID: ipv6GroupD, Name: "IPv6-D"},
|
||||
} {
|
||||
require.NoError(t, manager.CreateGroup(ctx, accountID, userID, g))
|
||||
}
|
||||
|
||||
_, err = manager.SavePolicy(ctx, accountID, userID, &types.Policy{
|
||||
Enabled: true,
|
||||
Rules: []*types.PolicyRule{{
|
||||
Enabled: true,
|
||||
Sources: []string{ipv6GroupA},
|
||||
Destinations: []string{ipv6GroupB},
|
||||
Bidirectional: true,
|
||||
Action: types.PolicyTrafficActionAccept,
|
||||
}},
|
||||
}, true)
|
||||
require.NoError(t, err)
|
||||
|
||||
// New accounts enable IPv6 for the All group; start from the requested groups.
|
||||
updateIPv6TestSettings(t, manager, accountID, func(s *types.Settings) {
|
||||
s.IPv6EnabledGroups = ipv6Groups
|
||||
})
|
||||
|
||||
tc := &ipv6AffectedTest{
|
||||
manager: manager,
|
||||
accountID: accountID,
|
||||
peer1: peer1,
|
||||
peer2: peer2,
|
||||
peer3: peer3,
|
||||
}
|
||||
tc.updMsg1 = updateManager.CreateChannel(ctx, peer1.ID)
|
||||
tc.updMsg2 = updateManager.CreateChannel(ctx, peer2.ID)
|
||||
tc.updMsg3 = updateManager.CreateChannel(ctx, peer3.ID)
|
||||
t.Cleanup(func() {
|
||||
updateManager.CloseChannel(ctx, peer1.ID)
|
||||
updateManager.CloseChannel(ctx, peer2.ID)
|
||||
updateManager.CloseChannel(ctx, peer3.ID)
|
||||
})
|
||||
|
||||
// The setup changes above dispatch asynchronously and can land after the
|
||||
// channels open, so drop them before the test acts.
|
||||
drainPeerUpdates(tc.updMsg1)
|
||||
drainPeerUpdates(tc.updMsg2)
|
||||
drainPeerUpdates(tc.updMsg3)
|
||||
|
||||
return tc
|
||||
}
|
||||
|
||||
// updateIPv6TestSettings applies mutate to a copy of the current settings, so only
|
||||
// the mutated fields differ from what is stored.
|
||||
func updateIPv6TestSettings(t *testing.T, manager *DefaultAccountManager, accountID string, mutate func(*types.Settings)) {
|
||||
t.Helper()
|
||||
ctx := context.Background()
|
||||
|
||||
current, err := manager.Store.GetAccountSettings(ctx, store.LockingStrengthNone, accountID)
|
||||
require.NoError(t, err)
|
||||
|
||||
updated := current.Copy()
|
||||
mutate(updated)
|
||||
|
||||
_, err = manager.UpdateAccountSettings(ctx, accountID, userID, updated)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
func (tc *ipv6AffectedTest) peerIPv6(t *testing.T, peerID string) netip.Addr {
|
||||
t.Helper()
|
||||
peer, err := tc.manager.Store.GetPeerByID(context.Background(), store.LockingStrengthNone, tc.accountID, peerID)
|
||||
require.NoError(t, err)
|
||||
return peer.IPv6
|
||||
}
|
||||
|
||||
func TestAffectedPeers_IPv6GroupEnabled_RefreshesOnlyReachablePeers(t *testing.T) {
|
||||
tc := setupIPv6AffectedTest(t, nil)
|
||||
|
||||
updateIPv6TestSettings(t, tc.manager, tc.accountID, func(s *types.Settings) {
|
||||
s.IPv6EnabledGroups = []string{ipv6GroupA}
|
||||
})
|
||||
require.True(t, tc.peerIPv6(t, tc.peer1.ID).IsValid(), "peer1 should get an IPv6 address")
|
||||
|
||||
peerShouldReceiveUpdate(t, tc.updMsg1)
|
||||
peerShouldReceiveUpdate(t, tc.updMsg2)
|
||||
peerShouldNotReceiveUpdate(t, tc.updMsg3)
|
||||
}
|
||||
|
||||
func TestAffectedPeers_IPv6GroupDisabled_RefreshesOnlyReachablePeers(t *testing.T) {
|
||||
tc := setupIPv6AffectedTest(t, []string{ipv6GroupA})
|
||||
require.True(t, tc.peerIPv6(t, tc.peer1.ID).IsValid(), "peer1 should start with an IPv6 address")
|
||||
|
||||
updateIPv6TestSettings(t, tc.manager, tc.accountID, func(s *types.Settings) {
|
||||
s.IPv6EnabledGroups = []string{}
|
||||
})
|
||||
require.False(t, tc.peerIPv6(t, tc.peer1.ID).IsValid(), "peer1 should lose its IPv6 address")
|
||||
|
||||
peerShouldReceiveUpdate(t, tc.updMsg1)
|
||||
peerShouldReceiveUpdate(t, tc.updMsg2)
|
||||
peerShouldNotReceiveUpdate(t, tc.updMsg3)
|
||||
}
|
||||
|
||||
// Widening the IPv6 range keeps peer addresses, but each holder's interface prefix
|
||||
// comes from the range, so holders refresh while peers that only reach them do not.
|
||||
func TestAffectedPeers_IPv6RangeWidened_RefreshesAddressHolders(t *testing.T) {
|
||||
tc := setupIPv6AffectedTest(t, []string{ipv6GroupA})
|
||||
oldIPv6 := tc.peerIPv6(t, tc.peer1.ID)
|
||||
require.True(t, oldIPv6.IsValid(), "peer1 should start with an IPv6 address")
|
||||
|
||||
// The range is allocated on the account network; settings may leave it empty.
|
||||
network, err := tc.manager.Store.GetAccountNetwork(context.Background(), store.LockingStrengthNone, tc.accountID)
|
||||
require.NoError(t, err)
|
||||
current := prefixFromIPNet(network.NetV6)
|
||||
require.True(t, current.IsValid(), "account should have an IPv6 range")
|
||||
widened := netip.PrefixFrom(current.Addr(), current.Bits()-8).Masked()
|
||||
|
||||
updateIPv6TestSettings(t, tc.manager, tc.accountID, func(s *types.Settings) {
|
||||
s.NetworkRangeV6 = widened
|
||||
})
|
||||
require.Equal(t, oldIPv6, tc.peerIPv6(t, tc.peer1.ID), "peer1 should keep its address inside the widened range")
|
||||
|
||||
peerShouldReceiveUpdate(t, tc.updMsg1)
|
||||
peerShouldNotReceiveUpdate(t, tc.updMsg2)
|
||||
peerShouldNotReceiveUpdate(t, tc.updMsg3)
|
||||
}
|
||||
|
||||
func TestAffectedPeers_IPv4RangeChange_RefreshesWholeAccount(t *testing.T) {
|
||||
tc := setupIPv6AffectedTest(t, nil)
|
||||
|
||||
updateIPv6TestSettings(t, tc.manager, tc.accountID, func(s *types.Settings) {
|
||||
s.NetworkRange = netip.MustParsePrefix("100.70.0.0/16")
|
||||
})
|
||||
|
||||
peerShouldReceiveUpdate(t, tc.updMsg1)
|
||||
peerShouldReceiveUpdate(t, tc.updMsg2)
|
||||
peerShouldReceiveUpdate(t, tc.updMsg3)
|
||||
}
|
||||
|
||||
func TestAffectedPeers_IPv6WithAccountWideChange_RefreshesWholeAccount(t *testing.T) {
|
||||
tc := setupIPv6AffectedTest(t, nil)
|
||||
|
||||
updateIPv6TestSettings(t, tc.manager, tc.accountID, func(s *types.Settings) {
|
||||
s.IPv6EnabledGroups = []string{ipv6GroupA}
|
||||
s.LazyConnectionEnabled = !s.LazyConnectionEnabled
|
||||
})
|
||||
|
||||
peerShouldReceiveUpdate(t, tc.updMsg1)
|
||||
peerShouldReceiveUpdate(t, tc.updMsg2)
|
||||
peerShouldReceiveUpdate(t, tc.updMsg3)
|
||||
}
|
||||
|
||||
// Joining an IPv6-enabled group that no policy references gives peer1 an address.
|
||||
// peer2 reaches peer1 through group A, not through the joined group, and must still
|
||||
// learn the new address.
|
||||
func TestAffectedPeers_GroupAddPeerIPv6_RefreshesPeersReachingThroughOtherGroups(t *testing.T) {
|
||||
tc := setupIPv6AffectedTest(t, []string{ipv6GroupD})
|
||||
|
||||
require.NoError(t, tc.manager.GroupAddPeer(context.Background(), tc.accountID, ipv6GroupD, tc.peer1.ID))
|
||||
require.True(t, tc.peerIPv6(t, tc.peer1.ID).IsValid(), "peer1 should get an IPv6 address")
|
||||
|
||||
peerShouldReceiveUpdate(t, tc.updMsg1)
|
||||
peerShouldReceiveUpdate(t, tc.updMsg2)
|
||||
peerShouldNotReceiveUpdate(t, tc.updMsg3)
|
||||
}
|
||||
|
||||
func TestAffectedPeers_UpdateGroupIPv6_RefreshesPeersReachingThroughOtherGroups(t *testing.T) {
|
||||
tc := setupIPv6AffectedTest(t, []string{ipv6GroupD})
|
||||
|
||||
require.NoError(t, tc.manager.UpdateGroup(context.Background(), tc.accountID, userID, &types.Group{
|
||||
ID: ipv6GroupD,
|
||||
Name: "IPv6-D",
|
||||
Peers: []string{tc.peer1.ID},
|
||||
}))
|
||||
require.True(t, tc.peerIPv6(t, tc.peer1.ID).IsValid(), "peer1 should get an IPv6 address")
|
||||
|
||||
peerShouldReceiveUpdate(t, tc.updMsg1)
|
||||
peerShouldReceiveUpdate(t, tc.updMsg2)
|
||||
peerShouldNotReceiveUpdate(t, tc.updMsg3)
|
||||
}
|
||||
|
||||
// Deleting an IPv6-enabled group removes its members' addresses after the
|
||||
// pre-delete snapshot was taken.
|
||||
func TestAffectedPeers_DeleteIPv6Group_RefreshesFormerMembersAndReachablePeers(t *testing.T) {
|
||||
tc := setupIPv6AffectedTest(t, []string{ipv6GroupD})
|
||||
ctx := context.Background()
|
||||
|
||||
require.NoError(t, tc.manager.GroupAddPeer(ctx, tc.accountID, ipv6GroupD, tc.peer1.ID))
|
||||
require.True(t, tc.peerIPv6(t, tc.peer1.ID).IsValid(), "peer1 should get an IPv6 address")
|
||||
drainPeerUpdates(tc.updMsg1)
|
||||
drainPeerUpdates(tc.updMsg2)
|
||||
drainPeerUpdates(tc.updMsg3)
|
||||
|
||||
require.NoError(t, tc.manager.DeleteGroup(ctx, tc.accountID, userID, ipv6GroupD))
|
||||
require.False(t, tc.peerIPv6(t, tc.peer1.ID).IsValid(), "peer1 should lose its IPv6 address")
|
||||
|
||||
peerShouldReceiveUpdate(t, tc.updMsg1)
|
||||
peerShouldReceiveUpdate(t, tc.updMsg2)
|
||||
peerShouldNotReceiveUpdate(t, tc.updMsg3)
|
||||
}
|
||||
@@ -108,11 +108,13 @@ func TestAffectedPeers_SaveUser_OnlyAffectedPeersUpdated(t *testing.T) {
|
||||
})
|
||||
|
||||
t.Run("auto group change reassigning IPv6 refreshes the changed peers and their observers", func(t *testing.T) {
|
||||
account, err := manager.Store.GetAccount(ctx, accountID)
|
||||
require.NoError(t, err)
|
||||
account.Settings.IPv6EnabledGroups = []string{"ug-v6"}
|
||||
require.NoError(t, manager.Store.SaveAccount(ctx, account))
|
||||
require.NoError(t, manager.CreateGroup(ctx, accountID, userID, &types.Group{ID: "ug-v6", Name: "ug-v6"}))
|
||||
// Apply through the settings API so the reconciliation that strips the other
|
||||
// peers' addresses happens here, leaving the target as the only peer the
|
||||
// user update reassigns.
|
||||
updateIPv6TestSettings(t, manager, accountID, func(s *types.Settings) {
|
||||
s.IPv6EnabledGroups = []string{"ug-v6"}
|
||||
})
|
||||
|
||||
drainPeerUpdates(updTarget)
|
||||
drainPeerUpdates(upd2)
|
||||
|
||||
+33
-13
@@ -166,9 +166,11 @@ func (am *DefaultAccountManager) UpdateGroup(ctx context.Context, accountID, use
|
||||
return err
|
||||
}
|
||||
|
||||
if err = am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, []string{newGroup.ID}); err != nil {
|
||||
ipv6Changed, err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, []string{newGroup.ID})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
change.ChangedPeerIDs = ipv6Changed
|
||||
|
||||
// A membership change does not alter which entities reference the group, so
|
||||
// the dependency walk runs once against the post-change snapshot. The new
|
||||
@@ -321,7 +323,7 @@ func (am *DefaultAccountManager) UpdateGroups(ctx context.Context, accountID, us
|
||||
var globalErr error
|
||||
for _, newGroup := range groups {
|
||||
change := affectedpeers.Change{ChangedGroupIDs: []string{newGroup.ID}}
|
||||
events, snap, err := am.updateSingleGroup(ctx, accountID, userID, newGroup, change)
|
||||
events, snap, change, err := am.updateSingleGroup(ctx, accountID, userID, newGroup, change)
|
||||
if err != nil {
|
||||
log.WithContext(ctx).Errorf("failed to update group %s: %v", newGroup.ID, err)
|
||||
if len(groups) == 1 {
|
||||
@@ -344,7 +346,7 @@ func (am *DefaultAccountManager) UpdateGroups(ctx context.Context, accountID, us
|
||||
return globalErr
|
||||
}
|
||||
|
||||
func (am *DefaultAccountManager) updateSingleGroup(ctx context.Context, accountID, userID string, newGroup *types.Group, change affectedpeers.Change) ([]func(), *affectedpeers.Snapshot, error) {
|
||||
func (am *DefaultAccountManager) updateSingleGroup(ctx context.Context, accountID, userID string, newGroup *types.Group, change affectedpeers.Change) ([]func(), *affectedpeers.Snapshot, affectedpeers.Change, error) {
|
||||
var events []func()
|
||||
var snap *affectedpeers.Snapshot
|
||||
err := am.Store.ExecuteInTransaction(ctx, func(transaction store.Store) error {
|
||||
@@ -364,9 +366,11 @@ func (am *DefaultAccountManager) updateSingleGroup(ctx context.Context, accountI
|
||||
return err
|
||||
}
|
||||
|
||||
if err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, []string{newGroup.ID}); err != nil {
|
||||
ipv6Changed, err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, []string{newGroup.ID})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
change.ChangedPeerIDs = ipv6Changed
|
||||
|
||||
if err := transaction.IncrementNetworkSerial(ctx, accountID); err != nil {
|
||||
return err
|
||||
@@ -377,7 +381,7 @@ func (am *DefaultAccountManager) updateSingleGroup(ctx context.Context, accountI
|
||||
snap, err = affectedpeers.Load(ctx, transaction, accountID, change)
|
||||
return err
|
||||
})
|
||||
return events, snap, err
|
||||
return events, snap, change, err
|
||||
}
|
||||
|
||||
// prepareGroupEvents prepares a list of event functions to be stored.
|
||||
@@ -480,8 +484,8 @@ func (am *DefaultAccountManager) DeleteGroups(ctx context.Context, accountID, us
|
||||
var allErrors error
|
||||
var groupIDsToDelete []string
|
||||
var deletedGroups []*types.Group
|
||||
var snap *affectedpeers.Snapshot
|
||||
var change affectedpeers.Change
|
||||
var snap, ipv6Snap *affectedpeers.Snapshot
|
||||
var change, ipv6Change affectedpeers.Change
|
||||
|
||||
extraSettings, err := am.settingsManager.GetExtraSettings(ctx, accountID)
|
||||
if err != nil {
|
||||
@@ -510,10 +514,20 @@ func (am *DefaultAccountManager) DeleteGroups(ctx context.Context, accountID, us
|
||||
return err
|
||||
}
|
||||
|
||||
if err = am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, groupIDsToDelete); err != nil {
|
||||
ipv6Changed, err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, groupIDsToDelete)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Members of a deleted IPv6-enabled group lose their address, which the
|
||||
// pre-delete snapshot cannot see, so they are resolved post-delete.
|
||||
if len(ipv6Changed) > 0 {
|
||||
ipv6Change = affectedpeers.Change{ChangedPeerIDs: ipv6Changed}
|
||||
if ipv6Snap, err = affectedpeers.Load(ctx, transaction, accountID, ipv6Change); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return transaction.IncrementNetworkSerial(ctx, accountID)
|
||||
})
|
||||
if err != nil {
|
||||
@@ -524,7 +538,7 @@ func (am *DefaultAccountManager) DeleteGroups(ctx context.Context, accountID, us
|
||||
am.StoreEvent(ctx, userID, group.ID, accountID, activity.GroupDeleted, group.EventMeta())
|
||||
}
|
||||
|
||||
am.ExpandAndUpdateAffected(ctx, accountID, snap, change)
|
||||
go am.dispatchAffected(ctx, accountID, []*affectedpeers.Snapshot{snap, ipv6Snap}, []affectedpeers.Change{change, ipv6Change})
|
||||
|
||||
return allErrors
|
||||
}
|
||||
@@ -564,11 +578,14 @@ func (am *DefaultAccountManager) GroupAddPeer(ctx context.Context, accountID, gr
|
||||
return err
|
||||
}
|
||||
|
||||
if err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, []string{groupID}); err != nil {
|
||||
ipv6Changed, err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, []string{groupID})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// A peer whose IPv6 address changed is visible to every peer that reaches it
|
||||
// through any of its groups, not only through this one.
|
||||
change.ChangedPeerIDs = ipv6Changed
|
||||
|
||||
var err error
|
||||
if snap, err = affectedpeers.Load(ctx, transaction, accountID, change); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -634,11 +651,14 @@ func (am *DefaultAccountManager) GroupDeletePeer(ctx context.Context, accountID,
|
||||
return err
|
||||
}
|
||||
|
||||
if err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, []string{groupID}); err != nil {
|
||||
ipv6Changed, err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, []string{groupID})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// A peer whose IPv6 address changed is visible to every peer that reaches it
|
||||
// through any of its groups, not only through this one.
|
||||
change.ChangedPeerIDs = ipv6Changed
|
||||
|
||||
var err error
|
||||
if snap, err = affectedpeers.Load(ctx, transaction, accountID, change); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -861,9 +861,11 @@ func (am *DefaultAccountManager) processUserUpdate(ctx context.Context, transact
|
||||
allGroupChanges := slices.Concat(removedGroups, addedGroups)
|
||||
change.LinkGroups = allGroupChanges
|
||||
|
||||
if err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, allGroupChanges); err != nil {
|
||||
ipv6Changed, err := am.reconcileIPv6ForGroupChanges(ctx, transaction, accountID, allGroupChanges)
|
||||
if err != nil {
|
||||
return change, nil, nil, nil, fmt.Errorf("reconcile IPv6 for group changes: %w", err)
|
||||
}
|
||||
change.ChangedPeerIDs = append(change.ChangedPeerIDs, ipv6Changed...)
|
||||
}
|
||||
|
||||
userEventsToAdd := am.prepareUserUpdateEvents(ctx, updatedUser.AccountID, initiatorUserId, oldUser, updatedUser, transferredOwnerRole, isNewUser, removedGroups, addedGroups, transaction)
|
||||
|
||||
Reference in New Issue
Block a user