diff --git a/weed/iam/integration/iam_manager.go b/weed/iam/integration/iam_manager.go index 4251587dc..45290c5e4 100644 --- a/weed/iam/integration/iam_manager.go +++ b/weed/iam/integration/iam_manager.go @@ -36,8 +36,15 @@ type IAMManager struct { // file, by ARN. With an in-memory store they are also written to the // store; a persistent store never holds them (see installOIDCProviderStore). staticOIDCProviders map[string]*OIDCProviderRecord - // cancelOIDCLoad stops a startup load still retrying against the store. + // oidcRetryMu guards the background refresh retry and which store is + // current: cancelOIDCLoad stops the retry, oidcRetryGen names the one + // running (0 when none) so at most one runs, and oidcRetryAgain records a + // refresh that failed while it ran, so the retry runs once more. + oidcRetryMu sync.Mutex cancelOIDCLoad context.CancelFunc + oidcRetryGen uint64 + oidcRetrySeq uint64 + oidcRetryAgain bool // oidcRefreshMu serializes refreshes from reading the store to handing // STS the result, so an older snapshot cannot replace a newer one. oidcRefreshMu sync.Mutex @@ -625,11 +632,12 @@ func (m *IAMManager) initOIDCProviderStore(config *IAMConfig) error { // providers so that an API call can shadow a bootstrap entry. Deleting the // stored provider brings the config-file one back. func (m *IAMManager) installOIDCProviderStore(store OIDCProviderStore, stsConfig *sts.STSConfig) { - if m.cancelOIDCLoad != nil { - m.cancelOIDCLoad() - m.cancelOIDCLoad = nil - } + // Cancel the old store's retry and switch stores in one step, so a failed + // refresh of the old store cannot start a retry after the cancel. + m.oidcRetryMu.Lock() + m.stopOIDCRetryLocked() m.oidcProviderStore = store + m.oidcRetryMu.Unlock() m.staticOIDCProviders = staticOIDCProviderRecords(stsConfig) if _, inMemory := store.(*MemoryOIDCProviderStore); inMemory { ctx := context.Background() @@ -646,16 +654,71 @@ func (m *IAMManager) installOIDCProviderStore(store OIDCProviderStore, stsConfig m.staticOIDCProviders = nil return } + // The metadata subscription only reports changes made from now on, so + // providers already in the store would stay unknown until one changes; a + // failed load is retried (RefreshOIDCProvidersFromStore). if err := m.RefreshOIDCProvidersFromStore(context.Background()); err != nil { - // The metadata subscription only reports changes made from now on, so - // providers already in the store would stay unknown until one changes. glog.Warningf("load OIDC providers from the store at startup: %v; retrying in the background", err) - ctx, cancel := context.WithCancel(context.Background()) - m.cancelOIDCLoad = cancel - go m.retryOIDCProviderLoad(ctx, store, oidcHydrateRetry) } } +// startOIDCRetry retries loading store in the background until it succeeds. +// A store that is no longer current gets no retry: nothing would cancel it, +// and its eventual success would replace the current store's providers. When +// a retry is already running, it is asked to run once more instead, because +// it may already have listed a snapshot older than this failure. +func (m *IAMManager) startOIDCRetry(store OIDCProviderStore) { + m.oidcRetryMu.Lock() + defer m.oidcRetryMu.Unlock() + if store != m.oidcProviderStore { + return + } + if m.oidcRetryGen != 0 { + m.oidcRetryAgain = true + return + } + ctx, cancel := context.WithCancel(context.Background()) + m.oidcRetrySeq++ + gen := m.oidcRetrySeq + m.cancelOIDCLoad, m.oidcRetryGen = cancel, gen + bounds := oidcHydrateRetry // read here, not in the goroutine: it outlives its caller + go func() { + defer cancel() + for { + m.retryOIDCProviderLoad(ctx, store, bounds) + m.oidcRetryMu.Lock() + if m.oidcRetryGen != gen { + m.oidcRetryMu.Unlock() + return // cancelled: another store was installed + } + if m.oidcRetryAgain && ctx.Err() == nil { + m.oidcRetryAgain = false + m.oidcRetryMu.Unlock() + continue + } + m.cancelOIDCLoad, m.oidcRetryGen, m.oidcRetryAgain = nil, 0, false + m.oidcRetryMu.Unlock() + return + } + }() +} + +// stopOIDCRetryLocked cancels a running retry, as installing another store +// must. The caller holds oidcRetryMu. +func (m *IAMManager) stopOIDCRetryLocked() { + if m.cancelOIDCLoad != nil { + m.cancelOIDCLoad() + } + m.cancelOIDCLoad, m.oidcRetryGen, m.oidcRetryAgain = nil, 0, false +} + +// currentOIDCProviderStore is the installed store, read under oidcRetryMu. +func (m *IAMManager) currentOIDCProviderStore() OIDCProviderStore { + m.oidcRetryMu.Lock() + defer m.oidcRetryMu.Unlock() + return m.oidcProviderStore +} + // staticOIDCProviderRecords describes the enabled OIDC providers of the IAM // config file as provider records. func staticOIDCProviderRecords(stsConfig *sts.STSConfig) map[string]*OIDCProviderRecord { @@ -697,6 +760,8 @@ var oidcHydrateRetry = struct{ initial, max time.Duration }{initial: time.Second // retryOIDCProviderLoad retries loading store, the store it was started for, // until it succeeds or ctx is cancelled because another store was installed. +// It calls refreshOIDCProvidersFrom, not RefreshOIDCProvidersFromStore, so it +// never schedules a retry of its own. func (m *IAMManager) retryOIDCProviderLoad(ctx context.Context, store OIDCProviderStore, bounds struct{ initial, max time.Duration }) { delay := bounds.initial for { @@ -735,8 +800,20 @@ func (m *IAMManager) refreshOIDCProvidersBestEffort(ctx context.Context, op, arn // the store is empty (clears the IAM-managed map). Records with empty URLs // or invalid configuration are logged and skipped so a single bad entry // does not stop the rest from refreshing. +// +// A refresh that fails keeps retrying in the background until the store +// answers. Every caller needs that: a metadata-subscription event reports each +// change ONCE, so a refresh that found the filer unreachable on it would leave +// a peer's new provider untrusted, or a deleted one trusted, until an unrelated +// later change; the refresh after a local IAM API mutation and the startup load +// have the same shape. At most one retry runs. func (m *IAMManager) RefreshOIDCProvidersFromStore(ctx context.Context) error { - return m.refreshOIDCProvidersFrom(ctx, m.oidcProviderStore) + store := m.currentOIDCProviderStore() + err := m.refreshOIDCProvidersFrom(ctx, store) + if err != nil && store != nil { + m.startOIDCRetry(store) + } + return err } // refreshOIDCProvidersFrom is RefreshOIDCProvidersFromStore for a given store. @@ -754,11 +831,15 @@ func (m *IAMManager) refreshOIDCProvidersFrom(ctx context.Context, store OIDCPro if err != nil { return fmt.Errorf("list OIDC providers: %w", err) } - // A startup retry is cancelled when another store is installed; its - // snapshot is of the old store and must not replace the new one's. + // A snapshot of a store that has since been replaced must not replace the + // current store's providers: a retry is cancelled when another store is + // installed, and a refresh may have listed the old store just before. if err := ctx.Err(); err != nil { return err } + if store != m.currentOIDCProviderStore() { + return fmt.Errorf("list OIDC providers: the store was replaced during the refresh") + } byIssuer := make(map[string][]sts.ScopedOIDCProvider, len(records)) for _, rec := range records { if rec == nil || rec.URL == "" { diff --git a/weed/iam/integration/oidc_provider_persist_test.go b/weed/iam/integration/oidc_provider_persist_test.go index 6904f232c..7190e2aa7 100644 --- a/weed/iam/integration/oidc_provider_persist_test.go +++ b/weed/iam/integration/oidc_provider_persist_test.go @@ -391,3 +391,171 @@ func TestAnOlderRefreshCannotRestoreADeletedProvider(t *testing.T) { require.NoError(t, deleteErr) assert.False(t, stsKnowsIssuer(t, mgr, persistTestAPIIssuer), "a refresh older than the deletion left the deleted provider trusted") } + +// A refresh that fails is retried until the store answers. A change event reports each mutation once, so +// a subscription refresh that found the filer unreachable — mid-restart, say — left the provider set stale +// until some unrelated later change; a peer's new provider stayed untrusted and a deleted one trusted. +func TestAFailedRefreshIsRetriedUntilTheStoreAnswers(t *testing.T) { + saved := oidcHydrateRetry + oidcHydrateRetry.initial, oidcHydrateRetry.max = time.Millisecond, 5*time.Millisecond + t.Cleanup(func() { oidcHydrateRetry = saved }) + + store := &unreachableThenReadyStore{MemoryOIDCProviderStore: NewMemoryOIDCProviderStore()} + mgr := startServer(t, store) + require.NoError(t, store.StoreProvider(context.Background(), "", &OIDCProviderRecord{ + ARN: arnOf(t, persistTestAPIIssuer), URL: persistTestAPIIssuer, ClientIDs: []string{"aud"}, + })) + store.mu.Lock() + store.failsLeft = 3 + store.mu.Unlock() + require.Error(t, mgr.RefreshOIDCProvidersFromStore(context.Background()), "precondition: the refresh fails") + + deadline := time.Now().Add(2 * time.Second) + for !stsKnowsIssuer(t, mgr, persistTestAPIIssuer) { + if time.Now().After(deadline) { + t.Fatal("a failed refresh was never retried: the stored provider stays untrusted until an unrelated change") + } + time.Sleep(5 * time.Millisecond) + } +} + +// Failures during an outage start ONE retry, not one per event: a filer that is down for a while produces a +// change event per mutation, and each would otherwise add a goroutine polling the same store. Once the store +// answers and the retry ends, a later failure starts a new one. +func TestFailedRefreshesShareOneRetry(t *testing.T) { + saved := oidcHydrateRetry + oidcHydrateRetry.initial, oidcHydrateRetry.max = time.Millisecond, time.Millisecond + t.Cleanup(func() { oidcHydrateRetry = saved }) + + store := &unreachableThenReadyStore{MemoryOIDCProviderStore: NewMemoryOIDCProviderStore()} + mgr := startServer(t, store) + store.mu.Lock() + store.failsLeft = 1 << 30 + store.mu.Unlock() + for range 20 { + require.Error(t, mgr.RefreshOIDCProvidersFromStore(context.Background())) + } + mgr.oidcRetryMu.Lock() + started := mgr.oidcRetrySeq + mgr.oidcRetryMu.Unlock() + assert.Equal(t, uint64(1), started, "twenty failed refreshes started more than one retry") + + store.mu.Lock() + store.failsLeft = 0 + store.mu.Unlock() + deadline := time.Now().Add(2 * time.Second) + for { + mgr.oidcRetryMu.Lock() + running := mgr.oidcRetryGen != 0 + mgr.oidcRetryMu.Unlock() + if !running { + break + } + if time.Now().After(deadline) { + t.Fatal("the retry never ended after the store answered") + } + time.Sleep(time.Millisecond) + } + store.mu.Lock() + store.failsLeft = 1 + store.mu.Unlock() + require.Error(t, mgr.RefreshOIDCProvidersFromStore(context.Background())) + mgr.oidcRetryMu.Lock() + assert.Equal(t, uint64(2), mgr.oidcRetrySeq, "a failure after the retry ended starts a new one") + mgr.oidcRetryMu.Unlock() +} + +// flakyCountingStore fails while fail is set and counts the lists that succeed. +type flakyCountingStore struct { + *MemoryOIDCProviderStore + mu sync.Mutex + fail bool + ok int +} + +func (s *flakyCountingStore) ListProviders(ctx context.Context, addr string) ([]*OIDCProviderRecord, error) { + s.mu.Lock() + if s.fail { + s.mu.Unlock() + return nil, errors.New("filer unavailable") + } + s.ok++ + s.mu.Unlock() + return s.MemoryOIDCProviderStore.ListProviders(ctx, addr) +} + +func (s *flakyCountingStore) set(fail bool) { + s.mu.Lock() + defer s.mu.Unlock() + s.fail = fail +} + +func (s *flakyCountingStore) successes() int { + s.mu.Lock() + defer s.mu.Unlock() + return s.ok +} + +// A refresh of store A that fails as store B is installed must not start a retry for A: nothing would cancel +// it, and when A answered it would replace B's providers (a removed one trusted again, B's own untrusted). +func TestAFailedRefreshOfASupersededStoreStartsNoRetry(t *testing.T) { + a := &persistentTestStore{NewMemoryOIDCProviderStore()} + mgr := startServer(t, a) + mgr.SetOIDCProviderStore(&persistentTestStore{NewMemoryOIDCProviderStore()}) + mgr.oidcRetryMu.Lock() + before := mgr.oidcRetrySeq + mgr.oidcRetryMu.Unlock() + + mgr.startOIDCRetry(a) // the late call of A's failed refresh, after B's install + + mgr.oidcRetryMu.Lock() + defer mgr.oidcRetryMu.Unlock() + assert.Equal(t, before, mgr.oidcRetrySeq, "a retry was started for the store that was replaced") + assert.Zero(t, mgr.oidcRetryGen) +} + +// A refresh that listed store A before B was installed must not hand STS A's providers afterwards. +func TestASupersededSnapshotIsNotApplied(t *testing.T) { + a := &persistentTestStore{NewMemoryOIDCProviderStore()} + require.NoError(t, a.StoreProvider(context.Background(), "", &OIDCProviderRecord{ + ARN: arnOf(t, persistTestAPIIssuer), URL: persistTestAPIIssuer, ClientIDs: []string{"aud"}, + })) + b := &persistentTestStore{NewMemoryOIDCProviderStore()} + require.NoError(t, b.StoreProvider(context.Background(), "", &OIDCProviderRecord{ + ARN: arnOf(t, persistTestStaticIssuer), URL: persistTestStaticIssuer, ClientIDs: []string{"aud"}, + })) + mgr := startServer(t, a) + mgr.SetOIDCProviderStore(b) + + assert.Error(t, mgr.refreshOIDCProvidersFrom(context.Background(), a), "a snapshot of the replaced store was applied") + assert.False(t, stsKnowsIssuer(t, mgr, persistTestAPIIssuer), "the replaced store's provider is trusted") + assert.True(t, stsKnowsIssuer(t, mgr, persistTestStaticIssuer), "the current store's provider is not") +} + +// A refresh that fails while a retry runs is not dropped: the retry may already have listed an older snapshot, +// so after its success it runs once more and picks up whatever the failed refresh would have seen. +func TestAFailureDuringARetryIsNotDropped(t *testing.T) { + saved := oidcHydrateRetry + oidcHydrateRetry.initial, oidcHydrateRetry.max = time.Millisecond, 2*time.Millisecond + t.Cleanup(func() { oidcHydrateRetry = saved }) + + store := &flakyCountingStore{MemoryOIDCProviderStore: NewMemoryOIDCProviderStore(), fail: true} + mgr := startServer(t, store) // the startup load fails: a retry is running + require.Error(t, mgr.RefreshOIDCProvidersFromStore(context.Background()), "a later refresh fails while it runs") + store.set(false) + + deadline := time.Now().Add(2 * time.Second) + for { + mgr.oidcRetryMu.Lock() + running := mgr.oidcRetryGen != 0 + mgr.oidcRetryMu.Unlock() + if !running { + break + } + if time.Now().After(deadline) { + t.Fatal("the retry never ended") + } + time.Sleep(time.Millisecond) + } + assert.GreaterOrEqual(t, store.successes(), 2, "the retry ended on its first success, dropping the refresh that failed while it ran") +}