fix(tags): collapse spelling-variant duplicate tags + prevent regrowth

Ten sam tag z różnych tubów tworzył osobne wpisy (Blowjob/Blow Job, Big Tits/
Bigtits, deep-throat/deepthroat, creampie-1/creampie, mojibake CJK 'Japanese'…)
bo resolve_tag brał slug/slugify(name) verbatim. ~11k nadmiarowych tagów.

- resolve_tag: przed utworzeniem nowego tagu szuka istniejącego kanonicznego po
  alnum-kluczu nazwy (_resolve_by_altkey) i reużywa go → warianty nie regenerują.
- Indeks funkcyjny ix_tags_name_altkey (migracja 0028) pod ten lookup.
- scripts/merge_altkey_tags.py: jednorazowy backfill istniejących duplikatów
  (kanoniczny = max scene_count, scene_tags/movie_tags/blacklisted_tags przepięte
  z dedupem na PK, dropy skasowane, taxonomy_counts odświeżone).
This commit is contained in:
goon-foss 2026-07-21 15:59:33 +02:00
parent f97369a764
commit 80ac621344
3 changed files with 195 additions and 1 deletions

View file

@ -0,0 +1,32 @@
"""tags: functional index on alnum-normalized name (variant collapse)
Revision ID: 0028_tags_name_altkey_index
Revises: 0027_playback_event_error_detail
Create Date: 2026-07-21
`resolve_tag` zwija warianty pisowni tego samego tagu (case/spacja/myślnik/sklejenie
+ mojibake) do istniejącego kanonicznego po alnum-kluczu nazwy. Lookup potrzebuje
indeksu funkcyjnego, inaczej seq-scan po ~180k tagach przy każdym miss-slug.
Idempotentne (IF NOT EXISTS): prod dostaje indeks ręcznym CREATE INDEX CONCURRENTLY
(bez locka) zanim/gdy ta migracja powstaje; guard chroni przed re-runem.
"""
from collections.abc import Sequence
from alembic import op
revision: str = "0028_tags_name_altkey_index"
down_revision: str | None = "0027_playback_event_error_detail"
branch_labels: str | Sequence[str] | None = None
depends_on: str | Sequence[str] | None = None
def upgrade() -> None:
op.execute(
"CREATE INDEX IF NOT EXISTS ix_tags_name_altkey "
"ON tags (regexp_replace(lower(btrim(name)), '[^a-z0-9]', '', 'g'))"
)
def downgrade() -> None:
op.execute("DROP INDEX IF EXISTS ix_tags_name_altkey")

View file

