2026-08-06 22:33:56 +12:00
|
|
|
package store
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"context"
|
|
|
|
|
"encoding/json"
|
|
|
|
|
"errors"
|
|
|
|
|
"fmt"
|
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
"github.com/jackc/pgx/v5"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
// ErrMediaRatingsNotFound means no durable external-rating response has been stored yet.
|
|
|
|
|
var ErrMediaRatingsNotFound = errors.New("store: media ratings not found")
|
|
|
|
|
|
|
|
|
|
// MediaRatings returns the raw provider response and the time it was fetched. Keeping the
|
|
|
|
|
// raw response lets API settings select different sources without invalidating this cache.
|
|
|
|
|
func (s *Store) MediaRatings(
|
|
|
|
|
ctx context.Context, mediaType, provider, providerID string,
|
|
|
|
|
) (json.RawMessage, time.Time, error) {
|
|
|
|
|
var ratings []byte
|
|
|
|
|
var fetchedAt time.Time
|
|
|
|
|
err := s.pool.QueryRow(ctx, `
|
|
|
|
|
SELECT ratings, fetched_at
|
|
|
|
|
FROM external_media_ratings
|
|
|
|
|
WHERE media_type = $1 AND provider = $2 AND provider_id = $3`,
|
|
|
|
|
mediaType, provider, providerID).Scan(&ratings, &fetchedAt)
|
|
|
|
|
if errors.Is(err, pgx.ErrNoRows) {
|
|
|
|
|
return nil, time.Time{}, ErrMediaRatingsNotFound
|
|
|
|
|
}
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, time.Time{}, fmt.Errorf("store: load media ratings: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return json.RawMessage(ratings), fetchedAt, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// RatingKey identifies one title at the external provider. The same key serves every
|
|
|
|
|
// television and every viewer, which is why ratings are stored per title rather than
|
|
|
|
|
// per Emby item — a series and each of its episodes share one row.
|
|
|
|
|
type RatingKey struct {
|
|
|
|
|
MediaType string // movie | show
|
|
|
|
|
Provider string // tmdb | imdb
|
|
|
|
|
ProviderID string
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (k RatingKey) valid() bool {
|
|
|
|
|
return k.MediaType != "" && k.Provider != "" && k.ProviderID != ""
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// MediaRatingsEntry is one stored provider response and its age.
|
|
|
|
|
type MediaRatingsEntry struct {
|
|
|
|
|
Ratings json.RawMessage
|
|
|
|
|
FetchedAt time.Time
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// MediaRatingsBatch reads many stored responses in one query. Rows attach ratings to
|
|
|
|
|
// every card they carry, so the per-title read of MediaRatings would otherwise mean a
|
|
|
|
|
// round trip per poster.
|
|
|
|
|
func (s *Store) MediaRatingsBatch(
|
|
|
|
|
ctx context.Context, keys []RatingKey,
|
|
|
|
|
) (map[RatingKey]MediaRatingsEntry, error) {
|
|
|
|
|
found := make(map[RatingKey]MediaRatingsEntry, len(keys))
|
|
|
|
|
mediaTypes, providers, providerIDs := ratingKeyColumns(keys)
|
|
|
|
|
if len(mediaTypes) == 0 {
|
|
|
|
|
return found, nil
|
|
|
|
|
}
|
|
|
|
|
rows, err := s.pool.Query(ctx, `
|
|
|
|
|
SELECT media_type, provider, provider_id, ratings, fetched_at
|
|
|
|
|
FROM external_media_ratings
|
|
|
|
|
WHERE (media_type, provider, provider_id) IN (
|
|
|
|
|
SELECT * FROM unnest($1::text[], $2::text[], $3::text[])
|
|
|
|
|
)`, mediaTypes, providers, providerIDs)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: load media ratings batch: %w", err)
|
|
|
|
|
}
|
|
|
|
|
defer rows.Close()
|
|
|
|
|
for rows.Next() {
|
|
|
|
|
var key RatingKey
|
|
|
|
|
var ratings []byte
|
|
|
|
|
var fetchedAt time.Time
|
|
|
|
|
if err := rows.Scan(&key.MediaType, &key.Provider, &key.ProviderID, &ratings, &fetchedAt); err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: scan media ratings batch: %w", err)
|
|
|
|
|
}
|
|
|
|
|
found[key] = MediaRatingsEntry{Ratings: json.RawMessage(ratings), FetchedAt: fetchedAt}
|
|
|
|
|
}
|
|
|
|
|
if err := rows.Err(); err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: load media ratings batch: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return found, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// ratingKeyColumns flattens keys into the three parallel arrays unnest expects, dropping
|
|
|
|
|
// incomplete ones so a single unresolved item cannot break the whole read.
|
|
|
|
|
func ratingKeyColumns(keys []RatingKey) ([]string, []string, []string) {
|
|
|
|
|
mediaTypes := make([]string, 0, len(keys))
|
|
|
|
|
providers := make([]string, 0, len(keys))
|
|
|
|
|
providerIDs := make([]string, 0, len(keys))
|
|
|
|
|
seen := make(map[RatingKey]bool, len(keys))
|
|
|
|
|
for _, key := range keys {
|
|
|
|
|
if !key.valid() || seen[key] {
|
|
|
|
|
continue
|
|
|
|
|
}
|
|
|
|
|
seen[key] = true
|
|
|
|
|
mediaTypes = append(mediaTypes, key.MediaType)
|
|
|
|
|
providers = append(providers, key.Provider)
|
|
|
|
|
providerIDs = append(providerIDs, key.ProviderID)
|
|
|
|
|
}
|
|
|
|
|
return mediaTypes, providers, providerIDs
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// SaveItemRatingRef remembers which external title an Emby item is, so later rows do not
|
|
|
|
|
// have to ask Emby for its ProviderIds again.
|
|
|
|
|
func (s *Store) SaveItemRatingRef(ctx context.Context, itemID string, key RatingKey) error {
|
|
|
|
|
if itemID == "" || !key.valid() {
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
_, err := s.pool.Exec(ctx, `
|
|
|
|
|
INSERT INTO item_rating_refs (item_id, media_type, provider, provider_id, updated_at)
|
|
|
|
|
VALUES ($1, $2, $3, $4, now())
|
|
|
|
|
ON CONFLICT (item_id) DO UPDATE SET
|
|
|
|
|
media_type = EXCLUDED.media_type,
|
|
|
|
|
provider = EXCLUDED.provider,
|
|
|
|
|
provider_id = EXCLUDED.provider_id,
|
|
|
|
|
updated_at = EXCLUDED.updated_at`,
|
|
|
|
|
itemID, key.MediaType, key.Provider, key.ProviderID)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return fmt.Errorf("store: save item rating ref: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// ItemRatingRefs resolves Emby item ids that have been looked up before.
|
|
|
|
|
func (s *Store) ItemRatingRefs(ctx context.Context, itemIDs []string) (map[string]RatingKey, error) {
|
|
|
|
|
found := make(map[string]RatingKey, len(itemIDs))
|
|
|
|
|
if len(itemIDs) == 0 {
|
|
|
|
|
return found, nil
|
|
|
|
|
}
|
|
|
|
|
rows, err := s.pool.Query(ctx, `
|
|
|
|
|
SELECT item_id, media_type, provider, provider_id
|
|
|
|
|
FROM item_rating_refs
|
|
|
|
|
WHERE item_id = ANY($1)`, itemIDs)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: load item rating refs: %w", err)
|
|
|
|
|
}
|
|
|
|
|
defer rows.Close()
|
|
|
|
|
for rows.Next() {
|
|
|
|
|
var itemID string
|
|
|
|
|
var key RatingKey
|
|
|
|
|
if err := rows.Scan(&itemID, &key.MediaType, &key.Provider, &key.ProviderID); err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: scan item rating refs: %w", err)
|
|
|
|
|
}
|
|
|
|
|
found[itemID] = key
|
|
|
|
|
}
|
|
|
|
|
if err := rows.Err(); err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: load item rating refs: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return found, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// LibraryProviderIDs reads the provider identifiers the library import stored, following
|
|
|
|
|
// an episode to its series — MDBList rates shows, never single episodes.
|
|
|
|
|
//
|
|
|
|
|
// The returned map is keyed by the *requested* item id and carries Emby's own
|
|
|
|
|
// ProviderIds object, so the caller applies the same provider precedence it uses for a
|
|
|
|
|
// live lookup rather than a second copy of that rule living in SQL.
|
|
|
|
|
func (s *Store) LibraryProviderIDs(
|
|
|
|
|
ctx context.Context, itemIDs []string,
|
|
|
|
|
) (map[string]LibraryProviderRef, error) {
|
|
|
|
|
found := make(map[string]LibraryProviderRef, len(itemIDs))
|
|
|
|
|
if len(itemIDs) == 0 {
|
|
|
|
|
return found, nil
|
|
|
|
|
}
|
|
|
|
|
rows, err := s.pool.Query(ctx, `
|
|
|
|
|
SELECT item.id,
|
|
|
|
|
item.type,
|
|
|
|
|
CASE WHEN item.type = 'Episode' THEN series.payload -> 'ProviderIds'
|
|
|
|
|
ELSE item.payload -> 'ProviderIds' END
|
|
|
|
|
FROM library_items item
|
|
|
|
|
LEFT JOIN library_items series
|
|
|
|
|
ON item.type = 'Episode' AND series.id = item.series_id
|
|
|
|
|
WHERE item.id = ANY($1)`, itemIDs)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: load library provider ids: %w", err)
|
|
|
|
|
}
|
|
|
|
|
defer rows.Close()
|
|
|
|
|
for rows.Next() {
|
|
|
|
|
var itemID, itemType string
|
|
|
|
|
var raw []byte
|
|
|
|
|
if err := rows.Scan(&itemID, &itemType, &raw); err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: scan library provider ids: %w", err)
|
|
|
|
|
}
|
|
|
|
|
ids := map[string]string{}
|
|
|
|
|
if len(raw) > 0 {
|
|
|
|
|
// Emby writes provider ids as strings, but a hand-edited or future payload
|
|
|
|
|
// need not: an undecodable object leaves the item unresolved rather than
|
|
|
|
|
// failing the whole read.
|
|
|
|
|
_ = json.Unmarshal(raw, &ids)
|
|
|
|
|
}
|
|
|
|
|
found[itemID] = LibraryProviderRef{Type: itemType, ProviderIDs: ids}
|
|
|
|
|
}
|
|
|
|
|
if err := rows.Err(); err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: load library provider ids: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return found, nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// LibraryProviderRef is one imported item's type and external identifiers.
|
|
|
|
|
type LibraryProviderRef struct {
|
|
|
|
|
Type string
|
|
|
|
|
ProviderIDs map[string]string
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// MediaRatingsStats reports how much of the durable cache exists and how much of it is
|
|
|
|
|
// old enough to be refreshed, which is the only way the admin page can say whether the
|
|
|
|
|
// integration is still spending external requests.
|
|
|
|
|
func (s *Store) MediaRatingsStats(ctx context.Context, staleBefore time.Time) (total, stale int, err error) {
|
|
|
|
|
err = s.pool.QueryRow(ctx, `
|
|
|
|
|
SELECT count(*), count(*) FILTER (WHERE fetched_at < $1)
|
|
|
|
|
FROM external_media_ratings`, staleBefore).Scan(&total, &stale)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return 0, 0, fmt.Errorf("store: media ratings stats: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return total, stale, nil
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-19 14:25:44 +12:00
|
|
|
// StaleRatingKeys is the oldest stored titles due to be re-checked, oldest first.
|
|
|
|
|
//
|
|
|
|
|
// Oldest first rather than by any measure of popularity, because the refresh is bounded
|
|
|
|
|
// per run: taking the oldest means every title comes round eventually, where taking the
|
|
|
|
|
// most-watched would leave the tail of the library permanently on scores from the year it
|
|
|
|
|
// was imported. The limit is what keeps a run's cost — and therefore the day's external
|
|
|
|
|
// allowance — a figure the operator can reason about.
|
|
|
|
|
func (s *Store) StaleRatingKeys(
|
|
|
|
|
ctx context.Context, staleBefore time.Time, limit int,
|
|
|
|
|
) ([]RatingKey, error) {
|
|
|
|
|
if limit <= 0 || limit > 500 {
|
|
|
|
|
limit = 100
|
|
|
|
|
}
|
|
|
|
|
rows, err := s.pool.Query(ctx, `
|
|
|
|
|
SELECT media_type, provider, provider_id
|
|
|
|
|
FROM external_media_ratings
|
|
|
|
|
WHERE fetched_at < $1
|
|
|
|
|
ORDER BY fetched_at ASC
|
|
|
|
|
LIMIT $2`, staleBefore, limit)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: stale rating keys: %w", err)
|
|
|
|
|
}
|
|
|
|
|
defer rows.Close()
|
|
|
|
|
keys := []RatingKey{}
|
|
|
|
|
for rows.Next() {
|
|
|
|
|
var key RatingKey
|
|
|
|
|
if err := rows.Scan(&key.MediaType, &key.Provider, &key.ProviderID); err != nil {
|
|
|
|
|
return nil, fmt.Errorf("store: scan stale rating key: %w", err)
|
|
|
|
|
}
|
|
|
|
|
keys = append(keys, key)
|
|
|
|
|
}
|
|
|
|
|
return keys, rows.Err()
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-06 22:33:56 +12:00
|
|
|
// SaveMediaRatings upserts a successfully loaded provider response.
|
|
|
|
|
func (s *Store) SaveMediaRatings(
|
|
|
|
|
ctx context.Context, mediaType, provider, providerID string, ratings json.RawMessage,
|
|
|
|
|
) error {
|
|
|
|
|
if !json.Valid(ratings) {
|
|
|
|
|
return errors.New("store: save media ratings: invalid JSON")
|
|
|
|
|
}
|
|
|
|
|
_, err := s.pool.Exec(ctx, `
|
|
|
|
|
INSERT INTO external_media_ratings (media_type, provider, provider_id, ratings, fetched_at)
|
|
|
|
|
VALUES ($1, $2, $3, $4, now())
|
|
|
|
|
ON CONFLICT (media_type, provider, provider_id) DO UPDATE SET
|
|
|
|
|
ratings = EXCLUDED.ratings,
|
|
|
|
|
fetched_at = EXCLUDED.fetched_at`,
|
|
|
|
|
mediaType, provider, providerID, ratings)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return fmt.Errorf("store: save media ratings: %w", err)
|
|
|
|
|
}
|
|
|
|
|
return nil
|
|
|
|
|
}
|