diff --git a/CHANGELOG.md b/CHANGELOG.md index dd279cfa131..fc7f2523e3f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/pkg/querier/tenantfederation/regex_resolver.go b/pkg/querier/tenantfederation/regex_resolver.go index 9fb2bc33aed..05f187a958d 100644 --- a/pkg/querier/tenantfederation/regex_resolver.go +++ b/pkg/querier/tenantfederation/regex_resolver.go @@ -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) @@ -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) { diff --git a/pkg/querier/tenantfederation/regex_resolver_test.go b/pkg/querier/tenantfederation/regex_resolver_test.go index 92285e835eb..9b92ccbba2d 100644 --- a/pkg/querier/tenantfederation/regex_resolver_test.go +++ b/pkg/querier/tenantfederation/regex_resolver_test.go @@ -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"}