From 238d03d0c6507044d9a1e153e2dc192803034436 Mon Sep 17 00:00:00 2001 From: jtrzupek Date: Thu, 2 Jul 2026 11:36:36 +0200 Subject: [PATCH] feat(movies): TPDB movie enrichment + dedup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- app/config.py | 10 ++ app/connectors/tpdb.py | 89 ++++++++++++ app/enrich/__init__.py | 1 + app/enrich/tpdb_movies.py | 273 +++++++++++++++++++++++++++++++++++++ app/resolve/movie_merge.py | 190 ++++++++++++++++++++++++++ app/scheduler/jobs.py | 34 +++++ app/scheduler/worker.py | 4 + 7 files changed, 601 insertions(+) create mode 100644 app/enrich/__init__.py create mode 100644 app/enrich/tpdb_movies.py create mode 100644 app/resolve/movie_merge.py diff --git a/app/config.py b/app/config.py index b8afd91..97d7ea1 100644 --- a/app/config.py +++ b/app/config.py @@ -136,6 +136,16 @@ class Settings(BaseSettings): ingest_watchdog_search_max_age_hours: int = Field( 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 # tags/performers/studios (hot-path /tags|/performers|/studios|/favorites czyta # gotową kolumnę zamiast agregować 6.3M scene_tags per-request). 3h cadence — diff --git a/app/connectors/tpdb.py b/app/connectors/tpdb.py index 5a80877..0412a5d 100644 --- a/app/connectors/tpdb.py +++ b/app/connectors/tpdb.py @@ -49,6 +49,7 @@ def _is_retryable_http_error(exc: BaseException) -> bool: from app.config import get_settings from app.connectors.base import ( BaseConnector, + RawMovie, RawPerformer, RawScene, RawStudio, @@ -195,6 +196,37 @@ class TPDBConnector(BaseConnector): first = data[0] 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= → 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/ → 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( self, 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 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, + ) diff --git a/app/enrich/__init__.py b/app/enrich/__init__.py new file mode 100644 index 0000000..88ea5d7 --- /dev/null +++ b/app/enrich/__init__.py @@ -0,0 +1 @@ +"""Post-ingest enrichment of canonical entities from authoritative metadata sources.""" diff --git a/app/enrich/tpdb_movies.py b/app/enrich/tpdb_movies.py new file mode 100644 index 0000000..a7d6454 --- /dev/null +++ b/app/enrich/tpdb_movies.py @@ -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=, 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 diff --git a/app/resolve/movie_merge.py b/app/resolve/movie_merge.py new file mode 100644 index 0000000..13b59a3 --- /dev/null +++ b/app/resolve/movie_merge.py @@ -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", + ) + ) diff --git a/app/scheduler/jobs.py b/app/scheduler/jobs.py index a8458e3..8aa6a69 100644 --- a/app/scheduler/jobs.py +++ b/app/scheduler/jobs.py @@ -349,6 +349,25 @@ def _job_performer_continuous(refresh_after_days: int) -> None: 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: """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"]) + 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"): sched.add_job( _job_thumb_asset_dedup, diff --git a/app/scheduler/worker.py b/app/scheduler/worker.py index f25bf93..2223a15 100644 --- a/app/scheduler/worker.py +++ b/app/scheduler/worker.py @@ -207,6 +207,10 @@ def run_forever() -> int: # Bulk-dedup performers — safety net dla duplikatów które resolver # pominął (np. freshporno scen przed fixem release_date). Run 12h. "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 # 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,