diff --git a/docs/superpowers/plans/2026-05-27-episode-catalog-performance.md b/docs/superpowers/plans/2026-05-27-episode-catalog-performance.md new file mode 100644 index 00000000..1f62291f --- /dev/null +++ b/docs/superpowers/plans/2026-05-27-episode-catalog-performance.md @@ -0,0 +1,523 @@ +# Episode Catalog Performance Plan + +Commands assume the repository root is the cwd. + +> **For agentic workers:** This is a discussion plan, not an execution script. Confirm the architecture decision points before implementing the durable catalog index work. + +## Goal + +Make `/api/v1/catalog?source=query&type=episode&library_id=` fast enough for large series libraries and concurrent users. The first page of common episode-library browse requests should avoid full-library scans, repeated media-file aggregation, and exact counts unless the caller explicitly needs them. + +Target behavior: + +- First-page episode browse returns in less than 800 ms API time for common sorts/filters on a library with about 1 million episodes. +- Hot SQL paths avoid per-request aggregation over `media_files` or all `episode_libraries`. +- Exact totals are not computed by default for expensive query shapes. +- Existing API semantics continue to work for web, Android, Apple, and Jellyfin-compatible clients unless explicitly versioned. + +## Current Findings + +Testing used the large Series library endpoint shape: + +```text +/api/v1/catalog?source=query&type=episode&library_id=2&limit=60&offset=0 +``` + +Library 2 currently has about 767k episode-library rows. The recent fixes improved title and date-added browse, but several sorts and filters still spend too much time in SQL, especially when exact totals are requested. + +### Sort timings + +With exact totals enabled: + +| Sort | API time | +| --- | ---: | +| `title asc` | 0.93s | +| `added_at desc` | 0.73s | +| `release_date desc` | 1.93s | +| `last_air_date desc` | 2.21s | +| `year desc` | 2.28s | +| `content_rating asc` | 1.21s | +| `runtime desc` | 1.97s | +| `rating_imdb desc` | 2.34s | +| `rating_tmdb desc` | 2.12s | +| `resolution desc` | 8.67s | +| `bitrate desc` | 5.53s | +| `progress desc` | 2.43s | +| `date_viewed desc` | 1.98s | +| `plays desc` | 2.06s | + +With `include_total=false`: + +| Sort | API time | +| --- | ---: | +| `title asc` | 0.39s | +| `added_at desc` | 0.29s | +| `release_date desc` | 1.50s | +| `last_air_date desc` | 1.84s | +| `year desc` | 1.45s | +| `runtime desc` | 1.55s | +| `rating_imdb desc` | 1.73s | +| `rating_tmdb desc` | 1.77s | +| `resolution desc` | 5.69s | +| `bitrate desc` | 5.09s | +| `date_viewed desc` | 2.59s | + +The count split helped, but the page query itself is still too slow for many sort families. + +### Filter timings + +With `sort=title&order=asc` and exact totals: + +| Filter | API time | Total | +| --- | ---: | ---: | +| none | 1.07s | 767k | +| `genre=Comedy` | 2.92s | 250k | +| `resolution=1080p` | 3.04s | 474k | +| `subtitle_language=en` | 6.57s | 524k | +| `dolby_vision=true` | 7.63s | 20k | +| `watched=true` | 16.32s | 11k | +| `watched=false` | 13.09s | 756k | +| `last_watched in_last 30d` | timed out at 30s | n/a | + +With `include_total=false`, most of those page queries drop below 1s, except `last_watched in_last 30d`, which still times out. This confirms two independent problems: + +- Exact totals are too expensive to run on every page 0 request. +- Some page plans start from the wrong side of the query and scan the whole library. + +## Root Causes + +1. The generic episode catalog path projects episodes into a `media_items`-shaped subquery and then asks one query builder to handle every sort/filter combination. This keeps code reusable, but it hides cheaper plans from PostgreSQL. + +2. Technical sorts and filters aggregate `media_files` per request: + +```text +media_files -> GROUP BY episode_id -> join all episode candidates -> sort +``` + +For `resolution` and `bitrate`, this means scanning and grouping hundreds of thousands of file rows before returning 60 items. + +3. Personalized sorts and filters left-join small user-state sets onto the entire episode library. For `date_viewed desc`, `plays desc`, and `last_watched`, the database sorts mostly-null rows from the whole library instead of starting with the few watched rows. + +4. Exact counts use the same broad filtered relation as the page query. This is acceptable for small result sets, but expensive for common filters like `watched=false`, `subtitle_language=en`, or genre filters that match hundreds of thousands of rows. + +5. Series-level episode filters and sorts duplicate series metadata across all episodes. A filter like `genre=Comedy` is logically a series filter, but the current episode projection evaluates it at episode scale. + +## Prototype Results + +The following SQL prototypes were tested against the same data shape to validate that the proposed plan is viable. + +### User-state-first plans + +Starting from watched/progress rows, then joining to episode library membership: + +| Query shape | SQL time | +| --- | ---: | +| `last_watched in_last 30d` | about 5 ms | +| `date_viewed desc` first page | about 282 ms | +| `watched=true` title page | about 261 ms | +| `watched=true` exact count | about 80 ms | + +This proves the `last_watched` timeout is a planner/source problem, not an unavoidable data-size problem. + +### Precomputed technical stats + +A temporary per-episode/per-library stats table was built with max resolution rank, max bitrate, HDR/Dolby Vision flags, and audio/subtitle language arrays. Build time was about 22s as a one-time backfill over the dev data; scanner maintenance would keep the permanent table updated incrementally. + +Using that temporary stats table: + +| Query shape | SQL time | +| --- | ---: | +| `subtitle_language=en` filter page | about 24 ms | +| `dolby_vision=true` filter page | about 128 ms | +| `bitrate desc` sort page | about 1 ms | +| `resolution desc` sort page | about 365 ms | + +This validates a durable browse index or stats table for technical fields. + +## Architecture Decision + +There are two viable paths. + +### Option A: Incremental Specialized Plans + +Add specific fast paths for technical stats and user-state filters while keeping the generic query builder as the main executor. + +Pros: + +- Lower implementation cost. +- Smaller schema change. +- Directly fixes the worst outliers: `resolution`, `bitrate`, `subtitle_language`, `dolby_vision`, `last_watched`, `watched`. + +Cons: + +- Leaves several episode metadata sorts around 1.5-2s. +- Keeps exact-count complexity spread through the generic executor. +- Each new slow query shape becomes another special case. + +### Option B: Durable Episode Browse Index + +Create a persistent per-library episode browse index table that stores the sort/filter keys needed to find page IDs quickly, then hydrate only the selected page rows from `episodes` and parent `media_items`. + +This is the recommended target if the server needs to handle hundreds of concurrent users. It turns page selection into indexed top-N scans over a narrow table and avoids repeated joins/aggregates over broad catalog tables. + +Proposed table shape: + +```sql +CREATE TABLE episode_catalog_entries ( + media_folder_id integer NOT NULL, + episode_id text NOT NULL, + series_id text NOT NULL, + sort_key text NOT NULL, + added_at timestamptz NOT NULL, + episode_air_date date, + year integer NOT NULL, + genres text[] NOT NULL, + studios text[] NOT NULL, + networks text[] NOT NULL, + countries text[] NOT NULL, + original_language text NOT NULL, + content_rating text NOT NULL, + content_rating_rank integer NOT NULL, + status text NOT NULL, + runtime integer NOT NULL, + rating_imdb numeric, + rating_tmdb numeric, + max_resolution_rank integer, + max_bitrate integer, + has_hdr boolean NOT NULL DEFAULT false, + has_dolby_vision boolean NOT NULL DEFAULT false, + audio_language_codes text[] NOT NULL DEFAULT '{}', + subtitle_language_codes text[] NOT NULL DEFAULT '{}', + updated_at timestamptz NOT NULL DEFAULT now(), + PRIMARY KEY (media_folder_id, episode_id) +); +``` + +The table should not become a second full metadata store unless measurements prove that hydration is too expensive. The first implementation can use it to pick ordered `episode_id` values, then join only those 60 IDs back to the existing episode projection for response shaping. + +Recommended indexes: + +```sql +CREATE INDEX idx_episode_catalog_entries_title +ON episode_catalog_entries (media_folder_id, sort_key, episode_id); + +CREATE INDEX idx_episode_catalog_entries_added +ON episode_catalog_entries (media_folder_id, added_at DESC, sort_key, episode_id); + +CREATE INDEX idx_episode_catalog_entries_air_date +ON episode_catalog_entries (media_folder_id, episode_air_date DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX idx_episode_catalog_entries_year +ON episode_catalog_entries (media_folder_id, year DESC, sort_key, episode_id); + +CREATE INDEX idx_episode_catalog_entries_runtime +ON episode_catalog_entries (media_folder_id, runtime DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX idx_episode_catalog_entries_imdb +ON episode_catalog_entries (media_folder_id, rating_imdb DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX idx_episode_catalog_entries_tmdb +ON episode_catalog_entries (media_folder_id, rating_tmdb DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX idx_episode_catalog_entries_resolution +ON episode_catalog_entries (media_folder_id, max_resolution_rank DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX idx_episode_catalog_entries_bitrate +ON episode_catalog_entries (media_folder_id, max_bitrate DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX idx_episode_catalog_entries_hdr +ON episode_catalog_entries (media_folder_id, sort_key, episode_id) +WHERE has_hdr; + +CREATE INDEX idx_episode_catalog_entries_dolby_vision +ON episode_catalog_entries (media_folder_id, sort_key, episode_id) +WHERE has_dolby_vision; +``` + +Evaluate either `btree_gin` multi-column GIN or separate GIN indexes for array filters: + +```sql +CREATE INDEX idx_episode_catalog_entries_genres_gin +ON episode_catalog_entries USING gin (genres); + +CREATE INDEX idx_episode_catalog_entries_audio_gin +ON episode_catalog_entries USING gin (audio_language_codes); + +CREATE INDEX idx_episode_catalog_entries_subtitle_gin +ON episode_catalog_entries USING gin (subtitle_language_codes); +``` + +## Query Plan Design + +Introduce a planner layer between `CatalogResolver` and `QueryExecutor`. + +```go +type CatalogPlan struct { + PageSQL string + PageArgs []any + CountSQL string + CountArgs []any + CountMode CatalogCountMode + SnapshotStrategy SnapshotStrategy + PlanName string +} +``` + +The planner chooses one of these sources: + +- `episode_catalog_entries` for normal episode library browse. +- `episode_user_state` CTE/source for watched/date-viewed/plays/last-watched shapes. +- Existing generic `QueryExecutor` fallback for unsupported combinations. + +The page query should use a narrow ID-first CTE: + +```sql +WITH page AS ( + SELECT ece.episode_id, ece.sort_key + FROM episode_catalog_entries ece + WHERE ece.media_folder_id = $1 + ORDER BY ece.sort_key ASC, ece.episode_id ASC + LIMIT $2 OFFSET $3 +) +SELECT ... +FROM page +JOIN episodes e ON e.content_id = page.episode_id +JOIN media_items si ON si.content_id = e.series_id +ORDER BY page.sort_key ASC, page.episode_id ASC; +``` + +Each concrete sort should carry the selected sort keys through the `page` CTE and use the same order in the final hydrated SELECT. Do not compute `row_number()` over the full candidate set just to preserve order; that would reintroduce a broad sort before `LIMIT`. + +For sort/filter shapes that can be satisfied entirely from `episode_catalog_entries`, counts become simple index-backed counts on the narrow table. For expensive or broad counts, the planner should return `total_exact=false` unless exact totals are explicitly requested. + +## User-State Plan + +Personalized sort/filter logic should start from user state when the requested result set is mostly watched/progress rows. + +Fast paths: + +- `last_watched` comparisons and `in_last`. +- `watched=true`. +- `in_progress=true`. +- `favorited=true`. +- `in_watchlist=true`. +- `date_viewed desc` first segment. +- `plays desc` first segment. +- `progress desc` first segment. + +The user-state source can start as request-time CTEs over `user_watch_history`, `user_watch_progress`, `user_favorites`, and `user_watchlist`. If concurrent load tests show those CTEs are still too expensive for heavy users, promote them to a maintained `profile_media_state` aggregate table. + +Special handling: + +- `watched=false` should not scan user state first because the result set is usually almost the whole library. Use the episode browse index with an anti-join against the small watched set, and compute exact count as `library_count - watched_count` when possible. +- `date_viewed desc` with deep offsets eventually reaches the unviewed segment. First implement the watched segment fast path and fall back only when the requested offset exceeds the watched count. +- `last_watched lt/lte` includes never-watched rows because the current SQL uses `-infinity`. Keep this behavior, but route `gt/gte/between/in_last` through user-state-first plans. + +## Count Strategy + +Exact totals are the largest remaining source of avoidable database load. The UI already understands `total_exact=false` and can estimate virtualized height. + +Plan: + +1. Change `LibraryBrowse` to call `useCatalogWindow(..., includeTotal: false)` for the first page unless a specific UI state truly needs an exact count. +2. Add backend count modes: + - `none`: return `total_exact=false` and one extra row for `has_more`. + - `fast`: exact count from a narrow indexed table or small user-state source. + - `cached`: count reused from a query hash and invalidated by scanner/user-state writes. + - `exact`: full exact count, only when requested. +3. For plain episode library counts, use `episode_catalog_entries` or `episode_libraries` directly. +4. For watched filters: + - `watched=true`: count from user-state source joined to library membership. + - `watched=false`: `library_total - watched_true_count` when the filter set permits it. +5. For technical filters, count from `episode_catalog_entries`. + +Do not run broad exact counts by default under web browse traffic. That path does not scale to hundreds of users. + +## Implementation Phases + +### Phase 0: Observability and Benchmark Harness + +- Add structured slow-query logging around catalog page and count execution: + - `source` + - `media_scope` + - `library_count` + - `sort` + - `filter_count` + - `include_total` + - `plan_name` + - page SQL duration + - count SQL duration +- Add a local benchmark helper under `scripts/` that exercises the sort/filter matrix without embedding credentials. +- Add a small `EXPLAIN (ANALYZE, BUFFERS)` note template for comparing plans. + +Verification: + +```bash +go test ./internal/catalog -count=1 +``` + +### Phase 1: Stop Exact Counts by Default in Library Browse + +- Change `web/src/pages/LibraryBrowse.tsx` to pass `includeTotal: false`. +- Keep `CatalogFiltersPanel` result count display compatible with estimated totals or suppress exact count text when `total_exact=false`. +- Confirm `ItemGrid` still virtualizes correctly using the existing estimated-total logic in `useCatalogWindow`. + +Verification: + +```bash +cd web && pnpm run lint +``` + +### Phase 2: Durable Episode Browse Index + +- Add migrations for `episode_catalog_entries`. +- Backfill from: + - `episode_libraries` + - `episodes` + - parent series rows in `media_items` + - active `media_files` +- Add a repository/service that can refresh entries for: + - one episode + - one series + - one library + - one changed media file +- Wire refresh calls into scanner and metadata writes where episode visibility or sort/filter keys change. +- Keep the existing `episode_libraries` table as the source of truth for membership; the new table is a maintained read model. + +Verification: + +```bash +go test ./internal/catalog ./internal/scanner -count=1 +``` + +### Phase 3: Episode Catalog Planner + +- Add a planner that recognizes episode library browse requests and emits ID-first page SQL against `episode_catalog_entries`. +- Keep the existing `QueryExecutor` as fallback. +- Support these first: + - no filters, all common sorts + - `genre`, `year`, `content_rating`, `status` + - `resolution`, `bitrate`, `audio_language`, `subtitle_language`, `hdr`, `dolby_vision` +- Add query-builder tests that assert the selected plan name and SQL shape for each supported sort/filter family. + +Verification: + +```bash +go test ./internal/catalog -count=1 +``` + +### Phase 4: User-State Planner + +- Add user-state-first plans for `last_watched`, `watched=true`, `in_progress=true`, `date_viewed desc`, `plays desc`, and `progress desc`. +- Add `watched=false` fast count using complement logic where safe. +- Preserve current hidden-history behavior. +- Preserve current semantics for never-watched rows on `last_watched lt/lte`. + +Verification: + +```bash +go test ./internal/catalog ./internal/userstore -count=1 +``` + +### Phase 5: Load Testing and Tuning + +- Deploy to dev. +- Re-run the sort/filter matrix with exact totals disabled and enabled. +- Run concurrent load against the high-traffic shapes: + - `title asc` + - `added_at desc` + - `release_date desc` + - `resolution desc` + - `bitrate desc` + - `subtitle_language=en` + - `dolby_vision=true` + - `date_viewed desc` + - `last_watched in_last 30d` +- Use PostgreSQL query plans and slow-query logs to tune indexes before considering the work complete. + +Target load result: + +- 50 concurrent catalog requests: p95 under 1s for indexed paths. +- No request over 5s for supported sort/filter shapes. +- No supported first-page request performs a full `media_files` aggregate. + +## Testing Matrix + +The benchmark harness should cover every executable sort and filter family. + +Sorts: + +- `title` +- `added_at` +- `release_date` +- `last_air_date` +- `year` +- `content_rating` +- `runtime` +- `rating_imdb` +- `rating_tmdb` +- `rating_rt_critic` +- `rating_rt_audience` +- `resolution` +- `bitrate` +- `progress` +- `date_viewed` +- `plays` + +Filters: + +- `type` +- `genre` +- `year` +- `rating_imdb` +- `studio` +- `network` +- `country` +- `original_language` +- `content_rating` +- `added_at` +- `release_date` +- `status` +- `actor` +- `director` +- `writer` +- `producer` +- `watched` +- `favorited` +- `in_watchlist` +- `in_progress` +- `last_watched` +- `resolution` +- `hdr` +- `dolby_vision` +- `bitrate` +- `audio_language` +- `subtitle_language` + +For each case, capture: + +- API time with exact totals. +- API time with `include_total=false`. +- page SQL time. +- count SQL time. +- rows scanned/aggregated from `EXPLAIN`. +- whether the planner used `episode_catalog_entries`, user state, or fallback. + +## Risks and Open Questions + +- The browse index is a read model. The source of truth stays in `episode_libraries`, `episodes`, `media_items`, and `media_files`, so refresh paths must be reliable and observable. +- Scanner and metadata updates can touch large series. Batch refreshes should be chunked and idempotent. +- Multi-library browse needs clear semantics for `added_at` and technical stats. The current single-library Series Library case should be optimized first. +- Exact snapshot semantics may conflict with cached counts. Page snapshots should remain stable; counts can be marked non-exact unless they come from the same snapshot-safe path. +- API response changes around count modes may require Android and Apple follow-up. Keeping the existing `total_exact=false` behavior avoids most client churn. +- Person filters may need fallback initially unless episode-level people data is available and indexed. +- Array GIN filters should be checked with real plans. If separate GIN indexes do too much post-filtering by library, use `btree_gin` multi-column indexes or add selective partial indexes for common libraries. +- Index maintenance, disk, and write amplification should be tracked explicitly. The btree indexes such as `idx_episode_catalog_entries_title`, `idx_episode_catalog_entries_added`, and `idx_episode_catalog_entries_air_date`, plus GIN indexes such as `idx_episode_catalog_entries_genres_gin`, trade scanner write cost and disk growth for browse latency. Measure `episode_catalog_entries` table and index size before rollout, and watch trigger update latency during large scanner runs. +- Backfill safety needs a production runbook. Treat the `episode_catalog_entries` backfill as idempotent work that can be chunked or retried with low-locking settings, capture a production-sized duration estimate before enabling it broadly, and document rollback to dropping the read model plus its triggers. +- Read-model health should be observable. Alert when `episode_catalog_entries.updated_at` lags source table changes beyond a stale threshold or refresh errors appear, fall back to the generic executor when health checks fail, and provide a remediation command to rebuild stale entries. + +## Recommended Decision + +Proceed with Option B in phases. The prototype timings show that denormalizing technical stats and starting user-state queries from user tables both work. A durable `episode_catalog_entries` read model generalizes those wins to normal episode metadata sorts too, which is the more scalable path for hundreds of users. + +Use Phase 1 as the immediate load reducer, then implement the durable browse index and planner behind the existing catalog API. Keep the generic executor as a fallback until the measured matrix shows the new planner covers the important sort/filter combinations. diff --git a/internal/catalog/episode_catalog_entries_executor.go b/internal/catalog/episode_catalog_entries_executor.go new file mode 100644 index 00000000..930f2b56 --- /dev/null +++ b/internal/catalog/episode_catalog_entries_executor.go @@ -0,0 +1,1024 @@ +package catalog + +import ( + "context" + "errors" + "fmt" + "strconv" + "strings" + "time" + + "github.com/jackc/pgx/v5/pgconn" + + "github.com/Silo-Server/silo-server/internal/models" +) + +type episodeCatalogEntryPageRow struct { + episodeID string + addedAt time.Time +} + +type episodeCatalogUserStatePlan struct { + entryDef QueryDefinition + source string + alias string + positiveSource bool + userClauses []string + userArgs []any + sortField string + sortOrder string + sortOnly bool + requiresFullPageGate bool + requireProgressRatio bool +} + +func (e *QueryExecutor) tryEpisodeCatalogUserStatePreviewPage( + ctx context.Context, + def QueryDefinition, + access AccessFilter, + limit int, + offset int, + includeTotal bool, +) ([]*models.MediaItem, int, bool, bool, error) { + def = def.Normalize() + effectiveScope := e.Scope + if effectiveScope == "" { + effectiveScope = def.MediaScope + } + if !isEpisodeCatalogScope(effectiveScope) { + return nil, 0, false, false, nil + } + if access.UserID == 0 || strings.TrimSpace(access.ProfileID) == "" { + return nil, 0, false, false, nil + } + + plan, ok, err := extractEpisodeCatalogUserStatePlan(def) + if err != nil || !ok { + return nil, 0, false, ok, err + } + if includeTotal && plan.sortOnly { + return nil, 0, false, false, nil + } + + libraryID, empty, ok := singleEpisodeCatalogLibraryID(def, access) + if !ok { + return nil, 0, false, false, nil + } + if empty { + return []*models.MediaItem{}, 0, false, true, nil + } + + if limit <= 0 { + limit = 20 + } + if offset < 0 { + offset = 0 + } + + ctes := []string{episodeCatalogUserStateCTE(plan)} + args := []any{access.UserID, access.ProfileID, libraryID} + argIdx := 4 + whereParts := []string{"ece.media_folder_id = $3"} + + if access.AllowedContentIDs != nil { + if len(access.AllowedContentIDs) == 0 { + return []*models.MediaItem{}, 0, false, true, nil + } + whereParts = append(whereParts, fmt.Sprintf("ece.episode_id = ANY($%d)", argIdx)) + args = append(args, access.AllowedContentIDs) + argIdx++ + } + + ApplySectionAccessFilter("ece", access, &whereParts, &args, &argIdx) + + if e.SnapshotAt != nil { + whereParts = append(whereParts, fmt.Sprintf("ece.episode_created_at <= $%d", argIdx)) + args = append(args, *e.SnapshotAt) + argIdx++ + } + + if prefix := strings.TrimSpace(access.NamePrefix); prefix != "" { + whereParts = append(whereParts, fmt.Sprintf( + "(ece.sort_key LIKE $%d ESCAPE '\\' OR LOWER(ece.title) LIKE $%d ESCAPE '\\')", + argIdx, + argIdx, + )) + args = append(args, escapePrefixForLike(prefix)+"%") + argIdx++ + } + + queryWhere, queryArgs, nextArgIdx, ok, err := buildEpisodeCatalogEntryQueryWhere(plan.entryDef, argIdx) + if err != nil || !ok { + return nil, 0, false, ok, err + } + if queryWhere != "" { + whereParts = append(whereParts, queryWhere) + args = append(args, queryArgs...) + argIdx = nextArgIdx + } + + userClauses := rebindUserStateClauses(plan.userClauses, argIdx) + whereParts = append(whereParts, userClauses...) + args = append(args, plan.userArgs...) + argIdx += len(plan.userArgs) + + fromClause := episodeCatalogUserStateFromClause(plan) + orderBy, ok := episodeCatalogUserStateOrderBy(plan) + if !ok { + return nil, 0, false, false, nil + } + whereClause := "WHERE " + strings.Join(whereParts, " AND ") + + pageLimit := limit + 1 + pageArgs := append([]any{}, args...) + limitArgIdx := argIdx + pageArgs = append(pageArgs, pageLimit) + offsetClause := "" + if offset > 0 { + offsetClause = fmt.Sprintf(" OFFSET $%d", limitArgIdx+1) + pageArgs = append(pageArgs, offset) + } + + pageSQL := fmt.Sprintf( + `WITH %s + SELECT ece.episode_id, ece.added_at + %s + %s + %s + LIMIT $%d%s`, + strings.Join(ctes, ",\n"), + fromClause, + whereClause, + orderBy, + limitArgIdx, + offsetClause, + ) + + pageRows, err := e.queryEpisodeCatalogEntryPage(ctx, pageSQL, pageArgs...) + if err != nil { + if episodeCatalogEntriesUnavailable(err) { + return nil, 0, false, false, nil + } + return nil, 0, false, true, err + } + if plan.requiresFullPageGate && len(pageRows) < pageLimit { + return nil, 0, false, false, nil + } + + hasMore := len(pageRows) > limit + if hasMore { + pageRows = pageRows[:limit] + } + + items, err := e.hydrateEpisodeCatalogEntryPage(ctx, pageRows) + if err != nil { + if episodeCatalogEntriesUnavailable(err) { + return nil, 0, false, false, nil + } + return nil, 0, false, true, err + } + + total := 0 + if includeTotal { + countSQL := fmt.Sprintf( + `WITH %s + SELECT COUNT(*) + %s + %s`, + strings.Join(ctes, ",\n"), + fromClause, + whereClause, + ) + if err := e.Pool.QueryRow(ctx, countSQL, args...).Scan(&total); err != nil { + if episodeCatalogEntriesUnavailable(err) { + return nil, 0, false, false, nil + } + return nil, 0, false, true, fmt.Errorf("counting episode user-state catalog entries: %w", err) + } + hasMore = total > offset+len(items) + } + + return items, total, hasMore, true, nil +} + +func (e *QueryExecutor) tryEpisodeCatalogEntriesPreviewPage( + ctx context.Context, + def QueryDefinition, + access AccessFilter, + limit int, + offset int, + includeTotal bool, +) ([]*models.MediaItem, int, bool, bool, error) { + def = def.Normalize() + effectiveScope := e.Scope + if effectiveScope == "" { + effectiveScope = def.MediaScope + } + if !isEpisodeCatalogScope(effectiveScope) { + return nil, 0, false, false, nil + } + + allowPersonalized := strings.TrimSpace(access.ProfileID) != "" + if err := def.ValidateWithOptions(allowPersonalized, allowPersonalized); err != nil { + return nil, 0, false, true, err + } + + libraryID, empty, ok := singleEpisodeCatalogLibraryID(def, access) + if !ok { + return nil, 0, false, false, nil + } + if empty { + return []*models.MediaItem{}, 0, false, true, nil + } + + if limit <= 0 { + limit = 20 + } + if offset < 0 { + offset = 0 + } + + whereParts := []string{"ece.media_folder_id = $1"} + args := []any{libraryID} + argIdx := 2 + + if access.AllowedContentIDs != nil { + if len(access.AllowedContentIDs) == 0 { + return []*models.MediaItem{}, 0, false, true, nil + } + whereParts = append(whereParts, fmt.Sprintf("ece.episode_id = ANY($%d)", argIdx)) + args = append(args, access.AllowedContentIDs) + argIdx++ + } + + ApplySectionAccessFilter("ece", access, &whereParts, &args, &argIdx) + + if e.SnapshotAt != nil { + whereParts = append(whereParts, fmt.Sprintf("ece.episode_created_at <= $%d", argIdx)) + args = append(args, *e.SnapshotAt) + argIdx++ + } + + if prefix := strings.TrimSpace(access.NamePrefix); prefix != "" { + whereParts = append(whereParts, fmt.Sprintf( + "(ece.sort_key LIKE $%d ESCAPE '\\' OR LOWER(ece.title) LIKE $%d ESCAPE '\\')", + argIdx, + argIdx, + )) + args = append(args, escapePrefixForLike(prefix)+"%") + argIdx++ + } + + queryWhere, queryArgs, nextArgIdx, ok, err := buildEpisodeCatalogEntryQueryWhere(def, argIdx) + if err != nil { + return nil, 0, false, true, err + } + if !ok { + return nil, 0, false, false, nil + } + if queryWhere != "" { + whereParts = append(whereParts, queryWhere) + args = append(args, queryArgs...) + argIdx = nextArgIdx + } + + orderBy, ok := episodeCatalogEntryOrderBy(def.Sort) + if !ok { + return nil, 0, false, false, nil + } + + whereClause := "WHERE " + strings.Join(whereParts, " AND ") + pageLimit := limit + 1 + pageArgs := append([]any{}, args...) + limitArgIdx := argIdx + pageArgs = append(pageArgs, pageLimit) + offsetClause := "" + if offset > 0 { + offsetClause = fmt.Sprintf(" OFFSET $%d", limitArgIdx+1) + pageArgs = append(pageArgs, offset) + } + + pageSQL := fmt.Sprintf( + `SELECT ece.episode_id, ece.added_at + FROM episode_catalog_entries ece + %s + %s + LIMIT $%d%s`, + whereClause, + orderBy, + limitArgIdx, + offsetClause, + ) + + pageRows, err := e.queryEpisodeCatalogEntryPage(ctx, pageSQL, pageArgs...) + if err != nil { + if episodeCatalogEntriesUnavailable(err) { + return nil, 0, false, false, nil + } + return nil, 0, false, true, err + } + + hasMore := len(pageRows) > limit + if hasMore { + pageRows = pageRows[:limit] + } + + items, err := e.hydrateEpisodeCatalogEntryPage(ctx, pageRows) + if err != nil { + if episodeCatalogEntriesUnavailable(err) { + return nil, 0, false, false, nil + } + return nil, 0, false, true, err + } + + total := 0 + if includeTotal { + countSQL := fmt.Sprintf( + `SELECT COUNT(*) + FROM episode_catalog_entries ece + %s`, + whereClause, + ) + if err := e.Pool.QueryRow(ctx, countSQL, args...).Scan(&total); err != nil { + if episodeCatalogEntriesUnavailable(err) { + return nil, 0, false, false, nil + } + return nil, 0, false, true, fmt.Errorf("counting episode catalog entries: %w", err) + } + hasMore = total > offset+len(items) + } + + return items, total, hasMore, true, nil +} + +func extractEpisodeCatalogUserStatePlan(def QueryDefinition) (episodeCatalogUserStatePlan, bool, error) { + plan := episodeCatalogUserStatePlan{ + entryDef: def, + sortField: NormalizeQuerySort(def.Sort).Field, + sortOrder: NormalizeQuerySort(def.Sort).Order, + } + switch plan.sortField { + case "date_viewed", "plays": + plan.source = "viewed" + plan.alias = "uv" + plan.positiveSource = true + plan.sortOnly = true + plan.requiresFullPageGate = true + case "progress": + plan.source = "progress" + plan.alias = "up" + plan.positiveSource = true + plan.sortOnly = true + plan.requiresFullPageGate = true + } + + if len(def.Groups) == 0 { + if plan.source == "" { + return episodeCatalogUserStatePlan{}, false, nil + } + plan.requireProgressRatio = plan.source == "progress" && plan.sortOnly + return plan, true, nil + } + if def.Match != "all" { + return episodeCatalogUserStatePlan{}, false, nil + } + + filteredGroups := make([]QueryGroup, 0, len(def.Groups)) + for _, group := range def.Groups { + if group.Match != "all" { + for _, rule := range group.Rules { + if episodeCatalogUserStateRule(rule) { + return episodeCatalogUserStatePlan{}, false, nil + } + } + filteredGroups = append(filteredGroups, group) + continue + } + + filteredRules := make([]QueryRule, 0, len(group.Rules)) + for _, rule := range group.Rules { + source, alias, positive, clause, args, ok, err := buildEpisodeCatalogUserStateRule(rule) + if err != nil { + return episodeCatalogUserStatePlan{}, true, err + } + if !ok { + filteredRules = append(filteredRules, rule) + continue + } + if plan.source != "" && plan.source != source { + return episodeCatalogUserStatePlan{}, false, nil + } + plan.source = source + plan.alias = alias + plan.positiveSource = positive + plan.userClauses = append(plan.userClauses, clause) + plan.userArgs = append(plan.userArgs, args...) + plan.sortOnly = false + plan.requiresFullPageGate = false + } + if len(filteredRules) > 0 { + group.Rules = filteredRules + filteredGroups = append(filteredGroups, group) + } + } + if plan.source == "" { + return episodeCatalogUserStatePlan{}, false, nil + } + plan.entryDef.Groups = filteredGroups + plan.requireProgressRatio = plan.source == "progress" && plan.sortOnly + return plan, true, nil +} + +func episodeCatalogUserStateRule(rule QueryRule) bool { + switch rule.Field { + case "watched", "in_progress", "last_watched": + return true + default: + return false + } +} + +func buildEpisodeCatalogUserStateRule(rule QueryRule) (string, string, bool, string, []any, bool, error) { + switch rule.Field { + case "watched": + value, ok := rule.Value.(bool) + if !ok { + return "", "", false, "", nil, true, fmt.Errorf("watched requires a boolean value") + } + if value { + return "viewed", "uv", true, "uv.episode_id IS NOT NULL", nil, true, nil + } + return "viewed", "uv", false, "uv.episode_id IS NULL", nil, true, nil + case "in_progress": + value, ok := rule.Value.(bool) + if !ok { + return "", "", false, "", nil, true, fmt.Errorf("in_progress requires a boolean value") + } + if value { + return "progress", "up", true, "up.episode_id IS NOT NULL", nil, true, nil + } + return "progress", "up", false, "up.episode_id IS NULL", nil, true, nil + case "last_watched": + switch rule.Op { + case "gt", "gte", "between", "in_last": + default: + return "", "", false, "", nil, false, nil + } + clause, args, err := buildEpisodeCatalogLastWatchedClause(rule) + if err != nil { + return "", "", false, "", nil, true, err + } + return "viewed", "uv", true, clause, args, true, nil + default: + return "", "", false, "", nil, false, nil + } +} + +func buildEpisodeCatalogLastWatchedClause(rule QueryRule) (string, []any, error) { + switch rule.Op { + case "gt", "gte": + operator := map[string]string{"gt": ">", "gte": ">="}[rule.Op] + return "uv.last_watched " + operator + " $%d::timestamptz", []any{rule.Value}, nil + case "between": + values, err := toBetweenValues(rule.Value) + if err != nil { + return "", nil, fmt.Errorf("between requires [min, max] array: %w", err) + } + return "uv.last_watched >= $%d::timestamptz AND uv.last_watched <= $%d::timestamptz", []any{values[0], values[1]}, nil + case "in_last": + duration, ok := rule.Value.(string) + if !ok { + return "", nil, fmt.Errorf("in_last requires a duration string like '30d'") + } + interval, err := parseDuration(duration) + if err != nil { + return "", nil, err + } + return "uv.last_watched >= NOW() - INTERVAL '" + interval + "'", nil, nil + default: + return "", nil, fmt.Errorf("unsupported last_watched operator %q", rule.Op) + } +} + +func rebindUserStateClauses(clauses []string, argIdx int) []string { + if len(clauses) == 0 { + return nil + } + out := make([]string, len(clauses)) + nextArg := argIdx + for i, clause := range clauses { + for strings.Contains(clause, "$%d") { + clause = strings.Replace(clause, "$%d", "$"+strconv.Itoa(nextArg), 1) + nextArg++ + } + out[i] = clause + } + return out +} + +func episodeCatalogUserStateCTE(plan episodeCatalogUserStatePlan) string { + switch plan.source { + case "progress": + progressRatioGate := "" + if plan.requireProgressRatio { + progressRatioGate = ` + AND COALESCE(uwp.duration_seconds, 0) > 0` + } + return fmt.Sprintf(`user_progress AS ( + SELECT uwp.media_item_id AS episode_id, + uwp.position_seconds::double precision / NULLIF(uwp.duration_seconds, 0) AS progress_ratio + FROM user_watch_progress uwp + LEFT JOIN user_history_hidden_items hhi + ON hhi.user_id = $1 + AND hhi.profile_id = $2 + AND hhi.media_item_id = uwp.media_item_id + WHERE uwp.user_id = $1 + AND uwp.profile_id = $2 + AND uwp.completed = FALSE + AND uwp.position_seconds > 0%s + AND (hhi.media_item_id IS NULL OR uwp.updated_at > hhi.hidden_before) + )`, progressRatioGate) + default: + return `user_viewed AS ( + SELECT src.episode_id, + MAX(src.last_watched) AS last_watched, + NULLIF(GREATEST( + COALESCE(MAX(src.history_play_count), 0), + COALESCE(MAX(src.progress_play_count), 0) + ), 0) AS play_count + FROM ( + SELECT uwh.media_item_id AS episode_id, + MAX(uwh.watched_at) AS last_watched, + COUNT(*)::integer AS history_play_count, + NULL::integer AS progress_play_count + FROM user_watch_history uwh + LEFT JOIN user_history_hidden_items hhi + ON hhi.user_id = $1 + AND hhi.profile_id = $2 + AND hhi.media_item_id = uwh.media_item_id + WHERE uwh.user_id = $1 + AND uwh.profile_id = $2 + AND uwh.completed = TRUE + AND (hhi.media_item_id IS NULL OR uwh.watched_at > hhi.hidden_before) + GROUP BY uwh.media_item_id + UNION ALL + SELECT uwp.media_item_id AS episode_id, + uwp.updated_at AS last_watched, + NULL::integer AS history_play_count, + 1 AS progress_play_count + FROM user_watch_progress uwp + LEFT JOIN user_history_hidden_items hhi + ON hhi.user_id = $1 + AND hhi.profile_id = $2 + AND hhi.media_item_id = uwp.media_item_id + WHERE uwp.user_id = $1 + AND uwp.profile_id = $2 + AND uwp.completed = TRUE + AND (hhi.media_item_id IS NULL OR uwp.updated_at > hhi.hidden_before) + ) src + GROUP BY src.episode_id + )` + } +} + +func episodeCatalogUserStateFromClause(plan episodeCatalogUserStatePlan) string { + tableName := "user_viewed" + if plan.source == "progress" { + tableName = "user_progress" + } + if plan.positiveSource { + return fmt.Sprintf( + "FROM %s %s JOIN episode_catalog_entries ece ON ece.episode_id = %s.episode_id", + tableName, + plan.alias, + plan.alias, + ) + } + return fmt.Sprintf( + "FROM episode_catalog_entries ece LEFT JOIN %s %s ON %s.episode_id = ece.episode_id", + tableName, + plan.alias, + plan.alias, + ) +} + +func episodeCatalogUserStateOrderBy(plan episodeCatalogUserStatePlan) (string, bool) { + dir := "DESC" + if plan.sortOrder == "asc" { + dir = "ASC" + } + switch plan.sortField { + case "date_viewed": + return fmt.Sprintf("ORDER BY uv.last_watched %s NULLS LAST, ece.sort_key ASC, ece.episode_id ASC", dir), true + case "plays": + return fmt.Sprintf("ORDER BY uv.play_count %s NULLS LAST, ece.sort_key ASC, ece.episode_id ASC", dir), true + case "progress": + return fmt.Sprintf("ORDER BY up.progress_ratio %s NULLS LAST, ece.sort_key ASC, ece.episode_id ASC", dir), true + default: + return episodeCatalogEntryOrderBy(QuerySort{Field: plan.sortField, Order: plan.sortOrder}) + } +} + +func (e *QueryExecutor) queryEpisodeCatalogEntryPage( + ctx context.Context, + sql string, + args ...any, +) ([]episodeCatalogEntryPageRow, error) { + rows, err := e.Pool.Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("querying episode catalog entry page: %w", err) + } + defer rows.Close() + + pageRows := make([]episodeCatalogEntryPageRow, 0) + for rows.Next() { + var row episodeCatalogEntryPageRow + if err := rows.Scan(&row.episodeID, &row.addedAt); err != nil { + return nil, fmt.Errorf("scanning episode catalog entry page: %w", err) + } + pageRows = append(pageRows, row) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("iterating episode catalog entry page: %w", err) + } + return pageRows, nil +} + +func (e *QueryExecutor) hydrateEpisodeCatalogEntryPage( + ctx context.Context, + pageRows []episodeCatalogEntryPageRow, +) ([]*models.MediaItem, error) { + if len(pageRows) == 0 { + return []*models.MediaItem{}, nil + } + + ids := make([]string, 0, len(pageRows)) + addedAtByID := make(map[string]time.Time, len(pageRows)) + for _, row := range pageRows { + ids = append(ids, row.episodeID) + addedAtByID[row.episodeID] = row.addedAt + } + + relation := fmt.Sprintf(episodeCatalogSelectBody, "e.content_id = ANY($1)") + sql := "SELECT " + qualifiedListItemColumns("mi") + " FROM " + relation + rows, err := e.Pool.Query(ctx, sql, ids) + if err != nil { + return nil, fmt.Errorf("hydrating episode catalog entry page: %w", err) + } + defer rows.Close() + + hydrated, err := scanItems(rows) + if err != nil { + return nil, err + } + + itemsByID := make(map[string]*models.MediaItem, len(hydrated)) + for _, item := range hydrated { + if item == nil { + continue + } + if addedAt, ok := addedAtByID[item.ContentID]; ok { + t := addedAt + item.AddedAt = &t + } else if item.AddedAt == nil && !item.CreatedAt.IsZero() { + t := item.CreatedAt + item.AddedAt = &t + } + itemsByID[item.ContentID] = item + } + + ordered := make([]*models.MediaItem, 0, len(pageRows)) + for _, row := range pageRows { + if item := itemsByID[row.episodeID]; item != nil { + ordered = append(ordered, item) + } + } + return ordered, nil +} + +func singleEpisodeCatalogLibraryID(def QueryDefinition, access AccessFilter) (int, bool, bool) { + libraryIDs := append([]int(nil), def.LibraryIDs...) + if access.AllowedLibraryIDs != nil { + if len(libraryIDs) == 0 { + libraryIDs = append([]int(nil), access.AllowedLibraryIDs...) + } else { + libraryIDs = intersectInts(libraryIDs, access.AllowedLibraryIDs) + } + } + if len(libraryIDs) == 0 { + if access.AllowedLibraryIDs != nil { + return 0, true, true + } + return 0, false, false + } + if len(libraryIDs) != 1 { + return 0, false, false + } + libraryID := libraryIDs[0] + if libraryID <= 0 { + return 0, true, true + } + for _, disabledID := range access.DisabledLibraryIDs { + if disabledID == libraryID { + return 0, true, true + } + } + return libraryID, false, true +} + +func buildEpisodeCatalogEntryQueryWhere(def QueryDefinition, argIdx int) (string, []any, int, bool, error) { + if len(def.Groups) == 0 { + return "", nil, argIdx, true, nil + } + + topJoiner := " AND " + if def.Match == "any" { + topJoiner = " OR " + } + + var args []any + var groupClauses []string + for _, group := range def.Groups { + clause, groupArgs, nextArgIdx, ok, err := buildEpisodeCatalogEntryGroupWhere(group, argIdx) + if err != nil || !ok { + return "", nil, argIdx, ok, err + } + argIdx = nextArgIdx + args = append(args, groupArgs...) + if clause != "" { + groupClauses = append(groupClauses, clause) + } + } + + switch len(groupClauses) { + case 0: + return "", args, argIdx, true, nil + case 1: + return groupClauses[0], args, argIdx, true, nil + default: + wrapped := make([]string, len(groupClauses)) + for i, clause := range groupClauses { + wrapped[i] = "(" + clause + ")" + } + return strings.Join(wrapped, topJoiner), args, argIdx, true, nil + } +} + +func buildEpisodeCatalogEntryGroupWhere(group QueryGroup, argIdx int) (string, []any, int, bool, error) { + if len(group.Rules) == 0 { + return "", nil, argIdx, true, nil + } + if group.Match == "all" && countSameFileTechnicalRules(group.Rules) > 1 { + return "", nil, argIdx, false, nil + } + + joiner := " AND " + if group.Match == "any" { + joiner = " OR " + } + + var args []any + var clauses []string + for _, rule := range group.Rules { + clause, ruleArgs, nextArgIdx, ok, err := buildEpisodeCatalogEntryRuleWhere(rule, argIdx) + if err != nil || !ok { + return "", nil, argIdx, ok, err + } + argIdx = nextArgIdx + args = append(args, ruleArgs...) + if clause != "" { + clauses = append(clauses, clause) + } + } + if len(clauses) == 0 { + return "", args, argIdx, true, nil + } + return strings.Join(clauses, joiner), args, argIdx, true, nil +} + +func buildEpisodeCatalogEntryRuleWhere(rule QueryRule, argIdx int) (string, []any, int, bool, error) { + switch rule.Field { + case "type": + return buildEpisodeCatalogTypeClause(rule, argIdx) + case "genre": + return buildEpisodeCatalogArrayClause("ece.genres", rule, argIdx) + case "studio": + return buildEpisodeCatalogArrayClause("ece.studios", rule, argIdx) + case "network": + return buildEpisodeCatalogArrayClause("ece.networks", rule, argIdx) + case "country": + return buildEpisodeCatalogArrayClause("ece.countries", rule, argIdx) + case "year": + return buildEpisodeCatalogComparisonClause("ece.year", rule, argIdx, "") + case "rating_imdb": + return buildEpisodeCatalogComparisonClause("ece.rating_imdb", rule, argIdx, "") + case "original_language": + return buildEpisodeCatalogEqualityClause("ece.original_language", rule, argIdx) + case "content_rating": + return buildEpisodeCatalogEqualityClause("ece.content_rating", rule, argIdx) + case "added_at": + return buildEpisodeCatalogComparisonClause("ece.episode_created_at", rule, argIdx, "timestamptz") + case "release_date": + return buildEpisodeCatalogComparisonClause("ece.episode_air_date", rule, argIdx, "date") + case "status": + return buildEpisodeCatalogEqualityClause("ece.status", rule, argIdx) + case "resolution": + return buildEpisodeCatalogResolutionClause(rule, argIdx) + case "hdr": + return buildEpisodeCatalogBoolPresenceClause("ece.has_hdr", "ece.has_non_hdr", rule, argIdx) + case "dolby_vision": + return buildEpisodeCatalogBoolPresenceClause("ece.has_dolby_vision", "ece.has_non_dolby_vision", rule, argIdx) + case "bitrate": + return buildEpisodeCatalogBitrateClause(rule, argIdx) + case "audio_language": + return buildEpisodeCatalogLanguageArrayClause("ece.audio_language_codes", rule, argIdx) + case "subtitle_language": + return buildEpisodeCatalogLanguageArrayClause("ece.subtitle_language_codes", rule, argIdx) + case "actor", "director", "writer", "producer", "watched", "favorited", "in_watchlist", "in_progress", "last_watched": + return "", nil, argIdx, false, nil + default: + return "", nil, argIdx, false, nil + } +} + +func buildEpisodeCatalogTypeClause(rule QueryRule, argIdx int) (string, []any, int, bool, error) { + value, ok := rule.Value.(string) + if !ok { + return "", nil, argIdx, true, fmt.Errorf("type requires a string value") + } + isEpisode := strings.EqualFold(strings.TrimSpace(value), "episode") + switch rule.Op { + case "is": + if isEpisode { + return "1 = 1", nil, argIdx, true, nil + } + return "1 = 0", nil, argIdx, true, nil + case "is_not": + if isEpisode { + return "1 = 0", nil, argIdx, true, nil + } + return "1 = 1", nil, argIdx, true, nil + default: + return "", nil, argIdx, true, fmt.Errorf("unsupported type operator %q", rule.Op) + } +} + +func buildEpisodeCatalogArrayClause(column string, rule QueryRule, argIdx int) (string, []any, int, bool, error) { + if rule.Op != "is" && rule.Op != "is_not" && rule.Op != "contains" { + return "", nil, argIdx, true, fmt.Errorf("unsupported array operator %q", rule.Op) + } + clause := fmt.Sprintf("%s @> ARRAY[$%d]::text[]", column, argIdx) + if rule.Op == "is_not" { + clause = "NOT (" + clause + ")" + } + return clause, []any{rule.Value}, argIdx + 1, true, nil +} + +func buildEpisodeCatalogLanguageArrayClause(column string, rule QueryRule, argIdx int) (string, []any, int, bool, error) { + value, ok := rule.Value.(string) + if !ok { + return "", nil, argIdx, true, fmt.Errorf("%s requires a string value", rule.Field) + } + rule.Value = strings.ToLower(strings.TrimSpace(value)) + return buildEpisodeCatalogArrayClause(column, rule, argIdx) +} + +func buildEpisodeCatalogResolutionClause(rule QueryRule, argIdx int) (string, []any, int, bool, error) { + value, ok := rule.Value.(string) + if !ok { + return "", nil, argIdx, true, fmt.Errorf("resolution requires a string value") + } + rule.Value = normalizeResolutionValue(value) + return buildEpisodeCatalogArrayClause("ece.resolution_codes", rule, argIdx) +} + +func buildEpisodeCatalogEqualityClause(column string, rule QueryRule, argIdx int) (string, []any, int, bool, error) { + if rule.Op != "is" && rule.Op != "is_not" { + return "", nil, argIdx, true, fmt.Errorf("unsupported equality operator %q", rule.Op) + } + clause := fmt.Sprintf("%s = $%d", column, argIdx) + if rule.Op == "is_not" { + clause = "NOT (" + clause + ")" + } + return clause, []any{rule.Value}, argIdx + 1, true, nil +} + +func buildEpisodeCatalogComparisonClause(column string, rule QueryRule, argIdx int, cast string) (string, []any, int, bool, error) { + placeholder := func(idx int) string { + if cast == "" { + return "$" + strconv.Itoa(idx) + } + return "$" + strconv.Itoa(idx) + "::" + cast + } + + switch rule.Op { + case "gt", "gte", "lt", "lte": + operator := map[string]string{"gt": ">", "gte": ">=", "lt": "<", "lte": "<="}[rule.Op] + return fmt.Sprintf("%s %s %s", column, operator, placeholder(argIdx)), []any{rule.Value}, argIdx + 1, true, nil + case "between": + values, err := toBetweenValues(rule.Value) + if err != nil { + return "", nil, argIdx, true, fmt.Errorf("between requires [min, max] array: %w", err) + } + return fmt.Sprintf("%s >= %s AND %s <= %s", column, placeholder(argIdx), column, placeholder(argIdx+1)), + []any{values[0], values[1]}, argIdx + 2, true, nil + case "is": + return fmt.Sprintf("%s = %s", column, placeholder(argIdx)), []any{rule.Value}, argIdx + 1, true, nil + case "is_not": + return fmt.Sprintf("NOT (%s = %s)", column, placeholder(argIdx)), []any{rule.Value}, argIdx + 1, true, nil + case "in_last": + duration, ok := rule.Value.(string) + if !ok { + return "", nil, argIdx, true, fmt.Errorf("in_last requires a duration string like '30d'") + } + interval, err := parseDuration(duration) + if err != nil { + return "", nil, argIdx, true, err + } + return fmt.Sprintf("%s >= NOW() - INTERVAL '%s'", column, interval), nil, argIdx, true, nil + default: + return "", nil, argIdx, true, fmt.Errorf("unsupported comparison operator %q", rule.Op) + } +} + +func buildEpisodeCatalogBoolPresenceClause(trueColumn, falseColumn string, rule QueryRule, argIdx int) (string, []any, int, bool, error) { + if rule.Op != "is" { + return "", nil, argIdx, true, fmt.Errorf("unsupported boolean operator %q", rule.Op) + } + value, ok := rule.Value.(bool) + if !ok { + return "", nil, argIdx, true, fmt.Errorf("%s requires a boolean value", rule.Field) + } + if value { + return trueColumn, nil, argIdx, true, nil + } + return falseColumn, nil, argIdx, true, nil +} + +func buildEpisodeCatalogBitrateClause(rule QueryRule, argIdx int) (string, []any, int, bool, error) { + switch rule.Op { + case "gt", "gte": + operator := map[string]string{"gt": ">", "gte": ">="}[rule.Op] + return fmt.Sprintf("ece.max_bitrate %s $%d", operator, argIdx), []any{rule.Value}, argIdx + 1, true, nil + case "lt", "lte": + operator := map[string]string{"lt": "<", "lte": "<="}[rule.Op] + return fmt.Sprintf("ece.min_bitrate %s $%d", operator, argIdx), []any{rule.Value}, argIdx + 1, true, nil + case "between": + return "", nil, argIdx, false, nil + default: + return "", nil, argIdx, true, fmt.Errorf("unsupported bitrate operator %q", rule.Op) + } +} + +func episodeCatalogEntryOrderBy(sortConfig QuerySort) (string, bool) { + sortConfig = NormalizeQuerySort(sortConfig) + dir := "DESC" + if sortConfig.Order == "asc" { + dir = "ASC" + } + + switch sortConfig.Field { + case "title": + return fmt.Sprintf("ORDER BY ece.sort_key %s, ece.episode_id ASC", dir), true + case "added_at": + return fmt.Sprintf("ORDER BY ece.added_at %s, ece.sort_key ASC, ece.episode_id ASC", dir), true + case "release_date", "last_air_date": + return fmt.Sprintf("ORDER BY ece.episode_air_date %s NULLS LAST, ece.sort_key ASC, ece.episode_id ASC", dir), true + case "year": + return fmt.Sprintf("ORDER BY ece.year %s, ece.sort_key ASC, ece.episode_id ASC", dir), true + case "content_rating": + return fmt.Sprintf("ORDER BY ece.content_rating_rank %s, ece.content_rating_label %s, ece.sort_key ASC, ece.episode_id ASC", dir, dir), true + case "runtime": + return fmt.Sprintf("ORDER BY ece.runtime %s NULLS LAST, ece.sort_key ASC, ece.episode_id ASC", dir), true + case "rating_imdb": + return fmt.Sprintf("ORDER BY ece.rating_imdb %s NULLS LAST, ece.sort_key ASC, ece.episode_id ASC", dir), true + case "rating_tmdb": + return fmt.Sprintf("ORDER BY ece.rating_tmdb %s NULLS LAST, ece.sort_key ASC, ece.episode_id ASC", dir), true + case "rating_rt_critic", "rating_rt_audience": + return "ORDER BY ece.sort_key ASC, ece.episode_id ASC", true + case "resolution": + return fmt.Sprintf("ORDER BY ece.max_resolution_rank %s NULLS LAST, ece.sort_key ASC, ece.episode_id ASC", dir), true + case "bitrate": + return fmt.Sprintf("ORDER BY ece.max_bitrate %s NULLS LAST, ece.sort_key ASC, ece.episode_id ASC", dir), true + default: + return "", false + } +} + +func countSameFileTechnicalRules(rules []QueryRule) int { + count := 0 + for _, rule := range rules { + if isSameFileTechnicalRule(rule) { + count++ + } + } + return count +} + +func episodeCatalogEntriesUnavailable(err error) bool { + var pgErr *pgconn.PgError + if !errors.As(err, &pgErr) { + return false + } + return pgErr.Code == "42P01" || pgErr.Code == "42883" +} diff --git a/internal/catalog/episode_catalog_entries_executor_test.go b/internal/catalog/episode_catalog_entries_executor_test.go new file mode 100644 index 00000000..f255f7ce --- /dev/null +++ b/internal/catalog/episode_catalog_entries_executor_test.go @@ -0,0 +1,206 @@ +package catalog + +import ( + "strings" + "testing" +) + +func TestEpisodeCatalogEntryOrderByUsesReadModelColumns(t *testing.T) { + tests := []struct { + name string + sort QuerySort + want []string + }{ + { + name: "resolution", + sort: QuerySort{Field: "resolution", Order: "desc"}, + want: []string{"ece.max_resolution_rank DESC NULLS LAST", "ece.sort_key ASC", "ece.episode_id ASC"}, + }, + { + name: "bitrate", + sort: QuerySort{Field: "bitrate", Order: "desc"}, + want: []string{"ece.max_bitrate DESC NULLS LAST", "ece.sort_key ASC", "ece.episode_id ASC"}, + }, + { + name: "date added", + sort: QuerySort{Field: "added_at", Order: "desc"}, + want: []string{"ece.added_at DESC", "ece.sort_key ASC", "ece.episode_id ASC"}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + orderBy, ok := episodeCatalogEntryOrderBy(tt.sort) + if !ok { + t.Fatalf("episodeCatalogEntryOrderBy(%+v) did not produce a plan", tt.sort) + } + for _, want := range tt.want { + if !strings.Contains(orderBy, want) { + t.Fatalf("ORDER BY %q does not contain %q", orderBy, want) + } + } + }) + } +} + +func TestBuildEpisodeCatalogEntryQueryWhereUsesReadModelFilters(t *testing.T) { + def := QueryDefinition{ + Match: "all", + Groups: []QueryGroup{{ + Match: "all", + Rules: []QueryRule{ + {Field: "genre", Op: "is", Value: "Comedy"}, + {Field: "content_rating", Op: "is", Value: "TV-14"}, + {Field: "year", Op: "between", Value: []any{2015, 2016}}, + }, + }}, + }.Normalize() + + where, args, nextArgIdx, ok, err := buildEpisodeCatalogEntryQueryWhere(def, 2) + if err != nil { + t.Fatalf("buildEpisodeCatalogEntryQueryWhere returned error: %v", err) + } + if !ok { + t.Fatal("buildEpisodeCatalogEntryQueryWhere unexpectedly fell back") + } + for _, want := range []string{ + "ece.genres @> ARRAY[$2]::text[]", + "ece.content_rating = $3", + "ece.year >= $4 AND ece.year <= $5", + } { + if !strings.Contains(where, want) { + t.Fatalf("WHERE %q does not contain %q", where, want) + } + } + if nextArgIdx != 6 { + t.Fatalf("nextArgIdx = %d, want 6", nextArgIdx) + } + if len(args) != 4 { + t.Fatalf("len(args) = %d, want 4", len(args)) + } +} + +func TestEpisodeCatalogAddedAtFilterUsesEpisodeCreatedAt(t *testing.T) { + def := QueryDefinition{ + Match: "all", + Groups: []QueryGroup{{ + Match: "all", + Rules: []QueryRule{ + {Field: "added_at", Op: "in_last", Value: "30d"}, + }, + }}, + }.Normalize() + + where, _, _, ok, err := buildEpisodeCatalogEntryQueryWhere(def, 2) + if err != nil { + t.Fatalf("buildEpisodeCatalogEntryQueryWhere returned error: %v", err) + } + if !ok { + t.Fatal("buildEpisodeCatalogEntryQueryWhere unexpectedly fell back") + } + if !strings.Contains(where, "ece.episode_created_at >= NOW() - INTERVAL '30 days'") { + t.Fatalf("expected added_at filter to use episode_created_at, got %q", where) + } + if strings.Contains(where, "ece.added_at") { + t.Fatalf("added_at filter must not use library first_seen_at, got %q", where) + } +} + +func TestEpisodeCatalogEntryQueryFallsBackForSameFileTechnicalAnd(t *testing.T) { + group := QueryGroup{ + Match: "all", + Rules: []QueryRule{ + {Field: "resolution", Op: "is", Value: "2160p"}, + {Field: "audio_language", Op: "is", Value: "en"}, + }, + } + + _, _, _, ok, err := buildEpisodeCatalogEntryGroupWhere(group, 1) + if err != nil { + t.Fatalf("buildEpisodeCatalogEntryGroupWhere returned error: %v", err) + } + if ok { + t.Fatal("same-file technical AND should fall back to the generic media_files EXISTS path") + } +} + +func TestEpisodeCatalogInProgressPlanKeepsUnknownDurationRows(t *testing.T) { + def := QueryDefinition{ + Match: "all", + Groups: []QueryGroup{{ + Match: "all", + Rules: []QueryRule{ + {Field: "in_progress", Op: "is", Value: true}, + }, + }}, + Sort: QuerySort{Field: "title", Order: "asc"}, + }.Normalize() + + plan, ok, err := extractEpisodeCatalogUserStatePlan(def) + if err != nil { + t.Fatalf("extractEpisodeCatalogUserStatePlan returned error: %v", err) + } + if !ok { + t.Fatal("expected in_progress filter to use a user-state plan") + } + sql := episodeCatalogUserStateCTE(plan) + if strings.Contains(sql, "COALESCE(uwp.duration_seconds, 0) > 0") { + t.Fatalf("in_progress filter must keep unknown-duration rows, got CTE:\n%s", sql) + } +} + +func TestEpisodeCatalogProgressSortRequiresProgressRatio(t *testing.T) { + def := QueryDefinition{ + Match: "all", + Groups: []QueryGroup{{ + Match: "all", + Rules: []QueryRule{ + {Field: "genre", Op: "is", Value: "Comedy"}, + }, + }}, + Sort: QuerySort{Field: "progress", Order: "desc"}, + }.Normalize() + + plan, ok, err := extractEpisodeCatalogUserStatePlan(def) + if err != nil { + t.Fatalf("extractEpisodeCatalogUserStatePlan returned error: %v", err) + } + if !ok { + t.Fatal("expected progress sort to use a user-state plan") + } + sql := episodeCatalogUserStateCTE(plan) + if !strings.Contains(sql, "COALESCE(uwp.duration_seconds, 0) > 0") { + t.Fatalf("progress sort must require a computable ratio before falling back, got CTE:\n%s", sql) + } +} + +func TestExtractEpisodeCatalogUserStatePlanForLastWatched(t *testing.T) { + def := QueryDefinition{ + Match: "all", + Groups: []QueryGroup{{ + Match: "all", + Rules: []QueryRule{ + {Field: "last_watched", Op: "in_last", Value: "30d"}, + {Field: "genre", Op: "is", Value: "Comedy"}, + }, + }}, + Sort: QuerySort{Field: "title", Order: "asc"}, + }.Normalize() + + plan, ok, err := extractEpisodeCatalogUserStatePlan(def) + if err != nil { + t.Fatalf("extractEpisodeCatalogUserStatePlan returned error: %v", err) + } + if !ok { + t.Fatal("expected last_watched filter to use a user-state plan") + } + if plan.source != "viewed" || plan.alias != "uv" || !plan.positiveSource { + t.Fatalf("unexpected user-state plan: %+v", plan) + } + if len(plan.userClauses) != 1 || !strings.Contains(plan.userClauses[0], "uv.last_watched") { + t.Fatalf("unexpected user clauses: %#v", plan.userClauses) + } + if len(plan.entryDef.Groups) != 1 || len(plan.entryDef.Groups[0].Rules) != 1 { + t.Fatalf("expected genre rule to remain in entry def, got %+v", plan.entryDef.Groups) + } +} diff --git a/internal/catalog/episode_catalog_source.go b/internal/catalog/episode_catalog_source.go index 3988672c..56b19e3a 100644 --- a/internal/catalog/episode_catalog_source.go +++ b/internal/catalog/episode_catalog_source.go @@ -12,6 +12,7 @@ const episodeCatalogSelectBody = `( 'episode'::text AS type, COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text) AS title, COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text) AS sort_title, + LOWER(COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text)) AS sort_key, COALESCE(NULLIF(BTRIM(e.default_metadata_language), ''), COALESCE(si.default_metadata_language, '')) AS default_metadata_language, ''::text AS original_title, COALESCE(si.year, EXTRACT(YEAR FROM e.air_date)::integer, 0) AS year, @@ -45,6 +46,7 @@ const episodeCatalogSelectBody = `( NULL::text AS last_air_date, e.air_date AS last_air_date_at, si.air_time, + COALESCE(si.show_status, '') AS show_status, si.matched_at, si.last_refreshed, si.refresh_failures, diff --git a/internal/catalog/episode_library_repo.go b/internal/catalog/episode_library_repo.go index 8b4e4d2a..f07500e7 100644 --- a/internal/catalog/episode_library_repo.go +++ b/internal/catalog/episode_library_repo.go @@ -17,9 +17,24 @@ func NewEpisodeLibraryRepository(pool *pgxpool.Pool) *EpisodeLibraryRepository { return &EpisodeLibraryRepository{pool: pool} } -// ReconcileFolderMembership removes episode memberships for episodes that no longer -// have any present files in the given folder. +// ReconcileFolderMembership restores missing episode memberships and removes +// memberships for episodes that no longer have any present files in the given +// folder. The returned count is the number of removed stale memberships. func (r *EpisodeLibraryRepository) ReconcileFolderMembership(ctx context.Context, folderID int) (int, error) { + if _, err := r.pool.Exec(ctx, ` + INSERT INTO episode_libraries (episode_id, media_folder_id, first_seen_at) + SELECT mf.episode_id, mf.media_folder_id, MIN(mf.created_at) + FROM media_files mf + JOIN episodes e ON e.content_id = mf.episode_id + WHERE mf.media_folder_id = $1 + AND mf.missing_since IS NULL + AND mf.episode_id IS NOT NULL + GROUP BY mf.episode_id, mf.media_folder_id + ON CONFLICT (episode_id, media_folder_id) DO NOTHING + `, folderID); err != nil { + return 0, fmt.Errorf("restoring episode library membership: %w", err) + } + tag, err := r.pool.Exec(ctx, ` DELETE FROM episode_libraries el WHERE el.media_folder_id = $1 diff --git a/internal/catalog/name_prefix_sql_test.go b/internal/catalog/name_prefix_sql_test.go index 629b8b6d..9f30bffa 100644 --- a/internal/catalog/name_prefix_sql_test.go +++ b/internal/catalog/name_prefix_sql_test.go @@ -23,7 +23,7 @@ func TestQueryExecutor_NamePrefix_PushedIntoWHERE(t *testing.T) { } // First arm: sort-key expression matching idx_media_items_sort_key. - if !strings.Contains(sql, "LOWER(COALESCE(NULLIF(BTRIM(mi.sort_title),''), mi.title)) LIKE") { + if !strings.Contains(sql, "LOWER(COALESCE(NULLIF(BTRIM(mi.sort_title), ''), mi.title)) LIKE") { t.Fatalf("expected sort-key LIKE arm matching idx_media_items_sort_key; got %q", sql) } // Second arm: LOWER(title) matching idx_media_items_search_exact_title; @@ -53,6 +53,31 @@ func TestQueryExecutor_NamePrefix_PushedIntoWHERE(t *testing.T) { } } +func TestQueryExecutor_NamePrefix_UsesEpisodeSortKeyForEpisodeScope(t *testing.T) { + exec := &QueryExecutor{} + access := AccessFilter{NamePrefix: "Pilot"} + + sql, args, err := exec.buildPreviewPageSQL( + QueryDefinition{MediaScope: "episode", LibraryIDs: []int{2}}, + access, + 20, + 0, + true, + ) + if err != nil { + t.Fatalf("buildPreviewPageSQL returned error: %v", err) + } + if !strings.Contains(sql, "mi.sort_key LIKE") { + t.Fatalf("expected episode prefix to use projected sort_key; got %q", sql) + } + if strings.Contains(sql, "BTRIM(mi.sort_title)") { + t.Fatalf("episode prefix should not recompute sort_title expression; got %q", sql) + } + if len(args) < 2 || args[1] != "pilot%" { + t.Fatalf("expected episode library args followed by prefix arg; got %v", args) + } +} + // TestBrowseFilters_NamePrefix_BothArmsSargable asserts that the dual-column // LIKE in BrowseRepository's WHERE clause uses an expression that matches // idx_media_items_sort_key (migration 102) on the first arm and diff --git a/internal/catalog/query_builder.go b/internal/catalog/query_builder.go index 2126aecd..f56a222f 100644 --- a/internal/catalog/query_builder.go +++ b/internal/catalog/query_builder.go @@ -234,10 +234,10 @@ func (qb *QueryBuilder) BuildSortPlan(sortConfig QuerySort) (QuerySortPlan, erro plan.OrderBy = qb.orderByExpr(qb.releaseDateSortExpr(), dir, true, titleExpr) return plan, nil case "added_at": - expr, joins, args := qb.addedAtSortPlan() + expr, joins, args, nullsLast := qb.addedAtSortPlan() plan.Joins = joins plan.Args = args - plan.OrderBy = qb.orderByExpr(expr, dir, len(joins) > 0, titleExpr) + plan.OrderBy = qb.orderByExpr(expr, dir, nullsLast, titleExpr) return plan, nil case "content_rating": rankExpr := qb.contentRatingRankExpr() @@ -1016,6 +1016,9 @@ func (qb *QueryBuilder) buildTimestampComparisonClause(column, op string, value } func (qb *QueryBuilder) normalizedTitleExpr() string { + if isEpisodeCatalogScope(qb.mediaScope) { + return fmt.Sprintf("%s.sort_key", qb.alias) + } return fmt.Sprintf( "LOWER(COALESCE(NULLIF(BTRIM(%s.sort_title), ''), %s.title))", qb.alias, @@ -1039,13 +1042,22 @@ func (qb *QueryBuilder) orderByExpr(expr, dir string, nullsLast bool, titleExpr return clause } -func (qb *QueryBuilder) addedAtSortPlan() (string, []string, []any) { +func (qb *QueryBuilder) addedAtSortPlan() (string, []string, []any, bool) { if len(qb.libraryIDs) == 0 { - return fmt.Sprintf("%s.created_at", qb.alias), nil, nil + return fmt.Sprintf("%s.created_at", qb.alias), nil, nil, false } placeholders, args := qb.consumeIntArgs(qb.libraryIDs) if isEpisodeCatalogScope(qb.mediaScope) { + if len(qb.libraryIDs) == 1 { + joinSQL := fmt.Sprintf( + `JOIN episode_libraries sort_added ON sort_added.episode_id = %s.content_id AND sort_added.media_folder_id = %s`, + qb.alias, + placeholders[0], + ) + return "sort_added.first_seen_at", []string{joinSQL}, args, false + } + joinSQL := fmt.Sprintf( `LEFT JOIN ( SELECT el.episode_id AS content_id, MIN(el.first_seen_at) AS added_at @@ -1056,7 +1068,7 @@ func (qb *QueryBuilder) addedAtSortPlan() (string, []string, []any) { strings.Join(placeholders, ", "), qb.alias, ) - return "sort_added.added_at", []string{joinSQL}, args + return "sort_added.added_at", []string{joinSQL}, args, true } joinSQL := fmt.Sprintf( @@ -1069,7 +1081,7 @@ func (qb *QueryBuilder) addedAtSortPlan() (string, []string, []any) { strings.Join(placeholders, ", "), qb.libraryContentExpr(), ) - return "sort_added.added_at", []string{joinSQL}, args + return "sort_added.added_at", []string{joinSQL}, args, true } func (qb *QueryBuilder) contentRatingRankExpr() string { diff --git a/internal/catalog/query_builder_test.go b/internal/catalog/query_builder_test.go index dd493de0..41400c14 100644 --- a/internal/catalog/query_builder_test.go +++ b/internal/catalog/query_builder_test.go @@ -74,6 +74,27 @@ func TestBuildSortClause_ReleaseDateUsesEpisodeAirDateForEpisodeScope(t *testing } } +func TestBuildSortClause_TitleUsesEpisodeSortKeyForEpisodeScope(t *testing.T) { + clause, args, err := NewQueryBuilder("mi"). + WithMediaScope("episode"). + BuildSortClause(QuerySort{ + Field: "title", + Order: "asc", + }) + if err != nil { + t.Fatalf("BuildSortClause returned error: %v", err) + } + if len(args) != 0 { + t.Fatalf("expected no args, got %v", args) + } + if !strings.Contains(clause, "ORDER BY mi.sort_key ASC, mi.content_id ASC") { + t.Fatalf("expected episode title sort to use indexed sort_key, got %q", clause) + } + if strings.Contains(clause, "BTRIM(mi.sort_title)") { + t.Fatalf("episode title sort should not recompute sort_title expression, got %q", clause) + } +} + func TestBuildSortClause_AddedAtUsesScopedFirstSeenAt(t *testing.T) { plan, err := NewQueryBuilder("mi"). WithLibraryScope([]int{3, 7}). @@ -201,15 +222,43 @@ func TestBuildSortPlan_AddedAtUsesEpisodeLibraryMembershipForEpisodeScope(t *tes if len(plan.Joins) != 1 { t.Fatalf("expected one join, got %v", plan.Joins) } - if !strings.Contains(plan.Joins[0], "FROM episode_libraries el") { + if !strings.Contains(plan.Joins[0], "JOIN episode_libraries sort_added") { t.Fatalf("expected episode added_at join to use episode_libraries, got %q", plan.Joins[0]) } - if !strings.Contains(plan.Joins[0], "GROUP BY el.episode_id") { - t.Fatalf("expected episode added_at join to group by episode_id, got %q", plan.Joins[0]) + if !strings.Contains(plan.Joins[0], "sort_added.media_folder_id = $1") { + t.Fatalf("expected direct episode added_at join to bind one library, got %q", plan.Joins[0]) } - if !strings.Contains(plan.Joins[0], "sort_added.content_id = mi.content_id") { + if !strings.Contains(plan.Joins[0], "sort_added.episode_id = mi.content_id") { t.Fatalf("expected episode added_at join to match episode content_id, got %q", plan.Joins[0]) } + if !strings.Contains(plan.OrderBy, "ORDER BY sort_added.first_seen_at DESC, mi.sort_key ASC, mi.content_id ASC") { + t.Fatalf("expected single-library episode added_at sort to use first_seen_at without NULLS LAST, got %q", plan.OrderBy) + } +} + +func TestBuildSortPlan_AddedAtAggregatesEpisodeLibrariesForMultiLibraryScope(t *testing.T) { + plan, err := NewQueryBuilder("mi"). + WithMediaScope("episode"). + WithLibraryScope([]int{6, 7}). + BuildSortPlan(QuerySort{Field: "added_at", Order: "desc"}) + if err != nil { + t.Fatalf("BuildSortPlan returned error: %v", err) + } + if len(plan.Args) != 2 || plan.Args[0] != 6 || plan.Args[1] != 7 { + t.Fatalf("expected scoped library args [6 7], got %v", plan.Args) + } + if len(plan.Joins) != 1 { + t.Fatalf("expected one join, got %v", plan.Joins) + } + if !strings.Contains(plan.Joins[0], "FROM episode_libraries el") { + t.Fatalf("expected episode added_at aggregate to use episode_libraries, got %q", plan.Joins[0]) + } + if !strings.Contains(plan.Joins[0], "GROUP BY el.episode_id") { + t.Fatalf("expected multi-library episode added_at join to group by episode_id, got %q", plan.Joins[0]) + } + if !strings.Contains(plan.OrderBy, "sort_added.added_at DESC NULLS LAST") { + t.Fatalf("expected multi-library episode added_at sort to keep NULLS LAST, got %q", plan.OrderBy) + } } func TestBuildSortPlan_FileSortUsesEpisodeIDsForEpisodeScope(t *testing.T) { diff --git a/internal/catalog/query_executor.go b/internal/catalog/query_executor.go index 6c0ea0a0..ce113cde 100644 --- a/internal/catalog/query_executor.go +++ b/internal/catalog/query_executor.go @@ -42,12 +42,34 @@ func (e *QueryExecutor) PreviewPage( return nil, 0, false, fmt.Errorf("query executor requires a database pool") } + if items, total, hasMore, ok, err := e.tryEpisodeCatalogUserStatePreviewPage( + ctx, + def, + access, + limit, + offset, + includeTotal, + ); ok || err != nil { + return items, total, hasMore, err + } + + if items, total, hasMore, ok, err := e.tryEpisodeCatalogEntriesPreviewPage( + ctx, + def, + access, + limit, + offset, + includeTotal, + ); ok || err != nil { + return items, total, hasMore, err + } + build, err := e.buildPreviewPagePlan(def, access, limit, offset) if err != nil { return nil, 0, false, err } - pagedSQL, pagedArgs := build.pagedSQL(includeTotal) + pagedSQL, pagedArgs := build.pagedSQL(false) rows, err := e.Pool.Query(ctx, pagedSQL, pagedArgs...) if err != nil { return nil, 0, false, fmt.Errorf("querying preview items: %w", err) @@ -58,14 +80,15 @@ func (e *QueryExecutor) PreviewPage( items []*models.MediaItem total int ) - if includeTotal { - items, total, err = scanItemsWithTotal(rows) - } else { - items, err = scanItems(rows) - } + items, err = scanItems(rows) if err != nil { return nil, 0, false, err } + hasMore := false + if len(items) > build.limit { + hasMore = true + items = items[:build.limit] + } // The preview path uses itemColumns which scans CreatedAt but not // AddedAt (set only by browse queries via MIN(mil.first_seen_at)). // Fall back to CreatedAt so the API response includes added_at. @@ -76,32 +99,15 @@ func (e *QueryExecutor) PreviewPage( } } - hasMore := false if includeTotal { - // COUNT(*) OVER () emits no rows when the data SELECT is empty, so - // total stays 0 even when the broader result set has matching rows - // (e.g. OFFSET past the last page). Re-query the count to give - // callers the real total. Skip when offset == 0 because in that - // case an empty page genuinely means total = 0. - // - // Use build.offset (normalized in buildPreviewPagePlan) rather than - // the raw offset parameter — a caller-supplied negative offset is - // floored to 0 in the plan, and the SQL uses the plan's value, so - // the fallback condition and HasMore must match. - if len(items) == 0 && build.offset > 0 { - countSQL, countArgs := build.countSQL() - if err := e.Pool.QueryRow(ctx, countSQL, countArgs...).Scan(&total); err != nil { - return nil, 0, false, fmt.Errorf("count fallback for empty page: %w", err) - } + countSQL, countArgs := build.countSQL() + if err := e.Pool.QueryRow(ctx, countSQL, countArgs...).Scan(&total); err != nil { + return nil, 0, false, fmt.Errorf("counting preview items: %w", err) } hasMore = total > build.offset+len(items) return items, total, hasMore, nil } - if len(items) > build.limit { - hasMore = true - items = items[:build.limit] - } return items, 0, hasMore, nil } @@ -121,6 +127,9 @@ type previewPagePlan struct { // fromClausePaged is the FROM clause for the paged query (includes any // sort-plan joins). fromClausePaged string + // fromClauseCount is the FROM clause for exact totals. It includes filter + // joins, but intentionally excludes sort-only joins. + fromClauseCount string whereClause string args []any orderBy string @@ -131,41 +140,33 @@ type previewPagePlan struct { } // countSQL renders a count-only query that returns the total number of rows -// matching the plan's WHERE clause, ignoring LIMIT/OFFSET/ORDER BY. Used as -// a fallback when pagedSQL(true) returned an empty page past offset 0: -// COUNT(*) OVER () emits no rows when the data SELECT is empty, so the -// caller would otherwise see total=0 even when the broader result set has -// matching rows. Wraps the inner query in `SELECT COUNT(*) FROM (...) sub` -// so any GROUP BY in the inner query is preserved (we count groups, matching -// what COUNT(*) OVER () would compute). +// matching the plan's WHERE clause, ignoring LIMIT/OFFSET/ORDER BY. PreviewPage +// uses it when callers need an exact total; keeping the count separate lets the +// data SELECT use top-N/index plans instead of forcing COUNT(*) OVER () across +// every matching row. Wraps the inner query in `SELECT COUNT(*) FROM (...) sub` +// so any GROUP BY in the inner query is preserved. // -// Bind cteArgs + args + sortArgs in the same order as pagedSQL. fromClausePaged -// embeds sort-plan join clauses (ORDER BY needs them to project the join -// columns), and those clauses reference $N placeholders for sortArgs; omitting -// sortArgs here would break sorts that need bound join args (added_at, -// progress, date_viewed, plays, resolution, bitrate). LIMIT/OFFSET are -// intentionally dropped — count is over the full filtered set. +// Bind cteArgs + args only. Sort-only joins and their args are intentionally +// excluded because ordering does not affect the filtered row count. func (p previewPagePlan) countSQL() (string, []any) { args := append([]any{}, p.cteArgs...) args = append(args, p.args...) - args = append(args, p.sortArgs...) withClause := "" if len(p.ctes) > 0 { withClause = "WITH " + strings.Join(p.ctes, ",\n") + "\n" } sql := fmt.Sprintf( "%sSELECT COUNT(*) FROM (SELECT 1 %s %s) sub", - withClause, p.fromClausePaged, p.whereClause, + withClause, p.fromClauseCount, p.whereClause, ) return sql, args } // pagedSQL renders the final paged SELECT and returns it together with the -// fully-bound arg list. When includeTotal is true the SELECT list is appended -// with COUNT(*) OVER () AS total_count so the caller can read the total from -// the first scanned row in a single round trip. When includeTotal is false we -// ask the database for one extra row to detect more pages without an exact -// count. +// fully-bound arg list. When includeTotal is false we ask the database for one +// extra row to detect more pages without an exact count. Exact totals are +// handled by countSQL instead of COUNT(*) OVER () so the page query can stop +// after the requested rows. func (p previewPagePlan) pagedSQL(includeTotal bool) (string, []any) { queryLimit := p.limit if !includeTotal { @@ -182,9 +183,6 @@ func (p previewPagePlan) pagedSQL(includeTotal bool) (string, []any) { args = append(args, p.offset) } selectList := qualifiedListItemColumns("mi") - if includeTotal { - selectList += ", COUNT(*) OVER () AS total_count" - } withClause := "" if len(p.ctes) > 0 { withClause = "WITH " + strings.Join(p.ctes, ",\n") + "\n" @@ -257,7 +255,7 @@ func (e *QueryExecutor) buildPreviewPagePlan( builder := NewQueryBuilder("mi"). WithArgIdx(len(baseArgs)+1). WithUserScope(access.UserID, access.ProfileID). - WithMediaScope(def.MediaScope). + WithMediaScope(effectiveScope). WithLibraryScope(libraryIDs) filterWhere, filterArgs, err := builder.Build(def) if err != nil { @@ -314,16 +312,15 @@ func (e *QueryExecutor) buildPreviewPagePlan( } if prefix := strings.TrimSpace(access.NamePrefix); prefix != "" { - // Dual-column OR matching browse.go and favorites_browse.go: items - // where a curated sort_title differs from title (e.g. title="The Office", - // sort_title="Office, The") would be silently lost on prefix="the" if - // we only checked the COALESCE'd sort-key expression. First arm uses - // the idx_media_items_sort_key expression (migration 102); second arm - // uses idx_media_items_search_exact_title on LOWER(title) (migration 001). - // Both arms are sargable; the planner can BitmapOr them. + // Dual-column OR matching browse.go and favorites_browse.go: items where + // a curated sort_title differs from title (e.g. title="The Office", + // sort_title="Office, The") would be silently lost on prefix="the" if we + // only checked the COALESCE'd sort-key expression. The first arm uses the + // scope-specific sort key; the second arm keeps literal title prefixes. + prefixSortExpr := builder.normalizedTitleExpr() conditions = append(conditions, fmt.Sprintf( - "(LOWER(COALESCE(NULLIF(BTRIM(mi.sort_title),''), mi.title)) LIKE $%d ESCAPE '\\' OR LOWER(mi.title) LIKE $%d ESCAPE '\\')", - argIdx, argIdx, + "(%s LIKE $%d ESCAPE '\\' OR LOWER(mi.title) LIKE $%d ESCAPE '\\')", + prefixSortExpr, argIdx, argIdx, )) args = append(args, escapePrefixForLike(prefix)+"%") argIdx++ @@ -334,6 +331,7 @@ func (e *QueryExecutor) buildPreviewPagePlan( whereClause = "WHERE " + strings.Join(conditions, " AND ") } fromClauseBase := "FROM " + baseRelation + fromClauseCount := fromClauseBase if limit <= 0 { limit = 20 @@ -363,9 +361,11 @@ func (e *QueryExecutor) buildPreviewPagePlan( ctes = []string{UserHistoryCTESQL(1)} fromClausePaged = rebindSQLPlaceholders(fromClausePaged, cteShift) + fromClauseCount = rebindSQLPlaceholders(fromClauseCount, cteShift) whereClause = rebindSQLPlaceholders(whereClause, cteShift) sortPlan.OrderBy = rebindSQLPlaceholders(sortPlan.OrderBy, cteShift) fromClausePaged += " LEFT JOIN user_last_watched uhist ON uhist.media_item_id = mi.content_id" + fromClauseCount += " LEFT JOIN user_last_watched uhist ON uhist.media_item_id = mi.content_id" limitArgIdx += cteShift } @@ -373,6 +373,7 @@ func (e *QueryExecutor) buildPreviewPagePlan( ctes: ctes, cteArgs: cteArgs, fromClausePaged: fromClausePaged, + fromClauseCount: fromClauseCount, whereClause: whereClause, args: args, orderBy: sortPlan.OrderBy, @@ -385,8 +386,9 @@ func (e *QueryExecutor) buildPreviewPagePlan( // buildPreviewPageSQL is a test-friendly facade over buildPreviewPagePlan that // returns the rendered paged SELECT plus the bound args. It performs no I/O. -// When includeTotal is true the emitted SELECT carries COUNT(*) OVER () as a -// total_count column so PreviewPage can avoid a separate count query. +// includeTotal only controls whether the page query fetches exactly limit rows +// or an extra row for has-more detection; exact totals are rendered separately +// by previewPagePlan.countSQL. func (e *QueryExecutor) buildPreviewPageSQL( def QueryDefinition, access AccessFilter, @@ -413,11 +415,10 @@ func escapePrefixForLike(s string) string { } // buildLibraryScopeJoin returns a WHERE clause that scopes the outer query -// to items whose membership in media_item_libraries (or episode_libraries for -// the episode catalog scope) matches the allow/deny lists. The clause is an -// EXISTS / NOT EXISTS semi-join that uses the (content_id, media_folder_id) -// PRIMARY KEY index directly without fanning out for items present in -// multiple libraries — Audit Pattern D (2026-05-01 §3 Pattern D). The prior +// to items whose library membership matches the allow/deny lists. The clause +// is an EXISTS / NOT EXISTS semi-join that uses membership indexes +// directly without fanning out for items present in multiple libraries — Audit +// Pattern D (2026-05-01 §3 Pattern D). The prior // shape wrapped the join in a SELECT DISTINCT subquery to defuse that // fanout; the DISTINCT was load-bearing because the PK is on the (content, // folder) PAIR, not on content alone. EXISTS is the canonical non-fanout diff --git a/internal/catalog/query_executor_test.go b/internal/catalog/query_executor_test.go index 128fdaeb..98881851 100644 --- a/internal/catalog/query_executor_test.go +++ b/internal/catalog/query_executor_test.go @@ -167,6 +167,64 @@ func TestEpisodeCatalogBaseRelationForLibraries_UsesEpisodeLibraries(t *testing. } } +func TestEpisodeCatalogProjectionIncludesSharedCatalogColumns(t *testing.T) { + sql, _, err := (&QueryExecutor{}).buildPreviewPageSQL( + QueryDefinition{ + MediaScope: "episode", + LibraryIDs: []int{2}, + Sort: QuerySort{Field: "title", Order: "asc"}, + }, + AccessFilter{}, + 20, + 0, + true, + ) + if err != nil { + t.Fatalf("buildPreviewPageSQL error: %v", err) + } + if !strings.Contains(sql, "COALESCE(si.show_status, '') AS show_status") { + t.Fatalf("expected episode projection to include show_status, got %s", sql) + } + if !strings.Contains(sql, "mi.show_status") { + t.Fatalf("expected outer catalog select to reference show_status, got %s", sql) + } + if !strings.Contains(sql, "LOWER(COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text)) AS sort_key") { + t.Fatalf("expected episode projection to include sort_key, got %s", sql) + } + if !strings.Contains(sql, "ORDER BY mi.sort_key ASC, mi.content_id ASC") { + t.Fatalf("expected episode title sort to use sort_key, got %s", sql) + } +} + +func TestEpisodeCatalogSingleLibraryAddedAtUsesDirectMembershipJoin(t *testing.T) { + sql, _, err := (&QueryExecutor{}).buildPreviewPageSQL( + QueryDefinition{ + MediaScope: "episode", + LibraryIDs: []int{2}, + Sort: QuerySort{Field: "added_at", Order: "desc"}, + }, + AccessFilter{}, + 20, + 0, + true, + ) + if err != nil { + t.Fatalf("buildPreviewPageSQL error: %v", err) + } + if !strings.Contains(sql, "JOIN episode_libraries sort_added") { + t.Fatalf("expected direct episode_libraries join for single-library added_at sort, got %s", sql) + } + if !strings.Contains(sql, "sort_added.media_folder_id = $2") { + t.Fatalf("expected direct added_at join to bind the single library, got %s", sql) + } + if strings.Contains(sql, "GROUP BY el.episode_id") { + t.Fatalf("single-library added_at sort should avoid aggregate membership join, got %s", sql) + } + if !strings.Contains(sql, "ORDER BY sort_added.first_seen_at DESC, mi.sort_key ASC, mi.content_id ASC") { + t.Fatalf("expected added_at order to use first_seen_at without NULLS LAST, got %s", sql) + } +} + func TestRebindSQLPlaceholders(t *testing.T) { got := rebindSQLPlaceholders("mi.created_at <= $1 AND mi.year >= $2", 3) want := "mi.created_at <= $4 AND mi.year >= $5" diff --git a/internal/catalog/window_count_test.go b/internal/catalog/window_count_test.go index 5f5c5f39..3ac01402 100644 --- a/internal/catalog/window_count_test.go +++ b/internal/catalog/window_count_test.go @@ -12,17 +12,18 @@ import ( // for use in arg-count assertions across countSQL tests. var placeholderRE = regexp.MustCompile(`\$(\d+)`) -// TestQueryExecutor_PreviewPage_UsesWindowCount asserts that buildPreviewPageSQL -// emits a single-pass paged SELECT that includes COUNT(*) OVER () so PreviewPage -// no longer needs a separate count query when includeTotal is true. -func TestQueryExecutor_PreviewPage_UsesWindowCount(t *testing.T) { +// TestQueryExecutor_PreviewPage_ExactTotalOmitsWindowCount asserts that +// buildPreviewPageSQL keeps the data SELECT free of COUNT(*) OVER (). PreviewPage +// runs the exact count separately so large ordered catalogs can use top-N/index +// plans for the page fetch. +func TestQueryExecutor_PreviewPage_ExactTotalOmitsWindowCount(t *testing.T) { exec := &QueryExecutor{Scope: "movie", BaseRelationSQL: "media_items mi"} sql, _, err := exec.buildPreviewPageSQL(QueryDefinition{}, AccessFilter{}, 20, 0, true /* includeTotal */) if err != nil { t.Fatalf("buildPreviewPageSQL error: %v", err) } - if !strings.Contains(sql, "COUNT(*) OVER ()") { - t.Fatalf("expected COUNT(*) OVER () for single-pass count; got:\n%s", sql) + if strings.Contains(sql, "COUNT(*) OVER ()") { + t.Fatalf("PreviewPage exact totals must omit COUNT(*) OVER (); got:\n%s", sql) } } @@ -74,16 +75,11 @@ func TestBrowseRepository_browse_SkipTotal_OmitsWindowCount(t *testing.T) { } } -// TestQueryExecutor_PreviewPage_CountSQL_BindsSortJoinArgs pins that -// previewPagePlan.countSQL binds the sort-plan join args (sortArgs) — not -// just cteArgs+args. fromClausePaged embeds sort-plan join clauses -// (added_at, progress, date_viewed, plays, resolution, bitrate all need -// LIBRARY-id-bound joins), so omitting sortArgs would leave bound -// placeholders inside the FROM clause unfilled and Postgres would error -// out with "missing argument" at the count fallback path. -// -// Regression guard for the post-perf-overhaul code review (macroscope High). -func TestQueryExecutor_PreviewPage_CountSQL_BindsSortJoinArgs(t *testing.T) { +// TestQueryExecutor_PreviewPage_CountSQL_OmitsSortOnlyJoins pins that exact +// totals do not carry ORDER BY-only joins. The page query still needs those +// joins for sorts such as added_at, but the count query only needs the base +// relation plus filter joins. +func TestQueryExecutor_PreviewPage_CountSQL_OmitsSortOnlyJoins(t *testing.T) { exec := &QueryExecutor{Scope: "movie", BaseRelationSQL: "media_items mi"} plan, err := exec.buildPreviewPagePlan( QueryDefinition{ @@ -102,15 +98,15 @@ func TestQueryExecutor_PreviewPage_CountSQL_BindsSortJoinArgs(t *testing.T) { sql, args := plan.countSQL() - // The sort-plan LEFT JOIN must be embedded in the count SQL — that's the - // reason sortArgs binding matters at all. - if !strings.Contains(sql, "sort_added") { - t.Fatalf("countSQL must include addedAtSortPlan's LEFT JOIN; got:\n%s", sql) + if strings.Contains(sql, "sort_added") { + t.Fatalf("countSQL must omit addedAtSortPlan's LEFT JOIN; got:\n%s", sql) + } + if len(args) != 1 { + t.Fatalf("countSQL must omit sortArgs and bind only the library-scope arg; got %v", args) } - // Verify args length covers every $N placeholder in the SQL. If we - // dropped sortArgs, the highest $N would exceed len(args) and Postgres - // would fail with "missing argument". + // Verify args still cover every $N placeholder after dropping sort-only + // joins and sortArgs. maxIdx := 0 for _, m := range placeholderRE.FindAllStringSubmatch(sql, -1) { idx, _ := strconv.Atoi(m[1]) @@ -122,23 +118,20 @@ func TestQueryExecutor_PreviewPage_CountSQL_BindsSortJoinArgs(t *testing.T) { t.Fatalf("expected at least one $N placeholder in countSQL; got:\n%s", sql) } if len(args) < maxIdx { - t.Fatalf("countSQL references $%d but only %d args bound (sortArgs likely dropped); sql:\n%s\nargs: %v", + t.Fatalf("countSQL references $%d but only %d args bound; sql:\n%s\nargs: %v", maxIdx, len(args), sql, args) } } // TestQueryExecutor_PreviewPage_CountSQL_OmitsLimitOffsetOrderBy pins the -// empty-page fallback SQL shape on previewPagePlan. When pagedSQL(true) -// returns an empty page past offset 0, the executor invokes countSQL() to -// recover the real total — COUNT(*) OVER () would otherwise emit no rows -// and leave the caller seeing total=0 even when broader matches exist. +// exact-total SQL shape on previewPagePlan. PreviewPage invokes countSQL() +// separately instead of making the data SELECT calculate COUNT(*) OVER (). // // The countSQL must: // - omit LIMIT/OFFSET (we want the unpaginated total) // - omit ORDER BY (irrelevant for a count, and may reference unbound args) // - wrap the inner FROM/WHERE in `SELECT COUNT(*) FROM (SELECT 1 ...) sub` -// so any GROUP BY in the inner query counts groups (not rows), matching -// what COUNT(*) OVER () would have computed. +// so any GROUP BY in the inner query counts groups (not rows). func TestQueryExecutor_PreviewPage_CountSQL_OmitsLimitOffsetOrderBy(t *testing.T) { exec := &QueryExecutor{Scope: "movie", BaseRelationSQL: "media_items mi"} plan, err := exec.buildPreviewPagePlan(QueryDefinition{}, AccessFilter{}, 20, 0) diff --git a/migrations/141_episode_title_sort_index.down.sql b/migrations/141_episode_title_sort_index.down.sql new file mode 100644 index 00000000..959a7406 --- /dev/null +++ b/migrations/141_episode_title_sort_index.down.sql @@ -0,0 +1 @@ +DROP INDEX IF EXISTS public.idx_episodes_sort_key_content; diff --git a/migrations/141_episode_title_sort_index.up.sql b/migrations/141_episode_title_sort_index.up.sql new file mode 100644 index 00000000..f249b2f4 --- /dev/null +++ b/migrations/141_episode_title_sort_index.up.sql @@ -0,0 +1,7 @@ +-- Episode library title browse. Matches episodeCatalogSelectBody's sort_key +-- expression so PostgreSQL can satisfy ORDER BY title from the episodes index. +CREATE INDEX IF NOT EXISTS idx_episodes_sort_key_content +ON public.episodes USING btree ( + LOWER(COALESCE(NULLIF(BTRIM(title), ''), 'Episode ' || episode_number::text)), + content_id +); diff --git a/migrations/142_episode_catalog_entries.down.sql b/migrations/142_episode_catalog_entries.down.sql new file mode 100644 index 00000000..e93f6816 --- /dev/null +++ b/migrations/142_episode_catalog_entries.down.sql @@ -0,0 +1,42 @@ +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_series ON public.media_items; +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_episodes ON public.episodes; +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_media_files_update ON public.media_files; +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_media_files_insert_delete ON public.media_files; +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_media_files ON public.media_files; +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_episode_libraries ON public.episode_libraries; + +DROP FUNCTION IF EXISTS public.episode_catalog_entries_series_trigger(); +DROP FUNCTION IF EXISTS public.episode_catalog_entries_episodes_trigger(); +DROP FUNCTION IF EXISTS public.episode_catalog_entries_media_files_trigger(); +DROP FUNCTION IF EXISTS public.episode_catalog_entries_episode_libraries_trigger(); +DROP FUNCTION IF EXISTS public.refresh_episode_catalog_entries_for_series(text); +DROP FUNCTION IF EXISTS public.refresh_episode_catalog_entries_for_episode(text); +DROP FUNCTION IF EXISTS public.refresh_episode_catalog_entry(text, integer); +DROP FUNCTION IF EXISTS public.episode_catalog_resolution_rank(text); +DROP FUNCTION IF EXISTS public.episode_catalog_normalized_resolution(text); +DROP FUNCTION IF EXISTS public.episode_catalog_rating_rank(text); + +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_subtitle_gin; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_audio_gin; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_resolution_codes_gin; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_countries_gin; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_networks_gin; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_studios_gin; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_genres_gin; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_non_dolby_vision; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_dolby_vision; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_non_hdr; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_hdr; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_status; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_bitrate; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_resolution; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_tmdb; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_imdb; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_runtime; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_content_rating; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_year; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_air_date; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_added; +DROP INDEX IF EXISTS public.idx_episode_catalog_entries_title; + +DROP TABLE IF EXISTS public.episode_catalog_entries; diff --git a/migrations/142_episode_catalog_entries.up.sql b/migrations/142_episode_catalog_entries.up.sql new file mode 100644 index 00000000..0020addb --- /dev/null +++ b/migrations/142_episode_catalog_entries.up.sql @@ -0,0 +1,674 @@ +CREATE TABLE IF NOT EXISTS public.episode_catalog_entries ( + media_folder_id integer NOT NULL, + episode_id text NOT NULL, + series_id text NOT NULL, + sort_key text NOT NULL, + title text NOT NULL, + added_at timestamp with time zone NOT NULL, + episode_air_date date, + year integer NOT NULL, + genres text[] NOT NULL DEFAULT '{}', + studios text[] NOT NULL DEFAULT '{}', + networks text[] NOT NULL DEFAULT '{}', + countries text[] NOT NULL DEFAULT '{}', + original_language text NOT NULL DEFAULT '', + content_rating text NOT NULL DEFAULT '', + content_rating_label text NOT NULL DEFAULT '~~~~', + content_rating_rank integer NOT NULL DEFAULT 2147483647, + status text NOT NULL DEFAULT 'matched', + runtime integer NOT NULL DEFAULT 0, + rating_imdb double precision, + rating_tmdb double precision, + max_resolution_rank integer, + resolution_codes text[] NOT NULL DEFAULT '{}', + max_bitrate integer, + min_bitrate integer, + has_hdr boolean NOT NULL DEFAULT false, + has_non_hdr boolean NOT NULL DEFAULT false, + has_dolby_vision boolean NOT NULL DEFAULT false, + has_non_dolby_vision boolean NOT NULL DEFAULT false, + audio_language_codes text[] NOT NULL DEFAULT '{}', + subtitle_language_codes text[] NOT NULL DEFAULT '{}', + episode_created_at timestamp with time zone NOT NULL, + updated_at timestamp with time zone NOT NULL DEFAULT now(), + CONSTRAINT episode_catalog_entries_pkey PRIMARY KEY (media_folder_id, episode_id), + CONSTRAINT episode_catalog_entries_episode_id_fkey FOREIGN KEY (episode_id) REFERENCES public.episodes(content_id) ON DELETE CASCADE, + CONSTRAINT episode_catalog_entries_media_folder_id_fkey FOREIGN KEY (media_folder_id) REFERENCES public.media_folders(id) ON DELETE CASCADE +); + +CREATE TEMP TABLE episode_catalog_entries_migration_clock AS +SELECT transaction_timestamp() AS started_at; + +CREATE OR REPLACE FUNCTION public.episode_catalog_rating_rank(rating text) +RETURNS integer +LANGUAGE sql +IMMUTABLE +AS $$ + SELECT CASE UPPER(NULLIF(BTRIM(rating), '')) + WHEN 'G' THEN 0 + WHEN 'TV-Y' THEN 0 + WHEN 'TV-G' THEN 0 + WHEN 'PG' THEN 1 + WHEN 'TV-Y7' THEN 1 + WHEN 'TV-PG' THEN 1 + WHEN 'PG-13' THEN 2 + WHEN 'TV-14' THEN 2 + WHEN 'R' THEN 3 + WHEN 'NC-17' THEN 3 + WHEN 'TV-MA' THEN 3 + ELSE 2147483647 + END +$$; + +CREATE OR REPLACE FUNCTION public.episode_catalog_normalized_resolution(resolution text) +RETURNS text +LANGUAGE sql +IMMUTABLE +AS $$ + SELECT CASE LOWER(NULLIF(BTRIM(resolution), '')) + WHEN '4k' THEN '2160p' + WHEN 'uhd' THEN '2160p' + ELSE LOWER(NULLIF(BTRIM(resolution), '')) + END +$$; + +CREATE OR REPLACE FUNCTION public.episode_catalog_resolution_rank(resolution text) +RETURNS integer +LANGUAGE sql +IMMUTABLE +AS $$ + SELECT CASE UPPER(public.episode_catalog_normalized_resolution(resolution)) + WHEN '480P' THEN 1 + WHEN '720P' THEN 2 + WHEN '1080P' THEN 3 + WHEN '2160P' THEN 4 + WHEN '4320P' THEN 5 + ELSE NULL + END +$$; + +CREATE OR REPLACE FUNCTION public.refresh_episode_catalog_entry(p_episode_id text, p_media_folder_id integer) +RETURNS void +LANGUAGE plpgsql +AS $$ +BEGIN + IF p_episode_id IS NULL OR p_media_folder_id IS NULL THEN + RETURN; + END IF; + + INSERT INTO public.episode_catalog_entries ( + media_folder_id, + episode_id, + series_id, + sort_key, + title, + added_at, + episode_air_date, + year, + genres, + studios, + networks, + countries, + original_language, + content_rating, + content_rating_label, + content_rating_rank, + status, + runtime, + rating_imdb, + rating_tmdb, + max_resolution_rank, + resolution_codes, + max_bitrate, + min_bitrate, + has_hdr, + has_non_hdr, + has_dolby_vision, + has_non_dolby_vision, + audio_language_codes, + subtitle_language_codes, + episode_created_at, + updated_at + ) + SELECT + el.media_folder_id, + e.content_id, + e.series_id, + LOWER(COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text)) AS sort_key, + COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text) AS title, + el.first_seen_at, + e.air_date, + COALESCE(si.year, EXTRACT(YEAR FROM e.air_date)::integer, 0) AS year, + COALESCE(si.genres, '{}'::text[]) AS genres, + COALESCE(si.studios, '{}'::text[]) AS studios, + COALESCE(si.networks, '{}'::text[]) AS networks, + COALESCE(si.countries, '{}'::text[]) AS countries, + COALESCE(si.original_language, '') AS original_language, + COALESCE(si.content_rating, '') AS content_rating, + LOWER(COALESCE(NULLIF(BTRIM(si.content_rating), ''), '~~~~')) AS content_rating_label, + public.episode_catalog_rating_rank(si.content_rating) AS content_rating_rank, + COALESCE(NULLIF(BTRIM(si.status), ''), 'matched') AS status, + COALESCE(NULLIF(e.runtime, 0), COALESCE(si.runtime, 0)) AS runtime, + e.rating_imdb, + e.rating_tmdb, + stats.max_resolution_rank, + COALESCE(stats.resolution_codes, '{}'::text[]) AS resolution_codes, + stats.max_bitrate, + stats.min_bitrate, + COALESCE(stats.has_hdr, false) AS has_hdr, + COALESCE(stats.has_non_hdr, false) AS has_non_hdr, + COALESCE(stats.has_dolby_vision, false) AS has_dolby_vision, + COALESCE(stats.has_non_dolby_vision, false) AS has_non_dolby_vision, + COALESCE(stats.audio_language_codes, '{}'::text[]) AS audio_language_codes, + COALESCE(stats.subtitle_language_codes, '{}'::text[]) AS subtitle_language_codes, + e.created_at, + NOW() + FROM public.episode_libraries el + JOIN public.episodes e ON e.content_id = el.episode_id + JOIN public.media_items si ON si.content_id = e.series_id + LEFT JOIN LATERAL ( + SELECT + MAX(public.episode_catalog_resolution_rank(mf.resolution)) AS max_resolution_rank, + ARRAY( + SELECT DISTINCT code + FROM public.media_files mf_res + CROSS JOIN LATERAL ( + SELECT public.episode_catalog_normalized_resolution(mf_res.resolution) AS code + ) normalized + WHERE mf_res.episode_id = e.content_id + AND mf_res.media_folder_id = el.media_folder_id + AND mf_res.missing_since IS NULL + AND normalized.code IS NOT NULL + ORDER BY code + ) AS resolution_codes, + MAX(mf.bitrate) FILTER (WHERE mf.bitrate IS NOT NULL AND mf.bitrate > 0) AS max_bitrate, + MIN(mf.bitrate) FILTER (WHERE mf.bitrate IS NOT NULL AND mf.bitrate > 0) AS min_bitrate, + BOOL_OR(mf.hdr IS TRUE) AS has_hdr, + BOOL_OR(mf.hdr IS FALSE) AS has_non_hdr, + BOOL_OR(EXISTS ( + SELECT 1 + FROM jsonb_array_elements(COALESCE(mf.video_tracks, '[]'::jsonb)) AS vt + WHERE NULLIF(BTRIM(vt->>'dolby_vision'), '') IS NOT NULL + )) AS has_dolby_vision, + BOOL_OR(NOT EXISTS ( + SELECT 1 + FROM jsonb_array_elements(COALESCE(mf.video_tracks, '[]'::jsonb)) AS vt + WHERE NULLIF(BTRIM(vt->>'dolby_vision'), '') IS NOT NULL + )) AS has_non_dolby_vision, + ARRAY( + SELECT DISTINCT LOWER(NULLIF(BTRIM(lang), '')) + FROM public.media_files mf_audio + CROSS JOIN LATERAL UNNEST(COALESCE(mf_audio.audio_language_codes, '{}'::text[])) AS lang + WHERE mf_audio.episode_id = e.content_id + AND mf_audio.media_folder_id = el.media_folder_id + AND mf_audio.missing_since IS NULL + AND NULLIF(BTRIM(lang), '') IS NOT NULL + ORDER BY LOWER(NULLIF(BTRIM(lang), '')) + ) AS audio_language_codes, + ARRAY( + SELECT DISTINCT lang_code + FROM ( + SELECT LOWER(NULLIF(BTRIM(lang), '')) AS lang_code + FROM public.media_files mf_sub + CROSS JOIN LATERAL UNNEST(COALESCE(mf_sub.subtitle_language_codes, '{}'::text[])) AS lang + WHERE mf_sub.episode_id = e.content_id + AND mf_sub.media_folder_id = el.media_folder_id + AND mf_sub.missing_since IS NULL + UNION + SELECT LOWER(NULLIF(BTRIM(track->>'language'), '')) AS lang_code + FROM public.media_files mf_ext + CROSS JOIN LATERAL jsonb_array_elements(COALESCE(mf_ext.external_subtitles, '[]'::jsonb)) AS track + WHERE mf_ext.episode_id = e.content_id + AND mf_ext.media_folder_id = el.media_folder_id + AND mf_ext.missing_since IS NULL + ) subtitle_codes + WHERE lang_code IS NOT NULL + ORDER BY lang_code + ) AS subtitle_language_codes + FROM public.media_files mf + WHERE mf.episode_id = e.content_id + AND mf.media_folder_id = el.media_folder_id + AND mf.missing_since IS NULL + ) stats ON TRUE + WHERE el.episode_id = p_episode_id + AND el.media_folder_id = p_media_folder_id + ON CONFLICT (media_folder_id, episode_id) DO UPDATE SET + series_id = EXCLUDED.series_id, + sort_key = EXCLUDED.sort_key, + title = EXCLUDED.title, + added_at = EXCLUDED.added_at, + episode_air_date = EXCLUDED.episode_air_date, + year = EXCLUDED.year, + genres = EXCLUDED.genres, + studios = EXCLUDED.studios, + networks = EXCLUDED.networks, + countries = EXCLUDED.countries, + original_language = EXCLUDED.original_language, + content_rating = EXCLUDED.content_rating, + content_rating_label = EXCLUDED.content_rating_label, + content_rating_rank = EXCLUDED.content_rating_rank, + status = EXCLUDED.status, + runtime = EXCLUDED.runtime, + rating_imdb = EXCLUDED.rating_imdb, + rating_tmdb = EXCLUDED.rating_tmdb, + max_resolution_rank = EXCLUDED.max_resolution_rank, + resolution_codes = EXCLUDED.resolution_codes, + max_bitrate = EXCLUDED.max_bitrate, + min_bitrate = EXCLUDED.min_bitrate, + has_hdr = EXCLUDED.has_hdr, + has_non_hdr = EXCLUDED.has_non_hdr, + has_dolby_vision = EXCLUDED.has_dolby_vision, + has_non_dolby_vision = EXCLUDED.has_non_dolby_vision, + audio_language_codes = EXCLUDED.audio_language_codes, + subtitle_language_codes = EXCLUDED.subtitle_language_codes, + episode_created_at = EXCLUDED.episode_created_at, + updated_at = NOW(); + + IF NOT FOUND THEN + DELETE FROM public.episode_catalog_entries + WHERE episode_id = p_episode_id + AND media_folder_id = p_media_folder_id; + END IF; +END; +$$; + +CREATE OR REPLACE FUNCTION public.refresh_episode_catalog_entries_for_episode(p_episode_id text) +RETURNS void +LANGUAGE sql +AS $$ + SELECT public.refresh_episode_catalog_entry(el.episode_id, el.media_folder_id) + FROM public.episode_libraries el + WHERE el.episode_id = p_episode_id; +$$; + +CREATE OR REPLACE FUNCTION public.refresh_episode_catalog_entries_for_series(p_series_id text) +RETURNS void +LANGUAGE sql +AS $$ + SELECT public.refresh_episode_catalog_entry(el.episode_id, el.media_folder_id) + FROM public.episode_libraries el + JOIN public.episodes e ON e.content_id = el.episode_id + WHERE e.series_id = p_series_id; +$$; + +CREATE OR REPLACE FUNCTION public.episode_catalog_entries_episode_libraries_trigger() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + IF TG_OP = 'DELETE' THEN + DELETE FROM public.episode_catalog_entries + WHERE episode_id = OLD.episode_id + AND media_folder_id = OLD.media_folder_id; + RETURN OLD; + END IF; + + PERFORM public.refresh_episode_catalog_entry(NEW.episode_id, NEW.media_folder_id); + IF TG_OP = 'UPDATE' + AND (OLD.episode_id IS DISTINCT FROM NEW.episode_id OR OLD.media_folder_id IS DISTINCT FROM NEW.media_folder_id) THEN + PERFORM public.refresh_episode_catalog_entry(OLD.episode_id, OLD.media_folder_id); + END IF; + RETURN NEW; +END; +$$; + +CREATE OR REPLACE FUNCTION public.episode_catalog_entries_media_files_trigger() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + IF TG_OP = 'DELETE' THEN + PERFORM public.refresh_episode_catalog_entry(OLD.episode_id, OLD.media_folder_id); + RETURN OLD; + END IF; + + IF TG_OP = 'UPDATE' THEN + PERFORM public.refresh_episode_catalog_entry(OLD.episode_id, OLD.media_folder_id); + END IF; + + PERFORM public.refresh_episode_catalog_entry(NEW.episode_id, NEW.media_folder_id); + RETURN NEW; +END; +$$; + +CREATE OR REPLACE FUNCTION public.episode_catalog_entries_episodes_trigger() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + IF TG_OP = 'DELETE' THEN + DELETE FROM public.episode_catalog_entries + WHERE episode_id = OLD.content_id; + RETURN OLD; + END IF; + + PERFORM public.refresh_episode_catalog_entries_for_episode(NEW.content_id); + IF TG_OP = 'UPDATE' AND OLD.content_id IS DISTINCT FROM NEW.content_id THEN + DELETE FROM public.episode_catalog_entries + WHERE episode_id = OLD.content_id; + END IF; + RETURN NEW; +END; +$$; + +CREATE OR REPLACE FUNCTION public.episode_catalog_entries_series_trigger() +RETURNS trigger +LANGUAGE plpgsql +AS $$ +BEGIN + IF TG_OP = 'DELETE' THEN + DELETE FROM public.episode_catalog_entries + WHERE series_id = OLD.content_id; + RETURN OLD; + END IF; + + IF COALESCE(NEW.type, '') = 'series' THEN + PERFORM public.refresh_episode_catalog_entries_for_series(NEW.content_id); + END IF; + RETURN NEW; +END; +$$; + +WITH active_episode_files AS ( + SELECT + mf.episode_id, + mf.media_folder_id, + public.episode_catalog_normalized_resolution(mf.resolution) AS normalized_resolution, + public.episode_catalog_resolution_rank(mf.resolution) AS resolution_rank, + mf.bitrate, + mf.hdr, + EXISTS ( + SELECT 1 + FROM jsonb_array_elements(COALESCE(mf.video_tracks, '[]'::jsonb)) AS vt + WHERE NULLIF(BTRIM(vt->>'dolby_vision'), '') IS NOT NULL + ) AS has_dolby_vision, + mf.audio_language_codes, + mf.subtitle_language_codes, + mf.external_subtitles + FROM public.media_files mf + WHERE mf.episode_id IS NOT NULL + AND mf.missing_since IS NULL +), +file_stats AS ( + SELECT + episode_id, + media_folder_id, + MAX(resolution_rank) AS max_resolution_rank, + ARRAY_AGG(DISTINCT normalized_resolution ORDER BY normalized_resolution) + FILTER (WHERE normalized_resolution IS NOT NULL) AS resolution_codes, + MAX(bitrate) FILTER (WHERE bitrate IS NOT NULL AND bitrate > 0) AS max_bitrate, + MIN(bitrate) FILTER (WHERE bitrate IS NOT NULL AND bitrate > 0) AS min_bitrate, + BOOL_OR(hdr IS TRUE) AS has_hdr, + BOOL_OR(hdr IS FALSE) AS has_non_hdr, + BOOL_OR(has_dolby_vision) AS has_dolby_vision, + BOOL_OR(NOT has_dolby_vision) AS has_non_dolby_vision + FROM active_episode_files + GROUP BY episode_id, media_folder_id +), +audio_stats AS ( + SELECT + episode_id, + media_folder_id, + ARRAY_AGG(DISTINCT lang_code ORDER BY lang_code) AS audio_language_codes + FROM ( + SELECT + aef.episode_id, + aef.media_folder_id, + LOWER(NULLIF(BTRIM(lang), '')) AS lang_code + FROM active_episode_files aef + CROSS JOIN LATERAL UNNEST(COALESCE(aef.audio_language_codes, '{}'::text[])) AS lang + ) codes + WHERE lang_code IS NOT NULL + GROUP BY episode_id, media_folder_id +), +subtitle_stats AS ( + SELECT + episode_id, + media_folder_id, + ARRAY_AGG(DISTINCT lang_code ORDER BY lang_code) AS subtitle_language_codes + FROM ( + SELECT + aef.episode_id, + aef.media_folder_id, + LOWER(NULLIF(BTRIM(lang), '')) AS lang_code + FROM active_episode_files aef + CROSS JOIN LATERAL UNNEST(COALESCE(aef.subtitle_language_codes, '{}'::text[])) AS lang + UNION + SELECT + aef.episode_id, + aef.media_folder_id, + LOWER(NULLIF(BTRIM(track->>'language'), '')) AS lang_code + FROM active_episode_files aef + CROSS JOIN LATERAL jsonb_array_elements(COALESCE(aef.external_subtitles, '[]'::jsonb)) AS track + ) codes + WHERE lang_code IS NOT NULL + GROUP BY episode_id, media_folder_id +) +INSERT INTO public.episode_catalog_entries ( + media_folder_id, + episode_id, + series_id, + sort_key, + title, + added_at, + episode_air_date, + year, + genres, + studios, + networks, + countries, + original_language, + content_rating, + content_rating_label, + content_rating_rank, + status, + runtime, + rating_imdb, + rating_tmdb, + max_resolution_rank, + resolution_codes, + max_bitrate, + min_bitrate, + has_hdr, + has_non_hdr, + has_dolby_vision, + has_non_dolby_vision, + audio_language_codes, + subtitle_language_codes, + episode_created_at, + updated_at +) +SELECT + el.media_folder_id, + e.content_id, + e.series_id, + LOWER(COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text)) AS sort_key, + COALESCE(NULLIF(BTRIM(e.title), ''), 'Episode ' || e.episode_number::text) AS title, + el.first_seen_at, + e.air_date, + COALESCE(si.year, EXTRACT(YEAR FROM e.air_date)::integer, 0) AS year, + COALESCE(si.genres, '{}'::text[]) AS genres, + COALESCE(si.studios, '{}'::text[]) AS studios, + COALESCE(si.networks, '{}'::text[]) AS networks, + COALESCE(si.countries, '{}'::text[]) AS countries, + COALESCE(si.original_language, '') AS original_language, + COALESCE(si.content_rating, '') AS content_rating, + LOWER(COALESCE(NULLIF(BTRIM(si.content_rating), ''), '~~~~')) AS content_rating_label, + public.episode_catalog_rating_rank(si.content_rating) AS content_rating_rank, + COALESCE(NULLIF(BTRIM(si.status), ''), 'matched') AS status, + COALESCE(NULLIF(e.runtime, 0), COALESCE(si.runtime, 0)) AS runtime, + e.rating_imdb, + e.rating_tmdb, + file_stats.max_resolution_rank, + COALESCE(file_stats.resolution_codes, '{}'::text[]) AS resolution_codes, + file_stats.max_bitrate, + file_stats.min_bitrate, + COALESCE(file_stats.has_hdr, false) AS has_hdr, + COALESCE(file_stats.has_non_hdr, false) AS has_non_hdr, + COALESCE(file_stats.has_dolby_vision, false) AS has_dolby_vision, + COALESCE(file_stats.has_non_dolby_vision, false) AS has_non_dolby_vision, + COALESCE(audio_stats.audio_language_codes, '{}'::text[]) AS audio_language_codes, + COALESCE(subtitle_stats.subtitle_language_codes, '{}'::text[]) AS subtitle_language_codes, + e.created_at, + NOW() +FROM public.episode_libraries el +JOIN public.episodes e ON e.content_id = el.episode_id +JOIN public.media_items si ON si.content_id = e.series_id +LEFT JOIN file_stats ON file_stats.episode_id = el.episode_id AND file_stats.media_folder_id = el.media_folder_id +LEFT JOIN audio_stats ON audio_stats.episode_id = el.episode_id AND audio_stats.media_folder_id = el.media_folder_id +LEFT JOIN subtitle_stats ON subtitle_stats.episode_id = el.episode_id AND subtitle_stats.media_folder_id = el.media_folder_id +ON CONFLICT (media_folder_id, episode_id) DO NOTHING; + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_title +ON public.episode_catalog_entries USING btree (media_folder_id, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_added +ON public.episode_catalog_entries USING btree (media_folder_id, added_at DESC, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_air_date +ON public.episode_catalog_entries USING btree (media_folder_id, episode_air_date DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_year +ON public.episode_catalog_entries USING btree (media_folder_id, year DESC, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_content_rating +ON public.episode_catalog_entries USING btree (media_folder_id, content_rating_rank, content_rating_label, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_runtime +ON public.episode_catalog_entries USING btree (media_folder_id, runtime DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_imdb +ON public.episode_catalog_entries USING btree (media_folder_id, rating_imdb DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_tmdb +ON public.episode_catalog_entries USING btree (media_folder_id, rating_tmdb DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_resolution +ON public.episode_catalog_entries USING btree (media_folder_id, max_resolution_rank DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_bitrate +ON public.episode_catalog_entries USING btree (media_folder_id, max_bitrate DESC NULLS LAST, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_status +ON public.episode_catalog_entries USING btree (media_folder_id, status, sort_key, episode_id); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_hdr +ON public.episode_catalog_entries USING btree (media_folder_id, sort_key, episode_id) +WHERE has_hdr; + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_non_hdr +ON public.episode_catalog_entries USING btree (media_folder_id, sort_key, episode_id) +WHERE has_non_hdr; + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_dolby_vision +ON public.episode_catalog_entries USING btree (media_folder_id, sort_key, episode_id) +WHERE has_dolby_vision; + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_non_dolby_vision +ON public.episode_catalog_entries USING btree (media_folder_id, sort_key, episode_id) +WHERE has_non_dolby_vision; + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_genres_gin +ON public.episode_catalog_entries USING gin (genres); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_studios_gin +ON public.episode_catalog_entries USING gin (studios); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_networks_gin +ON public.episode_catalog_entries USING gin (networks); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_countries_gin +ON public.episode_catalog_entries USING gin (countries); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_resolution_codes_gin +ON public.episode_catalog_entries USING gin (resolution_codes); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_audio_gin +ON public.episode_catalog_entries USING gin (audio_language_codes); + +CREATE INDEX IF NOT EXISTS idx_episode_catalog_entries_subtitle_gin +ON public.episode_catalog_entries USING gin (subtitle_language_codes); + +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_episode_libraries ON public.episode_libraries; +CREATE TRIGGER trg_episode_catalog_entries_episode_libraries +AFTER INSERT OR UPDATE OR DELETE ON public.episode_libraries +FOR EACH ROW EXECUTE FUNCTION public.episode_catalog_entries_episode_libraries_trigger(); + +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_media_files ON public.media_files; +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_media_files_insert_delete ON public.media_files; +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_media_files_update ON public.media_files; +CREATE TRIGGER trg_episode_catalog_entries_media_files_insert_delete +AFTER INSERT OR DELETE ON public.media_files +FOR EACH ROW EXECUTE FUNCTION public.episode_catalog_entries_media_files_trigger(); + +CREATE TRIGGER trg_episode_catalog_entries_media_files_update +AFTER UPDATE ON public.media_files +FOR EACH ROW +WHEN ( + OLD.episode_id IS DISTINCT FROM NEW.episode_id OR + OLD.media_folder_id IS DISTINCT FROM NEW.media_folder_id OR + OLD.missing_since IS DISTINCT FROM NEW.missing_since OR + OLD.resolution IS DISTINCT FROM NEW.resolution OR + OLD.bitrate IS DISTINCT FROM NEW.bitrate OR + OLD.hdr IS DISTINCT FROM NEW.hdr OR + OLD.video_tracks IS DISTINCT FROM NEW.video_tracks OR + OLD.external_subtitles IS DISTINCT FROM NEW.external_subtitles OR + OLD.audio_language_codes IS DISTINCT FROM NEW.audio_language_codes OR + OLD.subtitle_language_codes IS DISTINCT FROM NEW.subtitle_language_codes +) +EXECUTE FUNCTION public.episode_catalog_entries_media_files_trigger(); + +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_episodes ON public.episodes; +CREATE TRIGGER trg_episode_catalog_entries_episodes +AFTER INSERT OR UPDATE OF content_id, series_id, title, episode_number, air_date, runtime, rating_imdb, rating_tmdb, still_path, still_thumbhash, created_at OR DELETE +ON public.episodes +FOR EACH ROW EXECUTE FUNCTION public.episode_catalog_entries_episodes_trigger(); + +DROP TRIGGER IF EXISTS trg_episode_catalog_entries_series ON public.media_items; +CREATE TRIGGER trg_episode_catalog_entries_series +AFTER UPDATE OF content_id, type, year, genres, studios, networks, countries, original_language, content_rating, status, runtime ON public.media_items +FOR EACH ROW +WHEN (OLD.type = 'series' OR NEW.type = 'series') +EXECUTE FUNCTION public.episode_catalog_entries_series_trigger(); + +WITH migration_clock AS ( + SELECT started_at FROM episode_catalog_entries_migration_clock LIMIT 1 +), +changed_entries AS ( + SELECT DISTINCT mf.episode_id, mf.media_folder_id + FROM public.media_files mf, migration_clock mc + WHERE mf.episode_id IS NOT NULL + AND mf.updated_at >= mc.started_at + UNION + SELECT DISTINCT el.episode_id, el.media_folder_id + FROM public.episode_libraries el, migration_clock mc + WHERE el.first_seen_at >= mc.started_at + UNION + SELECT DISTINCT el.episode_id, el.media_folder_id + FROM public.episode_libraries el + JOIN public.episodes e ON e.content_id = el.episode_id + CROSS JOIN migration_clock mc + WHERE e.updated_at >= mc.started_at + UNION + SELECT DISTINCT el.episode_id, el.media_folder_id + FROM public.episode_libraries el + JOIN public.episodes e ON e.content_id = el.episode_id + JOIN public.media_items si ON si.content_id = e.series_id + CROSS JOIN migration_clock mc + WHERE si.updated_at >= mc.started_at + AND si.type = 'series' + UNION + SELECT DISTINCT ece.episode_id, ece.media_folder_id + FROM public.episode_catalog_entries ece + WHERE NOT EXISTS ( + SELECT 1 + FROM public.episode_libraries el + WHERE el.episode_id = ece.episode_id + AND el.media_folder_id = ece.media_folder_id + ) +) +SELECT COUNT(*) +FROM ( + SELECT public.refresh_episode_catalog_entry(episode_id, media_folder_id) + FROM changed_entries +) refreshed; diff --git a/scripts/benchmark_episode_catalog.py b/scripts/benchmark_episode_catalog.py new file mode 100644 index 00000000..d6ff54d0 --- /dev/null +++ b/scripts/benchmark_episode_catalog.py @@ -0,0 +1,293 @@ +#!/usr/bin/env python3 +"""Benchmark episode catalog sort/filter shapes without embedding credentials. + +Required environment: + SILO_BASE_URL Example: https://silo.example.com + SILO_BEARER_TOKEN API access token + SILO_PROFILE_ID Profile id used for personalized filters/sorts + +Optional environment: + SILO_LIBRARY_ID Defaults to 2 +""" + +from __future__ import annotations + +import argparse +import csv +import json +import os +import statistics +import sys +import time +import urllib.error +import urllib.parse +import urllib.request +from concurrent.futures import ThreadPoolExecutor, as_completed +from dataclasses import dataclass +from typing import Iterable + + +def group_rule(field: str, op: str, value: str) -> dict[str, str]: + return { + "groups[0][match]": "all", + "groups[0][rules][0][field]": field, + "groups[0][rules][0][op]": op, + "groups[0][rules][0][value]": value, + } + + +DEFAULT_SORT_CASES = [ + ("sort:title", {"sort": "title", "order": "asc"}), + ("sort:added_at", {"sort": "added_at", "order": "desc"}), + ("sort:release_date", {"sort": "release_date", "order": "desc"}), + ("sort:last_air_date", {"sort": "last_air_date", "order": "desc"}), + ("sort:year", {"sort": "year", "order": "desc"}), + ("sort:content_rating", {"sort": "content_rating", "order": "asc"}), + ("sort:runtime", {"sort": "runtime", "order": "desc"}), + ("sort:rating_imdb", {"sort": "rating_imdb", "order": "desc"}), + ("sort:rating_tmdb", {"sort": "rating_tmdb", "order": "desc"}), + ("sort:rating_rt_critic", {"sort": "rating_rt_critic", "order": "desc"}), + ("sort:rating_rt_audience", {"sort": "rating_rt_audience", "order": "desc"}), + ("sort:resolution", {"sort": "resolution", "order": "desc"}), + ("sort:bitrate", {"sort": "bitrate", "order": "desc"}), + ("sort:progress", {"sort": "progress", "order": "desc"}), + ("sort:date_viewed", {"sort": "date_viewed", "order": "desc"}), + ("sort:plays", {"sort": "plays", "order": "desc"}), +] + +DEFAULT_FILTER_CASES = [ + ("filter:genre", {"genre": "Comedy", "sort": "title", "order": "asc"}), + ("filter:year", {"year_min": "2020", "year_max": "2024", "sort": "year", "order": "desc"}), + ("filter:content_rating", {"content_rating": "tv-pg", "sort": "title", "order": "asc"}), + ( + "filter:resolution", + group_rule("resolution", "is", "1080p") | {"sort": "title", "order": "asc"}, + ), + ("filter:hdr", group_rule("hdr", "is", "true") | {"sort": "title", "order": "asc"}), + ( + "filter:dolby_vision", + group_rule("dolby_vision", "is", "true") | {"sort": "title", "order": "asc"}, + ), + ( + "filter:bitrate", + group_rule("bitrate", "gte", "8000000") | {"sort": "bitrate", "order": "desc"}, + ), + ( + "filter:audio_language", + group_rule("audio_language", "is", "en") | {"sort": "title", "order": "asc"}, + ), + ( + "filter:subtitle_language", + group_rule("subtitle_language", "is", "en") | {"sort": "title", "order": "asc"}, + ), + ( + "filter:watched_true", + group_rule("watched", "is", "true") | {"sort": "title", "order": "asc"}, + ), + ( + "filter:in_progress_true", + group_rule("in_progress", "is", "true") | {"sort": "progress", "order": "desc"}, + ), + ( + "filter:last_watched_30d", + group_rule("last_watched", "in_last", "30d") | {"sort": "date_viewed", "order": "desc"}, + ), +] + + +@dataclass(frozen=True) +class Config: + base_url: str + token: str + profile_id: str + library_id: str + limit: int + offset: int + timeout: float + + +@dataclass(frozen=True) +class CaseResult: + case: str + include_total: bool + status: int + elapsed_ms: float + total: int | None + total_exact: bool | None + has_more: bool | None + item_count: int | None + error: str + + +def getenv_required(name: str) -> str: + value = os.environ.get(name, "").strip() + if not value: + raise SystemExit(f"{name} is required") + return value + + +def validate_base_url(base_url: str) -> None: + parsed = urllib.parse.urlparse(base_url) + if parsed.scheme not in {"http", "https"} or not parsed.netloc: + raise ValueError("SILO_BASE_URL must be an http or https URL") + + +def build_url(config: Config, params: dict[str, str], include_total: bool) -> str: + validate_base_url(config.base_url) + query = { + "source": "query", + "type": "episode", + "library_id": config.library_id, + "limit": str(config.limit), + "offset": str(config.offset), + **params, + } + if not include_total: + query["include_total"] = "false" + api_url = urllib.parse.urljoin(config.base_url.rstrip("/") + "/", "api/v1/catalog") + return api_url + "?" + urllib.parse.urlencode(query) + + +def fetch_case(config: Config, case: str, params: dict[str, str], include_total: bool) -> CaseResult: + url = build_url(config, params, include_total) + request = urllib.request.Request( + url, + headers={ + "Accept": "application/json", + "Authorization": f"Bearer {config.token}", + "X-Profile-Id": config.profile_id, + }, + ) + started = time.perf_counter() + status = 0 + try: + with urllib.request.urlopen(request, timeout=config.timeout) as response: + status = response.status + body = response.read() + elapsed_ms = (time.perf_counter() - started) * 1000 + try: + payload = json.loads(body) + items = payload.get("items") + except (json.JSONDecodeError, AttributeError, TypeError) as err: + return CaseResult( + case, + include_total, + status, + elapsed_ms, + None, + None, + None, + None, + str(err), + ) + return CaseResult( + case=case, + include_total=include_total, + status=status, + elapsed_ms=elapsed_ms, + total=payload.get("total"), + total_exact=payload.get("total_exact"), + has_more=payload.get("has_more"), + item_count=len(items) if isinstance(items, list) else None, + error="", + ) + except urllib.error.HTTPError as err: + elapsed_ms = (time.perf_counter() - started) * 1000 + return CaseResult(case, include_total, err.code, elapsed_ms, None, None, None, None, err.reason) + except Exception as err: # noqa: BLE001 - benchmark output should capture failures. + elapsed_ms = (time.perf_counter() - started) * 1000 + return CaseResult(case, include_total, 0, elapsed_ms, None, None, None, None, str(err)) + + +def result_rows(results: Iterable[CaseResult]) -> list[dict[str, object]]: + rows = [] + for result in results: + rows.append( + { + "case": result.case, + "include_total": str(result.include_total).lower(), + "status": result.status, + "elapsed_ms": f"{result.elapsed_ms:.1f}", + "total": "" if result.total is None else result.total, + "total_exact": "" if result.total_exact is None else str(result.total_exact).lower(), + "has_more": "" if result.has_more is None else str(result.has_more).lower(), + "item_count": "" if result.item_count is None else result.item_count, + "error": result.error, + } + ) + return rows + + +def summarize(results: list[CaseResult]) -> None: + successful = [result.elapsed_ms for result in results if result.status == 200] + if not successful: + return + print( + f"# successful={len(successful)} p50_ms={statistics.median(successful):.1f} " + f"max_ms={max(successful):.1f}", + file=sys.stderr, + ) + + +def main() -> int: + parser = argparse.ArgumentParser(description="Benchmark episode catalog sort/filter cases.") + parser.add_argument("--limit", type=int, default=60) + parser.add_argument("--offset", type=int, default=0) + parser.add_argument("--timeout", type=float, default=30.0) + parser.add_argument("--repeat", type=int, default=1) + parser.add_argument("--concurrency", type=int, default=1) + parser.add_argument("--include-exact", action="store_true", help="Also run cases with exact totals enabled.") + args = parser.parse_args() + + config = Config( + base_url=getenv_required("SILO_BASE_URL"), + token=getenv_required("SILO_BEARER_TOKEN"), + profile_id=getenv_required("SILO_PROFILE_ID"), + library_id=os.environ.get("SILO_LIBRARY_ID", "2"), + limit=args.limit, + offset=args.offset, + timeout=args.timeout, + ) + try: + validate_base_url(config.base_url) + except ValueError as err: + raise SystemExit(str(err)) from err + + cases = DEFAULT_SORT_CASES + DEFAULT_FILTER_CASES + include_total_values = [False, True] if args.include_exact else [False] + jobs = [ + (name, params, include_total) + for _ in range(args.repeat) + for include_total in include_total_values + for name, params in cases + ] + + results: list[CaseResult] = [] + with ThreadPoolExecutor(max_workers=max(1, args.concurrency)) as executor: + futures = [executor.submit(fetch_case, config, *job) for job in jobs] + for future in as_completed(futures): + results.append(future.result()) + + results.sort(key=lambda row: (row.case, row.include_total, row.elapsed_ms)) + writer = csv.DictWriter( + sys.stdout, + fieldnames=[ + "case", + "include_total", + "status", + "elapsed_ms", + "total", + "total_exact", + "has_more", + "item_count", + "error", + ], + ) + writer.writeheader() + writer.writerows(result_rows(results)) + summarize(results) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/web/src/pages/LibraryBrowse.test.tsx b/web/src/pages/LibraryBrowse.test.tsx index 8318fdc3..c48a4d58 100644 --- a/web/src/pages/LibraryBrowse.test.tsx +++ b/web/src/pages/LibraryBrowse.test.tsx @@ -69,7 +69,7 @@ describe("LibraryBrowse", () => { expect(mocks.useCatalogWindow).toHaveBeenCalled(); }); - it("shows the filtered item count from the catalog result", () => { + it("does not show the estimated item count as an exact result count", () => { mocks.useCatalogWindow.mockReturnValue({ data: { totalItems: 1234, @@ -95,10 +95,11 @@ describe("LibraryBrowse", () => { />, ); - expect(markup).toContain("1,234 items"); + expect(markup).toContain("Filters"); + expect(markup).not.toContain("1,234 items"); }); - it("shows a loading state for the item count while catalog data loads", () => { + it("does not show an item count loading state when exact totals are skipped", () => { mocks.useCatalogWindow.mockReturnValue({ data: { totalItems: 0, @@ -124,7 +125,8 @@ describe("LibraryBrowse", () => { />, ); - expect(markup).toContain("Loading item count"); + expect(markup).toContain("Filters"); + expect(markup).not.toContain("Loading item count"); expect(markup).not.toContain("0 items"); }); @@ -152,7 +154,7 @@ describe("LibraryBrowse", () => { }), }), expect.objectContaining({ - includeTotal: true, + includeTotal: false, }), ); }); diff --git a/web/src/pages/LibraryBrowse.tsx b/web/src/pages/LibraryBrowse.tsx index 752e8779..c0cd4316 100644 --- a/web/src/pages/LibraryBrowse.tsx +++ b/web/src/pages/LibraryBrowse.tsx @@ -38,10 +38,6 @@ function getLibrarySortRelevanceScope( return "all"; } -function formatItemCount(count: number): string { - return `${count.toLocaleString()} item${count === 1 ? "" : "s"}`; -} - export default function LibraryBrowse({ libraryId, libraryType, @@ -106,13 +102,12 @@ export default function LibraryBrowse({ const catalogQuery = useCatalogWindow(state, { limit, - includeTotal: true, + includeTotal: false, visibleRange, }); const totalItems = catalogQuery.data?.totalItems ?? 0; const pages = catalogQuery.data?.pages ?? new Map(); const isLoading = catalogQuery.isLoading; - const itemCountLabel = formatItemCount(totalItems); return (
@@ -146,8 +141,6 @@ export default function LibraryBrowse({ allowPersonalizedFilters allowPersonalizedSorts sortRelevanceScope={sortRelevanceScope} - resultCountLabel={itemCountLabel} - resultCountLoading={isLoading} />