Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,7 @@
* [BUGFIX] Alertmanager: Reject the global `mattermost_webhook_url_file` setting in per-tenant configs, consistent with every other global `*_file` setting. #7768
* [BUGFIX] Alertmanager: Tighten per-tenant config validation to reject additional file-based settings. #7767
* [BUGFIX] Querier: Fix panic (`index out of range [-1]`) in the active request tracker when truncating a `match[]`/`query` value made entirely of invalid UTF-8 continuation bytes. The backwards scan for a rune boundary now stops at index 0 instead of underflowing. #7743
* [BUGFIX] Tenant Federation: Fix regex tenant federation resolving to an empty user list right after startup. #7811

## 1.21.1 2026-06-04

Expand Down
55 changes: 37 additions & 18 deletions pkg/querier/tenantfederation/regex_resolver.go
Original file line number Diff line number Diff line change
Expand Up @@ -106,11 +106,22 @@ func NewRegexResolver(cfg users.UsersScannerConfig, tenantFederationCfg Config,
Help: "Number of discovered users.",
})

r.Service = services.NewBasicService(nil, r.running, nil)
r.Service = services.NewBasicService(r.starting, r.running, nil)

return r, nil
}

// starting discovers users before the service becomes running. Otherwise queries
// served right after the startup would be matched against an empty user list until
// the first sync happens.
func (r *RegexResolver) starting(ctx context.Context) error {
if err := r.updateUsers(ctx); err != nil {
return errors.Wrap(err, "failed to discover users from bucket")
}

return nil
}

func (r *RegexResolver) running(ctx context.Context) error {
level.Info(r.logger).Log("msg", "regex-resolver started")
ticker := time.NewTicker(r.userSyncInterval)
Expand All @@ -121,29 +132,37 @@ func (r *RegexResolver) running(ctx context.Context) error {
case <-ctx.Done():
return ctx.Err()
case <-ticker.C:
// Active and deleting users are considered
active, deleting, _, err := r.userScanner.ScanUsers(ctx)
if err != nil {
if err := r.updateUsers(ctx); err != nil {
level.Error(r.logger).Log("msg", "failed to discover users from bucket", "err", err)
continue
}
}
}
}

newUsers := append(active, deleting...)
sort.Strings(newUsers)
func (r *RegexResolver) updateUsers(ctx context.Context) error {
// Active and deleting users are considered
active, deleting, _, err := r.userScanner.ScanUsers(ctx)
if err != nil {
return err
}

r.Lock()
changed := !slices.Equal(r.knownUsers, newUsers)
r.knownUsers = newUsers
if changed && r.matchedCache != nil {
// Reset the cache when the set of available users has changed.
r.matchedCache.Purge()
r.matchedCacheSize.Set(0)
}
r.Unlock()
r.lastUpdateUserRun.SetToCurrentTime()
r.discoveredUsers.Set(float64(len(active) + len(deleting)))
}
newUsers := append(active, deleting...)
sort.Strings(newUsers)

r.Lock()
changed := !slices.Equal(r.knownUsers, newUsers)
r.knownUsers = newUsers
if changed && r.matchedCache != nil {
// Reset the cache when the set of available users has changed.
r.matchedCache.Purge()
r.matchedCacheSize.Set(0)
}
r.Unlock()
r.lastUpdateUserRun.SetToCurrentTime()
r.discoveredUsers.Set(float64(len(active) + len(deleting)))

return nil
}

func (r *RegexResolver) TenantID(ctx context.Context) (string, error) {
Expand Down
51 changes: 51 additions & 0 deletions pkg/querier/tenantfederation/regex_resolver_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -266,6 +266,57 @@ func Test_RegexResolver_Cache(t *testing.T) {
}
}

func Test_RegexResolver_InitialSyncBeforeRunning(t *testing.T) {
reg := prometheus.NewRegistry()
existingTenants := []string{"user-1", "user-2"}
bucketClient := &bucket.ClientMock{}
bucketClient.MockIter("", existingTenants, nil)
bucketClient.MockIter("__markers__", []string{}, nil)
for _, tenant := range existingTenants {
bucketClient.MockExists(users.GetGlobalDeletionMarkPath(tenant), false, nil)
bucketClient.MockExists(users.GetLocalDeletionMarkPath(tenant), false, nil)
}

bucketClientFactory := func(ctx context.Context) (objstore.InstrumentedBucket, error) {
return bucketClient, nil
}

usersScannerConfig := users.UsersScannerConfig{Strategy: users.UserScanStrategyList}
// The sync interval is long enough to make sure the ticker never fires during the test,
// so users can only be discovered by the initial sync done before the service is running.
tenantFederationConfig := Config{UserSyncInterval: time.Hour, MaxTenant: 0, RegexCacheSize: 10}
regexResolver, err := NewRegexResolver(usersScannerConfig, tenantFederationConfig, reg, bucketClientFactory, log.NewNopLogger())
require.NoError(t, err)

require.NoError(t, services.StartAndAwaitRunning(context.Background(), regexResolver))
defer services.StopAndAwaitTerminated(context.Background(), regexResolver) //nolint:errcheck

ctx := user.InjectOrgID(context.Background(), "user-.+")
orgIDs, err := regexResolver.TenantIDs(ctx)
require.NoError(t, err)
require.Equal(t, existingTenants, orgIDs)
require.Equal(t, float64(len(existingTenants)), testutil.ToFloat64(regexResolver.discoveredUsers))
require.Greater(t, testutil.ToFloat64(regexResolver.lastUpdateUserRun), float64(0))
}

func Test_RegexResolver_InitialSyncFailure(t *testing.T) {
reg := prometheus.NewRegistry()
bucketClient := &bucket.ClientMock{}
bucketClient.MockIter("", nil, errors.New("failed to iterate"))

bucketClientFactory := func(ctx context.Context) (objstore.InstrumentedBucket, error) {
return bucketClient, nil
}

usersScannerConfig := users.UsersScannerConfig{Strategy: users.UserScanStrategyList}
tenantFederationConfig := Config{UserSyncInterval: time.Hour, MaxTenant: 0, RegexCacheSize: 10}
regexResolver, err := NewRegexResolver(usersScannerConfig, tenantFederationConfig, reg, bucketClientFactory, log.NewNopLogger())
require.NoError(t, err)

// The service must fail to start rather than serving queries against an empty user list.
require.Error(t, services.StartAndAwaitRunning(context.Background(), regexResolver))
}

func Test_RegexResolver_CacheInvalidation(t *testing.T) {
reg := prometheus.NewRegistry()
initialTenants := []string{"user-1", "user-2"}
Expand Down
Loading