mirror of
https://github.com/netbirdio/netbird.git
synced 2026-08-24 16:41:30 +02:00
feat(agentnetwork): allocate subdomains transactionally, retrying on conflict
Allocation read a per-cluster set of taken labels, picked one, and wrote it later. That had three defects: the set was per-cluster, which is wrong once labels must be unique across a shared zone; the read and the write were not atomic; and on pool exhaustion it appended the first four characters of the account ID with no retry and no uniqueness check -- and those four characters are constant for accounts created within roughly the same hour, so two such accounts could be handed the same label. Allocation now picks a label and inserts it inside a transaction, retrying with a fresh label when the database rejects a duplicate, and failing loudly when the attempt budget is exhausted. A fresh transaction per attempt is required rather than incidental: on PostgreSQL a failed statement poisons the enclosing transaction, so a single transaction wrapping the loop would fail every attempt after the first. Because the settings primary key is the account ID, a concurrent bootstrap for the same account fails on the primary key rather than the subdomain index. That is indistinguishable from a label collision by message, so the loop re-reads by account before retrying and returns the winner's row -- the same answer the sequential path gives.
This commit is contained in:
261
management/internals/modules/agentnetwork/allocate_test.go
Normal file
261
management/internals/modules/agentnetwork/allocate_test.go
Normal file
@@ -0,0 +1,261 @@
|
||||
package agentnetwork
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"runtime"
|
||||
"testing"
|
||||
|
||||
"github.com/golang/mock/gomock"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/netbirdio/netbird/management/internals/modules/agentnetwork/labelgen"
|
||||
"github.com/netbirdio/netbird/management/internals/modules/agentnetwork/types"
|
||||
"github.com/netbirdio/netbird/management/server/store"
|
||||
nbtypes "github.com/netbirdio/netbird/management/server/types"
|
||||
"github.com/netbirdio/netbird/shared/management/status"
|
||||
)
|
||||
|
||||
// TestIsUniqueConstraintError_RecognisesAllThreeDialects — the allocator's
|
||||
// retry loop hinges on this. A missed dialect turns a retryable collision into
|
||||
// a hard provider-create failure.
|
||||
func TestIsUniqueConstraintError_RecognisesAllThreeDialects(t *testing.T) {
|
||||
for name, err := range map[string]error{
|
||||
"postgres": errors.New(`ERROR: duplicate key value violates unique constraint (SQLSTATE 23505)`),
|
||||
"mysql": errors.New(`Error 1062 (23000): Duplicate entry 'brave-otter'`),
|
||||
"sqlite": errors.New(`UNIQUE constraint failed: agent_network_settings.subdomain`),
|
||||
} {
|
||||
assert.True(t, isUniqueConstraintError(err), "%s violation must be recognised", name)
|
||||
}
|
||||
|
||||
assert.False(t, isUniqueConstraintError(errors.New("connection refused")),
|
||||
"unrelated errors must not be treated as retryable collisions")
|
||||
}
|
||||
|
||||
// newAllocatorTestStore wires a real sqlite store, mirroring the pattern in
|
||||
// provider_bootstrap_test.go's bootstrapFixture. The allocator tests exercise
|
||||
// bootstrapSettingsIfNeeded directly against a managerImpl built in-package,
|
||||
// so no permissions manager or account manager is needed.
|
||||
func newAllocatorTestStore(t *testing.T) store.Store {
|
||||
t.Helper()
|
||||
if runtime.GOOS == "windows" {
|
||||
t.Skip("sqlite store not properly supported on Windows yet")
|
||||
}
|
||||
t.Setenv("NETBIRD_STORE_ENGINE", string(nbtypes.SqliteStoreEngine))
|
||||
|
||||
st, cleanUp, err := store.NewTestStoreFromSQL(context.Background(), "", t.TempDir())
|
||||
require.NoError(t, err, "test store setup must succeed")
|
||||
t.Cleanup(cleanUp)
|
||||
return st
|
||||
}
|
||||
|
||||
// TestBootstrapSettings_StampsZoneAndTupleLabel — new rows must carry the
|
||||
// configured zone and a tuple label, which together give the tenant a
|
||||
// placement-independent address.
|
||||
func TestBootstrapSettings_StampsZoneAndTupleLabel(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
st := newAllocatorTestStore(t)
|
||||
|
||||
m := &managerImpl{
|
||||
store: st,
|
||||
zone: "gateway.example",
|
||||
labelRng: rand.New(rand.NewSource(1)),
|
||||
}
|
||||
|
||||
settings, err := m.bootstrapSettingsIfNeeded(ctx, "account1", "cluster1.example.com")
|
||||
require.NoError(t, err, "bootstrap must succeed")
|
||||
require.NotNil(t, settings)
|
||||
|
||||
assert.Equal(t, "gateway.example", settings.Zone, "new row must carry the configured zone")
|
||||
assert.Equal(t, "cluster1.example.com", settings.Cluster)
|
||||
assert.Contains(t, settings.Subdomain, "-", "subdomain must be an adjective-noun tuple label")
|
||||
assert.Equal(t, "account1", settings.AccountID)
|
||||
assert.Equal(t, settings.Subdomain+".gateway.example", settings.Endpoint(),
|
||||
"endpoint must be placement-independent, hanging off the zone rather than the cluster")
|
||||
|
||||
persisted, err := st.GetAgentNetworkSettings(ctx, store.LockingStrengthNone, "account1")
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, settings.Subdomain, persisted.Subdomain, "returned settings must match the persisted row")
|
||||
assert.Equal(t, "gateway.example", persisted.Zone)
|
||||
}
|
||||
|
||||
// TestBootstrapSettings_RetriesOnCollision forces a duplicate by pre-inserting
|
||||
// a row whose subdomain matches the next label the seeded rng will draw, then
|
||||
// asserts allocation still succeeds with a different label and that no error
|
||||
// escapes.
|
||||
func TestBootstrapSettings_RetriesOnCollision(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
st := newAllocatorTestStore(t)
|
||||
|
||||
const seed = 7
|
||||
|
||||
// Precompute the label a freshly seeded rng will draw first, without
|
||||
// disturbing the rng the manager will actually use.
|
||||
predictor := rand.New(rand.NewSource(seed))
|
||||
firstDraw := labelgen.PickTuple(predictor)
|
||||
require.NotEmpty(t, firstDraw, "test precondition: label pools must be non-empty")
|
||||
|
||||
// Pre-insert a colliding row on a different account so the allocator's
|
||||
// first attempt hits the unique index and must retry.
|
||||
require.NoError(t, st.CreateAgentNetworkSettings(ctx, &types.Settings{
|
||||
AccountID: "other-account",
|
||||
Cluster: "cluster1.example.com",
|
||||
Subdomain: firstDraw,
|
||||
}), "seeding the colliding row must succeed")
|
||||
|
||||
m := &managerImpl{
|
||||
store: st,
|
||||
labelRng: rand.New(rand.NewSource(seed)),
|
||||
}
|
||||
|
||||
settings, err := m.bootstrapSettingsIfNeeded(ctx, "account1", "cluster1.example.com")
|
||||
require.NoError(t, err, "allocation must succeed after retrying past the collision")
|
||||
require.NotNil(t, settings)
|
||||
assert.NotEqual(t, firstDraw, settings.Subdomain,
|
||||
"the retried allocation must not reuse the already-taken label")
|
||||
}
|
||||
|
||||
// TestBootstrapSettings_IsIdempotent — calling twice for one account returns
|
||||
// the existing row unchanged (the early-return path), and does NOT
|
||||
// re-allocate.
|
||||
func TestBootstrapSettings_IsIdempotent(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
st := newAllocatorTestStore(t)
|
||||
|
||||
m := &managerImpl{
|
||||
store: st,
|
||||
labelRng: rand.New(rand.NewSource(3)),
|
||||
}
|
||||
|
||||
first, err := m.bootstrapSettingsIfNeeded(ctx, "account1", "cluster1.example.com")
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, first)
|
||||
|
||||
second, err := m.bootstrapSettingsIfNeeded(ctx, "account1", "cluster2.example.com")
|
||||
require.NoError(t, err, "second call must not error")
|
||||
require.NotNil(t, second)
|
||||
|
||||
assert.Equal(t, first.Subdomain, second.Subdomain, "second call must return the existing subdomain unchanged")
|
||||
assert.Equal(t, first.Cluster, second.Cluster, "second call must not repin the cluster to the new hint")
|
||||
|
||||
all, err := st.GetAllAgentNetworkSettings(ctx, store.LockingStrengthNone)
|
||||
require.NoError(t, err)
|
||||
var forAccount int
|
||||
for _, s := range all {
|
||||
if s.AccountID == "account1" {
|
||||
forAccount++
|
||||
}
|
||||
}
|
||||
assert.Equal(t, 1, forAccount, "exactly one row must exist for the account; no re-allocation")
|
||||
}
|
||||
|
||||
// TestBootstrapSettings_FailsAfterExhaustingAttempts — the retry loop's
|
||||
// failure mode. maxSubdomainAllocationAttempts consecutive collisions must
|
||||
// surface an error rather than inserting a duplicate, silently succeeding, or
|
||||
// looping forever.
|
||||
//
|
||||
// Seed 11 was checked to produce maxSubdomainAllocationAttempts distinct
|
||||
// labels from labelgen.PickTuple; a seed that repeated a label would leave
|
||||
// fewer than maxAttempts rows pre-inserted and the allocator would succeed on
|
||||
// the repeat instead of exhausting.
|
||||
func TestBootstrapSettings_FailsAfterExhaustingAttempts(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
st := newAllocatorTestStore(t)
|
||||
|
||||
const seed = 11
|
||||
predictor := rand.New(rand.NewSource(seed))
|
||||
seen := make(map[string]struct{}, maxSubdomainAllocationAttempts)
|
||||
for i := 0; i < maxSubdomainAllocationAttempts; i++ {
|
||||
label := labelgen.PickTuple(predictor)
|
||||
_, dup := seen[label]
|
||||
require.False(t, dup, "test precondition: seed %d must draw %d distinct labels, got a repeat %q at draw %d", seed, maxSubdomainAllocationAttempts, label, i)
|
||||
seen[label] = struct{}{}
|
||||
|
||||
require.NoError(t, st.CreateAgentNetworkSettings(ctx, &types.Settings{
|
||||
AccountID: fmt.Sprintf("squatter-%d", i),
|
||||
Cluster: "cluster1.example.com",
|
||||
Subdomain: label,
|
||||
}), "seeding colliding row %d must succeed", i)
|
||||
}
|
||||
|
||||
m := &managerImpl{
|
||||
store: st,
|
||||
labelRng: rand.New(rand.NewSource(seed)),
|
||||
}
|
||||
|
||||
settings, err := m.bootstrapSettingsIfNeeded(ctx, "account1", "cluster1.example.com")
|
||||
require.Error(t, err, "exhausting every attempt to a collision must not silently succeed")
|
||||
assert.Nil(t, settings, "no settings row may be returned on failure")
|
||||
assert.Contains(t, err.Error(), "attempts exhausted")
|
||||
|
||||
_, err = st.GetAgentNetworkSettings(ctx, store.LockingStrengthNone, "account1")
|
||||
assert.Error(t, err, "no settings row must be persisted for the account when allocation fails")
|
||||
}
|
||||
|
||||
// TestBootstrapSettings_ConcurrentBootstrapReturnsWinnersRow covers the
|
||||
// same-account race: Settings' primary key is AccountID, and the
|
||||
// existence pre-check in bootstrapSettingsIfNeeded runs outside the
|
||||
// transaction, so two concurrent first-provider creates for the same
|
||||
// account can both observe NotFound and both proceed to allocate. The
|
||||
// loser's INSERT then fails on the primary key rather than the subdomain
|
||||
// unique index — a string isUniqueConstraintError still recognises — and
|
||||
// must not be treated as a label collision to retry past; it must
|
||||
// re-read and return the winner's row.
|
||||
//
|
||||
// This is scripted against a gomock store rather than driven by real
|
||||
// goroutines against the sqlite test store: NewTestStoreFromSQL caps the
|
||||
// pool at a single open connection (see its startup log,
|
||||
// "max open db connections to 1"), which serialises statement execution
|
||||
// enough that reliably forcing the exact interleaving this test needs —
|
||||
// both pre-checks observing NotFound before either INSERT lands — would
|
||||
// depend on goroutine scheduling rather than the store, making a
|
||||
// real-goroutine version flaky rather than deterministic. Scripting the
|
||||
// exact sequence (pre-check miss, PK-shaped insert failure, re-read hit)
|
||||
// through a MockStore exercises the same re-read branch precisely and
|
||||
// deterministically.
|
||||
func TestBootstrapSettings_ConcurrentBootstrapReturnsWinnersRow(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
ctrl := gomock.NewController(t)
|
||||
mockStore := store.NewMockStore(ctrl)
|
||||
|
||||
winner := &types.Settings{
|
||||
AccountID: "account1",
|
||||
Cluster: "cluster1.example.com",
|
||||
Subdomain: "brave-otter",
|
||||
}
|
||||
|
||||
gomock.InOrder(
|
||||
// The pre-check: no row yet, so this bootstrap proceeds to allocate.
|
||||
mockStore.EXPECT().
|
||||
GetAgentNetworkSettings(gomock.Any(), store.LockingStrengthNone, "account1").
|
||||
Return(nil, status.Errorf(status.NotFound, "agent network settings not found")),
|
||||
// The insert loses the race. The message shape is the sqlite wording
|
||||
// for a primary-key violation on account_id (not the subdomain
|
||||
// index), per the review finding this test locks down.
|
||||
mockStore.EXPECT().
|
||||
ExecuteInTransaction(gomock.Any(), gomock.Any()).
|
||||
DoAndReturn(func(_ context.Context, f func(store.Store) error) error {
|
||||
return f(mockStore)
|
||||
}),
|
||||
// The re-read after the PK conflict finds the concurrent winner's row.
|
||||
mockStore.EXPECT().
|
||||
GetAgentNetworkSettings(gomock.Any(), store.LockingStrengthNone, "account1").
|
||||
Return(winner, nil),
|
||||
)
|
||||
mockStore.EXPECT().
|
||||
CreateAgentNetworkSettings(gomock.Any(), gomock.Any()).
|
||||
Return(errors.New("UNIQUE constraint failed: agent_network_settings.account_id"))
|
||||
|
||||
m := &managerImpl{
|
||||
store: mockStore,
|
||||
labelRng: rand.New(rand.NewSource(9)),
|
||||
}
|
||||
|
||||
settings, err := m.bootstrapSettingsIfNeeded(ctx, "account1", "cluster1.example.com")
|
||||
require.NoError(t, err, "losing the same-account race must not surface as an error")
|
||||
require.NotNil(t, settings)
|
||||
assert.Same(t, winner, settings, "the loser must return the concurrent winner's row, not retry past it")
|
||||
}
|
||||
@@ -133,7 +133,7 @@ type managerImpl struct {
|
||||
reconcileMu sync.Mutex
|
||||
reconcileCache map[string]map[string]*proto.ProxyMapping
|
||||
|
||||
// labelRngMu guards labelRng. PickUnique consumes math/rand.Source
|
||||
// labelRngMu guards labelRng. PickTuple consumes math/rand.Source
|
||||
// state; concurrent provider creates would otherwise race.
|
||||
labelRngMu sync.Mutex
|
||||
labelRng *rand.Rand
|
||||
@@ -309,6 +309,22 @@ func (m *managerImpl) DeleteProvider(ctx context.Context, accountID, userID, pro
|
||||
return nil
|
||||
}
|
||||
|
||||
// isUniqueConstraintError reports whether err is a duplicate-key rejection.
|
||||
//
|
||||
// The equivalent helper in management/server is unexported, so it cannot be
|
||||
// reused from here; this is a deliberate duplicate rather than a new dependency
|
||||
// on that package for a single three-line matcher. Keep the two in sync if a
|
||||
// dialect is added.
|
||||
func isUniqueConstraintError(err error) bool {
|
||||
if err == nil {
|
||||
return false
|
||||
}
|
||||
msg := err.Error()
|
||||
return strings.Contains(msg, "(SQLSTATE 23505)") || // postgres
|
||||
strings.Contains(msg, "Error 1062 (23000)") || // mysql
|
||||
strings.Contains(msg, "UNIQUE constraint failed") // sqlite
|
||||
}
|
||||
|
||||
func pluralize(n int, singular, plural string) string {
|
||||
if n == 1 {
|
||||
return singular
|
||||
@@ -632,12 +648,6 @@ func (m *managerImpl) GetSettings(ctx context.Context, accountID, userID string)
|
||||
return m.store.GetAgentNetworkSettings(ctx, store.LockingStrengthNone, accountID)
|
||||
}
|
||||
|
||||
// bootstrapSettingsIfNeeded creates the per-account agent-network
|
||||
// settings row when missing. The cluster comes from the create-time
|
||||
// hint the dashboard sends (auto-picked from the active cluster list);
|
||||
// the subdomain is picked from the curated wordlist avoiding
|
||||
// collisions on the same cluster. Idempotent: if a row already exists
|
||||
// it is returned untouched and the hint is ignored.
|
||||
// requireSettingsBootstrapPermission gates the one-time settings bootstrap a
|
||||
// first provider create performs. Pinning the account's cluster and subdomain
|
||||
// is a settings write, so it needs the settings permission on top of the
|
||||
@@ -654,6 +664,15 @@ func (m *managerImpl) requireSettingsBootstrapPermission(ctx context.Context, ac
|
||||
return m.requirePermission(ctx, accountID, userID, modules.AgentNetworkSettings, operations.Create)
|
||||
}
|
||||
|
||||
// maxSubdomainAllocationAttempts bounds the allocate-and-insert retry loop in
|
||||
// bootstrapSettingsIfNeeded. Package-level (rather than function-local) so
|
||||
// tests can assert on the exhaustion path without duplicating the literal.
|
||||
const maxSubdomainAllocationAttempts = 10
|
||||
|
||||
// bootstrapSettingsIfNeeded creates the per-account agent-network settings
|
||||
// row when missing, allocating a subdomain unique across the whole zone.
|
||||
// Idempotent: if a row already exists it is returned untouched and the
|
||||
// cluster hint is ignored.
|
||||
func (m *managerImpl) bootstrapSettingsIfNeeded(ctx context.Context, accountID, providerCluster string) (*types.Settings, error) {
|
||||
if accountID == "" {
|
||||
return nil, fmt.Errorf("bootstrap settings: account id is required")
|
||||
@@ -671,40 +690,66 @@ func (m *managerImpl) bootstrapSettingsIfNeeded(ctx context.Context, accountID,
|
||||
return nil, fmt.Errorf("get agent network settings: %w", err)
|
||||
}
|
||||
|
||||
siblings, err := m.store.GetAgentNetworkSettingsByCluster(ctx, store.LockingStrengthNone, providerCluster)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("list agent network settings on cluster: %w", err)
|
||||
}
|
||||
taken := make(map[string]struct{}, len(siblings))
|
||||
for _, s := range siblings {
|
||||
taken[s.Subdomain] = struct{}{}
|
||||
}
|
||||
|
||||
suffix := accountID
|
||||
if len(suffix) > 4 {
|
||||
suffix = suffix[:4]
|
||||
}
|
||||
|
||||
m.labelRngMu.Lock()
|
||||
subdomain := labelgen.PickUnique(m.labelRng, taken, suffix)
|
||||
m.labelRngMu.Unlock()
|
||||
|
||||
// Allocate a subdomain and insert in one transaction, retrying on a unique
|
||||
// violation. This replaces a read-then-write over a pre-computed "taken"
|
||||
// set, which had three defects: the set was per-cluster (wrong once the
|
||||
// endpoint hangs off a shared zone), the read and the write were not
|
||||
// atomic, and the exhaustion fallback appended accountID[:4] — constant for
|
||||
// every account created within the same ~68 minutes — with no retry and no
|
||||
// uniqueness check, so two such accounts could be handed the same label.
|
||||
now := time.Now().UTC()
|
||||
settings := &types.Settings{
|
||||
AccountID: accountID,
|
||||
Cluster: providerCluster,
|
||||
Subdomain: subdomain,
|
||||
// Logs on by default; usage is collected regardless. Retention bounds
|
||||
// how long full log rows are kept.
|
||||
AccountID: accountID,
|
||||
Cluster: providerCluster,
|
||||
Zone: m.zone,
|
||||
EnableLogCollection: true,
|
||||
AccessLogRetentionDays: types.DefaultAccessLogRetentionDays,
|
||||
CreatedAt: now,
|
||||
UpdatedAt: now,
|
||||
}
|
||||
if err := m.store.SaveAgentNetworkSettings(ctx, settings); err != nil {
|
||||
return nil, fmt.Errorf("save agent network settings: %w", err)
|
||||
|
||||
for attempt := 1; attempt <= maxSubdomainAllocationAttempts; attempt++ {
|
||||
m.labelRngMu.Lock()
|
||||
settings.Subdomain = labelgen.PickTuple(m.labelRng)
|
||||
m.labelRngMu.Unlock()
|
||||
|
||||
if settings.Subdomain == "" {
|
||||
// Only reachable if either word pool were emptied; a database
|
||||
// insert of an empty subdomain would collide with the unique
|
||||
// index in a confusing way and produce a broken endpoint like
|
||||
// ".gateway.example". Fail loudly instead of looping or inserting.
|
||||
return nil, fmt.Errorf(
|
||||
"allocate agent network subdomain for account %s: label generator returned an empty label",
|
||||
accountID)
|
||||
}
|
||||
|
||||
err := m.store.ExecuteInTransaction(ctx, func(transaction store.Store) error {
|
||||
return transaction.CreateAgentNetworkSettings(ctx, settings)
|
||||
})
|
||||
if err == nil {
|
||||
return settings, nil
|
||||
}
|
||||
if isUniqueConstraintError(err) {
|
||||
// A concurrent bootstrap for this account may have won the race: the
|
||||
// pre-check above is outside the transaction, and the settings PK is
|
||||
// account_id, so the loser's insert fails on the primary key rather
|
||||
// than the subdomain index. Re-read before assuming the label was
|
||||
// taken, so a same-account race resolves immediately instead of
|
||||
// burning every remaining attempt on the same primary-key conflict.
|
||||
if existing, getErr := m.store.GetAgentNetworkSettings(ctx, store.LockingStrengthNone, accountID); getErr == nil {
|
||||
return existing, nil
|
||||
}
|
||||
log.WithContext(ctx).Tracef(
|
||||
"agent-network subdomain %q taken, retrying (attempt %d/%d)",
|
||||
settings.Subdomain, attempt, maxSubdomainAllocationAttempts)
|
||||
continue
|
||||
}
|
||||
return nil, fmt.Errorf("create agent network settings: %w", err)
|
||||
}
|
||||
return settings, nil
|
||||
|
||||
return nil, fmt.Errorf(
|
||||
"allocate agent network subdomain for account %s: %d attempts exhausted",
|
||||
accountID, maxSubdomainAllocationAttempts)
|
||||
}
|
||||
|
||||
// ListConsumption returns every consumption row recorded for the
|
||||
|
||||
Reference in New Issue
Block a user