diff --git a/management/internals/modules/peers/ephemeral/manager/ephemeral.go b/management/internals/modules/peers/ephemeral/manager/ephemeral.go index 4a8683aa0..05c816e69 100644 --- a/management/internals/modules/peers/ephemeral/manager/ephemeral.go +++ b/management/internals/modules/peers/ephemeral/manager/ephemeral.go @@ -220,13 +220,13 @@ func (e *EphemeralManager) cleanup(ctx context.Context) { if err != nil { log.WithContext(ctx).Errorf("failed to delete ephemeral peers: %s", err) e.metrics.CountCleanupError() - continue } if len(skipped) > 0 { - // A skipped peer could not be deleted yet (still connected in the - // store, or seen too recently), which says nothing about whether a - // disconnect will ever be observed for it again. Schedule another - // attempt instead of dropping it, or it is never collected. + // A skipped peer was not deleted: it is still connected in the + // store, was seen too recently, or its deletion failed. None of + // that says whether a disconnect will ever be observed for it + // again, so schedule another attempt instead of dropping it, or + // it is never collected. log.WithContext(ctx).Debugf("cleanup: requeueing %d skipped ephemeral peers for account %s: %s", len(skipped), accountID, skipped) e.requeuePeers(ctx, accountID, skipped) } diff --git a/management/internals/modules/peers/ephemeral/manager/ephemeral_test.go b/management/internals/modules/peers/ephemeral/manager/ephemeral_test.go index 97b19e360..fc49c9d39 100644 --- a/management/internals/modules/peers/ephemeral/manager/ephemeral_test.go +++ b/management/internals/modules/peers/ephemeral/manager/ephemeral_test.go @@ -2,6 +2,7 @@ package manager import ( "context" + "errors" "fmt" "sync" "testing" @@ -320,6 +321,50 @@ func TestCleanupRequeuesVetoedPeers(t *testing.T) { assert.Len(t, mockStore.account.Peers, 0, "vetoed peer should be retried and deleted once the veto clears") } +// TestCleanupRequeuesFailedDeletes covers a peer whose deletion attempt errors +// (DeletePeers reports it as skipped alongside the aggregate error): it must be +// scheduled for another attempt rather than dropped from the list, or a +// transient store failure leaks it forever. +func TestCleanupRequeuesFailedDeletes(t *testing.T) { + t.Cleanup(func() { + timeNow = time.Now + }) + startTime := time.Now() + timeNow = func() time.Time { + return startTime + } + + mockStore := &MockStore{} + seedPeers(mockStore, 0, 1) + ctrl := gomock.NewController(t) + peersManager := peers.NewMockManager(ctrl) + + // The first attempt fails, the second deletes the peer. + first := peersManager.EXPECT(). + DeletePeers(gomock.Any(), gomock.Any(), []string{"ephemeral_peer_0"}, gomock.Any(), true). + Return([]string{"ephemeral_peer_0"}, errors.New("transient store failure")) + peersManager.EXPECT(). + DeletePeers(gomock.Any(), gomock.Any(), []string{"ephemeral_peer_0"}, gomock.Any(), true). + After(first). + DoAndReturn(func(_ context.Context, _ string, peerIDs []string, _ string, _ bool) ([]string, error) { + for _, peerID := range peerIDs { + delete(mockStore.account.Peers, peerID) + } + return nil, nil + }) + + mgr := NewEphemeralManager(mockStore, peersManager) + mgr.loadEphemeralPeers(context.Background()) + + startTime = startTime.Add(ephemeral.EphemeralLifeTime + time.Second) + mgr.cleanup(context.Background()) + assert.Len(t, mockStore.account.Peers, 1, "peer whose deletion failed must still exist") + + startTime = startTime.Add(ephemeral.EphemeralLifeTime + time.Second) + mgr.cleanup(context.Background()) + assert.Len(t, mockStore.account.Peers, 0, "peer whose deletion failed should be retried and deleted") +} + func seedPeers(store *MockStore, numberOfPeers int, numberOfEphemeralPeers int) { store.account = newAccountWithId(context.Background(), "my account", "", "", false) diff --git a/management/internals/modules/peers/manager.go b/management/internals/modules/peers/manager.go index d2be347c3..83615f926 100644 --- a/management/internals/modules/peers/manager.go +++ b/management/internals/modules/peers/manager.go @@ -8,9 +8,11 @@ import ( "net" "time" + "github.com/hashicorp/go-multierror" "github.com/rs/xid" log "github.com/sirupsen/logrus" + nberrors "github.com/netbirdio/netbird/client/errors" "github.com/netbirdio/netbird/management/internals/controllers/network_map" "github.com/netbirdio/netbird/management/internals/modules/peers/ephemeral" "github.com/netbirdio/netbird/management/server/account" @@ -31,9 +33,10 @@ type Manager interface { GetAllPeers(ctx context.Context, accountID, userID string) ([]*peer.Peer, error) GetPeersByGroupIDs(ctx context.Context, accountID string, groupsIDs []string) ([]*peer.Peer, error) // DeletePeers removes the given peers along with their group memberships and - // policies. With checkConnected, a peer that is still connected or was seen - // too recently is left in place and returned in skipped, so the caller can - // retry it later. + // policies. Every peer that was not deleted is returned in skipped, so the + // caller can retry it later: with checkConnected, a peer that is still + // connected or was seen too recently is left in place, and a peer whose + // deletion failed is skipped with the failure aggregated into err. DeletePeers(ctx context.Context, accountID string, peerIDs []string, userID string, checkConnected bool) (skipped []string, err error) SetNetworkMapController(networkMapController network_map.Controller) SetIntegratedPeerValidator(integratedPeerValidator integrated_validator.IntegratedValidator) @@ -148,16 +151,18 @@ const ( func (m *managerImpl) DeletePeers(ctx context.Context, accountID string, peerIDs []string, userID string, checkConnected bool) ([]string, error) { settings, err := m.store.GetAccountSettings(ctx, store.LockingStrengthNone, accountID) if err != nil { - return nil, err + return peerIDs, err } dnsDomain := m.networkMapController.GetDNSDomain(settings) var skipped []string + var merr *multierror.Error deletedAny := false for _, peerID := range peerIDs { outcome, events, err := m.deleteSinglePeer(ctx, accountID, peerID, userID, checkConnected, dnsDomain) if err != nil { - log.WithContext(ctx).Errorf("DeletePeers: failed to delete peer %s: %v", peerID, err) + merr = multierror.Append(merr, fmt.Errorf("delete peer %s: %w", peerID, err)) + skipped = append(skipped, peerID) continue } @@ -177,7 +182,7 @@ func (m *managerImpl) DeletePeers(ctx context.Context, accountID string, peerIDs m.accountManager.UpdateAccountPeers(ctx, accountID, types.UpdateReason{Resource: types.UpdateResourcePeer, Operation: types.UpdateOperationDelete}) } - return skipped, nil + return skipped, nberrors.FormatErrorOrNil(merr) } // deleteSinglePeer deletes one peer along with its group memberships and @@ -228,7 +233,7 @@ func (m *managerImpl) deleteSinglePeer(ctx context.Context, accountID, peerID, u // store once the transaction commits. func (m *managerImpl) deletePeerObjects(ctx context.Context, transaction store.Store, accountID, userID, dnsDomain string, p *peer.Peer) ([]func(), error) { if err := transaction.RemovePeerFromAllGroups(ctx, p.ID); err != nil { - return nil, fmt.Errorf("failed to remove peer %s from groups", p.ID) + return nil, fmt.Errorf("remove peer %s from groups: %w", p.ID, err) } var eventsToStore []func()