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/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. ✓ +``` 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..afb05454 --- /dev/null +++ b/docs/superpowers/specs/2026-05-29-trending-discover-persistent-snapshot-design.md @@ -0,0 +1,208 @@ +# 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 `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 + 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. +``` diff --git a/internal/api/router.go b/internal/api/router.go index 75e38dfa..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,6 +883,18 @@ func NewRouter(deps Dependencies) chi.Router { } } + // 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, libraryCollectionService, diff --git a/internal/sections/fetcher.go b/internal/sections/fetcher.go index 9a771d7e..279f809f 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 @@ -67,9 +73,16 @@ 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 + + // 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 + candidateGroup singleflight.Group // Clock returns the current time. Defaults to recipes.RealClock{}. // Tests inject recipes.FixedClock for deterministic seasonal/editorial behavior. @@ -1070,6 +1083,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 +2349,126 @@ 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" +} + +// 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 +// 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, window := canonicalTrendingKey(p.Source, p.Window) + + limit := s.ItemLimit + if limit <= 0 { + limit = 20 + } + + orderedIDs, err := f.loadTrendingDiscoverContentIDs(ctx, source, window) + 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. + ordered := orderMediaItems(items, orderedIDs) + if len(ordered) > limit { + ordered = ordered[:limit] + } + return ordered, len(ordered), 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 + } + snap, found, err := f.TrendingSnapshots.Get(ctx, source, window) + if err != nil { + return nil, err + } + if !found { + return nil, nil + } + return snap.ContentIDs, 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 { + // 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" + switch e.mediaType { + case "movie": + lookup = movieLookup + case "tv": + lookup = seriesLookup + default: + continue + } + 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/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 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/trending_read_test.go b/internal/sections/trending_read_test.go new file mode 100644 index 00000000..5743405d --- /dev/null +++ b/internal/sections/trending_read_test.go @@ -0,0 +1,65 @@ +package sections + +import ( + "context" + "errors" + "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) + } +} + +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 new file mode 100644 index 00000000..88a00892 --- /dev/null +++ b/internal/sections/trending_refresher.go @@ -0,0 +1,277 @@ +package sections + +import ( + "context" + "encoding/json" + "errors" + "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 == sourceTrakt { + 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", errors.Join(movieErr, showErr)) + } + // 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 { + 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 +} + +// 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) { + if r.Resolver == nil { + return nil, fmt.Errorf("trending_discover: external ID resolver not configured") + } + var movieBatch, seriesBatch catalog.ExternalIDBatch + for _, e := range entries { + 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) + } + 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..2f9bef86 --- /dev/null +++ b/internal/sections/trending_refresher_test.go @@ -0,0 +1,293 @@ +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 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 +} + +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 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"), + 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) + } +} diff --git a/internal/sections/trending_snapshot.go b/internal/sections/trending_snapshot.go new file mode 100644 index 00000000..d755e6ef --- /dev/null +++ b/internal/sections/trending_snapshot.go @@ -0,0 +1,156 @@ +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 +} + +// 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 != sourceTrakt { + source = sourceTMDB + } + if source == sourceTrakt { + return sourceTrakt, windowWeek + } + if window != windowDay { + window = windowWeek + } + return sourceTMDB, 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) + } + } +} 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, } diff --git a/internal/taskmanager/tasks/refresh_trending_discover.go b/internal/taskmanager/tasks/refresh_trending_discover.go new file mode 100644 index 00000000..3d9a3941 --- /dev/null +++ b/internal/taskmanager/tasks/refresh_trending_discover.go @@ -0,0 +1,65 @@ +package tasks + +import ( + "context" + "encoding/json" + "errors" + "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") + + 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) + } + if resultData != nil { + progress.SetResultData(resultData) + } + + progress.Report(100, "Trending discover refresh complete") + return nil +} diff --git a/migrations/167_trending_discover_snapshots.down.sql b/migrations/167_trending_discover_snapshots.down.sql new file mode 100644 index 00000000..72a6d6d4 --- /dev/null +++ b/migrations/167_trending_discover_snapshots.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS public.trending_discover_snapshots; diff --git a/migrations/167_trending_discover_snapshots.up.sql b/migrations/167_trending_discover_snapshots.up.sql new file mode 100644 index 00000000..706d5c83 --- /dev/null +++ b/migrations/167_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) +);