mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-01 19:19:07 +02:00
Requeue ephemeral peers whose deletion attempt failed
This commit is contained in:
@@ -220,13 +220,13 @@ func (e *EphemeralManager) cleanup(ctx context.Context) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
log.WithContext(ctx).Errorf("failed to delete ephemeral peers: %s", err)
|
log.WithContext(ctx).Errorf("failed to delete ephemeral peers: %s", err)
|
||||||
e.metrics.CountCleanupError()
|
e.metrics.CountCleanupError()
|
||||||
continue
|
|
||||||
}
|
}
|
||||||
if len(skipped) > 0 {
|
if len(skipped) > 0 {
|
||||||
// A skipped peer could not be deleted yet (still connected in the
|
// A skipped peer was not deleted: it is still connected in the
|
||||||
// store, or seen too recently), which says nothing about whether a
|
// store, was seen too recently, or its deletion failed. None of
|
||||||
// disconnect will ever be observed for it again. Schedule another
|
// that says whether a disconnect will ever be observed for it
|
||||||
// attempt instead of dropping it, or it is never collected.
|
// 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)
|
log.WithContext(ctx).Debugf("cleanup: requeueing %d skipped ephemeral peers for account %s: %s", len(skipped), accountID, skipped)
|
||||||
e.requeuePeers(ctx, accountID, skipped)
|
e.requeuePeers(ctx, accountID, skipped)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package manager
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"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")
|
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) {
|
func seedPeers(store *MockStore, numberOfPeers int, numberOfEphemeralPeers int) {
|
||||||
store.account = newAccountWithId(context.Background(), "my account", "", "", false)
|
store.account = newAccountWithId(context.Background(), "my account", "", "", false)
|
||||||
|
|
||||||
|
|||||||
@@ -8,9 +8,11 @@ import (
|
|||||||
"net"
|
"net"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/hashicorp/go-multierror"
|
||||||
"github.com/rs/xid"
|
"github.com/rs/xid"
|
||||||
log "github.com/sirupsen/logrus"
|
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/controllers/network_map"
|
||||||
"github.com/netbirdio/netbird/management/internals/modules/peers/ephemeral"
|
"github.com/netbirdio/netbird/management/internals/modules/peers/ephemeral"
|
||||||
"github.com/netbirdio/netbird/management/server/account"
|
"github.com/netbirdio/netbird/management/server/account"
|
||||||
@@ -31,9 +33,10 @@ type Manager interface {
|
|||||||
GetAllPeers(ctx context.Context, accountID, userID string) ([]*peer.Peer, error)
|
GetAllPeers(ctx context.Context, accountID, userID string) ([]*peer.Peer, error)
|
||||||
GetPeersByGroupIDs(ctx context.Context, accountID string, groupsIDs []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
|
// DeletePeers removes the given peers along with their group memberships and
|
||||||
// policies. With checkConnected, a peer that is still connected or was seen
|
// policies. Every peer that was not deleted is returned in skipped, so the
|
||||||
// too recently is left in place and returned in skipped, so the caller can
|
// caller can retry it later: with checkConnected, a peer that is still
|
||||||
// retry it later.
|
// 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)
|
DeletePeers(ctx context.Context, accountID string, peerIDs []string, userID string, checkConnected bool) (skipped []string, err error)
|
||||||
SetNetworkMapController(networkMapController network_map.Controller)
|
SetNetworkMapController(networkMapController network_map.Controller)
|
||||||
SetIntegratedPeerValidator(integratedPeerValidator integrated_validator.IntegratedValidator)
|
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) {
|
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)
|
settings, err := m.store.GetAccountSettings(ctx, store.LockingStrengthNone, accountID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return peerIDs, err
|
||||||
}
|
}
|
||||||
dnsDomain := m.networkMapController.GetDNSDomain(settings)
|
dnsDomain := m.networkMapController.GetDNSDomain(settings)
|
||||||
|
|
||||||
var skipped []string
|
var skipped []string
|
||||||
|
var merr *multierror.Error
|
||||||
deletedAny := false
|
deletedAny := false
|
||||||
for _, peerID := range peerIDs {
|
for _, peerID := range peerIDs {
|
||||||
outcome, events, err := m.deleteSinglePeer(ctx, accountID, peerID, userID, checkConnected, dnsDomain)
|
outcome, events, err := m.deleteSinglePeer(ctx, accountID, peerID, userID, checkConnected, dnsDomain)
|
||||||
if err != nil {
|
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
|
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})
|
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
|
// 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.
|
// store once the transaction commits.
|
||||||
func (m *managerImpl) deletePeerObjects(ctx context.Context, transaction store.Store, accountID, userID, dnsDomain string, p *peer.Peer) ([]func(), error) {
|
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 {
|
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()
|
var eventsToStore []func()
|
||||||
|
|||||||
Reference in New Issue
Block a user