diff --git a/internal/registries/popularity_github.go b/internal/registries/popularity_github.go index 0ad691b67..6cbaf2828 100644 --- a/internal/registries/popularity_github.go +++ b/internal/registries/popularity_github.go @@ -9,6 +9,7 @@ import ( "net/http" "net/url" "os" + "sort" "strconv" "strings" "sync" @@ -145,6 +146,9 @@ type githubStarsProvider struct { mu sync.Mutex entries map[string]*starsEntry store *popularityStore // nil = memory only + // storeMu serialises store writes; bbolt calls never run under p.mu. + // Lock order: storeMu -> mu. + storeMu sync.Mutex queue chan string queued map[string]chan struct{} // in-flight+queued dedup; closed on completion @@ -232,12 +236,33 @@ func NewGitHubStarsProvider(opts PopularityOptions) *githubStarsProvider { cancel: cancel, } if store != nil { + // Eager preload: memory is the source of truth from here on, so the + // FR-008 cap applies to what is on disk and Lookup never touches bbolt. p.entries = store.all() - p.mu.Lock() - for len(p.entries) > githubMaxCacheKeys { - p.evictIfNeededLocked() + if excess := len(p.entries) - githubMaxCacheKeys; excess > 0 { + type cachedKey struct { + key string + fetchedAt time.Time + } + oldest := make([]cachedKey, 0, len(p.entries)) + for key, entry := range p.entries { + oldest = append(oldest, cachedKey{key: key, fetchedAt: entry.FetchedAt}) + } + sort.Slice(oldest, func(i, j int) bool { + if oldest[i].fetchedAt.Equal(oldest[j].fetchedAt) { + return oldest[i].key < oldest[j].key + } + return oldest[i].fetchedAt.Before(oldest[j].fetchedAt) + }) + evicted := make([]string, excess) + for i := 0; i < excess; i++ { + evicted[i] = oldest[i].key + delete(p.entries, oldest[i].key) + } + if err := store.deleteMany(evicted); err != nil { + logger.Warn("catalog popularity: failed to evict excess entries", zap.Int("count", len(evicted)), zap.Error(err)) + } } - p.mu.Unlock() } if !p.disabled { p.startWorkers() @@ -292,20 +317,11 @@ func (p *githubStarsProvider) fresh(e *starsEntry) bool { return p.now().Sub(e.FetchedAt) < e.ttl() } -// getEntryLocked returns key's entry, lazily loading it from the store on a -// memory miss (plan.md: "Entries are loaded into memory lazily on the first -// Lookup miss"). Caller must hold p.mu. +// getEntryLocked returns key's in-memory entry, or nil. The constructor +// preloads every persisted entry, so memory is authoritative and this never +// reads bbolt (keeping bbolt I/O out from under p.mu). Caller must hold p.mu. func (p *githubStarsProvider) getEntryLocked(key string) *starsEntry { - if e, ok := p.entries[key]; ok { - return e - } - if p.store != nil { - if e, ok := p.store.get(key); ok { - p.entries[key] = e - return e - } - } - return nil + return p.entries[key] } // Lookup implements PopularityProvider. @@ -597,27 +613,26 @@ func (p *githubStarsProvider) applyResult(key string, prev *starsEntry, status, p.mu.Lock() p.entries[key] = entry - p.evictIfNeededLocked() + evictedKey := "" + if len(p.entries) > githubMaxCacheKeys { + evictedKey = p.evictOldestLocked() + } p.mu.Unlock() - if p.store != nil { - if err := p.store.put(key, entry); err != nil { - p.logger.Warn("catalog popularity: failed to persist entry", zap.String("key", key), zap.Error(err)) - } - } + // bbolt writes (an fsync each) happen outside p.mu so Lookup/enqueue + // never wait behind them. + p.persist(key, entry, evictedKey) if fetchErr == nil { p.applyBreaker(status, headers) } } -// evictIfNeededLocked drops the single oldest-FetchedAt entry once the cache -// exceeds its cap (FR-008). Caller must hold p.mu. Only ever one entry over -// cap at a time, since insertion happens one key at a time. -func (p *githubStarsProvider) evictIfNeededLocked() { - if len(p.entries) <= githubMaxCacheKeys { - return - } +// evictOldestLocked removes the oldest-FetchedAt entry from memory and +// returns its key ("" if the cache is empty). The caller is responsible for +// removing it from the store via persist AFTER releasing p.mu (FR-008). +// Caller must hold p.mu. +func (p *githubStarsProvider) evictOldestLocked() string { oldestKey := "" var oldestTime time.Time first := true @@ -626,13 +641,40 @@ func (p *githubStarsProvider) evictIfNeededLocked() { oldestKey, oldestTime, first = k, e.FetchedAt, false } } - if oldestKey == "" { + if oldestKey != "" { + delete(p.entries, oldestKey) + } + return oldestKey +} + +// persist mirrors an applyResult outcome to the store: it writes entry under +// key (skipped when key is "" or memory no longer holds this exact entry, +// i.e. it was evicted or superseded meanwhile) and deletes evictedKey +// (skipped when a newer fetch re-added it). Store writes are serialised by +// storeMu and the memory checks run after acquiring it, so a delayed put can +// never resurrect an evicted key and a delayed delete can never drop a +// re-added one. p.mu is only held for the brief check, never across bbolt. +func (p *githubStarsProvider) persist(key string, entry *starsEntry, evictedKey string) { + if p.store == nil { return } - delete(p.entries, oldestKey) - if p.store != nil { - if err := p.store.delete(oldestKey); err != nil { - p.logger.Warn("catalog popularity: failed to evict entry", zap.String("key", oldestKey), zap.Error(err)) + p.storeMu.Lock() + defer p.storeMu.Unlock() + + p.mu.Lock() + writeKey := key != "" && p.entries[key] == entry + _, evictedBack := p.entries[evictedKey] + deleteKey := evictedKey != "" && !evictedBack + p.mu.Unlock() + + if writeKey { + if err := p.store.put(key, entry); err != nil { + p.logger.Warn("catalog popularity: failed to persist entry", zap.String("key", key), zap.Error(err)) + } + } + if deleteKey { + if err := p.store.delete(evictedKey); err != nil { + p.logger.Warn("catalog popularity: failed to evict entry", zap.String("key", evictedKey), zap.Error(err)) } } } diff --git a/internal/registries/popularity_github_test.go b/internal/registries/popularity_github_test.go index 4d723b012..57440b5c0 100644 --- a/internal/registries/popularity_github_test.go +++ b/internal/registries/popularity_github_test.go @@ -6,6 +6,7 @@ import ( "net/http" "net/http/httptest" "strconv" + "strings" "sync" "sync/atomic" "testing" @@ -40,8 +41,10 @@ func TestGitHubStarsProvider_Fetch200StoresStarsAndSendsExactHeaders(t *testing. if got := gotHeaders.Get("X-GitHub-Api-Version"); got != "2022-11-28" { t.Errorf("X-GitHub-Api-Version header = %q", got) } - if got := gotHeaders.Get("User-Agent"); got == "" { - t.Error("expected a non-empty User-Agent") + // Compare against the exact versioned value: Go's default User-Agent is + // non-empty, so a bare non-empty check passes without the required header. + if got, want := gotHeaders.Get("User-Agent"), registryUserAgent(); got != want { + t.Errorf("User-Agent = %q, want %q", got, want) } if got := gotHeaders.Get("Authorization"); got != "" { t.Errorf("expected no Authorization header without MCPPROXY_GITHUB_TOKEN, got %q", got) @@ -533,3 +536,240 @@ func TestGitHubStarsProvider_KillSwitchNeverFetches(t *testing.T) { t.Fatalf("expected the kill switch to prevent any outbound request, got %d", got) } } + +// --- #1410: provider-level budget, TTL, env, breaker coverage --------------- + +// testClock is a mutable injected clock safe to read from worker goroutines. +type testClock struct { + mu sync.Mutex + t time.Time +} + +func newTestClock() *testClock { return &testClock{t: time.Unix(1_800_000_000, 0)} } + +func (c *testClock) Now() time.Time { + c.mu.Lock() + defer c.mu.Unlock() + return c.t +} + +func (c *testClock) Advance(d time.Duration) { + c.mu.Lock() + c.t = c.t.Add(d) + c.mu.Unlock() +} + +func countingServer(t *testing.T, handler http.HandlerFunc) (*httptest.Server, *int32) { + t.Helper() + var n int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + atomic.AddInt32(&n, 1) + handler(w, r) + })) + t.Cleanup(srv.Close) + return srv, &n +} + +func okStars(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(`{"stargazers_count": 7}`)) +} + +// TestGitHubStarsProvider_RollingBudgetGatesRequests proves the provider +// really stops issuing requests once the rolling budget is spent (not just +// that the rollingBudget helper counts), and resumes after the window rolls. +func TestGitHubStarsProvider_RollingBudgetGatesRequests(t *testing.T) { + srv, reqs := countingServer(t, okStars) + defer SetGitHubAPIBaseForTest(srv.URL)() + p := NewGitHubStarsProvider(PopularityOptions{}) + defer p.Close() + clock := newTestClock() + p.now = clock.Now + p.budget = newRollingBudget(2) + + p.Resolve(context.Background(), []string{"o/a", "o/b", "o/c", "o/d", "o/e"}, 2*time.Second) + if got := atomic.LoadInt32(reqs); got != 2 { + t.Fatalf("requests with a budget of 2 = %d, want exactly 2", got) + } + + // Budget spent: a further Resolve bails out without any request. + p.Resolve(context.Background(), []string{"o/f"}, 200*time.Millisecond) + if got := atomic.LoadInt32(reqs); got != 2 { + t.Fatalf("requests after the budget was spent = %d, want still 2", got) + } + + clock.Advance(61 * time.Minute) + p.Resolve(context.Background(), []string{"o/f"}, 2*time.Second) + if got := atomic.LoadInt32(reqs); got != 3 { + t.Fatalf("requests after the window rolled = %d, want 3", got) + } +} + +// TestGitHubStarsProvider_FreshEntriesDoNotRequeue pins that a fresh positive +// entry and a fresh negative (404) entry are not re-requested before their +// TTL, and that they are once it passes. +func TestGitHubStarsProvider_FreshEntriesDoNotRequeue(t *testing.T) { + srv, reqs := countingServer(t, func(w http.ResponseWriter, r *http.Request) { + if strings.HasSuffix(r.URL.Path, "/gone") { + w.WriteHeader(http.StatusNotFound) + return + } + okStars(w, r) + }) + defer SetGitHubAPIBaseForTest(srv.URL)() + p := NewGitHubStarsProvider(PopularityOptions{}) + defer p.Close() + clock := newTestClock() + p.now = clock.Now + + keys := []string{"o/live", "o/gone"} + p.Resolve(context.Background(), keys, 2*time.Second) + if got := atomic.LoadInt32(reqs); got != 2 { + t.Fatalf("initial requests = %d, want 2", got) + } + if _, state := p.Lookup("o/live"); state != LookupFresh { + t.Fatalf("o/live state = %d, want Fresh", state) + } + if _, state := p.Lookup("o/gone"); state != LookupNegative { + t.Fatalf("o/gone state = %d, want Negative", state) + } + + clock.Advance(githubFreshTTL - time.Minute) + for _, k := range keys { + if ch := p.enqueue(k); ch != nil { + t.Fatalf("enqueue(%s) admitted a fresh entry", k) + } + } + p.Resolve(context.Background(), keys, 200*time.Millisecond) + if got := atomic.LoadInt32(reqs); got != 2 { + t.Fatalf("requests before TTL expiry = %d, want still 2", got) + } + + clock.Advance(2 * time.Minute) + p.Resolve(context.Background(), keys, 2*time.Second) + if got := atomic.LoadInt32(reqs); got != 4 { + t.Fatalf("requests after TTL expiry = %d, want 4", got) + } +} + +// TestGitHubStarsProvider_GenericGitHubTokenIgnored pins FR-011: with only +// the generic GITHUB_TOKEN set, no Authorization header is sent and the +// unauthenticated budget is selected. +func TestGitHubStarsProvider_GenericGitHubTokenIgnored(t *testing.T) { + var gotAuth string + srv, _ := countingServer(t, func(w http.ResponseWriter, r *http.Request) { + gotAuth = r.Header.Get("Authorization") + okStars(w, r) + }) + defer SetGitHubAPIBaseForTest(srv.URL)() + t.Setenv("GITHUB_TOKEN", "generic-token") + t.Setenv("MCPPROXY_GITHUB_TOKEN", "") + p := NewGitHubStarsProvider(PopularityOptions{}) + defer p.Close() + + p.Resolve(context.Background(), []string{"o/r"}, 2*time.Second) + if _, state := p.Lookup("o/r"); state != LookupFresh { + t.Fatalf("expected the fetch to complete, got state=%d", state) + } + if gotAuth != "" { + t.Fatalf("Authorization = %q, want none when only GITHUB_TOKEN is set", gotAuth) + } + if p.budget.limit != githubRateLimitUnauth { + t.Fatalf("budget limit = %d, want the unauthenticated %d", p.budget.limit, githubRateLimitUnauth) + } +} + +// TestGitHubStarsProvider_ApplyBreakerTable covers the FR-009(d) pause +// computation and recovery once the injected clock passes the pause. +func TestGitHubStarsProvider_ApplyBreakerTable(t *testing.T) { + t.Setenv("MCPPROXY_CATALOG_POPULARITY", "false") + base := time.Unix(1_800_000_000, 0) + reset := base.Add(30 * time.Minute) + resetStr := strconv.FormatInt(reset.Unix(), 10) + httpDate := base.Add(90 * time.Second).UTC().Format(http.TimeFormat) + + tests := []struct { + name string + status int + headers rateLimitHeaders + wantPause time.Duration // 0 = not paused + }{ + {"429 Retry-After seconds", 429, rateLimitHeaders{retryAfter: "120"}, 120 * time.Second}, + {"429 Retry-After HTTP-date", 429, rateLimitHeaders{retryAfter: httpDate}, 90 * time.Second}, + {"429 Retry-After wins over reset", 429, rateLimitHeaders{retryAfter: "10", reset: resetStr}, 10 * time.Second}, + {"403 reset only", 403, rateLimitHeaders{reset: resetStr}, 30 * time.Minute}, + {"403 no headers uses default", 403, rateLimitHeaders{}, githubBreakerPause}, + {"429 bad Retry-After falls back to default", 429, rateLimitHeaders{retryAfter: "soon"}, githubBreakerPause}, + {"200 remaining at threshold", 200, rateLimitHeaders{remaining: "5", reset: resetStr}, 30 * time.Minute}, + {"200 remaining at threshold no reset", 200, rateLimitHeaders{remaining: "5"}, githubBreakerPause}, + {"200 remaining above threshold", 200, rateLimitHeaders{remaining: "6", reset: resetStr}, 0}, + {"200 no rate-limit headers", 200, rateLimitHeaders{}, 0}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + p := NewGitHubStarsProvider(PopularityOptions{}) + defer p.Close() + p.now = func() time.Time { return base } + + p.applyBreaker(tc.status, tc.headers) + + if tc.wantPause == 0 { + if p.pausedNow(base) { + t.Fatal("breaker engaged, want it idle") + } + return + } + if !p.pausedNow(base) { + t.Fatal("breaker idle, want it engaged") + } + if !p.pausedNow(base.Add(tc.wantPause - time.Second)) { + t.Errorf("breaker released before the %s pause elapsed", tc.wantPause) + } + if p.pausedNow(base.Add(tc.wantPause)) { + t.Errorf("breaker still engaged once the %s pause elapsed", tc.wantPause) + } + }) + } +} + +// TestGitHubStarsProvider_429PausesThenRecovers is the end-to-end 429 case: a +// 429 with Retry-After blocks further requests until the injected clock +// passes the pause, after which fetching resumes. +func TestGitHubStarsProvider_429PausesThenRecovers(t *testing.T) { + var limited atomic.Bool + limited.Store(true) + srv, reqs := countingServer(t, func(w http.ResponseWriter, r *http.Request) { + if limited.Load() { + w.Header().Set("Retry-After", "120") + w.WriteHeader(http.StatusTooManyRequests) + return + } + okStars(w, r) + }) + defer SetGitHubAPIBaseForTest(srv.URL)() + p := NewGitHubStarsProvider(PopularityOptions{}) + defer p.Close() + clock := newTestClock() + p.now = clock.Now + + p.Resolve(context.Background(), []string{"o/r1"}, 2*time.Second) + if got := atomic.LoadInt32(reqs); got != 1 { + t.Fatalf("requests before the pause = %d, want 1", got) + } + + clock.Advance(60 * time.Second) + p.Resolve(context.Background(), []string{"o/r2"}, 200*time.Millisecond) + if got := atomic.LoadInt32(reqs); got != 1 { + t.Fatalf("requests during the Retry-After pause = %d, want still 1", got) + } + + limited.Store(false) + clock.Advance(61 * time.Second) // 121s total: past Retry-After + p.Resolve(context.Background(), []string{"o/r2"}, 2*time.Second) + if got := atomic.LoadInt32(reqs); got != 2 { + t.Fatalf("requests after the pause elapsed = %d, want 2", got) + } + if stars, state := p.Lookup("o/r2"); state != LookupFresh || stars != 7 { + t.Fatalf("expected recovery to fetch o/r2 (Fresh/7), got state=%d stars=%d", state, stars) + } +} diff --git a/internal/registries/popularity_store.go b/internal/registries/popularity_store.go index 12e4efe52..2a6e9e8ee 100644 --- a/internal/registries/popularity_store.go +++ b/internal/registries/popularity_store.go @@ -113,3 +113,24 @@ func (s *popularityStore) delete(key string) error { return b.Delete([]byte(key)) }) } + +// deleteMany removes a set of keys in one bbolt write transaction. Startup +// uses this when trimming a legacy bucket that exceeds the in-memory cap, so +// a large cleanup incurs one fsync instead of one per evicted record. +func (s *popularityStore) deleteMany(keys []string) error { + if len(keys) == 0 { + return nil + } + return s.db.Update(func(tx *bbolt.Tx) error { + b := tx.Bucket([]byte(popularityBucketName)) + if b == nil { + return nil + } + for _, key := range keys { + if err := b.Delete([]byte(key)); err != nil { + return err + } + } + return nil + }) +} diff --git a/internal/registries/popularity_store_test.go b/internal/registries/popularity_store_test.go index ca93fb462..33254a95e 100644 --- a/internal/registries/popularity_store_test.go +++ b/internal/registries/popularity_store_test.go @@ -2,7 +2,9 @@ package registries import ( "encoding/json" + "fmt" "path/filepath" + "sync" "testing" "time" @@ -62,11 +64,11 @@ func TestPopularityStore_Delete(t *testing.T) { } } -// TestGitHubStarsProvider_LazyLoadFromStore pins that a provider restart -// (fresh githubStarsProvider over the SAME bbolt db) sees a previously -// fetched entry via Lookup with NO fetch — the lazy-load-on-miss path -// (plan.md data model) survives a restart. -func TestGitHubStarsProvider_LazyLoadFromStore(t *testing.T) { +// TestGitHubStarsProvider_LoadsPersistedEntriesOnStartup pins that a provider +// restart (fresh githubStarsProvider over the SAME bbolt db) eagerly preloads +// every persisted entry, so Lookup serves a previously fetched value with NO +// fetch and never reads bbolt afterwards. +func TestGitHubStarsProvider_LoadsPersistedEntriesOnStartup(t *testing.T) { db := openTempPopularityDB(t) store, err := newPopularityStore(db) if err != nil { @@ -88,6 +90,15 @@ func TestGitHubStarsProvider_LazyLoadFromStore(t *testing.T) { if stars != 77 { t.Fatalf("expected stars=77 from the pre-seeded store, got %d", stars) } + + // The preload is a one-time snapshot: an entry written to the store after + // construction is not visible (memory is authoritative). + if err := store.put("o/late", &starsEntry{Stars: 5, Status: 200, FetchedAt: fetchedAt}); err != nil { + t.Fatalf("late put: %v", err) + } + if _, state := provider.Lookup("o/late"); state != LookupAbsent { + t.Fatalf("expected a post-construction store write to stay invisible, got state=%d", state) + } } // TestGitHubStarsProvider_CapEviction pins FR-008's 5000-key cap: inserting @@ -137,7 +148,10 @@ func TestGitHubStarsProvider_CapHoldsAcrossRestart(t *testing.T) { if _, err := newPopularityStore(db); err != nil { t.Fatalf("newPopularityStore: %v", err) } - const over = 3 + // This is large enough to exercise trimming a legacy bucket with many + // surplus records without making the test itself expensive. The provider + // should remove the oldest records in one pass and one store transaction. + const over = 5_000 base := time.Now().Add(-time.Hour) if err := db.Update(func(tx *bbolt.Tx) error { b := tx.Bucket([]byte(popularityBucketName)) @@ -171,3 +185,185 @@ func TestGitHubStarsProvider_CapHoldsAcrossRestart(t *testing.T) { t.Errorf("newest key should survive, got state %d", state) } } + +// holdStoreWriter blocks every other bbolt writer (store.put/delete) until +// the returned release func is called. +func holdStoreWriter(t *testing.T, db *bbolt.DB) (release func()) { + t.Helper() + held := make(chan struct{}) + done := make(chan struct{}) + rel := make(chan struct{}) + go func() { + defer close(done) + _ = db.Update(func(_ *bbolt.Tx) error { + close(held) + <-rel + return nil + }) + }() + <-held + var once sync.Once + release = func() { once.Do(func() { close(rel); <-done }) } + t.Cleanup(release) + return release +} + +// TestGitHubStarsProvider_EvictionDoesNotBlockLookupOnStore pins that cap +// eviction never holds the provider mutex across the bbolt delete: with the +// store's writer lock held (so the delete cannot finish), Lookup must still +// return promptly and observe the eviction in memory. +func TestGitHubStarsProvider_EvictionDoesNotBlockLookupOnStore(t *testing.T) { + db := openTempPopularityDB(t) + t.Setenv("MCPPROXY_CATALOG_POPULARITY", "false") + provider := NewGitHubStarsProvider(PopularityOptions{DB: db}) + defer provider.Close() + + base := time.Now().Add(-time.Hour) + provider.mu.Lock() + for i := 0; i < githubMaxCacheKeys; i++ { + provider.entries[keyForIndex(i)] = &starsEntry{Stars: 1, Status: 200, FetchedAt: base.Add(time.Duration(i) * time.Second)} + } + provider.mu.Unlock() + // Mirror the to-be-evicted entry into the store so the delete has work. + if err := provider.store.put(keyForIndex(0), provider.entries[keyForIndex(0)]); err != nil { + t.Fatalf("seed put: %v", err) + } + + release := holdStoreWriter(t, db) + applied := make(chan struct{}) + go func() { + defer close(applied) + provider.applyResult("new/key", nil, 200, 5, "", rateLimitHeaders{}, nil) + }() + + oldest := keyForIndex(0) + deadline := time.After(3 * time.Second) + observed := make(chan struct{}) + go func() { + defer close(observed) + for { + if _, state := provider.Lookup(oldest); state == LookupAbsent { + return + } + time.Sleep(time.Millisecond) + } + }() + select { + case <-observed: + case <-deadline: + t.Fatal("Lookup blocked (or eviction never became visible) while the store delete was pending") + } + select { + case <-applied: + t.Fatal("applyResult finished although the store writer lock was held") + default: + } + + release() + <-applied + if _, ok := provider.store.get(oldest); ok { + t.Fatal("expected the evicted entry to be deleted from the store once it unblocked") + } + if _, ok := provider.store.get("new/key"); !ok { + t.Fatal("expected the new entry to be persisted") + } + provider.mu.Lock() + count := len(provider.entries) + provider.mu.Unlock() + if count != githubMaxCacheKeys { + t.Fatalf("entry count = %d, want %d", count, githubMaxCacheKeys) + } +} + +// TestGitHubStarsProvider_PersistGuardsAgainstEvictionRaces pins the ordering +// guards: a delayed put never resurrects an evicted key, and a delayed delete +// never drops a key that a newer fetch re-added. +func TestGitHubStarsProvider_PersistGuardsAgainstEvictionRaces(t *testing.T) { + db := openTempPopularityDB(t) + t.Setenv("MCPPROXY_CATALOG_POPULARITY", "false") + provider := NewGitHubStarsProvider(PopularityOptions{DB: db}) + defer provider.Close() + + // Stale put: the entry was evicted from memory before its put ran. + stale := &starsEntry{Stars: 1, Status: 200, FetchedAt: time.Now()} + provider.persist("gone/key", stale, "") + if _, ok := provider.store.get("gone/key"); ok { + t.Fatal("a put for an entry no longer in memory must be skipped") + } + + // Stale delete: the key was re-added before the delete ran. + fresh := &starsEntry{Stars: 2, Status: 200, FetchedAt: time.Now()} + provider.mu.Lock() + provider.entries["back/key"] = fresh + provider.mu.Unlock() + provider.persist("back/key", fresh, "") + provider.persist("", nil, "back/key") + if _, ok := provider.store.get("back/key"); !ok { + t.Fatal("a delete for a key that is back in memory must be skipped") + } +} + +// TestGitHubStarsProvider_ConcurrentApplyAndLookupUnderCap is a -race +// regression test: concurrent applyResult (evicting past the cap) and Lookup +// keep memory at the cap and leave the store consistent with memory. +func TestGitHubStarsProvider_ConcurrentApplyAndLookupUnderCap(t *testing.T) { + db := openTempPopularityDB(t) + t.Setenv("MCPPROXY_CATALOG_POPULARITY", "false") + provider := NewGitHubStarsProvider(PopularityOptions{DB: db}) + defer provider.Close() + + base := time.Now().Add(-time.Hour) + provider.mu.Lock() + for i := 0; i < githubMaxCacheKeys; i++ { + provider.entries[keyForIndex(i)] = &starsEntry{Stars: 1, Status: 200, FetchedAt: base.Add(time.Duration(i) * time.Second)} + } + provider.mu.Unlock() + + var wg sync.WaitGroup + stop := make(chan struct{}) + for g := 0; g < 4; g++ { + wg.Add(1) + go func() { + defer wg.Done() + for { + select { + case <-stop: + return + default: + provider.Lookup(keyForIndex(0)) + } + } + }() + } + var writers sync.WaitGroup + for g := 0; g < 4; g++ { + writers.Add(1) + go func(g int) { + defer writers.Done() + for i := 0; i < 10; i++ { + provider.applyResult(fmt.Sprintf("w%d/repo%d", g, i), nil, 200, 5, "", rateLimitHeaders{}, nil) + } + }(g) + } + writers.Wait() + close(stop) + wg.Wait() + + provider.mu.Lock() + count := len(provider.entries) + mem := make(map[string]struct{}, count) + for k := range provider.entries { + mem[k] = struct{}{} + } + provider.mu.Unlock() + if count != githubMaxCacheKeys { + t.Fatalf("entry count = %d, want %d", count, githubMaxCacheKeys) + } + // Only the 40 written keys were persisted (the seed was memory-only), so + // every persisted key must still be in memory: no resurrected evictions. + for k := range provider.store.all() { + if _, ok := mem[k]; !ok { + t.Errorf("store holds %q which is no longer in memory", k) + } + } +}