Merge pull request #23 from Silo-Server/feat/trending-section

feat(sections): persist trending_discover via background-refreshed snapshot
This commit is contained in:
Quick
2026-05-29 11:35:05 -04:00
committed by GitHub
17 changed files with 2733 additions and 3 deletions
+17
View File
@@ -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))
}
File diff suppressed because it is too large Load Diff
@@ -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.
```
+17
View File
@@ -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,
+138 -3
View File
@@ -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 {
@@ -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{})
}
+24
View File
@@ -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
@@ -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)
}
}
+65
View File
@@ -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)
}
}
+277
View File
@@ -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
}
@@ -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)
}
}
+156
View File
@@ -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()
}
@@ -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)
}
}
}
+3
View File
@@ -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,
}
@@ -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
}
@@ -0,0 +1 @@
DROP TABLE IF EXISTS public.trending_discover_snapshots;
@@ -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)
);