mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-07 22:19:08 +02:00
[management] fix the concurrency handling within the update message tests (#8000)
This commit is contained in:
@@ -13,6 +13,7 @@ import (
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"testing/synctest"
|
||||
"time"
|
||||
|
||||
"github.com/prometheus/client_golang/prometheus/push"
|
||||
@@ -40,6 +41,7 @@ import (
|
||||
"github.com/netbirdio/netbird/management/internals/modules/zones"
|
||||
networkmapdb "github.com/netbirdio/netbird/management/internals/network_map_db"
|
||||
networkmapdbfactory "github.com/netbirdio/netbird/management/internals/network_map_db/factory"
|
||||
networkmap_sqlite "github.com/netbirdio/netbird/management/internals/network_map_db/sqlite"
|
||||
"github.com/netbirdio/netbird/management/internals/server/config"
|
||||
nbgrpc "github.com/netbirdio/netbird/management/internals/shared/grpc"
|
||||
nbAccount "github.com/netbirdio/netbird/management/server/account"
|
||||
@@ -1297,7 +1299,9 @@ func TestAccountManager_AddPeerWithUserID(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestAccountManager_NetworkUpdates_SaveGroup(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_SaveGroup(t)
|
||||
runPeerUpdateTest(t, func(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_SaveGroup(t)
|
||||
})
|
||||
}
|
||||
|
||||
func testAccountManager_NetworkUpdates_SaveGroup(t *testing.T) {
|
||||
@@ -1330,6 +1334,8 @@ func testAccountManager_NetworkUpdates_SaveGroup(t *testing.T) {
|
||||
updMsg := updateManager.CreateChannel(context.Background(), peer1.ID)
|
||||
defer updateManager.CloseChannel(context.Background(), peer1.ID)
|
||||
|
||||
settleAffectedUpdates(updMsg)
|
||||
|
||||
wg := sync.WaitGroup{}
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
@@ -1352,7 +1358,9 @@ func testAccountManager_NetworkUpdates_SaveGroup(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestAccountManager_NetworkUpdates_DeletePolicy(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_DeletePolicy(t)
|
||||
runPeerUpdateTest(t, func(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_DeletePolicy(t)
|
||||
})
|
||||
}
|
||||
|
||||
func testAccountManager_NetworkUpdates_DeletePolicy(t *testing.T) {
|
||||
@@ -1361,13 +1369,7 @@ func testAccountManager_NetworkUpdates_DeletePolicy(t *testing.T) {
|
||||
updMsg := updateManager.CreateChannel(context.Background(), peer1.ID)
|
||||
defer updateManager.CloseChannel(context.Background(), peer1.ID)
|
||||
|
||||
// Ensure that we do not receive an update message before the policy is deleted
|
||||
time.Sleep(time.Second)
|
||||
select {
|
||||
case <-updMsg:
|
||||
t.Logf("received addPeer update message before policy deletion")
|
||||
default:
|
||||
}
|
||||
settleAffectedUpdates(updMsg)
|
||||
|
||||
wg := sync.WaitGroup{}
|
||||
wg.Add(1)
|
||||
@@ -1390,7 +1392,9 @@ func testAccountManager_NetworkUpdates_DeletePolicy(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestAccountManager_NetworkUpdates_SavePolicy(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_SavePolicy(t)
|
||||
runPeerUpdateTest(t, func(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_SavePolicy(t)
|
||||
})
|
||||
}
|
||||
|
||||
func testAccountManager_NetworkUpdates_SavePolicy(t *testing.T) {
|
||||
@@ -1410,6 +1414,8 @@ func testAccountManager_NetworkUpdates_SavePolicy(t *testing.T) {
|
||||
updMsg := updateManager.CreateChannel(context.Background(), peer1.ID)
|
||||
defer updateManager.CloseChannel(context.Background(), peer1.ID)
|
||||
|
||||
settleAffectedUpdates(updMsg)
|
||||
|
||||
wg := sync.WaitGroup{}
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
@@ -1443,7 +1449,9 @@ func testAccountManager_NetworkUpdates_SavePolicy(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestAccountManager_NetworkUpdates_DeletePeer(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_DeletePeer(t)
|
||||
runPeerUpdateTest(t, func(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_DeletePeer(t)
|
||||
})
|
||||
}
|
||||
|
||||
func testAccountManager_NetworkUpdates_DeletePeer(t *testing.T) {
|
||||
@@ -1482,6 +1490,8 @@ func testAccountManager_NetworkUpdates_DeletePeer(t *testing.T) {
|
||||
updMsg := updateManager.CreateChannel(context.Background(), peer1.ID)
|
||||
defer updateManager.CloseChannel(context.Background(), peer1.ID)
|
||||
|
||||
settleAffectedUpdates(updMsg)
|
||||
|
||||
wg := sync.WaitGroup{}
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
@@ -1503,7 +1513,9 @@ func testAccountManager_NetworkUpdates_DeletePeer(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestAccountManager_NetworkUpdates_DeleteGroup(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_DeleteGroup(t)
|
||||
runPeerUpdateTest(t, func(t *testing.T) {
|
||||
testAccountManager_NetworkUpdates_DeleteGroup(t)
|
||||
})
|
||||
}
|
||||
|
||||
func testAccountManager_NetworkUpdates_DeleteGroup(t *testing.T) {
|
||||
@@ -1550,6 +1562,8 @@ func testAccountManager_NetworkUpdates_DeleteGroup(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
settleAffectedUpdates(updMsg)
|
||||
|
||||
wg := sync.WaitGroup{}
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
@@ -3582,6 +3596,9 @@ func createManagerWithNetworkMapStore(t testing.TB) (*DefaultAccountManager, *up
|
||||
|
||||
nmdataStore, err := networkmapdbfactory.NewNetworkMapDBStore(context.Background(), types.SqliteStoreEngine, dataDir, MockIntegratedValidator{}, newSettingsMockManager(t))
|
||||
require.NoError(t, err)
|
||||
sqliteStore, ok := nmdataStore.Store.(*networkmap_sqlite.SqliteStore)
|
||||
require.True(t, ok, "network map store is %T, want *networkmap_sqlite.SqliteStore", nmdataStore.Store)
|
||||
t.Cleanup(func() { assert.NoError(t, sqliteStore.Db.Close()) })
|
||||
|
||||
manager, updateManager, err := buildTestManager(t, store, nmdataStore)
|
||||
require.NoError(t, err)
|
||||
@@ -3636,7 +3653,8 @@ func buildTestManager(t testing.TB, store store.Store, nmdataStore *networkmapdb
|
||||
Return(nil).
|
||||
AnyTimes()
|
||||
|
||||
cacheStore, err := cache.NewStore(ctx, 100*time.Millisecond, 300*time.Millisecond, 100)
|
||||
// The go-cache janitor only stops via a GC finalizer and would outlive synctest bubbles.
|
||||
cacheStore, err := cache.NewStore(ctx, 100*time.Millisecond, 0, 100)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
@@ -3648,8 +3666,11 @@ func buildTestManager(t testing.TB, store store.Store, nmdataStore *networkmapdb
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
cacheManager := manager.cacheManager
|
||||
t.Cleanup(func() { assert.NoError(t, cacheManager.Close()) })
|
||||
|
||||
proxyGrpcServer := nbgrpc.NewProxyServiceServer(nil, nil, nil, nbgrpc.ProxyOIDCConfig{}, peersManager, nil, nil, proxyManager, nil)
|
||||
t.Cleanup(proxyGrpcServer.Close)
|
||||
proxyController, err := proxymanager.NewGRPCController(proxyGrpcServer, noop.Meter{})
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
@@ -3743,6 +3764,33 @@ func setupNetworkMapTest(t *testing.T) (*DefaultAccountManager, *update_channel.
|
||||
// when the channel delivers.
|
||||
const peerUpdateTimeout = 5 * time.Second
|
||||
|
||||
// peerUpdateSettleTime bounds how far settleAffectedUpdates advances the fake clock. It must exceed
|
||||
// the account request and peer update buffer intervals.
|
||||
const peerUpdateSettleTime = time.Second
|
||||
|
||||
// runPeerUpdateTest runs f inside synctest.Test, so the peer update helpers observe every background
|
||||
// goroutine of the test, and lets the updates still in flight when f returns finish before the bubble
|
||||
// ends, since the bubble's clock stops with it.
|
||||
func runPeerUpdateTest(t *testing.T, f func(t *testing.T)) {
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
defer settleAffectedUpdates()
|
||||
f(t)
|
||||
})
|
||||
}
|
||||
|
||||
// settleAffectedUpdates runs the synctest bubble's fake clock past every update buffer interval until
|
||||
// all goroutines are blocked, then discards the updates already delivered to chans, so the next
|
||||
// assertion only observes updates from the action under test. It must be called inside synctest.Test.
|
||||
func settleAffectedUpdates(chans ...<-chan *network_map.UpdateMessage) {
|
||||
time.Sleep(peerUpdateSettleTime)
|
||||
synctest.Wait()
|
||||
for _, ch := range chans {
|
||||
for len(ch) > 0 {
|
||||
<-ch
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func drainPeerUpdates(ch <-chan *network_map.UpdateMessage) {
|
||||
for {
|
||||
select {
|
||||
@@ -3756,6 +3804,19 @@ func drainPeerUpdates(ch <-chan *network_map.UpdateMessage) {
|
||||
}
|
||||
}
|
||||
|
||||
// step runs f as one named stage of a test inside synctest.Test, which forbids t.Run, and names the
|
||||
// stage when it fails the test. Stages share t, so a fatal failure ends the remaining stages as well.
|
||||
func step(t *testing.T, name string, f func(t *testing.T)) {
|
||||
t.Helper()
|
||||
failedBefore := t.Failed()
|
||||
defer func() {
|
||||
if !failedBefore && t.Failed() {
|
||||
t.Logf("step %q failed", name)
|
||||
}
|
||||
}()
|
||||
f(t)
|
||||
}
|
||||
|
||||
func peerShouldNotReceiveUpdate(t *testing.T, updateMessage <-chan *network_map.UpdateMessage) {
|
||||
t.Helper()
|
||||
select {
|
||||
|
||||
Reference in New Issue
Block a user