feat(movies): TPDB movie enrichment + dedup
Enrich existing movies (from paradisehill/dooplay, which mostly lack cast) with metadata from TPDB's /movies API: cast, categories (tags), studio, director + a canonical TPDB UUID for dedup. Chosen over IAFD after a source-comparison research pass — IAFD has strong cast/studio but ZERO categories, while TPDB /movies has ~11 tags/movie, cast, studio, director, a canonical UUID (+ sparse phash), is already an integrated API (no scraping/anti-bot), and covers ~75-85% of our western-DVD-feature catalog. Enrichment only ever augments EXISTING movies and never creates new ones (TPDB has no playback, so a standalone TPDB movie would be unplayable). Writes to movie_performers / movie_tags / movie.studio_id, which the movies API + mobile detail already render, so no schema/API/UI change is needed. - connectors/tpdb.py: search_movies() + fetch_movie() + _parse_movie() reusing the existing _parse_studio/_parse_performer/_parse_tag. - enrich/tpdb_movies.py: match our movie to a TPDB /movies result by token_sort_ratio on normalized titles (sort, not set, to reject the short-title-subset trap "Fantasies" -> "Tara's Fetish Fantasies") with a +/-2yr guard; then attach cast/tags/studio/director. Incoming performers deduped by external_id to avoid the performer_external_refs PK clash. - resolve/movie_merge.py: merge_movies() mirror of scene_merge; two of our movies mapping to the same TPDB UUID are the same film -> merge. - scheduler: _job_tpdb_movie_enrich every 6h, batch 200, prioritizing playable movies missing cast/studio. Verified on a 150-movie batch: 119 enriched, 4 deduped, 26 no-match, 0 errors; matched titles/studios spot-checked correct (Big Butts Drive Me Nuts 4 -> 33 tags, Seinfeld #2 -> 10 cast/17 tags, German BB Video titles -> categories+studio). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
parent
cb0b843f48
commit
b9898ba592
7 changed files with 601 additions and 0 deletions
|
|
@ -136,6 +136,16 @@ class Settings(BaseSettings):
|
||||||
ingest_watchdog_search_max_age_hours: int = Field(
|
ingest_watchdog_search_max_age_hours: int = Field(
|
||||||
default=168, validation_alias="GOON_INGEST_WATCHDOG_SEARCH_MAX_AGE_HOURS"
|
default=168, validation_alias="GOON_INGEST_WATCHDOG_SEARCH_MAX_AGE_HOURS"
|
||||||
)
|
)
|
||||||
|
# TPDB movie enrichment — wzbogaca filmy obsadą + kategoriami (tagi) + studiem +
|
||||||
|
# reżyserem, plus kanoniczny TPDB UUID do dedupu mirrorów. paradisehill (primary)
|
||||||
|
# prawie nie ma obsady, więc to domyka największą lukę zakładki movies. Batch
|
||||||
|
# najświeższych GRYWALNYCH filmów bez obsady/studia per run; 6h cadence. 0 = off.
|
||||||
|
sched_tpdb_movie_enrich_hours: int = Field(
|
||||||
|
default=6, validation_alias="GOON_SCHED_TPDB_MOVIE_ENRICH_HOURS"
|
||||||
|
)
|
||||||
|
tpdb_movie_enrich_batch: int = Field(
|
||||||
|
default=200, validation_alias="GOON_TPDB_MOVIE_ENRICH_BATCH"
|
||||||
|
)
|
||||||
# Taxonomy scene_count refresh — przelicza denormalizowane liczniki scen na
|
# Taxonomy scene_count refresh — przelicza denormalizowane liczniki scen na
|
||||||
# tags/performers/studios (hot-path /tags|/performers|/studios|/favorites czyta
|
# tags/performers/studios (hot-path /tags|/performers|/studios|/favorites czyta
|
||||||
# gotową kolumnę zamiast agregować 6.3M scene_tags per-request). 3h cadence —
|
# gotową kolumnę zamiast agregować 6.3M scene_tags per-request). 3h cadence —
|
||||||
|
|
|
||||||
|
|
@ -49,6 +49,7 @@ def _is_retryable_http_error(exc: BaseException) -> bool:
|
||||||
from app.config import get_settings
|
from app.config import get_settings
|
||||||
from app.connectors.base import (
|
from app.connectors.base import (
|
||||||
BaseConnector,
|
BaseConnector,
|
||||||
|
RawMovie,
|
||||||
RawPerformer,
|
RawPerformer,
|
||||||
RawScene,
|
RawScene,
|
||||||
RawStudio,
|
RawStudio,
|
||||||
|
|
@ -195,6 +196,37 @@ class TPDBConnector(BaseConnector):
|
||||||
first = data[0]
|
first = data[0]
|
||||||
return str(first.get("id")) if first.get("id") else None
|
return str(first.get("id")) if first.get("id") else None
|
||||||
|
|
||||||
|
def search_movies(self, query: str, *, per_page: int = 10) -> list[dict[str, Any]]:
|
||||||
|
"""GET /movies?q=<query> → surowe payloady filmów (list endpoint jest już
|
||||||
|
hydrated: tags + performers + site embedded, bez detail-fetcha).
|
||||||
|
|
||||||
|
Używane do enrichmentu movies: wyszukujemy kandydatów po tytule, scorujemy
|
||||||
|
i mapujemy najlepszy przez `_parse_movie`. Pusty list na błąd/brak."""
|
||||||
|
if not query.strip():
|
||||||
|
return []
|
||||||
|
with self._client() as client:
|
||||||
|
try:
|
||||||
|
payload = self._get(client, "/movies", {"q": query, "per_page": per_page})
|
||||||
|
except httpx.HTTPStatusError as e:
|
||||||
|
log.warning("tpdb /movies q=%s failed: %s", query, e)
|
||||||
|
return []
|
||||||
|
return payload.get("data") or []
|
||||||
|
|
||||||
|
def fetch_movie(self, movie_uuid: str) -> dict[str, Any] | None:
|
||||||
|
"""GET /movies/<uuid> → pojedynczy film (detail). None gdy 404/błąd.
|
||||||
|
|
||||||
|
NB: list endpoint jest już hydrated, więc detail rzadko potrzebny — trzymamy
|
||||||
|
dla spójności / gdyby TPDB kiedyś przeniósł część pól tylko do detalu."""
|
||||||
|
with self._client() as client:
|
||||||
|
try:
|
||||||
|
payload = self._get(client, f"/movies/{movie_uuid}", {})
|
||||||
|
except httpx.HTTPStatusError as e:
|
||||||
|
if e.response.status_code == 404:
|
||||||
|
return None
|
||||||
|
log.warning("tpdb /movies/%s failed: %s", movie_uuid, e)
|
||||||
|
return None
|
||||||
|
return payload.get("data")
|
||||||
|
|
||||||
def _paginate_scenes(
|
def _paginate_scenes(
|
||||||
self,
|
self,
|
||||||
params: dict[str, Any],
|
params: dict[str, Any],
|
||||||
|
|
@ -327,3 +359,60 @@ def _parse_scene(raw: dict[str, Any]) -> RawScene | None:
|
||||||
fingerprints=[], # TPDB nie publikuje pHashy w głównym endpoint
|
fingerprints=[], # TPDB nie publikuje pHashy w głównym endpoint
|
||||||
raw=raw,
|
raw=raw,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _parse_movie(raw: dict[str, Any]) -> RawMovie | None:
|
||||||
|
"""TPDB movie payload → RawMovie. Reużywa _parse_studio/_parse_performer/_parse_tag
|
||||||
|
(movies mają identyczny kształt site/performers/tags co sceny)."""
|
||||||
|
external_id = raw.get("id")
|
||||||
|
title = raw.get("title")
|
||||||
|
if not external_id or not title:
|
||||||
|
return None
|
||||||
|
|
||||||
|
performers: list[RawPerformer] = []
|
||||||
|
for p in raw.get("performers") or []:
|
||||||
|
parsed = _parse_performer(p)
|
||||||
|
if parsed is not None:
|
||||||
|
performers.append(parsed)
|
||||||
|
tags: list[RawTag] = []
|
||||||
|
for t in raw.get("tags") or []:
|
||||||
|
parsed_t = _parse_tag(t)
|
||||||
|
if parsed_t is not None:
|
||||||
|
tags.append(parsed_t)
|
||||||
|
|
||||||
|
directors = raw.get("directors") or []
|
||||||
|
director = ", ".join(d.get("name") for d in directors if isinstance(d, dict) and d.get("name")) or None
|
||||||
|
|
||||||
|
d = _parse_date(raw.get("date"))
|
||||||
|
# duration bywa 1s-placeholderem na movies — poniżej 60s traktujemy jak brak.
|
||||||
|
dur = raw.get("duration")
|
||||||
|
duration_sec = int(dur) if dur and int(dur) > 60 else None
|
||||||
|
|
||||||
|
def _img(field: str) -> str | None:
|
||||||
|
v = raw.get(field)
|
||||||
|
if isinstance(v, str) and v:
|
||||||
|
return v
|
||||||
|
if isinstance(v, dict):
|
||||||
|
return v.get("large") or v.get("full") or v.get("medium") or None
|
||||||
|
return None
|
||||||
|
|
||||||
|
poster = _img("poster") or _img("posters")
|
||||||
|
backdrop = _img("background") or _img("backdrop")
|
||||||
|
|
||||||
|
return RawMovie(
|
||||||
|
external_id=str(external_id),
|
||||||
|
title=title,
|
||||||
|
description=raw.get("description"),
|
||||||
|
release_year=d.year if d else None,
|
||||||
|
release_date=d,
|
||||||
|
duration_sec=duration_sec,
|
||||||
|
director=director,
|
||||||
|
rating=float(raw["rating"]) if raw.get("rating") not in (None, "") else None,
|
||||||
|
poster_url=poster,
|
||||||
|
backdrop_url=backdrop,
|
||||||
|
url=raw.get("url"),
|
||||||
|
studio=_parse_studio(raw.get("site")),
|
||||||
|
performers=performers,
|
||||||
|
tags=tags,
|
||||||
|
raw=raw,
|
||||||
|
)
|
||||||
|
|
|
||||||
1
app/enrich/__init__.py
Normal file
1
app/enrich/__init__.py
Normal file
|
|
@ -0,0 +1 @@
|
||||||
|
"""Post-ingest enrichment of canonical entities from authoritative metadata sources."""
|
||||||
273
app/enrich/tpdb_movies.py
Normal file
273
app/enrich/tpdb_movies.py
Normal file
|
|
@ -0,0 +1,273 @@
|
||||||
|
"""TPDB movie enrichment + dedup.
|
||||||
|
|
||||||
|
TPDB `/movies` jest naszym kanonicznym źródłem metadanych filmów: obsada, kategorie
|
||||||
|
(tagi), studio, reżyser, rok — plus stabilny UUID do dedupu. paradisehill (primary
|
||||||
|
movie source) prawie nie ma obsady, więc TPDB wypełnia największą lukę.
|
||||||
|
|
||||||
|
Zasada: TPDB TYLKO wzbogaca ISTNIEJĄCE filmy (z paradisehill/dooplay) i dedupuje,
|
||||||
|
NIGDY nie tworzy nowych — TPDB nie ma playbacku, więc nowy film byłby niegrywalny.
|
||||||
|
|
||||||
|
Flow per film:
|
||||||
|
1. skip, jeśli film ma już TPDB movie ref (movie_external_refs jest tylko-movie,
|
||||||
|
więc dowolny ref z tpdb source = już wzbogacony),
|
||||||
|
2. search TPDB /movies?q=<tytuł>, wybierz najlepszego kandydata (token-set na tytule
|
||||||
|
+ guard roku ±2, próg `min_title`),
|
||||||
|
3. DEDUP: jeśli ten TPDB UUID jest już przypięty do INNEGO naszego filmu → to ten
|
||||||
|
sam film (mirror) → merge_movies (keep = starszy created_at),
|
||||||
|
4. wzbogać ocalały film: obsada (resolve_performer → MoviePerformer), tagi
|
||||||
|
(resolve_tag → MovieTag source=tpdb), studio (fill studio_id), reżyser/rok/
|
||||||
|
poster/rating (fill-only), przypnij MovieExternalRef(tpdb, uuid).
|
||||||
|
|
||||||
|
Match jest zachowawczy (próg 0.90 + guard roku), bo TPDB search zwraca dużo (np.
|
||||||
|
"Pirates" → gay Knightbreeders przed feature) — bez scoringu wzięlibyśmy zły film.
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
import uuid
|
||||||
|
from datetime import UTC, datetime
|
||||||
|
|
||||||
|
from rapidfuzz import fuzz
|
||||||
|
from sqlalchemy import func, select
|
||||||
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
from app.connectors.tpdb import TPDBConnector, _parse_movie
|
||||||
|
from app.models.movie import Movie, MovieExternalRef, MoviePerformer
|
||||||
|
from app.models.movie_playback_source import MoviePlaybackSource
|
||||||
|
from app.normalize.movies import normalize_movie
|
||||||
|
from app.normalize.text import normalize
|
||||||
|
from app.resolve.movie_merge import merge_movies
|
||||||
|
from app.resolve.movie_resolver import _sync_performers, _sync_tags
|
||||||
|
from app.resolve.performer_resolver import resolve_performer
|
||||||
|
from app.resolve.studio_resolver import resolve_studio
|
||||||
|
|
||||||
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
_DASH = str.maketrans({"–": "-", "—": "-", "‑": "-"})
|
||||||
|
|
||||||
|
|
||||||
|
def _cand_year(raw: dict) -> int | None:
|
||||||
|
d = raw.get("date")
|
||||||
|
if d and len(str(d)) >= 4:
|
||||||
|
try:
|
||||||
|
return int(str(d)[:4])
|
||||||
|
except ValueError:
|
||||||
|
return None
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _best_match(connector: TPDBConnector, movie: Movie, *, min_title: float) -> dict | None:
|
||||||
|
"""Najlepszy TPDB movie payload dla naszego filmu, albo None.
|
||||||
|
|
||||||
|
`token_sort_ratio` (a NIE token_set) na znormalizowanych tytułach: odporny na
|
||||||
|
inwersję ("Title, The" ↔ "The Title") i kolejność, ALE penalizuje różnicę długości,
|
||||||
|
więc krótki generyczny tytuł nie łapie dłuższego nadzbioru (bug: "Fantasies" →
|
||||||
|
"Tara's Fetish Fantasies", bo token_set nagradza podzbiór). Precyzja > recall:
|
||||||
|
lepiej pominąć niż wzbogacić zły film. Guard roku ±2 (inny rok = inna edycja)."""
|
||||||
|
query = (movie.title or "").translate(_DASH).strip()[:60]
|
||||||
|
if not query:
|
||||||
|
return None
|
||||||
|
my_norm = normalize(movie.title)
|
||||||
|
best: dict | None = None
|
||||||
|
best_score = 0.0
|
||||||
|
for raw in connector.search_movies(query, per_page=10):
|
||||||
|
cand_title = raw.get("title")
|
||||||
|
if not cand_title:
|
||||||
|
continue
|
||||||
|
score = fuzz.token_sort_ratio(my_norm, normalize(cand_title)) / 100.0
|
||||||
|
if score <= best_score:
|
||||||
|
continue
|
||||||
|
cy = _cand_year(raw)
|
||||||
|
if movie.release_year and cy and abs(movie.release_year - cy) > 2:
|
||||||
|
continue # guard: inny rok → prawdopodobnie inny film (Taxi 2 ≠ Taxi Violeur 2)
|
||||||
|
best_score = score
|
||||||
|
best = raw
|
||||||
|
return best if best is not None and best_score >= min_title else None
|
||||||
|
|
||||||
|
|
||||||
|
def _fill_scalar_fields(movie: Movie, norm) -> None:
|
||||||
|
"""Fill-only — nie nadpisujemy pól ustawionych przez primary (paradisehill/mirror)."""
|
||||||
|
if norm.director and not movie.director:
|
||||||
|
movie.director = norm.director
|
||||||
|
if norm.release_year and not movie.release_year:
|
||||||
|
movie.release_year = norm.release_year
|
||||||
|
if norm.release_date and not movie.release_date:
|
||||||
|
movie.release_date = norm.release_date
|
||||||
|
if norm.duration_sec and not movie.duration_sec:
|
||||||
|
movie.duration_sec = norm.duration_sec
|
||||||
|
if norm.description and not movie.description:
|
||||||
|
movie.description = norm.description
|
||||||
|
if norm.poster_url and not movie.poster_url:
|
||||||
|
movie.poster_url = norm.poster_url
|
||||||
|
if norm.backdrop_url and not movie.backdrop_url:
|
||||||
|
movie.backdrop_url = norm.backdrop_url
|
||||||
|
if norm.rating is not None and movie.rating is None:
|
||||||
|
movie.rating = norm.rating
|
||||||
|
|
||||||
|
|
||||||
|
def enrich_movie(
|
||||||
|
session: Session,
|
||||||
|
movie: Movie,
|
||||||
|
*,
|
||||||
|
connector: TPDBConnector,
|
||||||
|
source_id: uuid.UUID,
|
||||||
|
min_title: float = 0.90,
|
||||||
|
) -> str:
|
||||||
|
"""Zwraca: 'skip' | 'no_match' | 'enriched' | 'merged'."""
|
||||||
|
already = session.execute(
|
||||||
|
select(MovieExternalRef.external_id).where(
|
||||||
|
MovieExternalRef.source_id == source_id,
|
||||||
|
MovieExternalRef.movie_id == movie.id,
|
||||||
|
)
|
||||||
|
).first()
|
||||||
|
if already is not None:
|
||||||
|
return "skip"
|
||||||
|
|
||||||
|
raw = _best_match(connector, movie, min_title=min_title)
|
||||||
|
if raw is None:
|
||||||
|
return "no_match"
|
||||||
|
ext_id = str(raw["id"])
|
||||||
|
|
||||||
|
outcome = "enriched"
|
||||||
|
# DEDUP: ten TPDB UUID już przypięty do innego naszego filmu → mirror tego samego.
|
||||||
|
other = session.execute(
|
||||||
|
select(MovieExternalRef).where(
|
||||||
|
MovieExternalRef.source_id == source_id,
|
||||||
|
MovieExternalRef.external_id == ext_id,
|
||||||
|
)
|
||||||
|
).scalar_one_or_none()
|
||||||
|
if other is not None and other.movie_id != movie.id:
|
||||||
|
m2 = session.get(Movie, other.movie_id)
|
||||||
|
if m2 is not None:
|
||||||
|
now = datetime.now(UTC)
|
||||||
|
keep, drop = (
|
||||||
|
(movie, m2)
|
||||||
|
if (movie.created_at or now) <= (m2.created_at or now)
|
||||||
|
else (m2, movie)
|
||||||
|
)
|
||||||
|
movie = merge_movies(
|
||||||
|
session, keep_id=keep.id, drop_id=drop.id, resolved_by="tpdb:dedup"
|
||||||
|
)
|
||||||
|
outcome = "merged"
|
||||||
|
|
||||||
|
rm = _parse_movie(raw)
|
||||||
|
if rm is None:
|
||||||
|
return "no_match"
|
||||||
|
norm = normalize_movie(rm)
|
||||||
|
|
||||||
|
if norm.studio is not None:
|
||||||
|
studio = resolve_studio(session, norm=norm.studio, source_id=source_id)
|
||||||
|
if studio is not None and not movie.studio_id:
|
||||||
|
movie.studio_id = studio.id
|
||||||
|
|
||||||
|
# Dedup wchodzących performerów po external_id (TPDB potrafi wylistować tego samego
|
||||||
|
# kanonicznego performera 2×: pod aliasem i kanonicznie) — bez tego resolve_performer
|
||||||
|
# próbuje wstawić performer_external_refs 2× → UniqueViolation (jak w scene_resolver).
|
||||||
|
resolved: list[tuple[uuid.UUID, str | None]] = []
|
||||||
|
seen_perf_keys: set[str] = set()
|
||||||
|
for p_norm in norm.performers:
|
||||||
|
key = p_norm.external_id or normalize(p_norm.name)
|
||||||
|
if key in seen_perf_keys:
|
||||||
|
continue
|
||||||
|
seen_perf_keys.add(key)
|
||||||
|
performer = resolve_performer(session, norm=p_norm, source_id=source_id)
|
||||||
|
resolved.append((performer.id, p_norm.as_alias_in_scene))
|
||||||
|
_sync_performers(session, movie_id=movie.id, resolved=resolved)
|
||||||
|
_sync_tags(session, movie_id=movie.id, norm=norm, source_id=source_id)
|
||||||
|
_fill_scalar_fields(movie, norm)
|
||||||
|
|
||||||
|
ref = session.execute(
|
||||||
|
select(MovieExternalRef).where(
|
||||||
|
MovieExternalRef.source_id == source_id,
|
||||||
|
MovieExternalRef.external_id == ext_id,
|
||||||
|
)
|
||||||
|
).scalar_one_or_none()
|
||||||
|
if ref is None:
|
||||||
|
session.add(
|
||||||
|
MovieExternalRef(
|
||||||
|
source_id=source_id,
|
||||||
|
external_id=ext_id,
|
||||||
|
movie_id=movie.id,
|
||||||
|
confidence=1.0,
|
||||||
|
url=rm.url,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
elif ref.movie_id != movie.id:
|
||||||
|
ref.movie_id = movie.id
|
||||||
|
|
||||||
|
return outcome
|
||||||
|
|
||||||
|
|
||||||
|
def _candidate_movies(session: Session, *, source_id: uuid.UUID, limit: int) -> list[Movie]:
|
||||||
|
"""Filmy jeszcze nie wzbogacone TPDB, priorytet: mają żywy playback (są grywalne)
|
||||||
|
i brak obsady LUB brak studia (największa luka). Reszta później."""
|
||||||
|
has_tpdb = (
|
||||||
|
select(MovieExternalRef.movie_id)
|
||||||
|
.where(MovieExternalRef.source_id == source_id)
|
||||||
|
.scalar_subquery()
|
||||||
|
)
|
||||||
|
has_playback = (
|
||||||
|
select(MoviePlaybackSource.movie_id)
|
||||||
|
.where(MoviePlaybackSource.dead_at.is_(None))
|
||||||
|
.scalar_subquery()
|
||||||
|
)
|
||||||
|
perf_count = (
|
||||||
|
select(func.count())
|
||||||
|
.select_from(MoviePerformer)
|
||||||
|
.where(MoviePerformer.movie_id == Movie.id)
|
||||||
|
.correlate(Movie)
|
||||||
|
.scalar_subquery()
|
||||||
|
)
|
||||||
|
stmt = (
|
||||||
|
select(Movie)
|
||||||
|
.where(
|
||||||
|
Movie.id.not_in(has_tpdb),
|
||||||
|
Movie.id.in_(has_playback),
|
||||||
|
)
|
||||||
|
.where((perf_count == 0) | (Movie.studio_id.is_(None)))
|
||||||
|
.order_by(Movie.created_at.desc())
|
||||||
|
.limit(limit)
|
||||||
|
)
|
||||||
|
return list(session.execute(stmt).scalars().all())
|
||||||
|
|
||||||
|
|
||||||
|
def run_tpdb_movie_enrich(
|
||||||
|
session_factory,
|
||||||
|
*,
|
||||||
|
limit: int = 200,
|
||||||
|
min_title: float = 0.90,
|
||||||
|
) -> dict[str, int]:
|
||||||
|
"""Batch: wzbogać do `limit` filmów. Commit per-film (jeden błąd nie cofa reszty).
|
||||||
|
`session_factory` = kontekstowy scope (app.db.session_scope)."""
|
||||||
|
from app.ingest import get_or_create_source
|
||||||
|
from app.models.source import SourceKind
|
||||||
|
|
||||||
|
connector = TPDBConnector()
|
||||||
|
counters = {"seen": 0, "enriched": 0, "merged": 0, "no_match": 0, "skip": 0, "errors": 0}
|
||||||
|
|
||||||
|
with session_factory() as session:
|
||||||
|
src = get_or_create_source(session, kind=SourceKind.tpdb, name="tpdb")
|
||||||
|
source_id = src.id
|
||||||
|
session.commit()
|
||||||
|
movies = _candidate_movies(session, source_id=source_id, limit=limit)
|
||||||
|
movie_ids = [m.id for m in movies]
|
||||||
|
|
||||||
|
log.info("tpdb-movie-enrich: %d candidate movies", len(movie_ids))
|
||||||
|
for mid in movie_ids:
|
||||||
|
counters["seen"] += 1
|
||||||
|
try:
|
||||||
|
with session_factory() as session:
|
||||||
|
movie = session.get(Movie, mid)
|
||||||
|
if movie is None:
|
||||||
|
continue
|
||||||
|
outcome = enrich_movie(
|
||||||
|
session, movie, connector=connector, source_id=source_id, min_title=min_title
|
||||||
|
)
|
||||||
|
session.commit()
|
||||||
|
counters[outcome] = counters.get(outcome, 0) + 1
|
||||||
|
except Exception as e: # pragma: no cover - defensywnie, jeden film nie wywala batcha
|
||||||
|
counters["errors"] += 1
|
||||||
|
log.warning("tpdb-movie-enrich failed for %s: %s", mid, e)
|
||||||
|
|
||||||
|
log.info("tpdb-movie-enrich done: %s", counters)
|
||||||
|
return counters
|
||||||
190
app/resolve/movie_merge.py
Normal file
190
app/resolve/movie_merge.py
Normal file
|
|
@ -0,0 +1,190 @@
|
||||||
|
"""Scalanie dwóch kanonicznych movies w jeden (dedup).
|
||||||
|
|
||||||
|
`keep_id` przejmuje od `drop_id`: external_refs, performers, tags, chapters,
|
||||||
|
playback_sources — z deduplikacją na kluczach złączeń. Potem `drop` Movie jest
|
||||||
|
usuwany. Mirror `scene_merge.merge_scenes` dla movies.
|
||||||
|
|
||||||
|
Używane przez TPDB enrichment: gdy dwa nasze filmy (mirrory paradisehill/dooplay)
|
||||||
|
mapują się na TEN SAM TPDB movie UUID, to ten sam film → merge. `MergeKind.movie`
|
||||||
|
już istnieje w merge_candidates.
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import logging
|
||||||
|
import uuid
|
||||||
|
from datetime import UTC, datetime
|
||||||
|
|
||||||
|
from sqlalchemy import or_, select, update
|
||||||
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
from app.models.merge_candidate import MergeCandidate, MergeKind, MergeStatus
|
||||||
|
from app.models.movie import (
|
||||||
|
Movie,
|
||||||
|
MovieChapter,
|
||||||
|
MovieExternalRef,
|
||||||
|
MoviePerformer,
|
||||||
|
MovieTag,
|
||||||
|
)
|
||||||
|
from app.models.movie_playback_source import MoviePlaybackSource
|
||||||
|
|
||||||
|
log = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
class MovieMergeError(Exception):
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
def merge_movies(
|
||||||
|
session: Session,
|
||||||
|
*,
|
||||||
|
keep_id: uuid.UUID,
|
||||||
|
drop_id: uuid.UUID,
|
||||||
|
resolved_by: str | None = None,
|
||||||
|
) -> Movie:
|
||||||
|
if keep_id == drop_id:
|
||||||
|
raise MovieMergeError("cannot merge movie into itself")
|
||||||
|
keep = session.get(Movie, keep_id)
|
||||||
|
drop = session.get(Movie, drop_id)
|
||||||
|
if keep is None or drop is None:
|
||||||
|
raise MovieMergeError("movie not found")
|
||||||
|
|
||||||
|
_move_external_refs(session, keep_id=keep_id, drop_id=drop_id)
|
||||||
|
_move_performers(session, keep_id=keep_id, drop_id=drop_id)
|
||||||
|
_move_tags(session, keep_id=keep_id, drop_id=drop_id)
|
||||||
|
_move_chapters(session, keep_id=keep_id, drop_id=drop_id)
|
||||||
|
_move_playback_sources(session, keep_id=keep_id, drop_id=drop_id)
|
||||||
|
_coalesce_canonical_fields(keep, drop)
|
||||||
|
|
||||||
|
session.delete(drop)
|
||||||
|
session.flush()
|
||||||
|
_close_pending_candidates(session, movie_id=drop_id, resolved_by=resolved_by)
|
||||||
|
log.info("merged movie %s ← %s", keep_id, drop_id)
|
||||||
|
return keep
|
||||||
|
|
||||||
|
|
||||||
|
# ---- helpery --------------------------------------------------------------
|
||||||
|
|
||||||
|
def _move_external_refs(session: Session, *, keep_id: uuid.UUID, drop_id: uuid.UUID) -> None:
|
||||||
|
for ref in session.execute(
|
||||||
|
select(MovieExternalRef).where(MovieExternalRef.movie_id == drop_id)
|
||||||
|
).scalars().all():
|
||||||
|
clash = session.execute(
|
||||||
|
select(MovieExternalRef).where(
|
||||||
|
MovieExternalRef.source_id == ref.source_id,
|
||||||
|
MovieExternalRef.external_id == ref.external_id,
|
||||||
|
MovieExternalRef.movie_id == keep_id,
|
||||||
|
)
|
||||||
|
).scalar_one_or_none()
|
||||||
|
if clash is not None:
|
||||||
|
session.delete(ref)
|
||||||
|
else:
|
||||||
|
ref.movie_id = keep_id
|
||||||
|
|
||||||
|
|
||||||
|
def _move_performers(session: Session, *, keep_id: uuid.UUID, drop_id: uuid.UUID) -> None:
|
||||||
|
for link in session.execute(
|
||||||
|
select(MoviePerformer).where(MoviePerformer.movie_id == drop_id)
|
||||||
|
).scalars().all():
|
||||||
|
clash = session.execute(
|
||||||
|
select(MoviePerformer).where(
|
||||||
|
MoviePerformer.movie_id == keep_id,
|
||||||
|
MoviePerformer.performer_id == link.performer_id,
|
||||||
|
)
|
||||||
|
).scalar_one_or_none()
|
||||||
|
if clash is not None:
|
||||||
|
if link.as_alias and not clash.as_alias:
|
||||||
|
clash.as_alias = link.as_alias
|
||||||
|
session.delete(link)
|
||||||
|
else:
|
||||||
|
link.movie_id = keep_id
|
||||||
|
|
||||||
|
|
||||||
|
def _move_tags(session: Session, *, keep_id: uuid.UUID, drop_id: uuid.UUID) -> None:
|
||||||
|
for link in session.execute(
|
||||||
|
select(MovieTag).where(MovieTag.movie_id == drop_id)
|
||||||
|
).scalars().all():
|
||||||
|
clash = session.execute(
|
||||||
|
select(MovieTag).where(
|
||||||
|
MovieTag.movie_id == keep_id, MovieTag.tag_id == link.tag_id
|
||||||
|
)
|
||||||
|
).scalar_one_or_none()
|
||||||
|
if clash is not None:
|
||||||
|
session.delete(link)
|
||||||
|
else:
|
||||||
|
link.movie_id = keep_id
|
||||||
|
|
||||||
|
|
||||||
|
def _move_chapters(session: Session, *, keep_id: uuid.UUID, drop_id: uuid.UUID) -> None:
|
||||||
|
"""Chaptery są lokalne per-movie (UQ movie_id, chapter_index). Kolizja indeksu z
|
||||||
|
keep → keep swoje (kasujemy drop's); brak kolizji → przepięcie."""
|
||||||
|
for ch in session.execute(
|
||||||
|
select(MovieChapter).where(MovieChapter.movie_id == drop_id)
|
||||||
|
).scalars().all():
|
||||||
|
clash = session.execute(
|
||||||
|
select(MovieChapter).where(
|
||||||
|
MovieChapter.movie_id == keep_id,
|
||||||
|
MovieChapter.chapter_index == ch.chapter_index,
|
||||||
|
)
|
||||||
|
).scalar_one_or_none()
|
||||||
|
if clash is not None:
|
||||||
|
session.delete(ch)
|
||||||
|
else:
|
||||||
|
ch.movie_id = keep_id
|
||||||
|
|
||||||
|
|
||||||
|
def _move_playback_sources(session: Session, *, keep_id: uuid.UUID, drop_id: uuid.UUID) -> None:
|
||||||
|
"""Unique (origin, page_url) jest GLOBALNY → drop i keep nie mogą współdzielić
|
||||||
|
tego samego źródła, więc samo przepięcie movie_id nie grozi kolizją."""
|
||||||
|
session.execute(
|
||||||
|
update(MoviePlaybackSource)
|
||||||
|
.where(MoviePlaybackSource.movie_id == drop_id)
|
||||||
|
.values(movie_id=keep_id)
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _coalesce_canonical_fields(keep: Movie, drop: Movie) -> None:
|
||||||
|
"""Wypełnij braki w `keep` polami z `drop` (bez nadpisywania ustawionych)."""
|
||||||
|
if not keep.description and drop.description:
|
||||||
|
keep.description = drop.description
|
||||||
|
if not keep.duration_sec and drop.duration_sec:
|
||||||
|
keep.duration_sec = drop.duration_sec
|
||||||
|
if not keep.director and drop.director:
|
||||||
|
keep.director = drop.director
|
||||||
|
if not keep.country and drop.country:
|
||||||
|
keep.country = drop.country
|
||||||
|
if not keep.release_date and drop.release_date:
|
||||||
|
keep.release_date = drop.release_date
|
||||||
|
if not keep.release_year and drop.release_year:
|
||||||
|
keep.release_year = drop.release_year
|
||||||
|
if not keep.studio_id and drop.studio_id:
|
||||||
|
keep.studio_id = drop.studio_id
|
||||||
|
if not keep.poster_url and drop.poster_url:
|
||||||
|
keep.poster_url = drop.poster_url
|
||||||
|
if not keep.backdrop_url and drop.backdrop_url:
|
||||||
|
keep.backdrop_url = drop.backdrop_url
|
||||||
|
if keep.rating is None and drop.rating is not None:
|
||||||
|
keep.rating = drop.rating
|
||||||
|
if drop.title and len(drop.title) > len(keep.title or ""):
|
||||||
|
keep.title = drop.title
|
||||||
|
keep.title_normalized = drop.title_normalized
|
||||||
|
# created_at = najwcześniejsze "first seen" (NEW badge correctness, jak w scenach).
|
||||||
|
if drop.created_at and keep.created_at and drop.created_at < keep.created_at:
|
||||||
|
keep.created_at = drop.created_at
|
||||||
|
|
||||||
|
|
||||||
|
def _close_pending_candidates(
|
||||||
|
session: Session, *, movie_id: uuid.UUID, resolved_by: str | None
|
||||||
|
) -> None:
|
||||||
|
session.execute(
|
||||||
|
update(MergeCandidate)
|
||||||
|
.where(
|
||||||
|
MergeCandidate.kind == MergeKind.movie,
|
||||||
|
MergeCandidate.status == MergeStatus.pending,
|
||||||
|
or_(MergeCandidate.left_id == movie_id, MergeCandidate.right_id == movie_id),
|
||||||
|
)
|
||||||
|
.values(
|
||||||
|
status=MergeStatus.rejected,
|
||||||
|
resolved_at=datetime.now(UTC),
|
||||||
|
resolved_by=resolved_by or "auto:movie_dropped",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
@ -349,6 +349,25 @@ def _job_performer_continuous(refresh_after_days: int) -> None:
|
||||||
log.exception("[scheduler] performer-continuous failed")
|
log.exception("[scheduler] performer-continuous failed")
|
||||||
|
|
||||||
|
|
||||||
|
def _job_tpdb_movie_enrich(batch: int) -> None:
|
||||||
|
"""Wzbogaca filmy metadanymi z TPDB /movies: obsada, kategorie (tagi), studio,
|
||||||
|
reżyser + kanoniczny UUID do dedupu mirrorów. Batch najświeższych grywalnych
|
||||||
|
filmów bez obsady/studia per run. paradisehill (primary) prawie nie ma obsady,
|
||||||
|
więc to domyka największą lukę. Patrz app/enrich/tpdb_movies.py."""
|
||||||
|
log.info("[scheduler] tpdb-movie-enrich starting (batch=%d)", batch)
|
||||||
|
try:
|
||||||
|
from app.db import session_scope
|
||||||
|
from app.enrich.tpdb_movies import run_tpdb_movie_enrich
|
||||||
|
|
||||||
|
_run_with_timeout(
|
||||||
|
lambda: run_tpdb_movie_enrich(session_scope, limit=batch),
|
||||||
|
label="tpdb-movie-enrich",
|
||||||
|
)
|
||||||
|
log.info("[scheduler] tpdb-movie-enrich done")
|
||||||
|
except Exception:
|
||||||
|
log.exception("[scheduler] tpdb-movie-enrich failed")
|
||||||
|
|
||||||
|
|
||||||
def build_scheduler(cfg: dict[str, Any]) -> BlockingScheduler:
|
def build_scheduler(cfg: dict[str, Any]) -> BlockingScheduler:
|
||||||
"""Buduje scheduler na podstawie cfg dictu.
|
"""Buduje scheduler na podstawie cfg dictu.
|
||||||
|
|
||||||
|
|
@ -426,6 +445,21 @@ def build_scheduler(cfg: dict[str, Any]) -> BlockingScheduler:
|
||||||
)
|
)
|
||||||
log.info("scheduler: bulk-dedup performers every %dh", cfg["bulk_dedup_hours"])
|
log.info("scheduler: bulk-dedup performers every %dh", cfg["bulk_dedup_hours"])
|
||||||
|
|
||||||
|
if cfg.get("tpdb_movie_enrich_hours"):
|
||||||
|
_tme_batch = cfg.get("tpdb_movie_enrich_batch") or 200
|
||||||
|
sched.add_job(
|
||||||
|
lambda: _job_tpdb_movie_enrich(_tme_batch),
|
||||||
|
IntervalTrigger(hours=cfg["tpdb_movie_enrich_hours"], start_date=INTERVAL_ANCHOR),
|
||||||
|
id="tpdb_movie_enrich",
|
||||||
|
replace_existing=True,
|
||||||
|
max_instances=1,
|
||||||
|
coalesce=True,
|
||||||
|
)
|
||||||
|
log.info(
|
||||||
|
"scheduler: tpdb-movie-enrich every %dh (batch=%d)",
|
||||||
|
cfg["tpdb_movie_enrich_hours"], _tme_batch,
|
||||||
|
)
|
||||||
|
|
||||||
if cfg.get("thumb_dedup_hours"):
|
if cfg.get("thumb_dedup_hours"):
|
||||||
sched.add_job(
|
sched.add_job(
|
||||||
_job_thumb_asset_dedup,
|
_job_thumb_asset_dedup,
|
||||||
|
|
|
||||||
|
|
@ -207,6 +207,10 @@ def run_forever() -> int:
|
||||||
# Bulk-dedup performers — safety net dla duplikatów które resolver
|
# Bulk-dedup performers — safety net dla duplikatów które resolver
|
||||||
# pominął (np. freshporno scen przed fixem release_date). Run 12h.
|
# pominął (np. freshporno scen przed fixem release_date). Run 12h.
|
||||||
"bulk_dedup_hours": getattr(settings, "sched_bulk_dedup_hours", 12) or None,
|
"bulk_dedup_hours": getattr(settings, "sched_bulk_dedup_hours", 12) or None,
|
||||||
|
# TPDB movie enrichment — obsada/kategorie/studio/reżyser + UUID do dedupu.
|
||||||
|
# paradisehill prawie nie ma obsady; batch grywalnych filmów bez obsady/studia.
|
||||||
|
"tpdb_movie_enrich_hours": getattr(settings, "sched_tpdb_movie_enrich_hours", 6) or None,
|
||||||
|
"tpdb_movie_enrich_batch": getattr(settings, "tpdb_movie_enrich_batch", 200),
|
||||||
# Thumb-asset dedup — hdporn.gg/fullmovies.xxx same-video-różne-tytuły (reports
|
# Thumb-asset dedup — hdporn.gg/fullmovies.xxx same-video-różne-tytuły (reports
|
||||||
# 205b17d9/5a2944cb). bulk_dedup tego nie łapie; dupy odrastają przy re-ingeście.
|
# 205b17d9/5a2944cb). bulk_dedup tego nie łapie; dupy odrastają przy re-ingeście.
|
||||||
"thumb_dedup_hours": getattr(settings, "sched_thumb_dedup_hours", 12) or None,
|
"thumb_dedup_hours": getattr(settings, "sched_thumb_dedup_hours", 12) or None,
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue