"""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=, 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