Skip to content
Merged
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
112 changes: 77 additions & 35 deletions internal/registries/popularity_github.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (
"net/http"
"net/url"
"os"
"sort"
"strconv"
"strings"
"sync"
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand All @@ -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))
}
}
}
Expand Down
Loading
Loading