From ba4fd9d9fd5c4318253b5061ebfcd48dc1fab4ec Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 09:25:17 -0400 Subject: [PATCH 01/14] feat(sections): add trending_discover home section MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A library-agnostic home section that surfaces external global trending (TMDB or Trakt, admin-selectable) mixing movies + series, matched to titles in the viewer's enabled libraries. TMDB uses /trending/all/{window} (natively mixed); Trakt merges trending movies + shows. Fetched live with a 1h in-process cache, so no background job or stored collection — and no per-library duplication. Appears in the admin section gallery via its recipe presets (TMDB Trending Today/This Week, Trakt Trending); featured -> hero via the existing flag. Co-Authored-By: Claude Opus 4.8 (1M context) --- internal/api/router.go | 6 + internal/sections/fetcher.go | 238 +++++++++++++++++- .../sections/recipes/trending_discover.go | 61 +++++ internal/sections/trending_discover_test.go | 67 +++++ internal/sections/types.go | 3 + 5 files changed, 372 insertions(+), 3 deletions(-) create mode 100644 internal/sections/recipes/trending_discover.go create mode 100644 internal/sections/trending_discover_test.go diff --git a/internal/api/router.go b/internal/api/router.go index 75e38dfa..9125eedf 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -878,6 +878,12 @@ func NewRouter(deps Dependencies) chi.Router { } } + // Wire external-trending fetchers into the section fetcher so the + // trending_discover home section can pull TMDB/Trakt trending. + sectionFetcher.ItemRepo = itemRepo + sectionFetcher.TMDBTrending = libraryCollectionService.TMDBCollections + sectionFetcher.TraktTrending = libraryCollectionService.TraktCollections + libraryCollectionHandler = handlers.NewLibraryCollectionHandler( libraryCollectionRepo, libraryCollectionService, diff --git a/internal/sections/fetcher.go b/internal/sections/fetcher.go index 9a771d7e..df0168e1 100644 --- a/internal/sections/fetcher.go +++ b/internal/sections/fetcher.go @@ -67,9 +67,18 @@ type Fetcher struct { RecommendationRepo *recommendations.Repo // retained for non-reader call sites RecommendationReader recommendationReader NextUpRepo *catalog.NextUpRepository - candidateCacheMu sync.Mutex - candidateCache *editorialCandidateCache - candidateGroup singleflight.Group + + // ItemRepo resolves external IDs (TMDB/Trakt) to library content IDs for + // the trending_discover section. Nil disables external-trending matching. + ItemRepo *catalog.ItemRepository + // TMDBTrending and TraktTrending fetch external global trending lists. Each + // is nil when that provider is not configured. + TMDBTrending catalog.TMDBCollectionFetcher + TraktTrending catalog.TraktCollectionFetcher + + candidateCacheMu sync.Mutex + candidateCache *editorialCandidateCache + candidateGroup singleflight.Group // Clock returns the current time. Defaults to recipes.RealClock{}. // Tests inject recipes.FixedClock for deterministic seasonal/editorial behavior. @@ -1070,6 +1079,8 @@ func (f *Fetcher) fetchSection(ctx context.Context, s ResolvedSection, libraryID return f.fetchNewToLibrary(ctx, s, libraryID, libraryIDs, filter) case SectionMostWatched: return f.fetchMostWatched(ctx, s, libraryID, libraryIDs, filter) + case SectionTrendingDiscover: + return f.fetchTrendingDiscover(ctx, s, libraryID, libraryIDs, filter) case SectionAdminCuratedList: return f.fetchAdminCuratedList(ctx, s, libraryID, libraryIDs, filter) default: @@ -2334,6 +2345,227 @@ func scanMediaItems(rows pgx.Rows) ([]*models.MediaItem, error) { return items, rows.Err() } +// trendingDiscoverEntry is a provider-agnostic external trending result. +type trendingDiscoverEntry struct { + tmdbID string + imdbID string + tvdbID string + mediaType string // "movie" | "tv" +} + +// fetchTrendingDiscover surfaces external global trending (TMDB or Trakt), +// mixing movies + series, matched to titles in the viewer's enabled libraries. +// The external fetch + ID resolution are cached briefly so it does not hit the +// upstream API on every home-page load. +func (f *Fetcher) fetchTrendingDiscover(ctx context.Context, s ResolvedSection, libraryID *int, libraryIDs []int, filter catalog.AccessFilter) ([]*models.MediaItem, int, error) { + var p recipes.TrendingDiscoverParams + if len(s.Config) > 0 { + _ = json.Unmarshal(s.Config, &p) + } + source := p.Source + if source != "trakt" { + source = "tmdb" + } + window := p.Window + if window != "day" { + window = "week" + } + + limit := s.ItemLimit + if limit <= 0 { + limit = 20 + } + // Over-fetch: library-only matching drops globally-trending titles the + // server does not own, so request more candidates than the display limit. + fetchLimit := limit * 5 + if fetchLimit < 50 { + fetchLimit = 50 + } + if fetchLimit > 200 { + fetchLimit = 200 + } + + orderedIDs, err := f.loadTrendingDiscoverContentIDs(ctx, source, window, fetchLimit) + if err != nil { + return nil, 0, err + } + if len(orderedIDs) == 0 { + return []*models.MediaItem{}, 0, nil + } + + items, err := f.fetchItemsByContentIDs(ctx, orderedIDs, libraryID, libraryIDs, filter) + if err != nil { + return nil, 0, err + } + + // Re-order to trending rank (fetchItemsByContentIDs returns DB order) and + // truncate to the section's display limit. + byID := make(map[string]*models.MediaItem, len(items)) + for _, item := range items { + byID[item.ContentID] = item + } + ordered := make([]*models.MediaItem, 0, len(orderedIDs)) + for _, id := range orderedIDs { + item, ok := byID[id] + if !ok { + continue + } + ordered = append(ordered, item) + if len(ordered) >= limit { + break + } + } + return ordered, len(ordered), nil +} + +// loadTrendingDiscoverContentIDs fetches the external trending list and resolves +// it to ordered library content IDs, cached for an hour per (source, window). +func (f *Fetcher) loadTrendingDiscoverContentIDs(ctx context.Context, source, window string, fetchLimit int) ([]string, error) { + cache := f.ensureEditorialCandidateCache() + cacheKey := fmt.Sprintf("trending_discover|%s|%s|%d", source, window, fetchLimit) + now := f.Clock.Now() + if cached, ok := cache.get(cacheKey, now); ok { + return cached, nil + } + + entries, err := f.fetchTrendingDiscoverEntries(ctx, source, window, fetchLimit) + if err != nil { + return nil, err + } + // Don't cache when the provider is unconfigured/empty, so a newly-configured + // provider takes effect immediately rather than after the TTL. + if len(entries) == 0 { + return []string{}, nil + } + + contentIDs, err := f.resolveTrendingDiscoverIDs(ctx, entries) + if err != nil { + return nil, err + } + cache.set(cacheKey, contentIDs, now.Add(time.Hour)) + return contentIDs, nil +} + +// fetchTrendingDiscoverEntries pulls the raw trending list from the configured +// provider. A nil/unconfigured provider yields an empty list (no error) so the +// section simply renders empty. +func (f *Fetcher) fetchTrendingDiscoverEntries(ctx context.Context, source, window string, fetchLimit int) ([]trendingDiscoverEntry, error) { + if source == "trakt" { + if f.TraktTrending == nil { + return nil, nil + } + // Trakt has no mixed endpoint; fetch movies + shows and concatenate. + movies, movieErr := f.TraktTrending.GetCollectionPreset(ctx, "trending", "movie", fetchLimit, "") + shows, showErr := f.TraktTrending.GetCollectionPreset(ctx, "trending", "tv", fetchLimit, "") + if movieErr != nil && showErr != nil { + return nil, fmt.Errorf("trakt trending: %v / %v", movieErr, showErr) + } + merged := append(append([]catalog.TraktCollectionEntry{}, movies...), shows...) + out := make([]trendingDiscoverEntry, 0, len(merged)) + for _, e := range merged { + entry := trendingDiscoverEntry{imdbID: e.IMDbID, mediaType: e.MediaType} + if e.TMDBID > 0 { + entry.tmdbID = strconv.Itoa(e.TMDBID) + } + if e.TVDBID > 0 { + entry.tvdbID = strconv.Itoa(e.TVDBID) + } + out = append(out, entry) + } + return out, nil + } + + // Default: TMDB /trending/all/{window} — natively mixed movies + series. + if f.TMDBTrending == nil { + return nil, nil + } + entries, err := f.TMDBTrending.GetCollectionPreset(ctx, "trending", "all", window, fetchLimit) + if err != nil { + return nil, err + } + out := make([]trendingDiscoverEntry, 0, len(entries)) + for _, e := range entries { + entry := trendingDiscoverEntry{imdbID: e.IMDbID, mediaType: e.MediaType} + if e.ID > 0 { + entry.tmdbID = strconv.Itoa(e.ID) + } + if e.TVDBID > 0 { + entry.tvdbID = strconv.Itoa(e.TVDBID) + } + out = append(out, entry) + } + return out, nil +} + +// resolveTrendingDiscoverIDs matches trending entries to library content IDs via +// two batched external-ID lookups (movies, series), preserving trending order. +func (f *Fetcher) resolveTrendingDiscoverIDs(ctx context.Context, entries []trendingDiscoverEntry) ([]string, error) { + if f.ItemRepo == nil { + return nil, fmt.Errorf("trending_discover: item repository not configured") + } + var movieBatch, seriesBatch catalog.ExternalIDBatch + for _, e := range entries { + batch := &movieBatch + if e.mediaType == "tv" { + batch = &seriesBatch + } + if e.tmdbID != "" { + batch.TMDBIDs = append(batch.TMDBIDs, e.tmdbID) + } + if e.imdbID != "" { + batch.IMDbIDs = append(batch.IMDbIDs, e.imdbID) + } + if e.tvdbID != "" { + batch.TVDBIDs = append(batch.TVDBIDs, e.tvdbID) + } + } + movieLookup, err := f.ItemRepo.GetByExternalIDs(ctx, movieBatch, "movie") + if err != nil { + return nil, err + } + seriesLookup, err := f.ItemRepo.GetByExternalIDs(ctx, seriesBatch, "series") + if err != nil { + return nil, err + } + return orderedTrendingContentIDs(entries, movieLookup, seriesLookup), nil +} + +// orderedTrendingContentIDs maps trending entries to library content IDs in +// trending order, preferring TVDB (series) > TMDB > IMDb, de-duplicated. +func orderedTrendingContentIDs(entries []trendingDiscoverEntry, movieLookup, seriesLookup *catalog.ExternalIDLookup) []string { + seen := make(map[string]struct{}, len(entries)) + out := make([]string, 0, len(entries)) + for _, e := range entries { + lookup := movieLookup + isSeries := e.mediaType == "tv" + if isSeries { + lookup = seriesLookup + } + if lookup == nil { + continue + } + var id string + if isSeries && e.tvdbID != "" { + id = lookup.ByTVDB[e.tvdbID] + } + if id == "" && e.tmdbID != "" { + id = lookup.ByTMDB[e.tmdbID] + } + if id == "" && e.imdbID != "" { + id = lookup.ByIMDb[e.imdbID] + } + if id == "" { + continue + } + if _, dup := seen[id]; dup { + continue + } + seen[id] = struct{}{} + out = append(out, id) + } + return out +} + func (f *Fetcher) fetchTrending(ctx context.Context, s ResolvedSection, libraryID *int, libraryIDs []int, filter catalog.AccessFilter) ([]*models.MediaItem, int, error) { var p recipes.TrendingParams if len(s.Config) > 0 { diff --git a/internal/sections/recipes/trending_discover.go b/internal/sections/recipes/trending_discover.go new file mode 100644 index 00000000..c5fba9cb --- /dev/null +++ b/internal/sections/recipes/trending_discover.go @@ -0,0 +1,61 @@ +package recipes + +import ( + "encoding/json" + "errors" + "time" +) + +// TrendingDiscoverParams configures the trending_discover section: external +// global trending pulled from a single source (TMDB or Trakt), mixing movies + +// series in one list, matched to titles already in the library. +type TrendingDiscoverParams struct { + Source string `json:"source"` // "tmdb" | "trakt" + Window string `json:"window"` // "day" | "week" (TMDB only; ignored by Trakt) +} + +type trendingDiscoverRecipe struct{} + +func (trendingDiscoverRecipe) Type() string { return "trending_discover" } +func (trendingDiscoverRecipe) NewParams() any { return &TrendingDiscoverParams{} } +func (trendingDiscoverRecipe) DefaultCacheTTL() time.Duration { return time.Hour } +func (trendingDiscoverRecipe) Resolve(rc ResolverContext) (ResolvedItems, error) { + return delegateResolve("trending_discover", rc) +} + +func (trendingDiscoverRecipe) Validate(raw json.RawMessage) error { + if len(raw) == 0 { + return nil + } + var p TrendingDiscoverParams + if err := json.Unmarshal(raw, &p); err != nil { + return err + } + switch p.Source { + case "", "tmdb", "trakt": + default: + return errors.New(`trending_discover: source must be "tmdb" or "trakt"`) + } + switch p.Window { + case "", "day", "week": + default: + return errors.New(`trending_discover: window must be "day" or "week"`) + } + return nil +} + +func (trendingDiscoverRecipe) Definition() RecipeDefinition { + return RecipeDefinition{ + Type: "trending_discover", + Category: CategorySocial, + Presets: []GalleryPreset{ + {Key: "tdisc_tmdb_day", DisplayName: "TMDB Trending Today", Icon: "🔥", DescriptionShort: "Today's trending movies & shows from TMDB, matched to your library.", DefaultParams: json.RawMessage(`{"source":"tmdb","window":"day"}`)}, + {Key: "tdisc_tmdb_week", DisplayName: "TMDB Trending This Week", Icon: "🔥", DescriptionShort: "This week's trending movies & shows from TMDB, matched to your library.", DefaultParams: json.RawMessage(`{"source":"tmdb","window":"week"}`)}, + {Key: "tdisc_trakt", DisplayName: "Trakt Trending", Icon: "📈", DescriptionShort: "Trending movies & shows on Trakt, matched to your library.", DefaultParams: json.RawMessage(`{"source":"trakt","window":"week"}`)}, + }, + } +} + +func init() { + Register(trendingDiscoverRecipe{}) +} diff --git a/internal/sections/trending_discover_test.go b/internal/sections/trending_discover_test.go new file mode 100644 index 00000000..778a922a --- /dev/null +++ b/internal/sections/trending_discover_test.go @@ -0,0 +1,67 @@ +package sections + +import ( + "testing" + + "github.com/Silo-Server/silo-server/internal/catalog" +) + +func TestOrderedTrendingContentIDs_PreservesOrderAndSeriesPrefersTVDB(t *testing.T) { + entries := []trendingDiscoverEntry{ + {tmdbID: "1", mediaType: "movie"}, + {tmdbID: "2", tvdbID: "20", mediaType: "tv"}, + {imdbID: "tt3", mediaType: "movie"}, + } + movieLookup := &catalog.ExternalIDLookup{ + ByTMDB: map[string]string{"1": "cm1"}, + ByIMDb: map[string]string{"tt3": "cm3"}, + ByTVDB: map[string]string{}, + } + seriesLookup := &catalog.ExternalIDLookup{ + ByTVDB: map[string]string{"20": "cs2"}, + ByTMDB: map[string]string{"2": "cs2_tmdb"}, // TVDB should win for series + ByIMDb: map[string]string{}, + } + got := orderedTrendingContentIDs(entries, movieLookup, seriesLookup) + want := []string{"cm1", "cs2", "cm3"} + if len(got) != len(want) { + t.Fatalf("got %v, want %v", got, want) + } + for i := range want { + if got[i] != want[i] { + t.Fatalf("pos %d = %q want %q (full %v)", i, got[i], want[i], got) + } + } +} + +func TestOrderedTrendingContentIDs_SkipsUnmatchedAndDedups(t *testing.T) { + entries := []trendingDiscoverEntry{ + {tmdbID: "1", mediaType: "movie"}, // matches cX + {tmdbID: "404", mediaType: "movie"}, // no match -> skipped + {imdbID: "ttX", mediaType: "movie"}, // also resolves to cX -> deduped + } + movieLookup := &catalog.ExternalIDLookup{ + ByTMDB: map[string]string{"1": "cX"}, + ByIMDb: map[string]string{"ttX": "cX"}, + ByTVDB: map[string]string{}, + } + got := orderedTrendingContentIDs(entries, movieLookup, &catalog.ExternalIDLookup{}) + if len(got) != 1 || got[0] != "cX" { + t.Fatalf("expected [cX], got %v", got) + } +} + +func TestOrderedTrendingContentIDs_MovieIgnoresTVDB(t *testing.T) { + entries := []trendingDiscoverEntry{ + {tvdbID: "50", tmdbID: "5", mediaType: "movie"}, + } + movieLookup := &catalog.ExternalIDLookup{ + ByTVDB: map[string]string{"50": "cTV"}, // must be ignored for movies + ByTMDB: map[string]string{"5": "cTMDB"}, + ByIMDb: map[string]string{}, + } + got := orderedTrendingContentIDs(entries, movieLookup, &catalog.ExternalIDLookup{}) + if len(got) != 1 || got[0] != "cTMDB" { + t.Fatalf("movie should match TMDB not TVDB; got %v", got) + } +} diff --git a/internal/sections/types.go b/internal/sections/types.go index 10c8d8ec..6506f9b3 100644 --- a/internal/sections/types.go +++ b/internal/sections/types.go @@ -40,6 +40,8 @@ const ( SectionNewToLibrary SectionType = "new_to_library" SectionMostWatched SectionType = "most_watched" + SectionTrendingDiscover SectionType = "trending_discover" + SectionAdminCuratedList SectionType = "admin_curated_list" ) @@ -71,6 +73,7 @@ var ValidSectionTypes = map[SectionType]bool{ SectionProfileActivityFeed: true, SectionNewToLibrary: true, SectionMostWatched: true, + SectionTrendingDiscover: true, SectionAdminCuratedList: true, } From 7992470cc3750910a073e18542cce1b33c51acde Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:25:31 -0400 Subject: [PATCH 02/14] docs: add trending_discover persistent snapshot design Replace the in-process 1h trending cache with a background-refreshed, persisted snapshot for reliability under upstream failure and sync-run observability. Co-Authored-By: Claude Opus 4.8 (1M context) --- ...ing-discover-persistent-snapshot-design.md | 206 ++++++++++++++++++ 1 file changed, 206 insertions(+) create mode 100644 docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md diff --git a/docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md b/docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md new file mode 100644 index 00000000..84479bba --- /dev/null +++ b/docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md @@ -0,0 +1,206 @@ +# Trending Discover — Persistent Snapshot Design + +Status: Approved (design) +Date: 2026-05-29 + +Commands assume the repository root is the cwd. + +## Problem + +The `trending_discover` home section pulls global trending (TMDB or Trakt), +matches it to titles in the viewer's libraries, and renders the result. Today +the external fetch + external-ID resolution are memoised in an **in-process +map** (`editorialCandidateCache`) keyed by `source|window|fetchLimit` with a +1-hour TTL and a `singleflight` group to collapse concurrent misses +(`internal/sections/fetcher.go`, `loadTrendingDiscoverContentIDs`). + +That cache has three weaknesses, two of which we explicitly want to fix: + +1. **No resilience to upstream failure.** At TTL expiry a request still blocks + on the TMDB/Trakt call, and a slow or down provider degrades the home page. + There is no concept of "serve the last good list." +2. **No observability.** Nothing records when trending last refreshed, whether + it succeeded, or how many trending titles matched the catalog. +3. (Secondary) It is per-process and volatile — lost on restart, and each + server instance fetches independently. + +## Goals + +- Reads of the trending section **never call the upstream provider** and never + block on it. They serve a persisted, last-good list. +- A persisted, inspectable record with refresh history (the same "synced + record" shape used by library collections). + +## Non-goals + +- Reusing the library-collection machinery. `library_collections.library_id` is + `NOT NULL` (collections are library-scoped); trending is deliberately + library-agnostic and appears once. Making collections library-agnostic is a + large, separate change and is out of scope. +- An admin UI for the snapshot. The persisted columns make this possible later, + but it is not part of this work. +- A shared/distributed cache (Redis). The DB is the single source of truth. + +## Approach + +Replace the lazy in-process cache with a **background-refreshed persistent +snapshot**. A scheduled task refreshes the trending list on an interval; the +read path only ever reads the persisted snapshot. This mirrors the +collection-sync *pattern* (`CollectionSyncScheduler` + `SyncCollectionsTask`) +without reusing its library-scoped storage. + +The change cleanly separates **read** from **refresh**: + +- The refresh path owns the upstream call and external-ID resolution. +- The read path owns only a primary-key snapshot lookup plus the existing + per-viewer access filtering. + +## Data model + +New migration `166_trending_discover_snapshots` (next free number is 166; 165 is +the current max). + +``` +trending_discover_snapshots + source text NOT NULL -- 'tmdb' | 'trakt' + window text NOT NULL -- 'day' | 'week' (trakt canonicalized to 'week') + content_ids text[] NOT NULL DEFAULT '{}' -- ordered, resolved to library catalog, capped at 200 + entry_count int NOT NULL DEFAULT 0 -- raw provider entries fetched (matched-vs-trending observability) + refreshed_at timestamptz -- last successful refresh + last_attempt_at timestamptz + last_status text -- 'ok' | 'empty' | 'error' + last_error text + PRIMARY KEY (source, window) +``` + +Notes: + +- One row per canonical `(source, window)` — at most a handful of rows. +- `content_ids` is catalog-matched but **viewer-agnostic**, exactly like the + value the current cache holds. Per-viewer access filtering still happens at + read time in `fetchItemsByContentIDs`, so a single row serves all viewers. +- The list is stored at the over-fetch cap (200). The read path truncates to the + section's `ItemLimit`, decoupling the snapshot from per-section limits (the old + `fetchLimit` cache-key dimension goes away). +- **Reliability invariant:** a failed refresh updates `last_attempt_at`, + `last_status`, and `last_error` but **never clears `content_ids`**. The last + good list keeps serving through an upstream outage. + +The paired `.down.sql` drops the table. + +## Components + +All new server code lives in `internal/`. + +### `TrendingSnapshotRepository` (new, `internal/sections/`) + +Thin repository over the new table: + +- `Get(ctx, source, window) (TrendingSnapshot, error)` — PK lookup. +- `Upsert(ctx, snapshot) error` — insert/update on `(source, window)`. +- `ListAll(ctx) ([]TrendingSnapshot, error)` — for future admin/inspection use + and tests. + +`Upsert` distinguishes a successful refresh (writes `content_ids`, +`refreshed_at`, `entry_count`, `last_status='ok'|'empty'`) from a failed one +(touches only `last_attempt_at`/`last_status='error'`/`last_error`, leaving +`content_ids` intact). + +### `TrendingRefresher` (new, `internal/sections/trending_refresher.go`) + +Owns the upstream call. Holds `TMDBTrending`, `TraktTrending`, `ItemRepo`, the +section repository (to enumerate used combos), and the snapshot repository. The +existing `fetchTrendingDiscoverEntries` and `resolveTrendingDiscoverIDs` logic +**moves here** from `fetcher.go`. + +`RunOnce(ctx) (json.RawMessage, error)`: + +1. Enumerate the distinct canonical `(source, window)` pairs from **enabled** + `trending_discover` sections across all scopes (parse each section's + `TrendingDiscoverParams`, apply the same source/window normalization the + section uses, collapse Trakt to `week`). +2. For each pair: fetch the cap-200 list from the provider, resolve external IDs + to library content IDs, and `Upsert`. +3. Return a JSON summary (e.g. `{combos, refreshed, empty, failed}`), mirroring + `CollectionSyncResult`. + +No used combos → no upstream calls (a dormant feature costs nothing). A +configured-but-empty provider yields `last_status='empty'` (no error noise). A +per-pair failure is recorded and does not abort the other pairs. + +### `RefreshTrendingDiscoverTask` (new, `internal/taskmanager/tasks/`) + +Mirrors `SyncCollectionsTask`: delegates to `TrendingRefresher.RunOnce` and +reports progress/result data. Default trigger: hourly interval (matches the +recipe's `DefaultCacheTTL` of 1h), plus a one-shot run at startup so the first +snapshot lands quickly rather than after a full interval. + +### `Fetcher` (modified, `internal/sections/fetcher.go`) + +- `loadTrendingDiscoverContentIDs` now reads the snapshot row via the snapshot + repository and returns its `content_ids`. It drops the upstream fetch, the + `singleflight` group, and the `editorialCandidateCache` usage **for trending** + (the editorial-candidate cache remains for the editorial sections that still + use it). +- `fetchTrendingDiscover` keeps its read-time responsibilities: access + filtering, re-ordering to trending rank, truncation to `ItemLimit`. It no + longer computes a `fetchLimit` for the read path. +- The `TMDBTrending`, `TraktTrending`, and trending-only use of `ItemRepo` + migrate out of the read path into the refresher. The fetcher gains the + snapshot repository dependency. + +### Wiring (`cmd/silo/main.go`, `internal/api/router.go`) + +- Construct `TrendingSnapshotRepository` and `TrendingRefresher` (wired with the + TMDB/Trakt fetchers, `ItemRepo`, section repo, snapshot repo). +- Register `RefreshTrendingDiscoverTask` alongside the other task registrations. +- Hand the snapshot repository to the section `Fetcher`; move the upstream + fetchers from the fetcher wiring to the refresher wiring. + +## Data flow + +**Refresh (background, hourly + startup):** +task → `Refresher.RunOnce` → enumerate distinct canonical `(source, window)` +from enabled `trending_discover` sections → per pair: fetch cap-200 → resolve +external IDs → `Upsert` snapshot. + +**Read (per home render):** +`fetchTrendingDiscover` → `loadTrendingDiscoverContentIDs` reads the snapshot +row → `fetchItemsByContentIDs` applies the viewer's access filter → re-order to +trending rank → truncate to `ItemLimit`. No upstream call. A single PK lookup on +a ≤6-row table — negligible next to the queries `FetchAll` already runs, so no +in-process read cache is added (it would re-muddy the single-source-of-truth and +is premature). + +## Edge behavior + +- **Before first sync:** snapshot absent → section renders empty (no blocking, + no upstream call). The startup task run closes this gap to roughly one fetch + duration. +- **Provider unconfigured but a section exists:** `last_status='empty'`, no + error. +- **Upstream failure on a pair that already has a snapshot:** prior + `content_ids` keep serving; only the attempt/status/error columns update. +- **Trakt `day` vs `week`:** collapse to a single `week` row (Trakt ignores the + window). + +## Testing + +- **Repository round-trip:** upsert then get; `content_ids` ordering preserved; + success vs failure upsert paths. +- **Refresher `RunOnce`:** fake `TMDBTrending`/`TraktTrending`, fake section + enumerator, in-memory `ItemRepo`. Assert ordered IDs upserted; `ok`/`empty`/ + `error` transitions; and the reliability invariant — a failed refresh + preserves prior `content_ids`. +- **Read path:** inject nil upstream fetchers into the fetcher and assert the + trending read still works from the snapshot (proves reads never call upstream). +- **Canonicalization:** Trakt `day` and `week` resolve to one row. + +## Risks / follow-ups + +- The data captured (`refreshed_at`, `last_status`, `last_error`, `entry_count`) + enables a future admin/inspection surface; not built here. +- If multiple sections request very different `ItemLimit`s, all share the cap-200 + list and truncate — intended, and strictly more data than the old per-limit + cache held. +``` From cc5fee095606e0ef3c2fe4f5d78252cc49debad6 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:37:27 -0400 Subject: [PATCH 03/14] docs: add trending_discover persistent snapshot implementation plan Co-Authored-By: Claude Opus 4.8 (1M context) --- ...9-trending-discover-persistent-snapshot.md | 1297 +++++++++++++++++ 1 file changed, 1297 insertions(+) create mode 100644 docs/superpowers/plans/2026-05-29-trending-discover-persistent-snapshot.md diff --git a/docs/superpowers/plans/2026-05-29-trending-discover-persistent-snapshot.md b/docs/superpowers/plans/2026-05-29-trending-discover-persistent-snapshot.md new file mode 100644 index 00000000..0c8084a4 --- /dev/null +++ b/docs/superpowers/plans/2026-05-29-trending-discover-persistent-snapshot.md @@ -0,0 +1,1297 @@ +# Trending Discover Persistent Snapshot Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Replace the in-process 1-hour cache behind the `trending_discover` home section with a background-refreshed, persisted snapshot so reads never call the upstream provider and refresh history is observable. + +**Architecture:** A new `trending_discover_snapshots` table holds one row per canonical `(source, window)` with the ordered, catalog-resolved content IDs. A scheduled task (`refresh_trending_discover`) discovers which combos are used by enabled sections, fetches TMDB/Trakt, resolves external IDs, and upserts. The section read path only reads the snapshot row; per-viewer access filtering stays at read time. + +**Tech Stack:** Go, pgx/pgxpool, PostgreSQL, the in-house `taskmanager` framework. Spec: `docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md`. + +**Baseline note:** The working tree already contains a clean refactor of the trending code (`newTrendingEntry`, `orderMediaItems`, tidied `fetchTrendingDiscoverEntries`/`resolveTrendingDiscoverIDs`, and a `singleflight` wrapper in `loadTrendingDiscoverContentIDs`). This plan assumes that working-tree state. The free helpers `trendingDiscoverEntry`, `newTrendingEntry`, and `orderedTrendingContentIDs` in `internal/sections/fetcher.go` are reused (not duplicated); the `*Fetcher` methods `fetchTrendingDiscoverEntries` / `resolveTrendingDiscoverIDs` and the fields `TMDBTrending` / `TraktTrending` / `ItemRepo` are removed in Task 6. + +**Conventions used below** +- Run all commands from the repository root. +- Lint: `make lint` (golangci-lint). Format: `gofmt -w ` before committing. +- Tests in `internal/sections` are pure unit tests — they construct `&Fetcher{...}` / structs directly with fakes and never open a DB pool. Follow that pattern. + +--- + +## Task 1: Migration — `trending_discover_snapshots` + +**Files:** +- Create: `migrations/166_trending_discover_snapshots.up.sql` +- Create: `migrations/166_trending_discover_snapshots.down.sql` + +Note: `166` is the next free number (current max on this branch is `165`). The column is named `time_window` (not `window`) to match the existing collection vocabulary (`source_config.time_window`) and avoid any keyword friction. + +- [ ] **Step 1: Write the up migration** + +Create `migrations/166_trending_discover_snapshots.up.sql`: + +```sql +-- Persisted snapshot of external global trending (TMDB / Trakt) for the +-- trending_discover home section. One row per canonical (source, time_window): +-- content_ids are already resolved to library catalog content IDs and ordered +-- by trending rank. A background task refreshes these rows; the section read +-- path only reads them, so a slow or down provider never blocks the home page. +CREATE TABLE public.trending_discover_snapshots ( + source text NOT NULL, -- 'tmdb' | 'trakt' + time_window text NOT NULL, -- 'day' | 'week' (trakt pinned to 'week') + content_ids text[] NOT NULL DEFAULT '{}'::text[], + entry_count integer NOT NULL DEFAULT 0, -- raw provider entries fetched + refreshed_at timestamptz, -- last successful refresh + last_attempt_at timestamptz, -- last attempt (success or failure) + last_status text NOT NULL DEFAULT '', -- 'ok' | 'empty' | 'error' + last_error text NOT NULL DEFAULT '', + PRIMARY KEY (source, time_window) +); +``` + +- [ ] **Step 2: Write the down migration** + +Create `migrations/166_trending_discover_snapshots.down.sql`: + +```sql +DROP TABLE IF EXISTS public.trending_discover_snapshots; +``` + +- [ ] **Step 3: Verify the SQL parses by applying it to the dev DB** + +Migrations are applied on backend startup. Apply manually to confirm the SQL is valid (dev DB role/db is `continuum`, per workspace notes): + +Run: `psql "$DATABASE_URL" -f migrations/166_trending_discover_snapshots.up.sql && psql "$DATABASE_URL" -f migrations/166_trending_discover_snapshots.down.sql` +Expected: `CREATE TABLE` then `DROP TABLE` with no errors. (Re-apply the up migration afterward, or let the next server start apply it, so later manual testing has the table.) + +If `psql`/`$DATABASE_URL` is not configured in this environment, skip this step — Task 8 starts the dev backend, which applies the migration through the normal embedded-migration path. + +- [ ] **Step 4: Commit** + +```bash +git add migrations/166_trending_discover_snapshots.up.sql migrations/166_trending_discover_snapshots.down.sql +git commit -m "feat(sections): add trending_discover_snapshots table" +``` + +--- + +## Task 2: Snapshot model, canonical key, and repository + +**Files:** +- Create: `internal/sections/trending_snapshot.go` +- Test: `internal/sections/trending_snapshot_test.go` + +This task adds the persisted model, the `canonicalTrendingKey` helper (shared by the read path, the refresher, and the repo), and the `TrendingSnapshotRepository`. Only `canonicalTrendingKey` is unit-tested (the repo is a thin pgx wrapper exercised end-to-end in Task 8, matching how other repos in this package are tested). + +- [ ] **Step 1: Write the failing test for `canonicalTrendingKey`** + +Create `internal/sections/trending_snapshot_test.go`: + +```go +package sections + +import "testing" + +func TestCanonicalTrendingKey(t *testing.T) { + t.Parallel() + cases := []struct { + src, win string + wantSrc, wantWin string + }{ + {"tmdb", "day", "tmdb", "day"}, + {"tmdb", "week", "tmdb", "week"}, + {"tmdb", "", "tmdb", "week"}, + {"", "day", "tmdb", "day"}, + {"", "", "tmdb", "week"}, + {"trakt", "day", "trakt", "week"}, + {"trakt", "week", "trakt", "week"}, + {"trakt", "", "trakt", "week"}, + {"bogus", "bogus", "tmdb", "week"}, + } + for _, c := range cases { + gotSrc, gotWin := canonicalTrendingKey(c.src, c.win) + if gotSrc != c.wantSrc || gotWin != c.wantWin { + t.Errorf("canonicalTrendingKey(%q, %q) = (%q, %q); want (%q, %q)", + c.src, c.win, gotSrc, gotWin, c.wantSrc, c.wantWin) + } + } +} +``` + +- [ ] **Step 2: Run the test to verify it fails** + +Run: `go test ./internal/sections/ -run TestCanonicalTrendingKey` +Expected: FAIL — `undefined: canonicalTrendingKey`. + +- [ ] **Step 3: Write `internal/sections/trending_snapshot.go`** + +```go +package sections + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" +) + +// TrendingSnapshot is the persisted result of one external-trending refresh for +// a canonical (Source, Window). ContentIDs are resolved to library catalog +// content IDs and ordered by trending rank. The list is viewer-agnostic; +// per-viewer access filtering happens at read time. +type TrendingSnapshot struct { + Source string + Window string + ContentIDs []string + EntryCount int + RefreshedAt *time.Time + LastAttemptAt *time.Time + LastStatus string + LastError string +} + +// canonicalTrendingKey normalizes a section's configured source/window into the +// snapshot key space. Source is "trakt" only when explicitly set; everything +// else collapses to "tmdb". Trakt ignores the time window, so it is pinned to +// "week" to avoid duplicate identical rows. For TMDB, "day" is honored only +// when explicitly set; anything else is "week". +func canonicalTrendingKey(source, window string) (string, string) { + if source != "trakt" { + source = "tmdb" + } + if source == "trakt" { + return "trakt", "week" + } + if window != "day" { + window = "week" + } + return "tmdb", window +} + +// TrendingSnapshotRepository persists and reads trending_discover_snapshots. +type TrendingSnapshotRepository struct { + pool *pgxpool.Pool +} + +// NewTrendingSnapshotRepository creates a new TrendingSnapshotRepository. +func NewTrendingSnapshotRepository(pool *pgxpool.Pool) *TrendingSnapshotRepository { + return &TrendingSnapshotRepository{pool: pool} +} + +// Get returns the snapshot for the canonical (source, window). found is false +// when no row exists yet (before the first refresh). +func (r *TrendingSnapshotRepository) Get(ctx context.Context, source, window string) (TrendingSnapshot, bool, error) { + source, window = canonicalTrendingKey(source, window) + row := r.pool.QueryRow(ctx, ` + SELECT source, time_window, content_ids, entry_count, + refreshed_at, last_attempt_at, last_status, last_error + FROM trending_discover_snapshots + WHERE source = $1 AND time_window = $2`, source, window) + + var s TrendingSnapshot + err := row.Scan(&s.Source, &s.Window, &s.ContentIDs, &s.EntryCount, + &s.RefreshedAt, &s.LastAttemptAt, &s.LastStatus, &s.LastError) + if errors.Is(err, pgx.ErrNoRows) { + return TrendingSnapshot{}, false, nil + } + if err != nil { + return TrendingSnapshot{}, false, fmt.Errorf("getting trending snapshot: %w", err) + } + return s, true, nil +} + +// SaveSuccess records a completed refresh, replacing the content list. status is +// "ok" when at least one entry matched the catalog and "empty" when the provider +// returned entries but none matched. Used only when the provider actually +// returned data; see RecordAttempt for the no-data / failure paths. +func (r *TrendingSnapshotRepository) SaveSuccess(ctx context.Context, source, window string, contentIDs []string, entryCount int, status string, at time.Time) error { + source, window = canonicalTrendingKey(source, window) + if contentIDs == nil { + contentIDs = []string{} + } + _, err := r.pool.Exec(ctx, ` + INSERT INTO trending_discover_snapshots + (source, time_window, content_ids, entry_count, refreshed_at, last_attempt_at, last_status, last_error) + VALUES ($1, $2, $3, $4, $5, $5, $6, '') + ON CONFLICT (source, time_window) DO UPDATE SET + content_ids = EXCLUDED.content_ids, + entry_count = EXCLUDED.entry_count, + refreshed_at = EXCLUDED.refreshed_at, + last_attempt_at = EXCLUDED.last_attempt_at, + last_status = EXCLUDED.last_status, + last_error = ''`, + source, window, contentIDs, entryCount, at, status) + if err != nil { + return fmt.Errorf("saving trending snapshot: %w", err) + } + return nil +} + +// RecordAttempt records an attempt that produced no new content (an upstream +// failure or an unconfigured/empty provider) WITHOUT clearing the last-good +// content_ids. status is "error" or "empty". If no row exists yet it inserts a +// placeholder so the attempt is still observable. +func (r *TrendingSnapshotRepository) RecordAttempt(ctx context.Context, source, window, status, message string, at time.Time) error { + source, window = canonicalTrendingKey(source, window) + _, err := r.pool.Exec(ctx, ` + INSERT INTO trending_discover_snapshots + (source, time_window, last_attempt_at, last_status, last_error) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (source, time_window) DO UPDATE SET + last_attempt_at = EXCLUDED.last_attempt_at, + last_status = EXCLUDED.last_status, + last_error = EXCLUDED.last_error`, + source, window, at, status, message) + if err != nil { + return fmt.Errorf("recording trending snapshot attempt: %w", err) + } + return nil +} + +// ListAll returns every snapshot row, ordered, for inspection and tests. +func (r *TrendingSnapshotRepository) ListAll(ctx context.Context) ([]TrendingSnapshot, error) { + rows, err := r.pool.Query(ctx, ` + SELECT source, time_window, content_ids, entry_count, + refreshed_at, last_attempt_at, last_status, last_error + FROM trending_discover_snapshots + ORDER BY source, time_window`) + if err != nil { + return nil, fmt.Errorf("listing trending snapshots: %w", err) + } + defer rows.Close() + + var out []TrendingSnapshot + for rows.Next() { + var s TrendingSnapshot + if err := rows.Scan(&s.Source, &s.Window, &s.ContentIDs, &s.EntryCount, + &s.RefreshedAt, &s.LastAttemptAt, &s.LastStatus, &s.LastError); err != nil { + return nil, fmt.Errorf("scanning trending snapshot: %w", err) + } + out = append(out, s) + } + return out, rows.Err() +} +``` + +- [ ] **Step 4: Run the test to verify it passes** + +Run: `go test ./internal/sections/ -run TestCanonicalTrendingKey` +Expected: PASS. + +- [ ] **Step 5: Verify the package compiles** + +Run: `go build ./internal/sections/` +Expected: no output (success). + +- [ ] **Step 6: Commit** + +```bash +gofmt -w internal/sections/trending_snapshot.go internal/sections/trending_snapshot_test.go +git add internal/sections/trending_snapshot.go internal/sections/trending_snapshot_test.go +git commit -m "feat(sections): add trending snapshot model and repository" +``` + +--- + +## Task 3: Enumerate used `(source, window)` combos + +**Files:** +- Modify: `internal/sections/repo.go` (add method near the other list methods, e.g. after `ListByScopeAll`) + +The refresher needs the config JSON of every enabled `trending_discover` section across all scopes/libraries. `repo.go` already imports `encoding/json` and `fmt` and the `Repository` has a `pool *pgxpool.Pool` field. + +- [ ] **Step 1: Add `ListTrendingDiscoverConfigs` to `internal/sections/repo.go`** + +Insert this method (place it after the `ListByScopeAll` method, around line 138): + +```go +// ListTrendingDiscoverConfigs returns the config JSON of every enabled +// trending_discover section across all scopes and libraries. The trending +// refresh task uses this to discover which (source, window) combinations need a +// snapshot, so dormant configs (no enabled sections) trigger zero upstream work. +func (r *Repository) ListTrendingDiscoverConfigs(ctx context.Context) ([]json.RawMessage, error) { + rows, err := r.pool.Query(ctx, ` + SELECT config FROM page_sections + WHERE section_type = $1 AND enabled = true`, string(SectionTrendingDiscover)) + if err != nil { + return nil, fmt.Errorf("listing trending_discover configs: %w", err) + } + defer rows.Close() + + var out []json.RawMessage + for rows.Next() { + var cfg json.RawMessage + if err := rows.Scan(&cfg); err != nil { + return nil, fmt.Errorf("scanning trending_discover config: %w", err) + } + out = append(out, cfg) + } + return out, rows.Err() +} +``` + +- [ ] **Step 2: Verify the package compiles** + +Run: `go build ./internal/sections/` +Expected: no output (success). + +- [ ] **Step 3: Commit** + +```bash +gofmt -w internal/sections/repo.go +git add internal/sections/repo.go +git commit -m "feat(sections): list enabled trending_discover section configs" +``` + +--- + +## Task 4: `TrendingRefresher` + +**Files:** +- Create: `internal/sections/trending_refresher.go` +- Test: `internal/sections/trending_refresher_test.go` + +The refresher owns the upstream call and external-ID resolution via consumer-side interfaces (so it is unit-testable with fakes). It reuses the free helpers `trendingDiscoverEntry`, `newTrendingEntry`, and `orderedTrendingContentIDs` that live in `fetcher.go` (same package). Its `fetchEntries` / `resolveIDs` bodies are copied from the soon-to-be-removed `*Fetcher` methods; those originals are deleted in Task 6 (transient duplication, resolved within this plan). + +- [ ] **Step 1: Write the failing tests** + +Create `internal/sections/trending_refresher_test.go`: + +```go +package sections + +import ( + "context" + "encoding/json" + "errors" + "testing" + "time" + + "github.com/Silo-Server/silo-server/internal/catalog" + "github.com/Silo-Server/silo-server/internal/sections/recipes" +) + +type fakeSectionLister struct { + configs []json.RawMessage + err error +} + +func (f fakeSectionLister) ListTrendingDiscoverConfigs(context.Context) ([]json.RawMessage, error) { + return f.configs, f.err +} + +type savedSnap struct { + contentIDs []string + entryCount int + status string +} + +type attemptRec struct { + status string + message string +} + +type fakeSnapshotStore struct { + saved map[string]savedSnap + attempts map[string]attemptRec +} + +func newFakeSnapshotStore() *fakeSnapshotStore { + return &fakeSnapshotStore{saved: map[string]savedSnap{}, attempts: map[string]attemptRec{}} +} + +func (f *fakeSnapshotStore) SaveSuccess(_ context.Context, source, window string, contentIDs []string, entryCount int, status string, _ time.Time) error { + f.saved[source+"|"+window] = savedSnap{contentIDs: contentIDs, entryCount: entryCount, status: status} + return nil +} + +func (f *fakeSnapshotStore) RecordAttempt(_ context.Context, source, window, status, message string, _ time.Time) error { + f.attempts[source+"|"+window] = attemptRec{status: status, message: message} + return nil +} + +type fakeTMDB struct { + entries []catalog.TMDBCollectionEntry + err error +} + +func (f fakeTMDB) GetCollectionPreset(context.Context, string, string, string, int) ([]catalog.TMDBCollectionEntry, error) { + return f.entries, f.err +} + +type fakeResolver struct { + byType map[string]*catalog.ExternalIDLookup +} + +func (f fakeResolver) GetByExternalIDs(_ context.Context, _ catalog.ExternalIDBatch, itemType string) (*catalog.ExternalIDLookup, error) { + if lk, ok := f.byType[itemType]; ok { + return lk, nil + } + return &catalog.ExternalIDLookup{ByTMDB: map[string]string{}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, nil +} + +func tmdbConfig(t *testing.T, source, window string) json.RawMessage { + t.Helper() + raw, err := json.Marshal(recipes.TrendingDiscoverParams{Source: source, Window: window}) + if err != nil { + t.Fatalf("marshal config: %v", err) + } + return raw +} + +func TestRefresherSavesOrderedContentIDs(t *testing.T) { + store := newFakeSnapshotStore() + r := &TrendingRefresher{ + Sections: fakeSectionLister{configs: []json.RawMessage{tmdbConfig(t, "tmdb", "week")}}, + Snapshots: store, + Resolver: fakeResolver{byType: map[string]*catalog.ExternalIDLookup{ + "movie": {ByTMDB: map[string]string{"10": "c-movie"}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, + "series": {ByTMDB: map[string]string{"20": "c-series"}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, + }}, + TMDBTrending: fakeTMDB{entries: []catalog.TMDBCollectionEntry{ + {ID: 10, MediaType: "movie"}, + {ID: 20, MediaType: "tv"}, + }}, + Clock: recipes.FixedClock(time.Date(2026, 5, 29, 12, 0, 0, 0, time.UTC)), + } + + data, err := r.RunOnce(context.Background()) + if err != nil { + t.Fatalf("RunOnce: %v", err) + } + + var result TrendingRefreshResult + if err := json.Unmarshal(data, &result); err != nil { + t.Fatalf("unmarshal result: %v", err) + } + if result.Combos != 1 || result.Refreshed != 1 || result.Failed != 0 || result.Empty != 0 { + t.Fatalf("result = %+v; want {Combos:1 Refreshed:1 Empty:0 Failed:0}", result) + } + + got := store.saved["tmdb|week"] + want := []string{"c-movie", "c-series"} + if len(got.contentIDs) != len(want) || got.contentIDs[0] != want[0] || got.contentIDs[1] != want[1] { + t.Fatalf("saved content IDs = %v; want %v", got.contentIDs, want) + } + if got.status != "ok" || got.entryCount != 2 { + t.Fatalf("saved snap = %+v; want status ok, entryCount 2", got) + } +} + +func TestRefresherFailurePreservesLastGood(t *testing.T) { + store := newFakeSnapshotStore() + r := &TrendingRefresher{ + Sections: fakeSectionLister{configs: []json.RawMessage{tmdbConfig(t, "tmdb", "week")}}, + Snapshots: store, + Resolver: fakeResolver{}, + TMDBTrending: fakeTMDB{err: errors.New("tmdb 503")}, + Clock: recipes.FixedClock(time.Date(2026, 5, 29, 12, 0, 0, 0, time.UTC)), + } + + data, err := r.RunOnce(context.Background()) + if err != nil { + t.Fatalf("RunOnce: %v", err) + } + + if _, ok := store.saved["tmdb|week"]; ok { + t.Fatal("SaveSuccess must not be called on fetch failure (would clear last-good)") + } + att, ok := store.attempts["tmdb|week"] + if !ok || att.status != "error" { + t.Fatalf("attempt = %+v, ok=%v; want status error", att, ok) + } + + var result TrendingRefreshResult + _ = json.Unmarshal(data, &result) + if result.Failed != 1 { + t.Fatalf("result.Failed = %d; want 1", result.Failed) + } +} + +func TestRefresherEmptyProviderPreservesLastGood(t *testing.T) { + store := newFakeSnapshotStore() + r := &TrendingRefresher{ + Sections: fakeSectionLister{configs: []json.RawMessage{tmdbConfig(t, "tmdb", "week")}}, + Snapshots: store, + Resolver: fakeResolver{}, + // TMDBTrending nil => provider unconfigured => empty entries, no error. + Clock: recipes.FixedClock(time.Date(2026, 5, 29, 12, 0, 0, 0, time.UTC)), + } + + data, err := r.RunOnce(context.Background()) + if err != nil { + t.Fatalf("RunOnce: %v", err) + } + if _, ok := store.saved["tmdb|week"]; ok { + t.Fatal("SaveSuccess must not be called when provider returns no entries") + } + att := store.attempts["tmdb|week"] + if att.status != "empty" { + t.Fatalf("attempt status = %q; want empty", att.status) + } + var result TrendingRefreshResult + _ = json.Unmarshal(data, &result) + if result.Empty != 1 { + t.Fatalf("result.Empty = %d; want 1", result.Empty) + } +} + +func TestDistinctTrendingCombosCollapsesTrakt(t *testing.T) { + configs := []json.RawMessage{ + tmdbConfig(t, "trakt", "day"), + tmdbConfig(t, "trakt", "week"), + tmdbConfig(t, "tmdb", "day"), + tmdbConfig(t, "tmdb", "day"), + } + got := distinctTrendingCombos(configs) + if len(got) != 2 { + t.Fatalf("distinctTrendingCombos len = %d (%+v); want 2", len(got), got) + } + seen := map[trendingCombo]bool{} + for _, c := range got { + seen[c] = true + } + if !seen[trendingCombo{"trakt", "week"}] || !seen[trendingCombo{"tmdb", "day"}] { + t.Fatalf("combos = %+v; want {trakt week} and {tmdb day}", got) + } +} +``` + +- [ ] **Step 2: Run the tests to verify they fail** + +Run: `go test ./internal/sections/ -run 'TestRefresher|TestDistinctTrendingCombos'` +Expected: FAIL — `undefined: TrendingRefresher`, `undefined: TrendingRefreshResult`, `undefined: distinctTrendingCombos`, `undefined: trendingCombo`. + +- [ ] **Step 3: Write `internal/sections/trending_refresher.go`** + +```go +package sections + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "time" + + "github.com/Silo-Server/silo-server/internal/catalog" + "github.com/Silo-Server/silo-server/internal/sections/recipes" +) + +// trendingFetchCap is the over-fetch size for each refresh. Library-only +// matching drops globally-trending titles the server does not own, so we fetch +// well beyond any section's display limit and store the matched, ordered list. +const trendingFetchCap = 200 + +// trendingSectionConfigLister enumerates enabled trending_discover section +// configs. Satisfied by *Repository. +type trendingSectionConfigLister interface { + ListTrendingDiscoverConfigs(ctx context.Context) ([]json.RawMessage, error) +} + +// trendingSnapshotStore is the write side of the snapshot table. Satisfied by +// *TrendingSnapshotRepository. +type trendingSnapshotStore interface { + SaveSuccess(ctx context.Context, source, window string, contentIDs []string, entryCount int, status string, at time.Time) error + RecordAttempt(ctx context.Context, source, window, status, message string, at time.Time) error +} + +// trendingExternalIDResolver resolves external IDs to library content IDs. +// Satisfied by *catalog.ItemRepository. +type trendingExternalIDResolver interface { + GetByExternalIDs(ctx context.Context, batch catalog.ExternalIDBatch, itemType string) (*catalog.ExternalIDLookup, error) +} + +// TrendingRefresher fetches external global trending (TMDB/Trakt), resolves it +// to library content IDs, and persists one snapshot per canonical +// (source, window). It is driven by a TaskManager task on an interval. +type TrendingRefresher struct { + Sections trendingSectionConfigLister + Snapshots trendingSnapshotStore + Resolver trendingExternalIDResolver + TMDBTrending catalog.TMDBCollectionFetcher + TraktTrending catalog.TraktCollectionFetcher + + // Clock defaults to recipes.RealClock{}. Tests inject recipes.FixedClock. + Clock recipes.Clock + logger *slog.Logger +} + +// NewTrendingRefresher creates a refresher with real-clock and default logger. +func NewTrendingRefresher( + sectionsRepo trendingSectionConfigLister, + snapshots trendingSnapshotStore, + resolver trendingExternalIDResolver, + tmdb catalog.TMDBCollectionFetcher, + trakt catalog.TraktCollectionFetcher, +) *TrendingRefresher { + return &TrendingRefresher{ + Sections: sectionsRepo, + Snapshots: snapshots, + Resolver: resolver, + TMDBTrending: tmdb, + TraktTrending: trakt, + Clock: recipes.RealClock{}, + logger: slog.Default(), + } +} + +func (r *TrendingRefresher) now() time.Time { + if r.Clock != nil { + return r.Clock.Now() + } + return time.Now() +} + +// TrendingRefreshResult is the JSON summary attached to the task execution. +type TrendingRefreshResult struct { + Combos int `json:"combos"` + Refreshed int `json:"refreshed"` + Empty int `json:"empty"` + Failed int `json:"failed"` +} + +type trendingCombo struct { + source string + window string +} + +// distinctTrendingCombos parses section configs and returns the deduplicated set +// of canonical (source, window) pairs that need a snapshot. +func distinctTrendingCombos(configs []json.RawMessage) []trendingCombo { + seen := make(map[trendingCombo]struct{}, len(configs)) + out := make([]trendingCombo, 0, len(configs)) + for _, raw := range configs { + var p recipes.TrendingDiscoverParams + if len(raw) > 0 { + _ = json.Unmarshal(raw, &p) + } + source, window := canonicalTrendingKey(p.Source, p.Window) + c := trendingCombo{source: source, window: window} + if _, ok := seen[c]; ok { + continue + } + seen[c] = struct{}{} + out = append(out, c) + } + return out +} + +// RunOnce refreshes every (source, window) used by an enabled trending_discover +// section. Per-combo failures are recorded and never abort the others. The JSON +// summary is suitable for task result data. +func (r *TrendingRefresher) RunOnce(ctx context.Context) (json.RawMessage, error) { + configs, err := r.Sections.ListTrendingDiscoverConfigs(ctx) + if err != nil { + return nil, fmt.Errorf("listing trending_discover sections: %w", err) + } + + combos := distinctTrendingCombos(configs) + result := TrendingRefreshResult{Combos: len(combos)} + for _, c := range combos { + switch r.refreshCombo(ctx, c.source, c.window) { + case "ok": + result.Refreshed++ + case "empty": + result.Empty++ + default: + result.Failed++ + } + } + + data, _ := json.Marshal(result) + return data, nil +} + +// refreshCombo refreshes a single canonical (source, window) and returns its +// outcome: "ok", "empty", or "error". A fetch failure or an unconfigured/empty +// provider preserves the last-good content list (RecordAttempt). When the +// provider returns entries, the list is replaced even if nothing matched the +// catalog ("empty" status with an empty list) — that genuinely reflects current +// trending having no library matches. +func (r *TrendingRefresher) refreshCombo(ctx context.Context, source, window string) string { + now := r.now() + + entries, err := r.fetchEntries(ctx, source, window, trendingFetchCap) + if err != nil { + r.logger.Error("trending refresh: fetch failed", "source", source, "window", window, "error", err) + _ = r.Snapshots.RecordAttempt(ctx, source, window, "error", err.Error(), now) + return "error" + } + if len(entries) == 0 { + // Provider unconfigured or returned nothing: keep last-good, mark empty. + _ = r.Snapshots.RecordAttempt(ctx, source, window, "empty", "", now) + return "empty" + } + + contentIDs, err := r.resolveIDs(ctx, entries) + if err != nil { + r.logger.Error("trending refresh: resolve failed", "source", source, "window", window, "error", err) + _ = r.Snapshots.RecordAttempt(ctx, source, window, "error", err.Error(), now) + return "error" + } + + status := "ok" + if len(contentIDs) == 0 { + status = "empty" + } + if err := r.Snapshots.SaveSuccess(ctx, source, window, contentIDs, len(entries), status, now); err != nil { + r.logger.Error("trending refresh: save failed", "source", source, "window", window, "error", err) + return "error" + } + return status +} + +// fetchEntries pulls the raw trending list from the configured provider. A +// nil/unconfigured provider yields an empty list (no error). +func (r *TrendingRefresher) fetchEntries(ctx context.Context, source, window string, fetchLimit int) ([]trendingDiscoverEntry, error) { + if source == "trakt" { + if r.TraktTrending == nil { + return nil, nil + } + movies, movieErr := r.TraktTrending.GetCollectionPreset(ctx, "trending", "movie", fetchLimit, "") + shows, showErr := r.TraktTrending.GetCollectionPreset(ctx, "trending", "tv", fetchLimit, "") + if movieErr != nil && showErr != nil { + return nil, fmt.Errorf("trakt trending: %v / %v", movieErr, showErr) + } + out := make([]trendingDiscoverEntry, 0, len(movies)+len(shows)) + for _, e := range movies { + out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) + } + for _, e := range shows { + out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) + } + return out, nil + } + + if r.TMDBTrending == nil { + return nil, nil + } + entries, err := r.TMDBTrending.GetCollectionPreset(ctx, "trending", "all", window, fetchLimit) + if err != nil { + return nil, err + } + out := make([]trendingDiscoverEntry, 0, len(entries)) + for _, e := range entries { + out = append(out, newTrendingEntry(e.ID, e.TVDBID, e.IMDbID, e.MediaType)) + } + return out, nil +} + +// resolveIDs matches trending entries to library content IDs via two batched +// external-ID lookups (movies, series), preserving trending order. +func (r *TrendingRefresher) resolveIDs(ctx context.Context, entries []trendingDiscoverEntry) ([]string, error) { + if r.Resolver == nil { + return nil, fmt.Errorf("trending_discover: external ID resolver not configured") + } + var movieBatch, seriesBatch catalog.ExternalIDBatch + for _, e := range entries { + batch := &movieBatch + if e.mediaType == "tv" { + batch = &seriesBatch + } + if e.tmdbID != "" { + batch.TMDBIDs = append(batch.TMDBIDs, e.tmdbID) + } + if e.imdbID != "" { + batch.IMDbIDs = append(batch.IMDbIDs, e.imdbID) + } + if e.tvdbID != "" { + batch.TVDBIDs = append(batch.TVDBIDs, e.tvdbID) + } + } + movieLookup, err := r.Resolver.GetByExternalIDs(ctx, movieBatch, "movie") + if err != nil { + return nil, err + } + seriesLookup, err := r.Resolver.GetByExternalIDs(ctx, seriesBatch, "series") + if err != nil { + return nil, err + } + return orderedTrendingContentIDs(entries, movieLookup, seriesLookup), nil +} +``` + +- [ ] **Step 4: Run the tests to verify they pass** + +Run: `go test ./internal/sections/ -run 'TestRefresher|TestDistinctTrendingCombos'` +Expected: PASS (all four tests). + +- [ ] **Step 5: Verify the package compiles** + +Run: `go build ./internal/sections/` +Expected: no output (success). Note: `fetcher.go` still defines its own `fetchTrendingDiscoverEntries`/`resolveTrendingDiscoverIDs` at this point — that is expected and removed in Task 6. + +- [ ] **Step 6: Commit** + +```bash +gofmt -w internal/sections/trending_refresher.go internal/sections/trending_refresher_test.go +git add internal/sections/trending_refresher.go internal/sections/trending_refresher_test.go +git commit -m "feat(sections): add trending refresher with persisted snapshots" +``` + +--- + +## Task 5: `RefreshTrendingDiscoverTask` + +**Files:** +- Create: `internal/taskmanager/tasks/refresh_trending_discover.go` + +Mirrors `internal/taskmanager/tasks/sync_collections.go`. Triggers on startup (so the first snapshot lands quickly) and hourly thereafter. + +- [ ] **Step 1: Write `internal/taskmanager/tasks/refresh_trending_discover.go`** + +```go +package tasks + +import ( + "context" + "encoding/json" + "fmt" + + "github.com/Silo-Server/silo-server/internal/taskmanager" +) + +// TrendingDiscoverRefresher runs a single pass of the trending refresh. +// Satisfied by *sections.TrendingRefresher. +type TrendingDiscoverRefresher interface { + RunOnce(ctx context.Context) (json.RawMessage, error) +} + +// RefreshTrendingDiscoverTask refreshes the persisted external-trending +// snapshots used by trending_discover home sections. +type RefreshTrendingDiscoverTask struct { + refresher TrendingDiscoverRefresher +} + +// NewRefreshTrendingDiscoverTask creates a new RefreshTrendingDiscoverTask. +func NewRefreshTrendingDiscoverTask(refresher TrendingDiscoverRefresher) *RefreshTrendingDiscoverTask { + return &RefreshTrendingDiscoverTask{refresher: refresher} +} + +func (t *RefreshTrendingDiscoverTask) Key() string { return "refresh_trending_discover" } +func (t *RefreshTrendingDiscoverTask) Name() string { return "Refresh Trending Discover" } +func (t *RefreshTrendingDiscoverTask) Description() string { + return "Refreshes the persisted external trending list (TMDB/Trakt) for the Trending Discover home section" +} + +func (t *RefreshTrendingDiscoverTask) Category() taskmanager.TaskCategory { + return taskmanager.TaskCategoryLibrary +} + +func (t *RefreshTrendingDiscoverTask) IsHidden() bool { return false } + +func (t *RefreshTrendingDiscoverTask) DefaultTriggers() []taskmanager.TriggerConfig { + return []taskmanager.TriggerConfig{ + {Type: taskmanager.TriggerTypeStartup}, + {Type: taskmanager.TriggerTypeInterval, IntervalMs: 60 * 60 * 1000}, // hourly + } +} + +func (t *RefreshTrendingDiscoverTask) Execute(ctx context.Context, progress taskmanager.ProgressReporter) error { + progress.Report(0, "Refreshing trending discover") + + resultData, err := t.refresher.RunOnce(ctx) + if err != nil { + return fmt.Errorf("trending discover refresh: %w", err) + } + if resultData != nil { + progress.SetResultData(resultData) + } + + progress.Report(100, "Trending discover refresh complete") + return nil +} +``` + +- [ ] **Step 2: Verify the package compiles** + +Run: `go build ./internal/taskmanager/...` +Expected: no output (success). + +- [ ] **Step 3: Commit** + +```bash +gofmt -w internal/taskmanager/tasks/refresh_trending_discover.go +git add internal/taskmanager/tasks/refresh_trending_discover.go +git commit -m "feat(tasks): add refresh_trending_discover task" +``` + +--- + +## Task 6: Rewire the Fetcher read path to the snapshot + +**Files:** +- Modify: `internal/sections/fetcher.go` +- Test: `internal/sections/trending_read_test.go` (create) + +This task: (a) adds the `TrendingSnapshots` read dependency to `Fetcher`; (b) removes the `TMDBTrending`, `TraktTrending`, and `ItemRepo` fields plus the `fetchTrendingDiscoverEntries` / `resolveTrendingDiscoverIDs` methods (now owned by the refresher); (c) rewrites `loadTrendingDiscoverContentIDs` to read the snapshot; (d) simplifies `fetchTrendingDiscover`. The free helpers `trendingDiscoverEntry`, `newTrendingEntry`, `orderedTrendingContentIDs` stay. + +- [ ] **Step 1: Write the failing read-path test** + +Create `internal/sections/trending_read_test.go`: + +```go +package sections + +import ( + "context" + "testing" +) + +type fakeSnapshotGetter struct { + snap TrendingSnapshot + found bool + err error +} + +func (f fakeSnapshotGetter) Get(context.Context, string, string) (TrendingSnapshot, bool, error) { + return f.snap, f.found, f.err +} + +func TestLoadTrendingDiscoverContentIDsReadsSnapshot(t *testing.T) { + f := &Fetcher{TrendingSnapshots: fakeSnapshotGetter{ + snap: TrendingSnapshot{ContentIDs: []string{"a", "b"}}, + found: true, + }} + ids, err := f.loadTrendingDiscoverContentIDs(context.Background(), "tmdb", "week") + if err != nil { + t.Fatalf("loadTrendingDiscoverContentIDs: %v", err) + } + if len(ids) != 2 || ids[0] != "a" || ids[1] != "b" { + t.Fatalf("ids = %v; want [a b]", ids) + } +} + +func TestLoadTrendingDiscoverContentIDsNilGetter(t *testing.T) { + f := &Fetcher{} + ids, err := f.loadTrendingDiscoverContentIDs(context.Background(), "tmdb", "week") + if err != nil { + t.Fatalf("loadTrendingDiscoverContentIDs: %v", err) + } + if ids != nil { + t.Fatalf("ids = %v; want nil for nil getter", ids) + } +} + +func TestLoadTrendingDiscoverContentIDsNotFound(t *testing.T) { + f := &Fetcher{TrendingSnapshots: fakeSnapshotGetter{found: false}} + ids, err := f.loadTrendingDiscoverContentIDs(context.Background(), "tmdb", "week") + if err != nil { + t.Fatalf("loadTrendingDiscoverContentIDs: %v", err) + } + if ids != nil { + t.Fatalf("ids = %v; want nil when no snapshot exists", ids) + } +} +``` + +- [ ] **Step 2: Run the test to verify it fails** + +Run: `go test ./internal/sections/ -run TestLoadTrendingDiscoverContentIDs` +Expected: FAIL — `unknown field 'TrendingSnapshots' in struct literal` / `loadTrendingDiscoverContentIDs` signature mismatch. + +- [ ] **Step 3: Add the snapshot reader interface and field to `Fetcher`** + +In `internal/sections/fetcher.go`, replace the trending dependency fields in the `Fetcher` struct. Find this block (around lines 71-77): + +```go + // ItemRepo resolves external IDs (TMDB/Trakt) to library content IDs for + // the trending_discover section. Nil disables external-trending matching. + ItemRepo *catalog.ItemRepository + // TMDBTrending and TraktTrending fetch external global trending lists. Each + // is nil when that provider is not configured. + TMDBTrending catalog.TMDBCollectionFetcher + TraktTrending catalog.TraktCollectionFetcher +``` + +Replace it with: + +```go + // TrendingSnapshots reads the persisted external-trending snapshots that + // back the trending_discover section. Nil renders that section empty. + // Snapshots are produced out-of-band by TrendingRefresher, so the read path + // never calls the upstream provider. + TrendingSnapshots trendingSnapshotGetter +``` + +Then add the interface declaration just above the `Fetcher` struct definition (above `// Fetcher runs section queries against the database.`): + +```go +// trendingSnapshotGetter is the read side of the trending snapshot table. +// Satisfied by *TrendingSnapshotRepository. +type trendingSnapshotGetter interface { + Get(ctx context.Context, source, window string) (TrendingSnapshot, bool, error) +} +``` + +- [ ] **Step 4: Rewrite `loadTrendingDiscoverContentIDs`** + +In `internal/sections/fetcher.go`, replace the entire `loadTrendingDiscoverContentIDs` method (the version that uses `ensureEditorialCandidateCache` / `candidateGroup.Do`) with: + +```go +// loadTrendingDiscoverContentIDs returns the persisted, catalog-resolved content +// IDs for the canonical (source, window). It reads only the snapshot table; the +// upstream fetch happens out-of-band in TrendingRefresher. Returns nil when no +// snapshot reader is configured or no snapshot exists yet. +func (f *Fetcher) loadTrendingDiscoverContentIDs(ctx context.Context, source, window string) ([]string, error) { + if f.TrendingSnapshots == nil { + return nil, nil + } + snap, found, err := f.TrendingSnapshots.Get(ctx, source, window) + if err != nil { + return nil, err + } + if !found { + return nil, nil + } + return snap.ContentIDs, nil +} +``` + +- [ ] **Step 5: Simplify `fetchTrendingDiscover`** + +In `internal/sections/fetcher.go`, in `fetchTrendingDiscover`, replace the source/window normalization and the `fetchLimit` block. Find: + +```go + source := p.Source + if source != "trakt" { + source = "tmdb" + } + window := p.Window + if window != "day" { + window = "week" + } + + limit := s.ItemLimit + if limit <= 0 { + limit = 20 + } + // Over-fetch: library-only matching drops globally-trending titles the + // server does not own, so request more candidates than the display limit. + fetchLimit := limit * 5 + if fetchLimit < 50 { + fetchLimit = 50 + } + if fetchLimit > 200 { + fetchLimit = 200 + } + + orderedIDs, err := f.loadTrendingDiscoverContentIDs(ctx, source, window, fetchLimit) +``` + +Replace with: + +```go + source, window := canonicalTrendingKey(p.Source, p.Window) + + limit := s.ItemLimit + if limit <= 0 { + limit = 20 + } + + orderedIDs, err := f.loadTrendingDiscoverContentIDs(ctx, source, window) +``` + +- [ ] **Step 6: Delete the now-orphaned fetcher methods** + +In `internal/sections/fetcher.go`, delete the two methods `func (f *Fetcher) fetchTrendingDiscoverEntries(...)` and `func (f *Fetcher) resolveTrendingDiscoverIDs(...)` in their entirety (their logic now lives on `TrendingRefresher`). Keep `trendingDiscoverEntry`, `newTrendingEntry`, and `orderedTrendingContentIDs`. + +- [ ] **Step 7: Build and fix any leftover references** + +Run: `go build ./internal/sections/` +Expected: success. If the compiler reports `catalog` imported and not used, confirm `catalog` is still referenced elsewhere in `fetcher.go` (it is — `catalog.AccessFilter` and others). Do not remove the import. If it reports the removed methods are still referenced, ensure Step 4/5 fully replaced the old call site. + +- [ ] **Step 8: Run the read-path tests** + +Run: `go test ./internal/sections/ -run TestLoadTrendingDiscoverContentIDs` +Expected: PASS (all three). + +- [ ] **Step 9: Run the full sections package test suite** + +Run: `go test ./internal/sections/` +Expected: PASS. + +- [ ] **Step 10: Commit** + +```bash +gofmt -w internal/sections/fetcher.go internal/sections/trending_read_test.go +git add internal/sections/fetcher.go internal/sections/trending_read_test.go +git commit -m "refactor(sections): read trending_discover from persisted snapshot" +``` + +--- + +## Task 7: Wiring — construct the refresher, register the task, set the reader + +**Files:** +- Modify: `internal/api/router.go` (around lines 883-885) +- Modify: `cmd/silo/main.go` (declare near line 1235, build near 1238-1244, register near 1287-1289) + +- [ ] **Step 1: Point the section Fetcher at the snapshot reader (router.go)** + +In `internal/api/router.go`, replace the three trending-wiring lines (currently lines 883-885): + +```go + sectionFetcher.ItemRepo = itemRepo + sectionFetcher.TMDBTrending = libraryCollectionService.TMDBCollections + sectionFetcher.TraktTrending = libraryCollectionService.TraktCollections +``` + +with: + +```go + // trending_discover reads its list from the persisted snapshot table; + // the upstream fetch happens out-of-band in the refresh task. + sectionFetcher.TrendingSnapshots = sections.NewTrendingSnapshotRepository(deps.DB) +``` + +(The surrounding comment "Wire external-trending fetchers into the section fetcher..." can be updated to "Wire the trending snapshot reader into the section fetcher...".) + +- [ ] **Step 2: Verify the router compiles** + +Run: `go build ./internal/api/` +Expected: success. If `itemRepo` becomes unused in router.go after this change, the build will report it; in that case keep `itemRepo` only if other code uses it (search `grep -n "itemRepo" internal/api/router.go`) — it is used elsewhere (e.g. the library collection handler at ~line 888), so no removal is needed. + +- [ ] **Step 3: Declare the refresher variable (main.go)** + +In `cmd/silo/main.go`, find the declaration near line 1235: + +```go + var collectionSyncScheduler *catalog.CollectionSyncScheduler +``` + +Add directly below it: + +```go + var trendingRefresher *sections.TrendingRefresher +``` + +- [ ] **Step 4: Build the refresher (main.go)** + +In `cmd/silo/main.go`, find where `collectionSyncScheduler` is assigned (around line 1244): + +```go + collectionSyncScheduler = catalog.NewCollectionSyncScheduler(collectionRepo, collectionService, slog.Default()) +``` + +Add directly below it: + +```go + trendingRefresher = sections.NewTrendingRefresher( + sectionRepo, + sections.NewTrendingSnapshotRepository(pool), + catalog.NewItemRepository(deps.DB), + collectionService.TMDBCollections, + collectionService.TraktCollections, + ) +``` + +(`sectionRepo` is in scope from line 1178; `collectionService` from line 1241. `collectionService.TraktCollections` may be nil — the refresher handles a nil provider by recording an "empty" attempt, so passing nil is safe.) + +- [ ] **Step 5: Register the task (main.go)** + +In `cmd/silo/main.go`, find the collection task registration (around line 1287-1289): + +```go + if collectionSyncScheduler != nil { + taskMgr.Register(tasks.NewSyncCollectionsTask(collectionSyncScheduler)) + } +``` + +Add directly below it: + +```go + if trendingRefresher != nil { + taskMgr.Register(tasks.NewRefreshTrendingDiscoverTask(trendingRefresher)) + } +``` + +- [ ] **Step 6: Verify imports and build the whole binary** + +`cmd/silo/main.go` already imports `github.com/Silo-Server/silo-server/internal/sections` (used for `sectionRepo`) and `.../internal/catalog` and `.../internal/taskmanager/tasks`. No new imports needed. + +Run: `go build ./...` +Expected: success (no output). + +- [ ] **Step 7: Commit** + +```bash +gofmt -w internal/api/router.go cmd/silo/main.go +git add internal/api/router.go cmd/silo/main.go +git commit -m "feat: wire trending refresh task and snapshot reader" +``` + +--- + +## Task 8: Full verification + +**Files:** none (verification only). + +- [ ] **Step 1: Build everything** + +Run: `go build ./...` +Expected: success. + +- [ ] **Step 2: Vet** + +Run: `go vet ./internal/sections/... ./internal/taskmanager/... ./internal/api/... ./cmd/...` +Expected: no findings. + +- [ ] **Step 3: Run the affected test packages** + +Run: `go test ./internal/sections/... ./internal/taskmanager/...` +Expected: PASS. + +- [ ] **Step 4: Lint** + +Run: `make lint` +Expected: no new findings in the files touched by this plan. (Pre-existing findings elsewhere are out of scope.) + +- [ ] **Step 5: End-to-end smoke test against dev** + +Start the dev backend (this applies migration 166 via the embedded-migration path) and confirm the task runs and persists a snapshot: + +Run: `make dev-backend` (in a separate shell), then once it is up: +- Trigger or wait for the `refresh_trending_discover` task (it has a startup trigger). +- Verify a row exists: `psql "$DATABASE_URL" -c "SELECT source, time_window, array_length(content_ids,1), entry_count, last_status, refreshed_at FROM trending_discover_snapshots;"` +Expected: at least one row per used `(source, window)` with `last_status` of `ok`/`empty` and a recent `refreshed_at` (for `ok`). If no `trending_discover` section is enabled, expect zero rows (and zero upstream calls) — enable one in the admin UI to exercise the path. +- Load a home page that includes the trending section and confirm it renders the same items as before. + +- [ ] **Step 6: Final commit (if any formatting/cleanup remains)** + +```bash +git add -A +git commit -m "chore(sections): finalize trending snapshot verification" || echo "nothing to commit" +``` + +--- + +## Self-Review + +**1. Spec coverage** +- Persisted snapshot table (`trending_discover_snapshots`) → Task 1. ✓ +- Reads never call upstream (snapshot read only) → Task 6 (`loadTrendingDiscoverContentIDs` reads `TrendingSnapshots.Get`; Fetcher loses all upstream fetcher fields). ✓ +- Background refresh task (hourly + startup) → Task 5 (`TriggerTypeStartup` + hourly interval). ✓ +- Enumerate used `(source, window)` from enabled sections; dormant feature = zero upstream calls → Task 3 + Task 4 (`ListTrendingDiscoverConfigs`, `distinctTrendingCombos`, `RunOnce` over combos). ✓ +- Reliability invariant: failed/empty refresh preserves last-good `content_ids` → Task 2 (`RecordAttempt` vs `SaveSuccess`) + Task 4 (`refreshCombo` routing) + tests `TestRefresherFailurePreservesLastGood`, `TestRefresherEmptyProviderPreservesLastGood`. ✓ +- Observability (refreshed_at/last_status/last_error/entry_count + task run summary) → Task 1 columns + Task 4 `TrendingRefreshResult` + Task 5 `SetResultData`. ✓ +- Cap-200 over-fetch stored, read truncates to ItemLimit → Task 4 (`trendingFetchCap`) + Task 6 (`fetchTrendingDiscover` truncation retained). ✓ +- Trakt day/week collapse to one row → Task 2 (`canonicalTrendingKey`) + test `TestDistinctTrendingCombosCollapsesTrakt`. ✓ +- First-boot: empty until first sync, startup kick → Task 5 startup trigger; Task 6 nil/not-found returns nil → empty render. ✓ +- Refresher lives in `internal/sections/` → Task 4. ✓ +- Wiring (construct refresher, register task, set reader, remove old trending wiring) → Task 7. ✓ + +**2. Placeholder scan:** No TBD/TODO; every code step contains full code; every test step has assertions and an expected result. ✓ + +**3. Type consistency:** +- `canonicalTrendingKey(source, window string) (string, string)` — defined Task 2, used Tasks 2/4/6. ✓ +- `TrendingSnapshot{Source, Window, ContentIDs, EntryCount, RefreshedAt, LastAttemptAt, LastStatus, LastError}` — Task 2, used Tasks 4/6 tests. ✓ +- Repo methods `Get(ctx, source, window) (TrendingSnapshot, bool, error)`, `SaveSuccess(ctx, source, window, contentIDs, entryCount, status, at)`, `RecordAttempt(ctx, source, window, status, message, at)`, `ListAll(ctx)` — Task 2; interfaces `trendingSnapshotStore` (Task 4) and `trendingSnapshotGetter` (Task 6) match these signatures exactly. ✓ +- `ListTrendingDiscoverConfigs(ctx) ([]json.RawMessage, error)` — Task 3; interface `trendingSectionConfigLister` (Task 4) matches. ✓ +- `GetByExternalIDs(ctx, catalog.ExternalIDBatch, string) (*catalog.ExternalIDLookup, error)` — interface `trendingExternalIDResolver` (Task 4) matches `*catalog.ItemRepository`. ✓ +- `TrendingRefresher.RunOnce(ctx) (json.RawMessage, error)` — Task 4; interface `TrendingDiscoverRefresher` (Task 5) matches. ✓ +- `NewTrendingRefresher(lister, store, resolver, tmdb, trakt)` arg order — Task 4 definition matches Task 7 call site. ✓ +``` From 1a74bdbb6ace8733c27a8a587e1f7ac055b14d70 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:40:13 -0400 Subject: [PATCH 04/14] refactor(sections): tidy trending_discover fetch and cache helpers Extract newTrendingEntry, reuse orderMediaItems, and collapse concurrent cache-miss loads with singleflight. Baseline for the persistent snapshot work. Co-Authored-By: Claude Opus 4.8 (1M context) --- internal/sections/fetcher.go | 99 +++++++++++++++++++----------------- 1 file changed, 52 insertions(+), 47 deletions(-) diff --git a/internal/sections/fetcher.go b/internal/sections/fetcher.go index df0168e1..6e87fab1 100644 --- a/internal/sections/fetcher.go +++ b/internal/sections/fetcher.go @@ -2353,6 +2353,19 @@ type trendingDiscoverEntry struct { mediaType string // "movie" | "tv" } +// newTrendingEntry normalizes a provider entry into a trendingDiscoverEntry, +// stringifying the numeric external IDs and dropping any that are unset (<= 0). +func newTrendingEntry(tmdbID, tvdbID int, imdbID, mediaType string) trendingDiscoverEntry { + e := trendingDiscoverEntry{imdbID: imdbID, mediaType: mediaType} + if tmdbID > 0 { + e.tmdbID = strconv.Itoa(tmdbID) + } + if tvdbID > 0 { + e.tvdbID = strconv.Itoa(tvdbID) + } + return e +} + // fetchTrendingDiscover surfaces external global trending (TMDB or Trakt), // mixing movies + series, matched to titles in the viewer's enabled libraries. // The external fetch + ID resolution are cached briefly so it does not hit the @@ -2400,20 +2413,9 @@ func (f *Fetcher) fetchTrendingDiscover(ctx context.Context, s ResolvedSection, // Re-order to trending rank (fetchItemsByContentIDs returns DB order) and // truncate to the section's display limit. - byID := make(map[string]*models.MediaItem, len(items)) - for _, item := range items { - byID[item.ContentID] = item - } - ordered := make([]*models.MediaItem, 0, len(orderedIDs)) - for _, id := range orderedIDs { - item, ok := byID[id] - if !ok { - continue - } - ordered = append(ordered, item) - if len(ordered) >= limit { - break - } + ordered := orderMediaItems(items, orderedIDs) + if len(ordered) > limit { + ordered = ordered[:limit] } return ordered, len(ordered), nil } @@ -2423,27 +2425,42 @@ func (f *Fetcher) fetchTrendingDiscover(ctx context.Context, s ResolvedSection, func (f *Fetcher) loadTrendingDiscoverContentIDs(ctx context.Context, source, window string, fetchLimit int) ([]string, error) { cache := f.ensureEditorialCandidateCache() cacheKey := fmt.Sprintf("trending_discover|%s|%s|%d", source, window, fetchLimit) - now := f.Clock.Now() - if cached, ok := cache.get(cacheKey, now); ok { + if cached, ok := cache.get(cacheKey, f.now()); ok { return cached, nil } - entries, err := f.fetchTrendingDiscoverEntries(ctx, source, window, fetchLimit) - if err != nil { - return nil, err - } - // Don't cache when the provider is unconfigured/empty, so a newly-configured - // provider takes effect immediately rather than after the TTL. - if len(entries) == 0 { - return []string{}, nil - } + // Collapse concurrent cache-miss loads (e.g. many home-page requests at TTL + // expiry) into a single upstream fetch + ID resolution, as + // cachedEditorialCandidates does for editorial sections. + value, err, _ := f.candidateGroup.Do(cacheKey, func() (any, error) { + now := f.now() + if cached, ok := cache.get(cacheKey, now); ok { + return cached, nil + } - contentIDs, err := f.resolveTrendingDiscoverIDs(ctx, entries) + entries, err := f.fetchTrendingDiscoverEntries(ctx, source, window, fetchLimit) + if err != nil { + return nil, err + } + // Don't cache when the provider is unconfigured/empty, so a + // newly-configured provider takes effect immediately rather than after + // the TTL. + if len(entries) == 0 { + return []string{}, nil + } + + contentIDs, err := f.resolveTrendingDiscoverIDs(ctx, entries) + if err != nil { + return nil, err + } + cache.set(cacheKey, contentIDs, now.Add(time.Hour)) + return contentIDs, nil + }) if err != nil { return nil, err } - cache.set(cacheKey, contentIDs, now.Add(time.Hour)) - return contentIDs, nil + contentIDs, _ := value.([]string) + return append([]string(nil), contentIDs...), nil } // fetchTrendingDiscoverEntries pulls the raw trending list from the configured @@ -2460,17 +2477,12 @@ func (f *Fetcher) fetchTrendingDiscoverEntries(ctx context.Context, source, wind if movieErr != nil && showErr != nil { return nil, fmt.Errorf("trakt trending: %v / %v", movieErr, showErr) } - merged := append(append([]catalog.TraktCollectionEntry{}, movies...), shows...) - out := make([]trendingDiscoverEntry, 0, len(merged)) - for _, e := range merged { - entry := trendingDiscoverEntry{imdbID: e.IMDbID, mediaType: e.MediaType} - if e.TMDBID > 0 { - entry.tmdbID = strconv.Itoa(e.TMDBID) - } - if e.TVDBID > 0 { - entry.tvdbID = strconv.Itoa(e.TVDBID) - } - out = append(out, entry) + out := make([]trendingDiscoverEntry, 0, len(movies)+len(shows)) + for _, e := range movies { + out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) + } + for _, e := range shows { + out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) } return out, nil } @@ -2485,14 +2497,7 @@ func (f *Fetcher) fetchTrendingDiscoverEntries(ctx context.Context, source, wind } out := make([]trendingDiscoverEntry, 0, len(entries)) for _, e := range entries { - entry := trendingDiscoverEntry{imdbID: e.IMDbID, mediaType: e.MediaType} - if e.ID > 0 { - entry.tmdbID = strconv.Itoa(e.ID) - } - if e.TVDBID > 0 { - entry.tvdbID = strconv.Itoa(e.TVDBID) - } - out = append(out, entry) + out = append(out, newTrendingEntry(e.ID, e.TVDBID, e.IMDbID, e.MediaType)) } return out, nil } From ca4cd230a47501f4d8aaaa4568a3fa59a43a9706 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:40:58 -0400 Subject: [PATCH 05/14] feat(sections): add trending_discover_snapshots table --- .../166_trending_discover_snapshots.down.sql | 1 + .../166_trending_discover_snapshots.up.sql | 16 ++++++++++++++++ 2 files changed, 17 insertions(+) create mode 100644 migrations/166_trending_discover_snapshots.down.sql create mode 100644 migrations/166_trending_discover_snapshots.up.sql diff --git a/migrations/166_trending_discover_snapshots.down.sql b/migrations/166_trending_discover_snapshots.down.sql new file mode 100644 index 00000000..72a6d6d4 --- /dev/null +++ b/migrations/166_trending_discover_snapshots.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS public.trending_discover_snapshots; diff --git a/migrations/166_trending_discover_snapshots.up.sql b/migrations/166_trending_discover_snapshots.up.sql new file mode 100644 index 00000000..706d5c83 --- /dev/null +++ b/migrations/166_trending_discover_snapshots.up.sql @@ -0,0 +1,16 @@ +-- Persisted snapshot of external global trending (TMDB / Trakt) for the +-- trending_discover home section. One row per canonical (source, time_window): +-- content_ids are already resolved to library catalog content IDs and ordered +-- by trending rank. A background task refreshes these rows; the section read +-- path only reads them, so a slow or down provider never blocks the home page. +CREATE TABLE public.trending_discover_snapshots ( + source text NOT NULL, -- 'tmdb' | 'trakt' + time_window text NOT NULL, -- 'day' | 'week' (trakt pinned to 'week') + content_ids text[] NOT NULL DEFAULT '{}'::text[], + entry_count integer NOT NULL DEFAULT 0, -- raw provider entries fetched + refreshed_at timestamptz, -- last successful refresh + last_attempt_at timestamptz, -- last attempt (success or failure) + last_status text NOT NULL DEFAULT '', -- 'ok' | 'empty' | 'error' + last_error text NOT NULL DEFAULT '', + PRIMARY KEY (source, time_window) +); From 94a588c7a3fa65e290eed5e4635d3a5af9715556 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:41:47 -0400 Subject: [PATCH 06/14] feat(sections): add trending snapshot model and repository --- internal/sections/trending_snapshot.go | 148 ++++++++++++++++++++ internal/sections/trending_snapshot_test.go | 28 ++++ 2 files changed, 176 insertions(+) create mode 100644 internal/sections/trending_snapshot.go create mode 100644 internal/sections/trending_snapshot_test.go diff --git a/internal/sections/trending_snapshot.go b/internal/sections/trending_snapshot.go new file mode 100644 index 00000000..033d3b19 --- /dev/null +++ b/internal/sections/trending_snapshot.go @@ -0,0 +1,148 @@ +package sections + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" +) + +// TrendingSnapshot is the persisted result of one external-trending refresh for +// a canonical (Source, Window). ContentIDs are resolved to library catalog +// content IDs and ordered by trending rank. The list is viewer-agnostic; +// per-viewer access filtering happens at read time. +type TrendingSnapshot struct { + Source string + Window string + ContentIDs []string + EntryCount int + RefreshedAt *time.Time + LastAttemptAt *time.Time + LastStatus string + LastError string +} + +// canonicalTrendingKey normalizes a section's configured source/window into the +// snapshot key space. Source is "trakt" only when explicitly set; everything +// else collapses to "tmdb". Trakt ignores the time window, so it is pinned to +// "week" to avoid duplicate identical rows. For TMDB, "day" is honored only +// when explicitly set; anything else is "week". +func canonicalTrendingKey(source, window string) (string, string) { + if source != "trakt" { + source = "tmdb" + } + if source == "trakt" { + return "trakt", "week" + } + if window != "day" { + window = "week" + } + return "tmdb", window +} + +// TrendingSnapshotRepository persists and reads trending_discover_snapshots. +type TrendingSnapshotRepository struct { + pool *pgxpool.Pool +} + +// NewTrendingSnapshotRepository creates a new TrendingSnapshotRepository. +func NewTrendingSnapshotRepository(pool *pgxpool.Pool) *TrendingSnapshotRepository { + return &TrendingSnapshotRepository{pool: pool} +} + +// Get returns the snapshot for the canonical (source, window). found is false +// when no row exists yet (before the first refresh). +func (r *TrendingSnapshotRepository) Get(ctx context.Context, source, window string) (TrendingSnapshot, bool, error) { + source, window = canonicalTrendingKey(source, window) + row := r.pool.QueryRow(ctx, ` + SELECT source, time_window, content_ids, entry_count, + refreshed_at, last_attempt_at, last_status, last_error + FROM trending_discover_snapshots + WHERE source = $1 AND time_window = $2`, source, window) + + var s TrendingSnapshot + err := row.Scan(&s.Source, &s.Window, &s.ContentIDs, &s.EntryCount, + &s.RefreshedAt, &s.LastAttemptAt, &s.LastStatus, &s.LastError) + if errors.Is(err, pgx.ErrNoRows) { + return TrendingSnapshot{}, false, nil + } + if err != nil { + return TrendingSnapshot{}, false, fmt.Errorf("getting trending snapshot: %w", err) + } + return s, true, nil +} + +// SaveSuccess records a completed refresh, replacing the content list. status is +// "ok" when at least one entry matched the catalog and "empty" when the provider +// returned entries but none matched. Used only when the provider actually +// returned data; see RecordAttempt for the no-data / failure paths. +func (r *TrendingSnapshotRepository) SaveSuccess(ctx context.Context, source, window string, contentIDs []string, entryCount int, status string, at time.Time) error { + source, window = canonicalTrendingKey(source, window) + if contentIDs == nil { + contentIDs = []string{} + } + _, err := r.pool.Exec(ctx, ` + INSERT INTO trending_discover_snapshots + (source, time_window, content_ids, entry_count, refreshed_at, last_attempt_at, last_status, last_error) + VALUES ($1, $2, $3, $4, $5, $5, $6, '') + ON CONFLICT (source, time_window) DO UPDATE SET + content_ids = EXCLUDED.content_ids, + entry_count = EXCLUDED.entry_count, + refreshed_at = EXCLUDED.refreshed_at, + last_attempt_at = EXCLUDED.last_attempt_at, + last_status = EXCLUDED.last_status, + last_error = ''`, + source, window, contentIDs, entryCount, at, status) + if err != nil { + return fmt.Errorf("saving trending snapshot: %w", err) + } + return nil +} + +// RecordAttempt records an attempt that produced no new content (an upstream +// failure or an unconfigured/empty provider) WITHOUT clearing the last-good +// content_ids. status is "error" or "empty". If no row exists yet it inserts a +// placeholder so the attempt is still observable. +func (r *TrendingSnapshotRepository) RecordAttempt(ctx context.Context, source, window, status, message string, at time.Time) error { + source, window = canonicalTrendingKey(source, window) + _, err := r.pool.Exec(ctx, ` + INSERT INTO trending_discover_snapshots + (source, time_window, last_attempt_at, last_status, last_error) + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT (source, time_window) DO UPDATE SET + last_attempt_at = EXCLUDED.last_attempt_at, + last_status = EXCLUDED.last_status, + last_error = EXCLUDED.last_error`, + source, window, at, status, message) + if err != nil { + return fmt.Errorf("recording trending snapshot attempt: %w", err) + } + return nil +} + +// ListAll returns every snapshot row, ordered, for inspection and tests. +func (r *TrendingSnapshotRepository) ListAll(ctx context.Context) ([]TrendingSnapshot, error) { + rows, err := r.pool.Query(ctx, ` + SELECT source, time_window, content_ids, entry_count, + refreshed_at, last_attempt_at, last_status, last_error + FROM trending_discover_snapshots + ORDER BY source, time_window`) + if err != nil { + return nil, fmt.Errorf("listing trending snapshots: %w", err) + } + defer rows.Close() + + var out []TrendingSnapshot + for rows.Next() { + var s TrendingSnapshot + if err := rows.Scan(&s.Source, &s.Window, &s.ContentIDs, &s.EntryCount, + &s.RefreshedAt, &s.LastAttemptAt, &s.LastStatus, &s.LastError); err != nil { + return nil, fmt.Errorf("scanning trending snapshot: %w", err) + } + out = append(out, s) + } + return out, rows.Err() +} diff --git a/internal/sections/trending_snapshot_test.go b/internal/sections/trending_snapshot_test.go new file mode 100644 index 00000000..cb9fb000 --- /dev/null +++ b/internal/sections/trending_snapshot_test.go @@ -0,0 +1,28 @@ +package sections + +import "testing" + +func TestCanonicalTrendingKey(t *testing.T) { + t.Parallel() + cases := []struct { + src, win string + wantSrc, wantWin string + }{ + {"tmdb", "day", "tmdb", "day"}, + {"tmdb", "week", "tmdb", "week"}, + {"tmdb", "", "tmdb", "week"}, + {"", "day", "tmdb", "day"}, + {"", "", "tmdb", "week"}, + {"trakt", "day", "trakt", "week"}, + {"trakt", "week", "trakt", "week"}, + {"trakt", "", "trakt", "week"}, + {"bogus", "bogus", "tmdb", "week"}, + } + for _, c := range cases { + gotSrc, gotWin := canonicalTrendingKey(c.src, c.win) + if gotSrc != c.wantSrc || gotWin != c.wantWin { + t.Errorf("canonicalTrendingKey(%q, %q) = (%q, %q); want (%q, %q)", + c.src, c.win, gotSrc, gotWin, c.wantSrc, c.wantWin) + } + } +} From 7d70bf0cf46c60f581a2d9f9ccd2bbd1b5e46dfd Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:42:09 -0400 Subject: [PATCH 07/14] feat(sections): list enabled trending_discover section configs --- internal/sections/repo.go | 24 ++++++++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/internal/sections/repo.go b/internal/sections/repo.go index 7fb4ca19..a0db33c5 100644 --- a/internal/sections/repo.go +++ b/internal/sections/repo.go @@ -135,6 +135,30 @@ func (r *Repository) ListByScopeAll(ctx context.Context, scope string, libraryID return scanSections(rows) } +// ListTrendingDiscoverConfigs returns the config JSON of every enabled +// trending_discover section across all scopes and libraries. The trending +// refresh task uses this to discover which (source, window) combinations need a +// snapshot, so dormant configs (no enabled sections) trigger zero upstream work. +func (r *Repository) ListTrendingDiscoverConfigs(ctx context.Context) ([]json.RawMessage, error) { + rows, err := r.pool.Query(ctx, ` + SELECT config FROM page_sections + WHERE section_type = $1 AND enabled = true`, string(SectionTrendingDiscover)) + if err != nil { + return nil, fmt.Errorf("listing trending_discover configs: %w", err) + } + defer rows.Close() + + var out []json.RawMessage + for rows.Next() { + var cfg json.RawMessage + if err := rows.Scan(&cfg); err != nil { + return nil, fmt.Errorf("scanning trending_discover config: %w", err) + } + out = append(out, cfg) + } + return out, rows.Err() +} + // Update modifies an existing section. func (r *Repository) Update(ctx context.Context, s *PageSection) error { query := `UPDATE page_sections SET From ff5f399e2c15371ed43ae53c9e4a7b316b02df6b Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:43:54 -0400 Subject: [PATCH 08/14] feat(sections): add trending refresher with persisted snapshots --- internal/sections/trending_refresher.go | 252 +++++++++++++++++++ internal/sections/trending_refresher_test.go | 197 +++++++++++++++ 2 files changed, 449 insertions(+) create mode 100644 internal/sections/trending_refresher.go create mode 100644 internal/sections/trending_refresher_test.go diff --git a/internal/sections/trending_refresher.go b/internal/sections/trending_refresher.go new file mode 100644 index 00000000..b5bc92d0 --- /dev/null +++ b/internal/sections/trending_refresher.go @@ -0,0 +1,252 @@ +package sections + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "time" + + "github.com/Silo-Server/silo-server/internal/catalog" + "github.com/Silo-Server/silo-server/internal/sections/recipes" +) + +// trendingFetchCap is the over-fetch size for each refresh. Library-only +// matching drops globally-trending titles the server does not own, so we fetch +// well beyond any section's display limit and store the matched, ordered list. +const trendingFetchCap = 200 + +// trendingSectionConfigLister enumerates enabled trending_discover section +// configs. Satisfied by *Repository. +type trendingSectionConfigLister interface { + ListTrendingDiscoverConfigs(ctx context.Context) ([]json.RawMessage, error) +} + +// trendingSnapshotStore is the write side of the snapshot table. Satisfied by +// *TrendingSnapshotRepository. +type trendingSnapshotStore interface { + SaveSuccess(ctx context.Context, source, window string, contentIDs []string, entryCount int, status string, at time.Time) error + RecordAttempt(ctx context.Context, source, window, status, message string, at time.Time) error +} + +// trendingExternalIDResolver resolves external IDs to library content IDs. +// Satisfied by *catalog.ItemRepository. +type trendingExternalIDResolver interface { + GetByExternalIDs(ctx context.Context, batch catalog.ExternalIDBatch, itemType string) (*catalog.ExternalIDLookup, error) +} + +// TrendingRefresher fetches external global trending (TMDB/Trakt), resolves it +// to library content IDs, and persists one snapshot per canonical +// (source, window). It is driven by a TaskManager task on an interval. +type TrendingRefresher struct { + Sections trendingSectionConfigLister + Snapshots trendingSnapshotStore + Resolver trendingExternalIDResolver + TMDBTrending catalog.TMDBCollectionFetcher + TraktTrending catalog.TraktCollectionFetcher + + // Clock defaults to recipes.RealClock{}. Tests inject recipes.FixedClock. + Clock recipes.Clock + logger *slog.Logger +} + +// NewTrendingRefresher creates a refresher with real-clock and default logger. +func NewTrendingRefresher( + sectionsRepo trendingSectionConfigLister, + snapshots trendingSnapshotStore, + resolver trendingExternalIDResolver, + tmdb catalog.TMDBCollectionFetcher, + trakt catalog.TraktCollectionFetcher, +) *TrendingRefresher { + return &TrendingRefresher{ + Sections: sectionsRepo, + Snapshots: snapshots, + Resolver: resolver, + TMDBTrending: tmdb, + TraktTrending: trakt, + Clock: recipes.RealClock{}, + logger: slog.Default(), + } +} + +func (r *TrendingRefresher) now() time.Time { + if r.Clock != nil { + return r.Clock.Now() + } + return time.Now() +} + +func (r *TrendingRefresher) log() *slog.Logger { + if r.logger != nil { + return r.logger + } + return slog.Default() +} + +// TrendingRefreshResult is the JSON summary attached to the task execution. +type TrendingRefreshResult struct { + Combos int `json:"combos"` + Refreshed int `json:"refreshed"` + Empty int `json:"empty"` + Failed int `json:"failed"` +} + +type trendingCombo struct { + source string + window string +} + +// distinctTrendingCombos parses section configs and returns the deduplicated set +// of canonical (source, window) pairs that need a snapshot. +func distinctTrendingCombos(configs []json.RawMessage) []trendingCombo { + seen := make(map[trendingCombo]struct{}, len(configs)) + out := make([]trendingCombo, 0, len(configs)) + for _, raw := range configs { + var p recipes.TrendingDiscoverParams + if len(raw) > 0 { + _ = json.Unmarshal(raw, &p) + } + source, window := canonicalTrendingKey(p.Source, p.Window) + c := trendingCombo{source: source, window: window} + if _, ok := seen[c]; ok { + continue + } + seen[c] = struct{}{} + out = append(out, c) + } + return out +} + +// RunOnce refreshes every (source, window) used by an enabled trending_discover +// section. Per-combo failures are recorded and never abort the others. The JSON +// summary is suitable for task result data. +func (r *TrendingRefresher) RunOnce(ctx context.Context) (json.RawMessage, error) { + configs, err := r.Sections.ListTrendingDiscoverConfigs(ctx) + if err != nil { + return nil, fmt.Errorf("listing trending_discover sections: %w", err) + } + + combos := distinctTrendingCombos(configs) + result := TrendingRefreshResult{Combos: len(combos)} + for _, c := range combos { + switch r.refreshCombo(ctx, c.source, c.window) { + case "ok": + result.Refreshed++ + case "empty": + result.Empty++ + default: + result.Failed++ + } + } + + data, _ := json.Marshal(result) + return data, nil +} + +// refreshCombo refreshes a single canonical (source, window) and returns its +// outcome: "ok", "empty", or "error". A fetch failure or an unconfigured/empty +// provider preserves the last-good content list (RecordAttempt). When the +// provider returns entries, the list is replaced even if nothing matched the +// catalog ("empty" status with an empty list) — that genuinely reflects current +// trending having no library matches. +func (r *TrendingRefresher) refreshCombo(ctx context.Context, source, window string) string { + now := r.now() + + entries, err := r.fetchEntries(ctx, source, window, trendingFetchCap) + if err != nil { + r.log().Error("trending refresh: fetch failed", "source", source, "window", window, "error", err) + _ = r.Snapshots.RecordAttempt(ctx, source, window, "error", err.Error(), now) + return "error" + } + if len(entries) == 0 { + // Provider unconfigured or returned nothing: keep last-good, mark empty. + _ = r.Snapshots.RecordAttempt(ctx, source, window, "empty", "", now) + return "empty" + } + + contentIDs, err := r.resolveIDs(ctx, entries) + if err != nil { + r.log().Error("trending refresh: resolve failed", "source", source, "window", window, "error", err) + _ = r.Snapshots.RecordAttempt(ctx, source, window, "error", err.Error(), now) + return "error" + } + + status := "ok" + if len(contentIDs) == 0 { + status = "empty" + } + if err := r.Snapshots.SaveSuccess(ctx, source, window, contentIDs, len(entries), status, now); err != nil { + r.log().Error("trending refresh: save failed", "source", source, "window", window, "error", err) + return "error" + } + return status +} + +// fetchEntries pulls the raw trending list from the configured provider. A +// nil/unconfigured provider yields an empty list (no error). +func (r *TrendingRefresher) fetchEntries(ctx context.Context, source, window string, fetchLimit int) ([]trendingDiscoverEntry, error) { + if source == "trakt" { + if r.TraktTrending == nil { + return nil, nil + } + movies, movieErr := r.TraktTrending.GetCollectionPreset(ctx, "trending", "movie", fetchLimit, "") + shows, showErr := r.TraktTrending.GetCollectionPreset(ctx, "trending", "tv", fetchLimit, "") + if movieErr != nil && showErr != nil { + return nil, fmt.Errorf("trakt trending: %v / %v", movieErr, showErr) + } + out := make([]trendingDiscoverEntry, 0, len(movies)+len(shows)) + for _, e := range movies { + out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) + } + for _, e := range shows { + out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) + } + return out, nil + } + + if r.TMDBTrending == nil { + return nil, nil + } + entries, err := r.TMDBTrending.GetCollectionPreset(ctx, "trending", "all", window, fetchLimit) + if err != nil { + return nil, err + } + out := make([]trendingDiscoverEntry, 0, len(entries)) + for _, e := range entries { + out = append(out, newTrendingEntry(e.ID, e.TVDBID, e.IMDbID, e.MediaType)) + } + return out, nil +} + +// resolveIDs matches trending entries to library content IDs via two batched +// external-ID lookups (movies, series), preserving trending order. +func (r *TrendingRefresher) resolveIDs(ctx context.Context, entries []trendingDiscoverEntry) ([]string, error) { + if r.Resolver == nil { + return nil, fmt.Errorf("trending_discover: external ID resolver not configured") + } + var movieBatch, seriesBatch catalog.ExternalIDBatch + for _, e := range entries { + batch := &movieBatch + if e.mediaType == "tv" { + batch = &seriesBatch + } + if e.tmdbID != "" { + batch.TMDBIDs = append(batch.TMDBIDs, e.tmdbID) + } + if e.imdbID != "" { + batch.IMDbIDs = append(batch.IMDbIDs, e.imdbID) + } + if e.tvdbID != "" { + batch.TVDBIDs = append(batch.TVDBIDs, e.tvdbID) + } + } + movieLookup, err := r.Resolver.GetByExternalIDs(ctx, movieBatch, "movie") + if err != nil { + return nil, err + } + seriesLookup, err := r.Resolver.GetByExternalIDs(ctx, seriesBatch, "series") + if err != nil { + return nil, err + } + return orderedTrendingContentIDs(entries, movieLookup, seriesLookup), nil +} diff --git a/internal/sections/trending_refresher_test.go b/internal/sections/trending_refresher_test.go new file mode 100644 index 00000000..4ef6bb89 --- /dev/null +++ b/internal/sections/trending_refresher_test.go @@ -0,0 +1,197 @@ +package sections + +import ( + "context" + "encoding/json" + "errors" + "testing" + "time" + + "github.com/Silo-Server/silo-server/internal/catalog" + "github.com/Silo-Server/silo-server/internal/sections/recipes" +) + +type fakeSectionLister struct { + configs []json.RawMessage + err error +} + +func (f fakeSectionLister) ListTrendingDiscoverConfigs(context.Context) ([]json.RawMessage, error) { + return f.configs, f.err +} + +type savedSnap struct { + contentIDs []string + entryCount int + status string +} + +type attemptRec struct { + status string + message string +} + +type fakeSnapshotStore struct { + saved map[string]savedSnap + attempts map[string]attemptRec +} + +func newFakeSnapshotStore() *fakeSnapshotStore { + return &fakeSnapshotStore{saved: map[string]savedSnap{}, attempts: map[string]attemptRec{}} +} + +func (f *fakeSnapshotStore) SaveSuccess(_ context.Context, source, window string, contentIDs []string, entryCount int, status string, _ time.Time) error { + f.saved[source+"|"+window] = savedSnap{contentIDs: contentIDs, entryCount: entryCount, status: status} + return nil +} + +func (f *fakeSnapshotStore) RecordAttempt(_ context.Context, source, window, status, message string, _ time.Time) error { + f.attempts[source+"|"+window] = attemptRec{status: status, message: message} + return nil +} + +type fakeTMDB struct { + entries []catalog.TMDBCollectionEntry + err error +} + +func (f fakeTMDB) GetCollectionPreset(context.Context, string, string, string, int) ([]catalog.TMDBCollectionEntry, error) { + return f.entries, f.err +} + +type fakeResolver struct { + byType map[string]*catalog.ExternalIDLookup +} + +func (f fakeResolver) GetByExternalIDs(_ context.Context, _ catalog.ExternalIDBatch, itemType string) (*catalog.ExternalIDLookup, error) { + if lk, ok := f.byType[itemType]; ok { + return lk, nil + } + return &catalog.ExternalIDLookup{ByTMDB: map[string]string{}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, nil +} + +func tmdbConfig(t *testing.T, source, window string) json.RawMessage { + t.Helper() + raw, err := json.Marshal(recipes.TrendingDiscoverParams{Source: source, Window: window}) + if err != nil { + t.Fatalf("marshal config: %v", err) + } + return raw +} + +func TestRefresherSavesOrderedContentIDs(t *testing.T) { + store := newFakeSnapshotStore() + r := &TrendingRefresher{ + Sections: fakeSectionLister{configs: []json.RawMessage{tmdbConfig(t, "tmdb", "week")}}, + Snapshots: store, + Resolver: fakeResolver{byType: map[string]*catalog.ExternalIDLookup{ + "movie": {ByTMDB: map[string]string{"10": "c-movie"}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, + "series": {ByTMDB: map[string]string{"20": "c-series"}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, + }}, + TMDBTrending: fakeTMDB{entries: []catalog.TMDBCollectionEntry{ + {ID: 10, MediaType: "movie"}, + {ID: 20, MediaType: "tv"}, + }}, + Clock: recipes.FixedClock(time.Date(2026, 5, 29, 12, 0, 0, 0, time.UTC)), + } + + data, err := r.RunOnce(context.Background()) + if err != nil { + t.Fatalf("RunOnce: %v", err) + } + + var result TrendingRefreshResult + if err := json.Unmarshal(data, &result); err != nil { + t.Fatalf("unmarshal result: %v", err) + } + if result.Combos != 1 || result.Refreshed != 1 || result.Failed != 0 || result.Empty != 0 { + t.Fatalf("result = %+v; want {Combos:1 Refreshed:1 Empty:0 Failed:0}", result) + } + + got := store.saved["tmdb|week"] + want := []string{"c-movie", "c-series"} + if len(got.contentIDs) != len(want) || got.contentIDs[0] != want[0] || got.contentIDs[1] != want[1] { + t.Fatalf("saved content IDs = %v; want %v", got.contentIDs, want) + } + if got.status != "ok" || got.entryCount != 2 { + t.Fatalf("saved snap = %+v; want status ok, entryCount 2", got) + } +} + +func TestRefresherFailurePreservesLastGood(t *testing.T) { + store := newFakeSnapshotStore() + r := &TrendingRefresher{ + Sections: fakeSectionLister{configs: []json.RawMessage{tmdbConfig(t, "tmdb", "week")}}, + Snapshots: store, + Resolver: fakeResolver{}, + TMDBTrending: fakeTMDB{err: errors.New("tmdb 503")}, + Clock: recipes.FixedClock(time.Date(2026, 5, 29, 12, 0, 0, 0, time.UTC)), + } + + data, err := r.RunOnce(context.Background()) + if err != nil { + t.Fatalf("RunOnce: %v", err) + } + + if _, ok := store.saved["tmdb|week"]; ok { + t.Fatal("SaveSuccess must not be called on fetch failure (would clear last-good)") + } + att, ok := store.attempts["tmdb|week"] + if !ok || att.status != "error" { + t.Fatalf("attempt = %+v, ok=%v; want status error", att, ok) + } + + var result TrendingRefreshResult + _ = json.Unmarshal(data, &result) + if result.Failed != 1 { + t.Fatalf("result.Failed = %d; want 1", result.Failed) + } +} + +func TestRefresherEmptyProviderPreservesLastGood(t *testing.T) { + store := newFakeSnapshotStore() + r := &TrendingRefresher{ + Sections: fakeSectionLister{configs: []json.RawMessage{tmdbConfig(t, "tmdb", "week")}}, + Snapshots: store, + Resolver: fakeResolver{}, + // TMDBTrending nil => provider unconfigured => empty entries, no error. + Clock: recipes.FixedClock(time.Date(2026, 5, 29, 12, 0, 0, 0, time.UTC)), + } + + data, err := r.RunOnce(context.Background()) + if err != nil { + t.Fatalf("RunOnce: %v", err) + } + if _, ok := store.saved["tmdb|week"]; ok { + t.Fatal("SaveSuccess must not be called when provider returns no entries") + } + att := store.attempts["tmdb|week"] + if att.status != "empty" { + t.Fatalf("attempt status = %q; want empty", att.status) + } + var result TrendingRefreshResult + _ = json.Unmarshal(data, &result) + if result.Empty != 1 { + t.Fatalf("result.Empty = %d; want 1", result.Empty) + } +} + +func TestDistinctTrendingCombosCollapsesTrakt(t *testing.T) { + configs := []json.RawMessage{ + tmdbConfig(t, "trakt", "day"), + tmdbConfig(t, "trakt", "week"), + tmdbConfig(t, "tmdb", "day"), + tmdbConfig(t, "tmdb", "day"), + } + got := distinctTrendingCombos(configs) + if len(got) != 2 { + t.Fatalf("distinctTrendingCombos len = %d (%+v); want 2", len(got), got) + } + seen := map[trendingCombo]bool{} + for _, c := range got { + seen[c] = true + } + if !seen[trendingCombo{"trakt", "week"}] || !seen[trendingCombo{"tmdb", "day"}] { + t.Fatalf("combos = %+v; want {trakt week} and {tmdb day}", got) + } +} From 3c3ef82ae817ec293a128c959c85a04fd3094d57 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:44:15 -0400 Subject: [PATCH 09/14] feat(tasks): add refresh_trending_discover task --- .../tasks/refresh_trending_discover.go | 60 +++++++++++++++++++ 1 file changed, 60 insertions(+) create mode 100644 internal/taskmanager/tasks/refresh_trending_discover.go diff --git a/internal/taskmanager/tasks/refresh_trending_discover.go b/internal/taskmanager/tasks/refresh_trending_discover.go new file mode 100644 index 00000000..7324d517 --- /dev/null +++ b/internal/taskmanager/tasks/refresh_trending_discover.go @@ -0,0 +1,60 @@ +package tasks + +import ( + "context" + "encoding/json" + "fmt" + + "github.com/Silo-Server/silo-server/internal/taskmanager" +) + +// TrendingDiscoverRefresher runs a single pass of the trending refresh. +// Satisfied by *sections.TrendingRefresher. +type TrendingDiscoverRefresher interface { + RunOnce(ctx context.Context) (json.RawMessage, error) +} + +// RefreshTrendingDiscoverTask refreshes the persisted external-trending +// snapshots used by trending_discover home sections. +type RefreshTrendingDiscoverTask struct { + refresher TrendingDiscoverRefresher +} + +// NewRefreshTrendingDiscoverTask creates a new RefreshTrendingDiscoverTask. +func NewRefreshTrendingDiscoverTask(refresher TrendingDiscoverRefresher) *RefreshTrendingDiscoverTask { + return &RefreshTrendingDiscoverTask{refresher: refresher} +} + +func (t *RefreshTrendingDiscoverTask) Key() string { return "refresh_trending_discover" } +func (t *RefreshTrendingDiscoverTask) Name() string { return "Refresh Trending Discover" } +func (t *RefreshTrendingDiscoverTask) Description() string { + return "Refreshes the persisted external trending list (TMDB/Trakt) for the Trending Discover home section" +} + +func (t *RefreshTrendingDiscoverTask) Category() taskmanager.TaskCategory { + return taskmanager.TaskCategoryLibrary +} + +func (t *RefreshTrendingDiscoverTask) IsHidden() bool { return false } + +func (t *RefreshTrendingDiscoverTask) DefaultTriggers() []taskmanager.TriggerConfig { + return []taskmanager.TriggerConfig{ + {Type: taskmanager.TriggerTypeStartup}, + {Type: taskmanager.TriggerTypeInterval, IntervalMs: 60 * 60 * 1000}, // hourly + } +} + +func (t *RefreshTrendingDiscoverTask) Execute(ctx context.Context, progress taskmanager.ProgressReporter) error { + progress.Report(0, "Refreshing trending discover") + + resultData, err := t.refresher.RunOnce(ctx) + if err != nil { + return fmt.Errorf("trending discover refresh: %w", err) + } + if resultData != nil { + progress.SetResultData(resultData) + } + + progress.Report(100, "Trending discover refresh complete") + return nil +} From 38723c8b8161d4fe43fdfb5b092abd3b7881efc1 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:46:01 -0400 Subject: [PATCH 10/14] refactor(sections): read trending_discover from persisted snapshot --- internal/sections/fetcher.go | 155 ++++-------------------- internal/sections/trending_read_test.go | 52 ++++++++ 2 files changed, 75 insertions(+), 132 deletions(-) create mode 100644 internal/sections/trending_read_test.go diff --git a/internal/sections/fetcher.go b/internal/sections/fetcher.go index 6e87fab1..80bf526c 100644 --- a/internal/sections/fetcher.go +++ b/internal/sections/fetcher.go @@ -59,6 +59,12 @@ type recommendationReader interface { GetTasteMatchRow(ctx context.Context, userID int, profileID, genre string, limit int, filter catalog.AccessFilter) (*recommendations.ForYouRow, error) } +// trendingSnapshotGetter is the read side of the trending snapshot table. +// Satisfied by *TrendingSnapshotRepository. +type trendingSnapshotGetter interface { + Get(ctx context.Context, source, window string) (TrendingSnapshot, bool, error) +} + // Fetcher runs section queries against the database. type Fetcher struct { pool *pgxpool.Pool @@ -68,13 +74,11 @@ type Fetcher struct { RecommendationReader recommendationReader NextUpRepo *catalog.NextUpRepository - // ItemRepo resolves external IDs (TMDB/Trakt) to library content IDs for - // the trending_discover section. Nil disables external-trending matching. - ItemRepo *catalog.ItemRepository - // TMDBTrending and TraktTrending fetch external global trending lists. Each - // is nil when that provider is not configured. - TMDBTrending catalog.TMDBCollectionFetcher - TraktTrending catalog.TraktCollectionFetcher + // TrendingSnapshots reads the persisted external-trending snapshots that + // back the trending_discover section. Nil renders that section empty. + // Snapshots are produced out-of-band by TrendingRefresher, so the read path + // never calls the upstream provider. + TrendingSnapshots trendingSnapshotGetter candidateCacheMu sync.Mutex candidateCache *editorialCandidateCache @@ -2375,30 +2379,14 @@ func (f *Fetcher) fetchTrendingDiscover(ctx context.Context, s ResolvedSection, if len(s.Config) > 0 { _ = json.Unmarshal(s.Config, &p) } - source := p.Source - if source != "trakt" { - source = "tmdb" - } - window := p.Window - if window != "day" { - window = "week" - } + source, window := canonicalTrendingKey(p.Source, p.Window) limit := s.ItemLimit if limit <= 0 { limit = 20 } - // Over-fetch: library-only matching drops globally-trending titles the - // server does not own, so request more candidates than the display limit. - fetchLimit := limit * 5 - if fetchLimit < 50 { - fetchLimit = 50 - } - if fetchLimit > 200 { - fetchLimit = 200 - } - orderedIDs, err := f.loadTrendingDiscoverContentIDs(ctx, source, window, fetchLimit) + orderedIDs, err := f.loadTrendingDiscoverContentIDs(ctx, source, window) if err != nil { return nil, 0, err } @@ -2420,119 +2408,22 @@ func (f *Fetcher) fetchTrendingDiscover(ctx context.Context, s ResolvedSection, return ordered, len(ordered), nil } -// loadTrendingDiscoverContentIDs fetches the external trending list and resolves -// it to ordered library content IDs, cached for an hour per (source, window). -func (f *Fetcher) loadTrendingDiscoverContentIDs(ctx context.Context, source, window string, fetchLimit int) ([]string, error) { - cache := f.ensureEditorialCandidateCache() - cacheKey := fmt.Sprintf("trending_discover|%s|%s|%d", source, window, fetchLimit) - if cached, ok := cache.get(cacheKey, f.now()); ok { - return cached, nil - } - - // Collapse concurrent cache-miss loads (e.g. many home-page requests at TTL - // expiry) into a single upstream fetch + ID resolution, as - // cachedEditorialCandidates does for editorial sections. - value, err, _ := f.candidateGroup.Do(cacheKey, func() (any, error) { - now := f.now() - if cached, ok := cache.get(cacheKey, now); ok { - return cached, nil - } - - entries, err := f.fetchTrendingDiscoverEntries(ctx, source, window, fetchLimit) - if err != nil { - return nil, err - } - // Don't cache when the provider is unconfigured/empty, so a - // newly-configured provider takes effect immediately rather than after - // the TTL. - if len(entries) == 0 { - return []string{}, nil - } - - contentIDs, err := f.resolveTrendingDiscoverIDs(ctx, entries) - if err != nil { - return nil, err - } - cache.set(cacheKey, contentIDs, now.Add(time.Hour)) - return contentIDs, nil - }) - if err != nil { - return nil, err - } - contentIDs, _ := value.([]string) - return append([]string(nil), contentIDs...), nil -} - -// fetchTrendingDiscoverEntries pulls the raw trending list from the configured -// provider. A nil/unconfigured provider yields an empty list (no error) so the -// section simply renders empty. -func (f *Fetcher) fetchTrendingDiscoverEntries(ctx context.Context, source, window string, fetchLimit int) ([]trendingDiscoverEntry, error) { - if source == "trakt" { - if f.TraktTrending == nil { - return nil, nil - } - // Trakt has no mixed endpoint; fetch movies + shows and concatenate. - movies, movieErr := f.TraktTrending.GetCollectionPreset(ctx, "trending", "movie", fetchLimit, "") - shows, showErr := f.TraktTrending.GetCollectionPreset(ctx, "trending", "tv", fetchLimit, "") - if movieErr != nil && showErr != nil { - return nil, fmt.Errorf("trakt trending: %v / %v", movieErr, showErr) - } - out := make([]trendingDiscoverEntry, 0, len(movies)+len(shows)) - for _, e := range movies { - out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) - } - for _, e := range shows { - out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) - } - return out, nil - } - - // Default: TMDB /trending/all/{window} — natively mixed movies + series. - if f.TMDBTrending == nil { +// loadTrendingDiscoverContentIDs returns the persisted, catalog-resolved content +// IDs for the canonical (source, window). It reads only the snapshot table; the +// upstream fetch happens out-of-band in TrendingRefresher. Returns nil when no +// snapshot reader is configured or no snapshot exists yet. +func (f *Fetcher) loadTrendingDiscoverContentIDs(ctx context.Context, source, window string) ([]string, error) { + if f.TrendingSnapshots == nil { return nil, nil } - entries, err := f.TMDBTrending.GetCollectionPreset(ctx, "trending", "all", window, fetchLimit) + snap, found, err := f.TrendingSnapshots.Get(ctx, source, window) if err != nil { return nil, err } - out := make([]trendingDiscoverEntry, 0, len(entries)) - for _, e := range entries { - out = append(out, newTrendingEntry(e.ID, e.TVDBID, e.IMDbID, e.MediaType)) + if !found { + return nil, nil } - return out, nil -} - -// resolveTrendingDiscoverIDs matches trending entries to library content IDs via -// two batched external-ID lookups (movies, series), preserving trending order. -func (f *Fetcher) resolveTrendingDiscoverIDs(ctx context.Context, entries []trendingDiscoverEntry) ([]string, error) { - if f.ItemRepo == nil { - return nil, fmt.Errorf("trending_discover: item repository not configured") - } - var movieBatch, seriesBatch catalog.ExternalIDBatch - for _, e := range entries { - batch := &movieBatch - if e.mediaType == "tv" { - batch = &seriesBatch - } - if e.tmdbID != "" { - batch.TMDBIDs = append(batch.TMDBIDs, e.tmdbID) - } - if e.imdbID != "" { - batch.IMDbIDs = append(batch.IMDbIDs, e.imdbID) - } - if e.tvdbID != "" { - batch.TVDBIDs = append(batch.TVDBIDs, e.tvdbID) - } - } - movieLookup, err := f.ItemRepo.GetByExternalIDs(ctx, movieBatch, "movie") - if err != nil { - return nil, err - } - seriesLookup, err := f.ItemRepo.GetByExternalIDs(ctx, seriesBatch, "series") - if err != nil { - return nil, err - } - return orderedTrendingContentIDs(entries, movieLookup, seriesLookup), nil + return snap.ContentIDs, nil } // orderedTrendingContentIDs maps trending entries to library content IDs in diff --git a/internal/sections/trending_read_test.go b/internal/sections/trending_read_test.go new file mode 100644 index 00000000..6c1e97a7 --- /dev/null +++ b/internal/sections/trending_read_test.go @@ -0,0 +1,52 @@ +package sections + +import ( + "context" + "testing" +) + +type fakeSnapshotGetter struct { + snap TrendingSnapshot + found bool + err error +} + +func (f fakeSnapshotGetter) Get(context.Context, string, string) (TrendingSnapshot, bool, error) { + return f.snap, f.found, f.err +} + +func TestLoadTrendingDiscoverContentIDsReadsSnapshot(t *testing.T) { + f := &Fetcher{TrendingSnapshots: fakeSnapshotGetter{ + snap: TrendingSnapshot{ContentIDs: []string{"a", "b"}}, + found: true, + }} + ids, err := f.loadTrendingDiscoverContentIDs(context.Background(), "tmdb", "week") + if err != nil { + t.Fatalf("loadTrendingDiscoverContentIDs: %v", err) + } + if len(ids) != 2 || ids[0] != "a" || ids[1] != "b" { + t.Fatalf("ids = %v; want [a b]", ids) + } +} + +func TestLoadTrendingDiscoverContentIDsNilGetter(t *testing.T) { + f := &Fetcher{} + ids, err := f.loadTrendingDiscoverContentIDs(context.Background(), "tmdb", "week") + if err != nil { + t.Fatalf("loadTrendingDiscoverContentIDs: %v", err) + } + if ids != nil { + t.Fatalf("ids = %v; want nil for nil getter", ids) + } +} + +func TestLoadTrendingDiscoverContentIDsNotFound(t *testing.T) { + f := &Fetcher{TrendingSnapshots: fakeSnapshotGetter{found: false}} + ids, err := f.loadTrendingDiscoverContentIDs(context.Background(), "tmdb", "week") + if err != nil { + t.Fatalf("loadTrendingDiscoverContentIDs: %v", err) + } + if ids != nil { + t.Fatalf("ids = %v; want nil when no snapshot exists", ids) + } +} From 9d03ad0b37ad5285dae23f4fbe727e820a3f1a39 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:49:01 -0400 Subject: [PATCH 11/14] feat: wire trending refresh task and snapshot reader --- cmd/silo/main.go | 17 +++++++++++++++++ internal/api/router.go | 21 ++++++++++++++++----- 2 files changed, 33 insertions(+), 5 deletions(-) diff --git a/cmd/silo/main.go b/cmd/silo/main.go index 1305fed6..907145fb 100644 --- a/cmd/silo/main.go +++ b/cmd/silo/main.go @@ -1234,6 +1234,7 @@ func main() { // Construct collection service for both the router and the collection sync scheduler. var collectionSyncScheduler *catalog.CollectionSyncScheduler var userCollectionScheduler *usercollections.Scheduler + var trendingRefresher *sections.TrendingRefresher if needsWorkers && deps.DB != nil { collectionRepo := catalog.NewLibraryCollectionRepository(deps.DB) collItemRepo := catalog.NewItemRepository(deps.DB) @@ -1243,6 +1244,19 @@ func main() { deps.CollectionService = collectionService collectionSyncScheduler = catalog.NewCollectionSyncScheduler(collectionRepo, collectionService, slog.Default()) + // The trending refresher reuses the section repo (to find used source/ + // window combos), a snapshot repo, an item repo (external-ID matching), + // and the TMDB fetcher. The Trakt fetcher needs settingsRepo and is + // propagated onto deps.TrendingRefresher later in router.go. + trendingRefresher = sections.NewTrendingRefresher( + sectionRepo, + sections.NewTrendingSnapshotRepository(pool), + catalog.NewItemRepository(deps.DB), + collectionService.TMDBCollections, + collectionService.TraktCollections, + ) + deps.TrendingRefresher = trendingRefresher + if deps.UserStoreProvider != nil { userSync := usercollections.NewService(deps.UserStoreProvider, collItemRepo, libraryItemRepo, nil, slog.Default()) userSync.TMDBCollections = collectionService.TMDBCollections @@ -1287,6 +1301,9 @@ func main() { if collectionSyncScheduler != nil { taskMgr.Register(tasks.NewSyncCollectionsTask(collectionSyncScheduler)) } + if trendingRefresher != nil { + taskMgr.Register(tasks.NewRefreshTrendingDiscoverTask(trendingRefresher)) + } if userCollectionScheduler != nil { taskMgr.Register(tasks.NewSyncUserCollectionsTask(userCollectionScheduler)) } diff --git a/internal/api/router.go b/internal/api/router.go index 9125eedf..9888ea68 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -141,6 +141,11 @@ type Dependencies struct { UserCollectionSync *usercollections.Service UserCollectionScheduler *usercollections.Scheduler + // TrendingRefresher refreshes the persisted trending_discover snapshots. + // Built in main.go with TMDB wired; its Trakt fetcher is propagated here in + // router.go once the Trakt adapter exists (mirrors UserCollectionSync). + TrendingRefresher *sections.TrendingRefresher + // MDBListClient is used by user-facing list discovery endpoints // (search/top). May be nil; the handlers report "not configured" in // that case rather than failing. @@ -878,11 +883,17 @@ func NewRouter(deps Dependencies) chi.Router { } } - // Wire external-trending fetchers into the section fetcher so the - // trending_discover home section can pull TMDB/Trakt trending. - sectionFetcher.ItemRepo = itemRepo - sectionFetcher.TMDBTrending = libraryCollectionService.TMDBCollections - sectionFetcher.TraktTrending = libraryCollectionService.TraktCollections + // Propagate the now-wired Trakt fetcher to the trending refresher (built + // in main.go with TMDB only, before the Trakt adapter existed). + if deps.TrendingRefresher != nil && deps.TrendingRefresher.TraktTrending == nil { + deps.TrendingRefresher.TraktTrending = libraryCollectionService.TraktCollections + } + + // Wire the trending snapshot reader into the section fetcher. The + // trending_discover home section reads its list from the persisted + // snapshot table; the upstream fetch happens out-of-band in the refresh + // task, so the read path never calls the provider. + sectionFetcher.TrendingSnapshots = sections.NewTrendingSnapshotRepository(deps.DB) libraryCollectionHandler = handlers.NewLibraryCollectionHandler( libraryCollectionRepo, From 4455c2cd4080b543972451af1ecab70d5d2dda97 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:51:54 -0400 Subject: [PATCH 12/14] chore(sections): satisfy lint (wrap trakt errors, lift source/window constants) --- internal/sections/trending_refresher.go | 4 ++-- internal/sections/trending_snapshot.go | 22 +++++++++++++++------- 2 files changed, 17 insertions(+), 9 deletions(-) diff --git a/internal/sections/trending_refresher.go b/internal/sections/trending_refresher.go index b5bc92d0..245aa859 100644 --- a/internal/sections/trending_refresher.go +++ b/internal/sections/trending_refresher.go @@ -185,14 +185,14 @@ func (r *TrendingRefresher) refreshCombo(ctx context.Context, source, window str // fetchEntries pulls the raw trending list from the configured provider. A // nil/unconfigured provider yields an empty list (no error). func (r *TrendingRefresher) fetchEntries(ctx context.Context, source, window string, fetchLimit int) ([]trendingDiscoverEntry, error) { - if source == "trakt" { + if source == sourceTrakt { if r.TraktTrending == nil { return nil, nil } movies, movieErr := r.TraktTrending.GetCollectionPreset(ctx, "trending", "movie", fetchLimit, "") shows, showErr := r.TraktTrending.GetCollectionPreset(ctx, "trending", "tv", fetchLimit, "") if movieErr != nil && showErr != nil { - return nil, fmt.Errorf("trakt trending: %v / %v", movieErr, showErr) + return nil, fmt.Errorf("trakt trending: %w / %w", movieErr, showErr) } out := make([]trendingDiscoverEntry, 0, len(movies)+len(shows)) for _, e := range movies { diff --git a/internal/sections/trending_snapshot.go b/internal/sections/trending_snapshot.go index 033d3b19..d755e6ef 100644 --- a/internal/sections/trending_snapshot.go +++ b/internal/sections/trending_snapshot.go @@ -25,22 +25,30 @@ type TrendingSnapshot struct { LastError string } +// Canonical trending source and window values used as snapshot keys. +const ( + sourceTMDB = "tmdb" + sourceTrakt = "trakt" + windowDay = "day" + windowWeek = "week" +) + // canonicalTrendingKey normalizes a section's configured source/window into the // snapshot key space. Source is "trakt" only when explicitly set; everything // else collapses to "tmdb". Trakt ignores the time window, so it is pinned to // "week" to avoid duplicate identical rows. For TMDB, "day" is honored only // when explicitly set; anything else is "week". func canonicalTrendingKey(source, window string) (string, string) { - if source != "trakt" { - source = "tmdb" + if source != sourceTrakt { + source = sourceTMDB } - if source == "trakt" { - return "trakt", "week" + if source == sourceTrakt { + return sourceTrakt, windowWeek } - if window != "day" { - window = "week" + if window != windowDay { + window = windowWeek } - return "tmdb", window + return sourceTMDB, window } // TrendingSnapshotRepository persists and reads trending_discover_snapshots. From c98f828c25932c11e0e12ee7582728fa426e8e15 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 10:57:57 -0400 Subject: [PATCH 13/14] fix(migrations): renumber trending_discover_snapshots 166 -> 167 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The shared dev DB already recorded version 166 (166_trending_blend_collection_type from another branch), so the integer-version migration runner silently skipped our 166 and the table was never created — the trending section errored out empty. 167 is the next free version. Co-Authored-By: Claude Opus 4.8 (1M context) --- ...26-05-29-trending-discover-persistent-snapshot-design.md | 6 ++++-- ...ts.down.sql => 167_trending_discover_snapshots.down.sql} | 0 ...pshots.up.sql => 167_trending_discover_snapshots.up.sql} | 0 3 files changed, 4 insertions(+), 2 deletions(-) rename migrations/{166_trending_discover_snapshots.down.sql => 167_trending_discover_snapshots.down.sql} (100%) rename migrations/{166_trending_discover_snapshots.up.sql => 167_trending_discover_snapshots.up.sql} (100%) diff --git a/docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md b/docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md index 84479bba..afb05454 100644 --- a/docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md +++ b/docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md @@ -57,8 +57,10 @@ The change cleanly separates **read** from **refresh**: ## Data model -New migration `166_trending_discover_snapshots` (next free number is 166; 165 is -the current max). +New migration `167_trending_discover_snapshots`. (Originally numbered 166, but +the shared dev DB had already recorded version 166 from another branch's +`166_trending_blend_collection_type`; the migration runner dedupes by integer +version, so 166 was silently skipped. Renumbered to 167, the next free version.) ``` trending_discover_snapshots diff --git a/migrations/166_trending_discover_snapshots.down.sql b/migrations/167_trending_discover_snapshots.down.sql similarity index 100% rename from migrations/166_trending_discover_snapshots.down.sql rename to migrations/167_trending_discover_snapshots.down.sql diff --git a/migrations/166_trending_discover_snapshots.up.sql b/migrations/167_trending_discover_snapshots.up.sql similarity index 100% rename from migrations/166_trending_discover_snapshots.up.sql rename to migrations/167_trending_discover_snapshots.up.sql From f10536165842f231577dc8fbf8eeb654d02c8d55 Mon Sep 17 00:00:00 2001 From: Quick <31828688+Quick104@users.noreply.github.com> Date: Fri, 29 May 2026 11:28:59 -0400 Subject: [PATCH 14/14] fix(sections): harden trending refresher per PR review - Interleave Trakt movies/shows by rank so the mixed row shows both types instead of burying all series past the display limit. - Treat any Trakt sub-fetch failure as fatal (errors.Join) so a partial result never overwrites the last-good snapshot with a media type missing. - Skip non-title entries (TMDB trending/all returns media_type "person") in both ID batching and ordering so they can't match an unrelated library title. - Guard the refresh task against a nil refresher. - Tests: person skip, Trakt interleave, Trakt partial-failure preserves last-good, snapshot read error propagation. Co-Authored-By: Claude Opus 4.8 (1M context) --- internal/sections/fetcher.go | 11 ++- internal/sections/trending_read_test.go | 13 +++ internal/sections/trending_refresher.go | 49 +++++++--- internal/sections/trending_refresher_test.go | 96 +++++++++++++++++++ .../tasks/refresh_trending_discover.go | 5 + 5 files changed, 160 insertions(+), 14 deletions(-) diff --git a/internal/sections/fetcher.go b/internal/sections/fetcher.go index 80bf526c..279f809f 100644 --- a/internal/sections/fetcher.go +++ b/internal/sections/fetcher.go @@ -2432,10 +2432,17 @@ func orderedTrendingContentIDs(entries []trendingDiscoverEntry, movieLookup, ser seen := make(map[string]struct{}, len(entries)) out := make([]string, 0, len(entries)) for _, e := range entries { - lookup := movieLookup + // Only movies and series are resolvable; skip anything else (TMDB + // trending/all also returns media_type "person"). + var lookup *catalog.ExternalIDLookup isSeries := e.mediaType == "tv" - if isSeries { + switch e.mediaType { + case "movie": + lookup = movieLookup + case "tv": lookup = seriesLookup + default: + continue } if lookup == nil { continue diff --git a/internal/sections/trending_read_test.go b/internal/sections/trending_read_test.go index 6c1e97a7..5743405d 100644 --- a/internal/sections/trending_read_test.go +++ b/internal/sections/trending_read_test.go @@ -2,6 +2,7 @@ package sections import ( "context" + "errors" "testing" ) @@ -50,3 +51,15 @@ func TestLoadTrendingDiscoverContentIDsNotFound(t *testing.T) { t.Fatalf("ids = %v; want nil when no snapshot exists", ids) } } + +func TestLoadTrendingDiscoverContentIDsPropagatesError(t *testing.T) { + boom := errors.New("boom") + f := &Fetcher{TrendingSnapshots: fakeSnapshotGetter{err: boom}} + ids, err := f.loadTrendingDiscoverContentIDs(context.Background(), "tmdb", "week") + if !errors.Is(err, boom) { + t.Fatalf("err = %v; want boom", err) + } + if ids != nil { + t.Fatalf("ids = %v; want nil on error", ids) + } +} diff --git a/internal/sections/trending_refresher.go b/internal/sections/trending_refresher.go index 245aa859..88a00892 100644 --- a/internal/sections/trending_refresher.go +++ b/internal/sections/trending_refresher.go @@ -3,6 +3,7 @@ package sections import ( "context" "encoding/json" + "errors" "fmt" "log/slog" "time" @@ -189,19 +190,18 @@ func (r *TrendingRefresher) fetchEntries(ctx context.Context, source, window str if r.TraktTrending == nil { return nil, nil } + // Trakt has no mixed endpoint; fetch movies + shows separately. Treat ANY + // failure as fatal so a partial result never overwrites the last-good + // snapshot with one media type missing. movies, movieErr := r.TraktTrending.GetCollectionPreset(ctx, "trending", "movie", fetchLimit, "") shows, showErr := r.TraktTrending.GetCollectionPreset(ctx, "trending", "tv", fetchLimit, "") - if movieErr != nil && showErr != nil { - return nil, fmt.Errorf("trakt trending: %w / %w", movieErr, showErr) + if movieErr != nil || showErr != nil { + return nil, fmt.Errorf("trakt trending: %w", errors.Join(movieErr, showErr)) } - out := make([]trendingDiscoverEntry, 0, len(movies)+len(shows)) - for _, e := range movies { - out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) - } - for _, e := range shows { - out = append(out, newTrendingEntry(e.TMDBID, e.TVDBID, e.IMDbID, e.MediaType)) - } - return out, nil + // Interleave by rank so the mixed row actually shows both movies and + // series; plain concatenation would bury all series past the display + // limit whenever enough movies match the library. + return interleaveTraktEntries(movies, shows), nil } if r.TMDBTrending == nil { @@ -218,6 +218,24 @@ func (r *TrendingRefresher) fetchEntries(ctx context.Context, source, window str return out, nil } +// interleaveTraktEntries alternates rank-ordered movies and shows so the mixed +// trending row surfaces both media types. Each input list is already in trending +// order; alternating preserves that order within each type while mixing them. +func interleaveTraktEntries(movies, shows []catalog.TraktCollectionEntry) []trendingDiscoverEntry { + out := make([]trendingDiscoverEntry, 0, len(movies)+len(shows)) + for i := 0; i < len(movies) || i < len(shows); i++ { + if i < len(movies) { + m := movies[i] + out = append(out, newTrendingEntry(m.TMDBID, m.TVDBID, m.IMDbID, m.MediaType)) + } + if i < len(shows) { + s := shows[i] + out = append(out, newTrendingEntry(s.TMDBID, s.TVDBID, s.IMDbID, s.MediaType)) + } + } + return out +} + // resolveIDs matches trending entries to library content IDs via two batched // external-ID lookups (movies, series), preserving trending order. func (r *TrendingRefresher) resolveIDs(ctx context.Context, entries []trendingDiscoverEntry) ([]string, error) { @@ -226,9 +244,16 @@ func (r *TrendingRefresher) resolveIDs(ctx context.Context, entries []trendingDi } var movieBatch, seriesBatch catalog.ExternalIDBatch for _, e := range entries { - batch := &movieBatch - if e.mediaType == "tv" { + var batch *catalog.ExternalIDBatch + switch e.mediaType { + case "movie": + batch = &movieBatch + case "tv": batch = &seriesBatch + default: + // Skip non-title entries (TMDB trending/all can return "person") so + // they never match an unrelated library title by shared external ID. + continue } if e.tmdbID != "" { batch.TMDBIDs = append(batch.TMDBIDs, e.tmdbID) diff --git a/internal/sections/trending_refresher_test.go b/internal/sections/trending_refresher_test.go index 4ef6bb89..2f9bef86 100644 --- a/internal/sections/trending_refresher_test.go +++ b/internal/sections/trending_refresher_test.go @@ -59,6 +59,18 @@ func (f fakeTMDB) GetCollectionPreset(context.Context, string, string, string, i return f.entries, f.err } +type fakeTrakt struct { + byMediaType map[string][]catalog.TraktCollectionEntry + errByType map[string]error +} + +func (f fakeTrakt) GetCollectionPreset(_ context.Context, _, mediaType string, _ int, _ string) ([]catalog.TraktCollectionEntry, error) { + if err := f.errByType[mediaType]; err != nil { + return nil, err + } + return f.byMediaType[mediaType], nil +} + type fakeResolver struct { byType map[string]*catalog.ExternalIDLookup } @@ -176,6 +188,90 @@ func TestRefresherEmptyProviderPreservesLastGood(t *testing.T) { } } +func TestRefresherSkipsPersonEntries(t *testing.T) { + store := newFakeSnapshotStore() + r := &TrendingRefresher{ + Sections: fakeSectionLister{configs: []json.RawMessage{tmdbConfig(t, "tmdb", "week")}}, + Snapshots: store, + Resolver: fakeResolver{byType: map[string]*catalog.ExternalIDLookup{ + // "99" is present in the movie lookup to simulate a person ID that + // collides with an unrelated library movie's TMDB ID. + "movie": {ByTMDB: map[string]string{"10": "c-movie", "99": "c-person-collision"}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, + "series": {ByTMDB: map[string]string{"20": "c-series"}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, + }}, + TMDBTrending: fakeTMDB{entries: []catalog.TMDBCollectionEntry{ + {ID: 10, MediaType: "movie"}, + {ID: 99, MediaType: "person"}, + {ID: 20, MediaType: "tv"}, + }}, + Clock: recipes.FixedClock(time.Date(2026, 5, 29, 12, 0, 0, 0, time.UTC)), + } + + if _, err := r.RunOnce(context.Background()); err != nil { + t.Fatalf("RunOnce: %v", err) + } + + got := store.saved["tmdb|week"].contentIDs + want := []string{"c-movie", "c-series"} + if len(got) != len(want) || got[0] != want[0] || got[1] != want[1] { + t.Fatalf("content IDs = %v; want %v (person entry must be skipped)", got, want) + } +} + +func TestRefresherTraktInterleavesMoviesAndShows(t *testing.T) { + store := newFakeSnapshotStore() + r := &TrendingRefresher{ + Sections: fakeSectionLister{configs: []json.RawMessage{tmdbConfig(t, "trakt", "week")}}, + Snapshots: store, + Resolver: fakeResolver{byType: map[string]*catalog.ExternalIDLookup{ + "movie": {ByTMDB: map[string]string{"1": "m1", "2": "m2"}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, + "series": {ByTMDB: map[string]string{"3": "s1"}, ByIMDb: map[string]string{}, ByTVDB: map[string]string{}}, + }}, + TraktTrending: fakeTrakt{byMediaType: map[string][]catalog.TraktCollectionEntry{ + "movie": {{TMDBID: 1, MediaType: "movie"}, {TMDBID: 2, MediaType: "movie"}}, + "tv": {{TMDBID: 3, MediaType: "tv"}}, + }}, + Clock: recipes.FixedClock(time.Date(2026, 5, 29, 12, 0, 0, 0, time.UTC)), + } + + if _, err := r.RunOnce(context.Background()); err != nil { + t.Fatalf("RunOnce: %v", err) + } + + // Interleaved order: movie[0], show[0], movie[1] => m1, s1, m2. A plain + // concat would have buried s1 after all movies. + got := store.saved["trakt|week"].contentIDs + want := []string{"m1", "s1", "m2"} + if len(got) != len(want) || got[0] != want[0] || got[1] != want[1] || got[2] != want[2] { + t.Fatalf("content IDs = %v; want %v (interleaved)", got, want) + } +} + +func TestRefresherTraktPartialFailurePreservesLastGood(t *testing.T) { + store := newFakeSnapshotStore() + r := &TrendingRefresher{ + Sections: fakeSectionLister{configs: []json.RawMessage{tmdbConfig(t, "trakt", "week")}}, + Snapshots: store, + Resolver: fakeResolver{}, + TraktTrending: fakeTrakt{ + byMediaType: map[string][]catalog.TraktCollectionEntry{"movie": {{TMDBID: 1, MediaType: "movie"}}}, + errByType: map[string]error{"tv": errors.New("trakt shows 500")}, + }, + Clock: recipes.FixedClock(time.Date(2026, 5, 29, 12, 0, 0, 0, time.UTC)), + } + + if _, err := r.RunOnce(context.Background()); err != nil { + t.Fatalf("RunOnce: %v", err) + } + + if _, ok := store.saved["trakt|week"]; ok { + t.Fatal("SaveSuccess must not run when one Trakt sub-fetch fails (would drop a media type)") + } + if store.attempts["trakt|week"].status != "error" { + t.Fatalf("attempt status = %q; want error", store.attempts["trakt|week"].status) + } +} + func TestDistinctTrendingCombosCollapsesTrakt(t *testing.T) { configs := []json.RawMessage{ tmdbConfig(t, "trakt", "day"), diff --git a/internal/taskmanager/tasks/refresh_trending_discover.go b/internal/taskmanager/tasks/refresh_trending_discover.go index 7324d517..3d9a3941 100644 --- a/internal/taskmanager/tasks/refresh_trending_discover.go +++ b/internal/taskmanager/tasks/refresh_trending_discover.go @@ -3,6 +3,7 @@ package tasks import ( "context" "encoding/json" + "errors" "fmt" "github.com/Silo-Server/silo-server/internal/taskmanager" @@ -47,6 +48,10 @@ func (t *RefreshTrendingDiscoverTask) DefaultTriggers() []taskmanager.TriggerCon func (t *RefreshTrendingDiscoverTask) Execute(ctx context.Context, progress taskmanager.ProgressReporter) error { progress.Report(0, "Refreshing trending discover") + if t.refresher == nil { + return errors.New("trending discover refresh: refresher not configured") + } + resultData, err := t.refresher.RunOnce(ctx) if err != nil { return fmt.Errorf("trending discover refresh: %w", err)