From 034ddd644c376bb45339011579e8943eb2abe3db Mon Sep 17 00:00:00 2001 From: goon-foss Date: Mon, 27 Jul 2026 15:04:10 +0200 Subject: [PATCH] feat(dedup): provenance przypisan performerow do scen (scene_performers.source_id) scene_performers nie zapisywalo, KTO przypisal performera do sceny. Przez to nie dalo sie odroznic obsady z kanonu od performera doklejonego przez tube-search, ktory dopasowuje wynik po JEDNYM tokenie zapytania. Stad przypadek "Rapture": 30 z 72 jej przypisan bylo falszywych (film "The Rapture" z Mimi Rogers, literowka "hymen raptured", tytuly scen) i trzeba bylo je przegladac recznie, bo nie bylo sygnalu. Kolumna trzyma NAJBARDZIEJ wiarygodne zrodlo, jakie przypisalo. Regula idzie tylko w gore: NULL wypelnia cokolwiek, kanon nadpisuje scraper, scraper NIGDY nie zamazuje kanonu. Inaczej dowolny pozniejszy tube-search zatarlby informacje, ze obsade potwierdzilo TPDB. To samo przy merge_scenes: przy kolizji przenosimy lepsza provenance na keepera, zeby scalanie nie degradowalo potwierdzonej obsady. Backfill wypelnil 4 147 128 wierszy (96,3 procent) wnioskowaniem, ktore jest pewne: scena z refami z DOKLADNIE jednego zrodla ma cala obsade stamtad. Sceny wielozrodlowe (67 tys.) zostaja z NULL-em, bo tam zgadywanie zepsuloby sens kolumny. Wsadowo, z commitem na paczke: pojedynczy UPDATE na 4 mln wierszy trzymalby locki na scene_performers i zablokowal ingest. Rozklad: tube-scraper 3 219 677, tpdb 747 540, stashdb 180 355. Co-Authored-By: Claude Opus 5 --- .../20260727_0026_scene_performer_source.py | 49 +++++++++++ app/models/scene.py | 10 +++ app/resolve/scene_merge.py | 19 ++++ app/resolve/scene_resolver.py | 31 +++++-- scripts/backfill_scene_performer_source.py | 88 +++++++++++++++++++ 5 files changed, 191 insertions(+), 6 deletions(-) create mode 100644 alembic/versions/20260727_0026_scene_performer_source.py create mode 100644 scripts/backfill_scene_performer_source.py diff --git a/alembic/versions/20260727_0026_scene_performer_source.py b/alembic/versions/20260727_0026_scene_performer_source.py new file mode 100644 index 0000000..0f317b0 --- /dev/null +++ b/alembic/versions/20260727_0026_scene_performer_source.py @@ -0,0 +1,49 @@ +"""provenance przypisań performerów do scen + +Revision ID: 0026_scene_performer_source +Revises: 0025_phash_blacklist +Create Date: 2026-07-27 + +`scene_performers` nie zapisywało, KTO przypisał performera do sceny. Przez to nie da +się odróżnić obsady pochodzącej z kanonu od performera doklejonego przez tube-search +(`_search_base` dopasowuje wynik po JEDNYM tokenie zapytania, więc „Rapture" łapała +każdą scenę ze słowem „rapture" — 30 z 72 jej przypisań było fałszywych i trzeba je +było przejrzeć ręcznie, bo nie było na to żadnego sygnału). + +Kolumna nullable, bo dla 4,3 mln istniejących wierszy źródła nie da się odtworzyć — +wypełnia się od teraz, przy kolejnych ingestach. +""" +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op +from sqlalchemy.dialects import postgresql + +revision: str = "0026_scene_performer_source" +down_revision: str | None = "0025_phash_blacklist" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.add_column( + "scene_performers", + sa.Column("source_id", postgresql.UUID(as_uuid=True), nullable=True), + ) + op.create_foreign_key( + "fk_scene_performers_source_id_sources", + "scene_performers", + "sources", + ["source_id"], + ["id"], + ondelete="SET NULL", + ) + op.create_index("ix_scene_performers_source_id", "scene_performers", ["source_id"]) + + +def downgrade() -> None: + op.drop_index("ix_scene_performers_source_id", table_name="scene_performers") + op.drop_constraint( + "fk_scene_performers_source_id_sources", "scene_performers", type_="foreignkey" + ) + op.drop_column("scene_performers", "source_id") diff --git a/app/models/scene.py b/app/models/scene.py index f09d435..a277bdd 100644 --- a/app/models/scene.py +++ b/app/models/scene.py @@ -102,6 +102,16 @@ class ScenePerformer(Base): role: Mapped[str | None] = mapped_column(String(64)) position: Mapped[int | None] = mapped_column(Integer) as_alias: Mapped[str | None] = mapped_column(String(256)) + #: Skąd wzięło się TO przypisanie. Bez tego nie da się odróżnić obsady z kanonu od + #: performera doklejonego przez luźny tube-search (ten dopasowuje wynik po JEDNYM + #: tokenie zapytania, więc „Rapture" łapała każdą scenę ze słowem „rapture" — + #: 30 z 72 jej scen było fałszywych i trzeba je było przeglądać ręcznie). + #: Trzymamy NAJBARDZIEJ wiarygodne źródło, jakie przypisało: kanoniczne + #: (tpdb/stashdb) nadpisuje scraperowe, nigdy odwrotnie. NULL = wiersz sprzed + #: migracji 0026, źródła nie da się odtworzyć wstecz. + source_id: Mapped[uuid.UUID | None] = mapped_column( + UUID(as_uuid=True), ForeignKey("sources.id", ondelete="SET NULL"), index=True + ) class SceneTag(Base): diff --git a/app/resolve/scene_merge.py b/app/resolve/scene_merge.py index 167f21a..f5651bf 100644 --- a/app/resolve/scene_merge.py +++ b/app/resolve/scene_merge.py @@ -29,10 +29,18 @@ from app.models.scene import ( ScenePerformer, SceneTag, ) +from app.models.source import Source, SourceKind log = logging.getLogger(__name__) +def _is_canonical(session: Session, source_id: uuid.UUID | None) -> bool: + if source_id is None: + return False + src = session.get(Source, source_id) + return src is not None and src.kind in (SourceKind.tpdb, SourceKind.stashdb) + + class MergeError(Exception): pass @@ -145,6 +153,17 @@ def _move_performers(session: Session, *, keep_id: uuid.UUID, drop_id: uuid.UUID if clash is not None: if link.as_alias and not clash.as_alias: clash.as_alias = link.as_alias + # Kasujemy wiersz drop-a, więc jego provenance przepadłaby razem z nim. + # Przenosimy ją, gdy jest lepsza: kanon > scraper > NULL. Inaczej merge + # potrafiłby zdegradować obsadę potwierdzoną przez TPDB do „nie wiadomo skąd". + if clash.source_id is None and link.source_id is not None: + clash.source_id = link.source_id + elif ( + link.source_id is not None + and _is_canonical(session, link.source_id) + and not _is_canonical(session, clash.source_id) + ): + clash.source_id = link.source_id session.delete(link) else: link.scene_id = keep_id diff --git a/app/resolve/scene_resolver.py b/app/resolve/scene_resolver.py index 388308c..60f086b 100644 --- a/app/resolve/scene_resolver.py +++ b/app/resolve/scene_resolver.py @@ -327,7 +327,7 @@ def resolve_scene( if decision == "auto": _update_scene_fields(best_scene, norm, studio_id=studio_id, source_kind=source_kind, session=session) _attach_external_ref(session, scene_id=best_scene.id, source_id=source_id, norm=norm) - _sync_performers(session, scene_id=best_scene.id, resolved=resolved_performers) + _sync_performers(session, scene_id=best_scene.id, resolved=resolved_performers, source_id=source_id) _sync_tags(session, scene_id=best_scene.id, norm=norm, source_id=source_id) _sync_fingerprints( session, scene_id=best_scene.id, norm=norm, source_id=source_id @@ -349,7 +349,7 @@ def resolve_scene( if decision == "review": new_scene = _create_canonical(session, norm=norm, studio_id=studio_id, backfill=backfill) _attach_external_ref(session, scene_id=new_scene.id, source_id=source_id, norm=norm) - _sync_performers(session, scene_id=new_scene.id, resolved=resolved_performers) + _sync_performers(session, scene_id=new_scene.id, resolved=resolved_performers, source_id=source_id) _sync_tags(session, scene_id=new_scene.id, norm=norm, source_id=source_id) _sync_fingerprints( session, scene_id=new_scene.id, norm=norm, source_id=source_id @@ -376,7 +376,7 @@ def resolve_scene( # Brak żadnego sensownego dopasowania → nowa kanoniczna new_scene = _create_canonical(session, norm=norm, studio_id=studio_id, backfill=backfill) _attach_external_ref(session, scene_id=new_scene.id, source_id=source_id, norm=norm) - _sync_performers(session, scene_id=new_scene.id, resolved=resolved_performers) + _sync_performers(session, scene_id=new_scene.id, resolved=resolved_performers, source_id=source_id) _sync_tags(session, scene_id=new_scene.id, norm=norm, source_id=source_id) _sync_fingerprints(session, scene_id=new_scene.id, norm=norm, source_id=source_id) _sync_playback_sources(session, scene_id=new_scene.id, norm=norm) @@ -542,17 +542,25 @@ def _sync_attached_entities( for p_norm in norm.performers: performer = resolve_performer(session, norm=p_norm, source_id=source_id) resolved.append((performer.id, p_norm.as_alias_in_scene)) - _sync_performers(session, scene_id=scene.id, resolved=resolved) + _sync_performers(session, scene_id=scene.id, resolved=resolved, source_id=source_id) _sync_tags(session, scene_id=scene.id, norm=norm, source_id=source_id) _sync_fingerprints(session, scene_id=scene.id, norm=norm, source_id=source_id) _sync_playback_sources(session, scene_id=scene.id, norm=norm) +def _is_canonical_source(session: Session, source_id: uuid.UUID | None) -> bool: + if source_id is None: + return False + src = session.get(Source, source_id) + return src is not None and src.kind in (SourceKind.tpdb, SourceKind.stashdb) + + def _sync_performers( session: Session, *, scene_id: uuid.UUID, resolved: list[tuple[uuid.UUID, str | None]], + source_id: uuid.UUID | None = None, ) -> None: # Deduplikuj — dwa różne aliasy tej samej osoby (np. "Aj Applegate" + "AJ Applegate") # przejdą przez resolve_performer zwracając ten sam Performer.id. Bez tej dedup @@ -581,10 +589,21 @@ def _sync_performers( performer_id=performer_id, position=position, as_alias=as_alias, + source_id=source_id, ) ) - elif as_alias and not existing.as_alias: - existing.as_alias = as_alias + else: + if as_alias and not existing.as_alias: + existing.as_alias = as_alias + # Provenance idzie tylko w górę: NULL wypełnia cokolwiek, kanon nadpisuje + # scraper, ale scraper NIGDY nie zamazuje kanonu. Inaczej dowolny późniejszy + # tube-search zatarłby informację, że obsadę potwierdziło TPDB. + if source_id is not None and existing.source_id != source_id: + if existing.source_id is None or ( + _is_canonical_source(session, source_id) + and not _is_canonical_source(session, existing.source_id) + ): + existing.source_id = source_id def _sync_tags( diff --git a/scripts/backfill_scene_performer_source.py b/scripts/backfill_scene_performer_source.py new file mode 100644 index 0000000..bdb0afd --- /dev/null +++ b/scripts/backfill_scene_performer_source.py @@ -0,0 +1,88 @@ +"""Backfill `scene_performers.source_id` dla wierszy sprzed migracji 0026. + +Wnioskowanie jest pewne tylko dla scen, które mają `scene_external_refs` z DOKŁADNIE +jednego źródła: skoro tylko ono kiedykolwiek dotknęło tę scenę, to cała jej obsada +stamtąd pochodzi. Takich scen jest 3,39 mln i pokrywają 4,15 mln przypisań, czyli ~96% +tabeli. Sceny wieloźródłowe (67 tys.) zostawiamy z NULL-em — tam nie da się powiedzieć, +które źródło przypisało którego performera, a zgadywanie zepsułoby cały sens kolumny. + +Granularność jest na poziomie ŹRÓDŁA, nie pojedynczego tuba: wszystkie direct-scrapery +dzielą jeden rekord `sources` (`SCRAPER_SOURCE_NAME`). To wystarcza do pytania, które +było nie do zadania wcześniej: „czy tę obsadę potwierdził kanon, czy dokleił ją +tube-search". + +Wsadowo i z osobną transakcją na paczkę — pojedynczy UPDATE na 4 mln wierszy trzymałby +locki na `scene_performers` przez cały czas i zablokował ingest (zdarzyło się przy +migracji 0026: porzucone zapytanie analityczne wstrzymało ALTER TABLE, a za nim ustawiła +się kolejka). + +Uruchomienie (kontener worker): + python -m scripts.backfill_scene_performer_source # dry-run + python -m scripts.backfill_scene_performer_source --yes # wykonaj +""" +from __future__ import annotations + +import argparse +import logging + +from sqlalchemy import text + +from app.db import session_scope + +log = logging.getLogger(__name__) + +_BATCH = 50_000 + +_COUNT_SQL = """ +SELECT count(*) FROM scene_performers sp +WHERE sp.source_id IS NULL + AND (SELECT count(DISTINCT r.source_id) FROM scene_external_refs r + WHERE r.scene_id = sp.scene_id) = 1 +""" + +# ctid-based batching: bez sztucznego klucza na (scene_id, performer_id) to najtańszy +# sposób na „weź N dowolnych jeszcze niewypełnionych". +_UPDATE_SQL = """ +UPDATE scene_performers sp +SET source_id = ( + SELECT (array_agg(DISTINCT r.source_id))[1] + FROM scene_external_refs r WHERE r.scene_id = sp.scene_id + ) +WHERE sp.ctid IN ( + SELECT sp2.ctid FROM scene_performers sp2 + WHERE sp2.source_id IS NULL + AND (SELECT count(DISTINCT r.source_id) FROM scene_external_refs r + WHERE r.scene_id = sp2.scene_id) = 1 + LIMIT :batch +) +""" + + +def main() -> None: + logging.basicConfig(level=logging.INFO, format="%(message)s") + ap = argparse.ArgumentParser(description=__doc__) + ap.add_argument("--yes", action="store_true", help="wykonaj (bez tego dry-run)") + ap.add_argument("--batch", type=int, default=_BATCH) + args = ap.parse_args() + + with session_scope() as s: + total = int(s.execute(text(_COUNT_SQL)).scalar_one()) + log.info("do wypełnienia: %d przypisań (sceny jednoźródłowe)", total) + if not args.yes: + log.info("(dry-run — uruchom z --yes)") + return + + done = 0 + while True: + with session_scope() as s: + n = int(s.execute(text(_UPDATE_SQL), {"batch": args.batch}).rowcount or 0) + s.commit() + if n == 0: + break + done += n + log.info(" wypełnione %d / %d", done, total) + log.info("GOTOWE: %d wierszy", done) + + +if __name__ == "__main__": + main()