@ -1,9 +1,10 @@
"""Resolver tagów. Tag identyfikuje slug (case-insensitive)."""
from __future__ import annotations
import re
import uuid
from sqlalchemy import select
from sqlalchemy import func, select
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.orm import Session
@ -11,6 +12,35 @@ from app.models.tag import Tag
from app.normalize.scenes import NormalizedTag
from app.normalize.text import slugify
_ALTKEY_RE = re.compile(r"[^a-z0-9]")
def _altkey(value: str) -> str:
"""Alnum-klucz nazwy tagu: lower + usunięcie wszystkiego poza [a-z0-9]. Zwija
warianty pisowni tego samego pojęcia (Blowjob/Blow Job/Bigtits, mojibake CJK
czysty rdzeń). Musi być zgodny z regexp w scripts/merge_altkey_tags.py + indeksem
ix_tags_name_altkey."""
return _ALTKEY_RE.sub("", value.strip().lower())
def _resolve_by_altkey(session: Session, name: str) -> Tag | None:
"""Znajdź istniejący kanoniczny tag o tym samym alnum-kluczu nazwy (największy
scene_count wygrywa). Zapobiega regeneracji duplikatów: nowy wariant "Deep Throat"
trafia do istniejącego "deepthroat" zamiast tworzyć osobny tag. Używa indeksu
funkcyjnego ix_tags_name_altkey."""
key = _altkey(name)
if not key:
return None
return (
session.execute(
select(Tag)
.where(func.regexp_replace(func.lower(func.btrim(Tag.name)), "[^a-z0-9]", "", "g") == key)
.order_by(Tag.scene_count.desc())
)
.scalars()
.first()
)
def _canonical_dup2_slug(session: Session, slug: str) -> str:
"""Kanonizuje numbered-duplicate slug `<base>2` → `<base>`.
@ -42,6 +72,12 @@ def resolve_tag(session: Session, *, norm: NormalizedTag) -> Tag | None:
return tag
name = (norm.name or "").strip() or slug.replace("-", " ").title()
# Prewencja regeneracji duplikatów: zanim utworzymy nowy tag, sprawdź czy istnieje
# kanoniczny o tym samym alnum-kluczu nazwy (wariant pisowni). Jeśli tak, reużyj go.
# (One-time backfill istniejących dup: scripts/merge_altkey_tags.py.)
existing = _resolve_by_altkey(session, name)
if existing is not None:
return existing
if len(name) > 120:
name = name[:120]
# Concurrent insert race: worker scraper + API enrich_tags_from_tube oba

View file

@ -0,0 +1,126 @@
"""Bulk-merge wariantów tego samego tagu różniących się tylko wielkością liter,
spacją, myślnikiem albo sklejeniem słów.
Kontekst: `resolve_tag` bierze slug ze źródła albo `slugify(name)` verbatim, więc
te same kategorie z różnych tubów tworzą OSOBNE tagi:
"Blowjob"/"Blow Job", "Big Tits"/"Bigtits", "Deepthroat"/"Deep Throat",
"Doggystyle"/"Doggy Style", "POV"x3. Przy scene_count>=30 to ~376 nadmiarowych tagów.
Klucz scalania = alnum-normalizacja nazwy: lower + usunięcie WSZYSTKICH nie-alfanum
znaków. Łączy TYLKO warianty pisowni (case/spacja/myślnik/sklejenie), NIE łączy
różnych pojęć (inny klucz) ani liczby poj./mn. (nie ścinamy 's'). Kanoniczny w
klastrze = tag o największym scene_count. Reszta scalona: scene_tags + movie_tags +
blacklisted_tags przepisane (dedup na PK), dropy skasowane, scene_count odświeżony.
Prewencja regeneracji: `app/resolve/tag_resolver.py` (_canonical_altkey_slug) +
indeks funkcyjny `ix_tags_altkey`.
Użycie:
python scripts/merge_altkey_tags.py [--dry-run] [--min-count N]
--min-count: scalaj tylko klastry, gdzie kanoniczny ma >= N scen (domyślnie 0 = wszystkie).
"""
from __future__ import annotations
import argparse
import logging
from sqlalchemy import text
from app.db import session_scope
log = logging.getLogger("merge_altkey_tags")
# drop→keep: w każdym klastrze o wspólnym alnum-kluczu kanoniczny = max scene_count
# (tie-break: id), reszta to dropy. Pomijamy pusty klucz (nazwy bez alfanum).
_DUP_MAP_SQL = """
WITH norm AS (
SELECT id, slug, name, scene_count,
regexp_replace(lower(btrim(name)), '[^a-z0-9]', '', 'g') AS k
FROM tags
),
ranked AS (
SELECT id, slug, scene_count, k,
first_value(id) OVER (PARTITION BY k ORDER BY scene_count DESC, id) AS keep_id
FROM norm
WHERE k <> ''
)
SELECT r.id AS drop_id, r.slug AS drop_slug, r.scene_count AS drop_cnt,
b.id AS keep_id, b.slug AS keep_slug, b.scene_count AS keep_cnt
FROM ranked r
JOIN tags b ON b.id = r.keep_id
WHERE r.id <> r.keep_id
AND b.scene_count >= :min_count
ORDER BY b.scene_count DESC
"""
def main() -> None:
ap = argparse.ArgumentParser()
ap.add_argument("--dry-run", action="store_true")
ap.add_argument("--min-count", type=int, default=0)
args = ap.parse_args()
logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s")
with session_scope() as s:
pairs = list(s.execute(text(_DUP_MAP_SQL), {"min_count": args.min_count}))
log.info("found %d drop tags across clusters (min_count=%d)", len(pairs), args.min_count)
for p in pairs[:80]:
log.info(" %-34s (%6d) -> %-30s (%6d)", p.drop_slug, p.drop_cnt, p.keep_slug, p.keep_cnt)
if len(pairs) > 80:
log.info(" ... (%d more)", len(pairs) - 80)
if not pairs:
return
s.execute(
text(
"CREATE TEMP TABLE _dup_map ON COMMIT DROP AS " + _DUP_MAP_SQL
),
{"min_count": args.min_count},
)
if args.dry_run:
n = s.execute(text("SELECT count(*) FROM scene_tags st JOIN _dup_map m ON st.tag_id=m.drop_id")).scalar()
nm = s.execute(text("SELECT count(*) FROM movie_tags mt JOIN _dup_map m ON mt.tag_id=m.drop_id")).scalar()
log.info("DRY-RUN: would touch %d scene_tags + %d movie_tags, drop %d tags", n, nm, len(pairs))
s.rollback()
return
r1 = s.execute(text("""
UPDATE scene_tags st SET tag_id = m.keep_id
FROM _dup_map m
WHERE st.tag_id = m.drop_id
AND NOT EXISTS (SELECT 1 FROM scene_tags k
WHERE k.scene_id = st.scene_id AND k.tag_id = m.keep_id)
"""))
log.info("scene_tags migrated: %d", r1.rowcount)
r2 = s.execute(text("""
UPDATE movie_tags mt SET tag_id = m.keep_id
FROM _dup_map m
WHERE mt.tag_id = m.drop_id
AND NOT EXISTS (SELECT 1 FROM movie_tags k
WHERE k.movie_id = mt.movie_id AND k.tag_id = m.keep_id)
"""))
log.info("movie_tags migrated: %d", r2.rowcount)
# blacklisted_tags PK = (device_id, tag_id) → przenieś ban z dropa na kanoniczny per device.
r3 = s.execute(text("""
INSERT INTO blacklisted_tags (device_id, tag_id)
SELECT bt.device_id, m.keep_id
FROM blacklisted_tags bt JOIN _dup_map m ON bt.tag_id = m.drop_id
ON CONFLICT DO NOTHING
"""))
if r3.rowcount:
log.info("blacklist refs moved: %d", r3.rowcount)
rd = s.execute(text("DELETE FROM tags WHERE id IN (SELECT drop_id FROM _dup_map)"))
log.info("dup tags deleted: %d", rd.rowcount)
s.commit()
from app.scheduler.taxonomy_counts import refresh_taxonomy_counts
changed = refresh_taxonomy_counts()
log.info("taxonomy counts refreshed: %s", changed)
if __name__ == "__main__":
main()