kink-discovery / backend /catalog.py
Perplexed7675's picture
Sync from kink_cli (Docker Space)
817e34f verified
Raw
History Blame Contribute Delete
38.7 kB
"""Backend catalog service."""
from __future__ import annotations
import json
import re
import sqlite3
import threading
import time
from collections import defaultdict
from pathlib import Path
from typing import Any
from sqlmodel import Session, select
from backend.image_policy import clip_min_similarity_catalog, strict_clip_catalog_enabled
from backend.scenarios import (
play_excluded_from_surfacing,
scenario_title_fields,
strip_profile_bucket_text_prefix,
title_surface_as_bundle,
title_surface_as_profile_bucket_noise,
)
from backend.scrape_artifacts import clean_scraped_catalog_text, clean_scraped_notes
from models import (
Alias,
Asset,
Definition,
FetlifeKinkMeta,
FetlifePictureRef,
Kink,
KinkExample,
KinkLemma,
KinkScenarioParent,
SimilarityEdge,
)
def _lemma_name_hash(name: str) -> str:
"""Fingerprint a kink name so the lemma cache invalidates when the name changes."""
import hashlib
return hashlib.sha1(name.encode("utf-8")).hexdigest()
def _load_cached_lemma_signatures(
self,
candidates: dict[str, dict[str, Any]],
) -> tuple[dict[str, str], list[tuple[str, str, str]]]:
"""Return (cached_lemmas_by_kid, empty_writebacks_list).
A cached row is only returned when its ``name_hash`` matches the current candidate name —
stale rows are silently ignored (caller will recompute and overwrite).
"""
try:
with self._sqlite() as conn:
rows = conn.execute("SELECT kink_id, name_hash, signature FROM kinklemma").fetchall()
except sqlite3.OperationalError:
return {}, []
cached: dict[str, str] = {}
for row in rows:
kid = row["kink_id"]
if kid not in candidates:
continue
if row["name_hash"] != _lemma_name_hash(str(candidates[kid].get("name", "") or "")):
continue
cached[kid] = row["signature"] or ""
return cached, []
def _persist_cached_lemma_signatures(
self,
rows: list[tuple[str, str, str]],
) -> None:
if not rows:
return
try:
with self._sqlite() as conn:
conn.executemany(
"INSERT INTO kinklemma (kink_id, name_hash, signature) VALUES (?, ?, ?) "
"ON CONFLICT(kink_id) DO UPDATE SET name_hash=excluded.name_hash, signature=excluded.signature",
rows,
)
conn.commit()
except sqlite3.OperationalError:
return
_IMAGE_TOKEN_RE = re.compile(r"[a-z0-9]+")
_IMAGE_RELEVANCE_STOPWORDS = {"and", "the", "play", "sex", "in", "with"}
_PRODUCT_FLAG_KEYS = (
"starter_eligible",
"starter_tier",
"starter_reason",
"shared_eligible",
"prompt_eligible",
"canonical_priority",
"image_trust_state",
"detail_summary",
)
def _image_relevance_from_tokens(kink_tokens: set[str], text_tokens: set[str], reuse_count: int) -> dict[str, Any]:
score = 1.0
reasons: list[str] = []
if reuse_count > 1:
penalty = min(0.08 * (reuse_count - 1), 0.55)
score -= penalty
reasons.append(f"reused:{reuse_count}")
if text_tokens and kink_tokens and not (kink_tokens & text_tokens):
score -= 0.18
reasons.append("caption_mismatch")
if len(text_tokens) >= 18:
score -= 0.18
reasons.append("tag_dense")
if not text_tokens:
score -= 0.08
reasons.append("no_caption")
score = max(0.0, round(score, 3))
if reuse_count > 12:
return {"score": score, "reason": "attachment_reused_too_often", "trusted": False}
if score < 0.42:
return {"score": score, "reason": ",".join(reasons) or "low_relevance", "trusted": False}
return {"score": score, "reason": ",".join(reasons), "trusted": True}
def _starter_sort_key(item: dict[str, Any]) -> tuple[int, float, float, float, str]:
return (
int(item.get("starter_tier_rank", 0) or 0),
float(item.get("source_backed_popularity", 0.0) or 0.0),
float(item.get("starter_score", 0.0) or 0.0),
float(item.get("popularity", 0.0) or 0.0),
str(item.get("name", "")).lower(),
)
def _catalog_invalidation_token(self) -> tuple[Any, ...]:
"""Bust catalog cache when CLIP file or policy env changes (no restart required)."""
p = self.path.parent / "clip_asset_scores.json"
try:
st = p.stat()
mtime_ns = st.st_mtime_ns
except OSError:
mtime_ns = -1
return (
mtime_ns,
clip_min_similarity_catalog(),
strict_clip_catalog_enabled(),
)
def _load_clip_asset_scores_file(path: Path) -> dict[str, dict[str, Any]]:
"""Load optional offline CLIP verdicts (see scripts/clip_align_kink_assets.py).
Each entry may include ``accepted: false`` to drop an image at catalog build time.
With ``KINK_STRICT_CLIP_IMAGES=1`` and a non-empty file, URLs missing from this map are
dropped too; see ``backend/image_policy.py``.
"""
if not path.is_file():
return {}
try:
raw = json.loads(path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
return {}
if not isinstance(raw, dict):
return {}
out: dict[str, dict[str, Any]] = {}
for url, payload in raw.items():
if not isinstance(url, str) or not isinstance(payload, dict):
continue
out[url] = payload
return out
def _invalidate_catalog_cache(self, *, force: bool = False) -> None:
now = time.monotonic()
if not force and now - self._last_cache_invalidation_at < 30.0:
return
self._last_cache_invalidation_at = now
self._catalog_cache = None
self._catalog_token = None
def _build_catalog_cache(self) -> dict[str, Any]:
type_rows = self._content_type_map()
if not type_rows:
type_rows = self._sync_content_types()
with Session(self.engine) as session:
kinks = session.exec(select(Kink).order_by(Kink.name)).all()
aliases = session.exec(select(Alias)).all()
definitions = session.exec(select(Definition)).all()
examples = session.exec(select(KinkExample)).all()
assets = session.exec(select(Asset)).all()
picture_refs = session.exec(select(FetlifePictureRef)).all()
fetlife_meta_rows = session.exec(select(FetlifeKinkMeta)).all()
scenario_parent_rows = session.exec(select(KinkScenarioParent)).all()
scenario_bridge_rows = session.exec(
select(SimilarityEdge).where(SimilarityEdge.similarity_type == "scenario_bridge")
).all()
alias_map: dict[str, list[str]] = defaultdict(list)
definition_map: dict[str, list[dict[str, str]]] = defaultdict(list)
example_map: dict[str, list[dict[str, str]]] = defaultdict(list)
asset_map: dict[str, list[dict[str, Any]]] = defaultdict(list)
raw_asset_count_by_kink: dict[str, int] = defaultdict(int)
inbound_counts: dict[str, int] = defaultdict(int)
outbound_counts: dict[str, int] = defaultdict(int)
fetlife_meta_map = {row.kink_id: row for row in fetlife_meta_rows}
parents_by_scenario: dict[str, list[dict[str, Any]]] = defaultdict(list)
children_by_parent: dict[str, list[dict[str, Any]]] = defaultdict(list)
seen_scenario_parent_pairs: set[tuple[str, str]] = set()
for row in scenario_parent_rows:
pair = (row.scenario_kink_id, row.parent_kink_id)
seen_scenario_parent_pairs.add(pair)
payload = {
"scenario_kink_id": row.scenario_kink_id,
"parent_kink_id": row.parent_kink_id,
"score": float(row.score),
"method": row.method,
"updated_at": row.updated_at,
}
parents_by_scenario[row.scenario_kink_id].append(payload)
children_by_parent[row.parent_kink_id].append(payload)
for row in scenario_bridge_rows:
pair = (row.left_kink_id, row.right_kink_id)
if pair in seen_scenario_parent_pairs:
continue
payload = {
"scenario_kink_id": row.left_kink_id,
"parent_kink_id": row.right_kink_id,
"score": float(row.score),
"method": row.method,
"updated_at": "",
}
parents_by_scenario[row.left_kink_id].append(payload)
children_by_parent[row.right_kink_id].append(payload)
for rows in parents_by_scenario.values():
rows.sort(key=lambda item: (-float(item["score"]), str(item["parent_kink_id"])))
for rows in children_by_parent.values():
rows.sort(key=lambda item: (-float(item["score"]), str(item["scenario_kink_id"])))
with self._sqlite() as conn:
for row in conn.execute(
"SELECT right_kink_id, COUNT(*) AS n FROM similarityedge GROUP BY right_kink_id"
).fetchall():
inbound_counts[row["right_kink_id"]] = int(row["n"])
for row in conn.execute(
"SELECT left_kink_id, COUNT(*) AS n FROM similarityedge GROUP BY left_kink_id"
).fetchall():
outbound_counts[row["left_kink_id"]] = int(row["n"])
for row in aliases:
alias_map[row.kink_id].append(row.alias)
for row in definitions:
definition_map[row.kink_id].append(
{
"text": row.text,
"source_id": row.source_id,
"source_url": row.source_url,
"license": row.license,
}
)
for row in examples:
example_map[row.kink_id].append({"text": row.text, "kind": row.kind})
fallback_definition_by_name: dict[str, str] = {}
for kink in kinks:
normalized_name = self._normalized_play_name(kink.name)
if not normalized_name:
continue
candidate_definition = strip_profile_bucket_text_prefix(kink.short_definition.strip())
candidate_definition = clean_scraped_catalog_text(kink.name, candidate_definition)
if not candidate_definition:
defs = definition_map.get(kink.id, [])
if defs:
candidate_definition = strip_profile_bucket_text_prefix(defs[0]["text"].strip())
candidate_definition = clean_scraped_catalog_text(kink.name, candidate_definition)
if not candidate_definition:
continue
incumbent = fallback_definition_by_name.get(normalized_name, "")
if not incumbent or len(candidate_definition) < len(incumbent):
fallback_definition_by_name[normalized_name] = candidate_definition
picture_ref_by_key: dict[tuple[str, str], FetlifePictureRef] = {}
reuse_counts: dict[str, int] = {}
for row in picture_refs:
picture_ref_by_key[(row.kink_id, row.attachment_id)] = row
reuse_counts.setdefault(row.attachment_id, 0)
for row in picture_refs:
reuse_counts[row.attachment_id] += 1
kink_image_tokens_by_id = {
kink.id: {
token
for token in _IMAGE_TOKEN_RE.findall(kink.name.lower())
if token not in _IMAGE_RELEVANCE_STOPWORDS
}
for kink in kinks
}
picture_text_tokens_by_key = {
key: set(
_IMAGE_TOKEN_RE.findall(
" ".join(
part
for part in (row.picture_title, row.caption_text)
if part
).lower()
)
)
for key, row in picture_ref_by_key.items()
}
clip_by_url = _load_clip_asset_scores_file(self.path.parent / "clip_asset_scores.json")
clip_catalog_min = clip_min_similarity_catalog()
# Prototype: KINK_STRICT_CLIP_IMAGES=1 + non-empty clip file → require CLIP row and accepted==true; tighten FL ref path.
enforce_clip_row = strict_clip_catalog_enabled() and bool(clip_by_url)
filtered_image_reasons_by_kink: dict[str, dict[str, int]] = defaultdict(lambda: defaultdict(int))
image_quality_by_kink: dict[str, dict[str, float]] = defaultdict(lambda: {"score_total": 0.0, "count": 0.0})
for row in assets:
raw_asset_count_by_kink[row.kink_id] += 1
attachment_id = ""
if row.source_id == "fetlife_fetish_pages":
attachment_id = Path(row.asset_url).stem
picture_ref = picture_ref_by_key.get((row.kink_id, attachment_id)) if attachment_id else None
reuse_count = reuse_counts.get(attachment_id, 1) if attachment_id else 1
if picture_ref:
relevance = _image_relevance_from_tokens(
kink_image_tokens_by_id[row.kink_id],
picture_text_tokens_by_key[(row.kink_id, attachment_id)],
reuse_count,
)
elif row.source_id == "fetlife_fetish_pages" and strict_clip_catalog_enabled():
relevance = {
"score": 0.28,
"reason": "fetlife_picture_ref_missing_strict",
"trusted": False,
}
else:
relevance = {"score": 0.72, "reason": "", "trusted": True}
if not relevance["trusted"]:
filtered_image_reasons_by_kink[row.kink_id][str(relevance["reason"])] += 1
continue
clip_meta = clip_by_url.get(row.asset_url)
clip_sim = None
if enforce_clip_row:
if clip_meta is None:
filtered_image_reasons_by_kink[row.kink_id]["clip_missing_row"] += 1
continue
if clip_meta.get("accepted") is not True:
filtered_image_reasons_by_kink[row.kink_id]["clip_not_accepted_strict"] += 1
continue
if clip_meta:
kid = str(clip_meta.get("kink_id", "") or "")
if kid and kid != row.kink_id:
filtered_image_reasons_by_kink[row.kink_id]["clip_kink_mismatch"] += 1
continue
try:
clip_sim = float(clip_meta.get("similarity"))
except (TypeError, ValueError):
clip_sim = None
if clip_sim is not None and clip_sim < clip_catalog_min:
filtered_image_reasons_by_kink[row.kink_id]["clip_below_catalog_min_similarity"] += 1
continue
if not enforce_clip_row and clip_meta.get("accepted") is False:
filtered_image_reasons_by_kink[row.kink_id]["clip_rejected_offline"] += 1
continue
image_quality_by_kink[row.kink_id]["score_total"] += float(relevance["score"])
image_quality_by_kink[row.kink_id]["count"] += 1.0
entry = {
"asset_url": row.asset_url,
"license": row.license,
"creator": picture_ref.creator if picture_ref and picture_ref.creator else row.creator,
"caption_text": picture_ref.caption_text if picture_ref else "",
"picture_page_url": picture_ref.picture_page_url if picture_ref else "",
"tag_spam_score": max(0, reuse_count - 1),
"attachment_reuse_count": reuse_count,
"image_relevance_score": relevance["score"],
"image_filter_reason": relevance["reason"],
"click_to_reveal": row.click_to_reveal,
"is_illustration": row.is_illustration,
"is_explicit": row.is_explicit,
}
if clip_sim is not None:
entry["clip_similarity"] = clip_sim
if clip_meta is not None and "accepted" in clip_meta:
entry["clip_accepted"] = clip_meta["accepted"]
asset_map[row.kink_id].append(entry)
for items in asset_map.values():
items.sort(
key=lambda item: (
item.get("tag_spam_score", 0),
0 if item.get("caption_text") else 1,
-float(item["clip_similarity"]) if item.get("clip_similarity") is not None else 0.0,
item.get("asset_url", ""),
)
)
detail_by_id: dict[str, dict[str, Any]] = {}
summary_by_id: dict[str, dict[str, Any]] = {}
search_docs: list[dict[str, Any]] = []
for kink in kinks:
defs = definition_map.get(kink.id, [])
raw_notes = kink.notes.strip()
display_notes = clean_scraped_notes(kink.id, kink.name, raw_notes)
definition = kink.short_definition.strip() or (defs[0]["text"].strip() if defs else "")
definition = strip_profile_bucket_text_prefix(definition)
definition = clean_scraped_catalog_text(kink.name, definition)
if not definition and kink.cluster == "fetlife_fetish":
definition = strip_profile_bucket_text_prefix(
fallback_definition_by_name.get(self._normalized_play_name(kink.name), "")
)
definition = clean_scraped_catalog_text(kink.name, definition)
if not definition:
from backend.curated_definitions import CURATED_DEFINITIONS
curated = CURATED_DEFINITIONS.get(kink.id, "").strip()
if curated:
definition = curated
examples_for_kink = example_map.get(kink.id, [])
example_texts = [
text
for row in examples_for_kink
if row["kind"] == "example"
for text in [clean_scraped_catalog_text(kink.name, row["text"])]
if text
]
visual_cues = [
text
for row in examples_for_kink
if row["kind"] == "visual_cue"
for text in [clean_scraped_catalog_text(kink.name, row["text"])]
if text
]
summary = self._display_summary(definition, display_notes, example_texts, kink.name)
explicit_popularity = self._extract_popularity_from_notes(raw_notes)
source_backed_popularity = fetlife_meta_map.get(kink.id).popularity if kink.id in fetlife_meta_map else 0.0
popularity = source_backed_popularity or explicit_popularity or (
inbound_counts.get(kink.id, 0) * 8
+ outbound_counts.get(kink.id, 0) * 4
+ len(example_texts)
+ len(asset_map.get(kink.id, [])) * 3
+ len(defs) * 2
)
base_payload = {
"id": kink.id,
"name": kink.name,
"cluster": kink.cluster,
"definition": definition,
"summary": summary,
"notes": display_notes,
}
type_row = type_rows.get(kink.id)
content_kind = type_row["content_kind"] if type_row else self._content_kind(base_payload)
type_evidence = type_row["evidence"] if type_row else ""
scenario_links = parents_by_scenario.get(kink.id, [])
scenario_parent_ids = [row["parent_kink_id"] for row in scenario_links]
is_scenario = bool(scenario_links)
scenario_child_count = len(children_by_parent.get(kink.id, []))
if content_kind == "play":
scenario_title_score, title_surface_flag = scenario_title_fields(kink.name)
else:
scenario_title_score, title_surface_flag = 0.0, False
title_bundle_flag = title_surface_as_bundle(kink.name)
profile_bucket_noise_flag = title_surface_as_profile_bucket_noise(kink.name)
from backend.direction_shapes import lookup_direction_shape
shape_info = lookup_direction_shape(kink.id, kink.name)
direction_shape = shape_info.get("shape", "action")
direction_role_a = shape_info.get("role_a", "")
direction_role_b = shape_info.get("role_b", "")
detail = {
"id": kink.id,
"name": kink.name,
"cluster": kink.cluster,
"content_kind": content_kind,
"type_evidence": type_evidence,
"definition": definition,
"summary": summary,
"notes": display_notes,
"risk_level": kink.risk_level,
"is_extreme": kink.is_extreme,
"popularity": popularity,
"source_backed_popularity": source_backed_popularity,
"fetlife_fetish_id": fetlife_meta_map.get(kink.id).fetish_id if kink.id in fetlife_meta_map else "",
"has_images": bool(asset_map.get(kink.id)),
"has_real_images": bool(fetlife_meta_map.get(kink.id).has_real_images) if kink.id in fetlife_meta_map else False,
"raw_asset_count": int(raw_asset_count_by_kink.get(kink.id, 0)),
"filtered_asset_count": len(asset_map.get(kink.id, [])),
"representative_image_count": len(asset_map.get(kink.id, [])),
"image_relevance_score": (
round(image_quality_by_kink[kink.id]["score_total"] / image_quality_by_kink[kink.id]["count"], 3)
if image_quality_by_kink[kink.id]["count"]
else 0.0
),
"image_filter_reasons": dict(filtered_image_reasons_by_kink.get(kink.id, {})),
"aliases": alias_map.get(kink.id, []),
"definitions": defs,
"examples": example_texts,
"visual_cues": visual_cues,
"assets": asset_map.get(kink.id, []),
"similar": [],
"similar_count": int(fetlife_meta_map.get(kink.id).similar_count) if kink.id in fetlife_meta_map else 0,
"is_scenario": is_scenario,
"scenario_parent_ids": scenario_parent_ids,
"scenario_parent_links": scenario_links,
"scenario_child_count": scenario_child_count,
"scenario_title_score": scenario_title_score,
"title_surface_as_scenario": title_surface_flag,
"title_surface_as_bundle": title_bundle_flag,
"title_surface_as_profile_bucket_noise": profile_bucket_noise_flag,
"direction_shape": direction_shape,
"direction_role_a": direction_role_a,
"direction_role_b": direction_role_b,
}
detail.update(self._derived_product_flags(detail))
if is_scenario or title_surface_flag:
detail.update(
{
"starter_eligible": False,
"starter_tier": "scenario",
"shared_eligible": False,
"prompt_eligible": False,
}
)
detail["starter_tier_rank"] = self._starter_tier_rank(detail)
detail["starter_score"] = self._starter_score(detail)
summary_payload = {
"id": kink.id,
"name": kink.name,
"cluster": kink.cluster,
"content_kind": content_kind,
"type_evidence": type_evidence,
"summary": summary,
"definition": definition,
"notes": display_notes,
"has_images": bool(asset_map.get(kink.id)),
"has_real_images": bool(fetlife_meta_map.get(kink.id).has_real_images) if kink.id in fetlife_meta_map else False,
"raw_asset_count": int(raw_asset_count_by_kink.get(kink.id, 0)),
"filtered_asset_count": len(asset_map.get(kink.id, [])),
"representative_image_count": len(asset_map.get(kink.id, [])),
"image_relevance_score": (
round(image_quality_by_kink[kink.id]["score_total"] / image_quality_by_kink[kink.id]["count"], 3)
if image_quality_by_kink[kink.id]["count"]
else 0.0
),
"popularity": popularity,
"source_backed_popularity": source_backed_popularity,
"similar_count": int(fetlife_meta_map.get(kink.id).similar_count) if kink.id in fetlife_meta_map else 0,
"is_scenario": is_scenario,
"scenario_parent_ids": scenario_parent_ids,
"scenario_child_count": scenario_child_count,
"scenario_title_score": scenario_title_score,
"title_surface_as_scenario": title_surface_flag,
"title_surface_as_bundle": title_bundle_flag,
"title_surface_as_profile_bucket_noise": profile_bucket_noise_flag,
"direction_shape": direction_shape,
"direction_role_a": direction_role_a,
"direction_role_b": direction_role_b,
}
for key in _PRODUCT_FLAG_KEYS:
summary_payload[key] = detail[key]
summary_payload["starter_tier_rank"] = detail["starter_tier_rank"]
summary_payload["starter_score"] = detail["starter_score"]
surface_excluded = play_excluded_from_surfacing(summary_payload)
detail["surface_excluded"] = surface_excluded
summary_payload["surface_excluded"] = surface_excluded
detail_by_id[kink.id] = detail
summary_by_id[kink.id] = summary_payload
search_docs.append(
{
"id": kink.id,
"name": kink.name,
"cluster": kink.cluster,
"content_kind": content_kind,
"type_evidence": type_evidence,
"aliases": alias_map.get(kink.id, []),
"summary": summary,
"definition": definition,
"examples": example_texts[:8],
"visual_cues": visual_cues[:8],
"popularity": popularity,
"source_backed_popularity": source_backed_popularity,
"is_scenario": is_scenario,
"scenario_parent_ids": scenario_parent_ids,
"scenario_child_count": scenario_child_count,
"scenario_title_score": scenario_title_score,
"title_surface_as_scenario": title_surface_flag,
"title_surface_as_bundle": title_bundle_flag,
"title_surface_as_profile_bucket_noise": profile_bucket_noise_flag,
"starter_eligible": summary_payload["starter_eligible"],
"starter_tier": summary_payload["starter_tier"],
"prompt_eligible": summary_payload["prompt_eligible"],
"shared_eligible": summary_payload["shared_eligible"],
"image_trust_state": summary_payload["image_trust_state"],
"detail_summary": summary_payload["detail_summary"],
"haystack": " ".join(
[
kink.name,
*alias_map.get(kink.id, []),
definition,
summary,
display_notes,
*example_texts[:8],
*visual_cues[:8],
]
).lower(),
}
)
self._mark_runtime_duplicates(detail_by_id, summary_by_id, search_docs)
self._maybe_rebuild_fts_index(search_docs)
play_exact_sets: defaultdict[str, set[str]] = defaultdict(set)
for doc in search_docs:
if doc["content_kind"] != "play":
continue
play_exact_sets[doc["name"].lower()].add(doc["id"])
for alias in doc["aliases"]:
play_exact_sets[str(alias).lower()].add(doc["id"])
play_exact_ids_by_lower = {k: sorted(v) for k, v in play_exact_sets.items()}
play_summaries = [
item
for item in summary_by_id.values()
if item["content_kind"] == "play" and not item["surface_excluded"]
]
play_summaries.sort(key=lambda item: (-item["popularity"], item["name"].lower()))
play_details = [detail_by_id[item["id"]] for item in play_summaries]
starter_summaries = [item for item in play_summaries if item["starter_eligible"]]
starter_summaries.sort(key=_starter_sort_key, reverse=True)
return {
"detail_by_id": detail_by_id,
"summary_by_id": summary_by_id,
"search_docs": search_docs,
"play_exact_ids_by_lower": play_exact_ids_by_lower,
"play_summaries": play_summaries,
"play_details": play_details,
"starter_summaries": starter_summaries,
}
def _mark_runtime_duplicates(
self,
detail_by_id: dict[str, dict[str, Any]],
summary_by_id: dict[str, dict[str, Any]],
search_docs: list[dict[str, Any]],
) -> None:
"""Group play kinks by (cluster, merge_signature) and (cluster, compact_form); within each
group with >= 2 members, keep the canonical (highest popularity, has-definition, shorter name)
and mark every other member ``surface_excluded=True``.
This dedupe is purely runtime — the SQLite catalog is not modified — so it lands on the next
deploy without needing a Hub dataset rebuild.
"""
from backend.kink_merge import compact_form, lemma_signature, merge_signature
candidates: dict[str, dict[str, Any]] = {}
for kid, summary in summary_by_id.items():
if summary.get("content_kind") != "play":
continue
if summary.get("surface_excluded"):
continue
candidates[kid] = summary
sig_groups: defaultdict[tuple[str, str], list[str]] = defaultdict(list)
compact_groups: defaultdict[tuple[str, str], list[str]] = defaultdict(list)
# Lemma-based grouping (spaCy ``en_core_web_sm`` lemmatization + POS-aware stopword filter).
# Catches morphological siblings ("Cuddling" / "Cuddles" / "Cuddly", "After sex cuddles" /
# "Cuddling after sex") that ``merge_signature``'s sorted-token-set misses. F1 0.80 on the
# hand-labeled set vs F1 0.29 for ``merge_signature`` alone — see scripts/merge_dedup_eval.py.
lemma_groups: defaultdict[tuple[str, str], list[str]] = defaultdict(list)
cached_lemmas, lemma_writebacks = _load_cached_lemma_signatures(self, candidates)
for kid, summary in candidates.items():
cluster = str(summary.get("cluster", "") or "")
name = str(summary.get("name", "") or "")
sig = merge_signature(name)
if sig:
sig_groups[(cluster, sig)].append(kid)
cf = compact_form(name)
if len(cf) >= 8:
compact_groups[(cluster, cf)].append(kid)
if kid in cached_lemmas:
lemma_sig = cached_lemmas[kid]
else:
lemma_sig = lemma_signature(name)
lemma_writebacks.append((kid, _lemma_name_hash(name), lemma_sig))
if lemma_sig:
lemma_groups[(cluster, lemma_sig)].append(kid)
_persist_cached_lemma_signatures(self, lemma_writebacks)
def _canonical_key(kid: str) -> tuple[int, int, float, int, str]:
"""Sort key — higher is better. Strict preference: has_definition first, then well-cased
Title-style name (so ``Anal Creampie`` wins over ``analcreampie``), then popularity, then
shorter name, then id stability."""
summary = candidates[kid]
detail = detail_by_id.get(kid, {})
name = str(summary.get("name", "") or "")
has_def = 1 if (detail.get("definition") or "").strip() else 0
# Heuristic: a name with at least one uppercase letter is better-cased than all-lowercase.
well_cased = 1 if any(ch.isupper() for ch in name) else 0
pop = float(summary.get("popularity", 0.0) or 0.0)
# Ties on the rest: prefer SHORTER name (negate length so larger key = shorter name).
return (has_def, well_cased, pop, -len(name), kid)
duplicate_ids: set[str] = set()
for groups in (sig_groups, compact_groups, lemma_groups):
for ids in groups.values():
if len(ids) < 2:
continue
canonical = max(ids, key=_canonical_key)
for kid in ids:
if kid != canonical:
duplicate_ids.add(kid)
if not duplicate_ids:
return
for kid in duplicate_ids:
if kid in detail_by_id:
detail_by_id[kid]["surface_excluded"] = True
detail_by_id[kid]["surface_excluded_reason"] = "runtime_duplicate"
if kid in summary_by_id:
summary_by_id[kid]["surface_excluded"] = True
summary_by_id[kid]["surface_excluded_reason"] = "runtime_duplicate"
duplicate_set = duplicate_ids
for doc in search_docs:
if doc.get("id") in duplicate_set:
doc["surface_excluded"] = True
def _rebuild_fts_index(self, search_docs: list[dict[str, Any]]) -> None:
with self._sqlite() as conn:
conn.execute("DELETE FROM play_fts")
conn.executemany(
"INSERT INTO play_fts (kink_id, content_kind, name, aliases, definition, examples) VALUES (?, ?, ?, ?, ?, ?)",
[
(
doc["id"],
doc["content_kind"],
doc["name"],
" ".join(doc["aliases"]),
doc["definition"],
" ".join(doc["examples"]),
)
for doc in search_docs
],
)
conn.commit()
def _maybe_rebuild_fts_index(self, search_docs: list[dict[str, Any]]) -> None:
try:
with self._sqlite() as conn:
row = conn.execute("SELECT COUNT(*) AS n FROM play_fts").fetchone()
existing_count = int(row["n"]) if row else 0
if existing_count == 0 or existing_count < max(1000, len(search_docs) // 2):
self._rebuild_fts_index(search_docs)
except sqlite3.OperationalError:
return
_CATALOG_PICKLE_VERSION = 1
def _catalog_pickle_path(self) -> Path:
return self.path.parent / f".{self.path.name}.catalog_v{_CATALOG_PICKLE_VERSION}.pkl"
def _catalog_pickle_key(self, token: tuple[Any, ...]) -> tuple[Any, ...]:
"""A content-derived key that survives side-table writes (e.g. kinklemma) but invalidates
when the underlying catalog changes. Uses SQL counters instead of file mtime because
_mark_runtime_duplicates writes back to kinklemma during the same build, bumping mtime.
Deliberately excludes ``token`` (CLIP env policy): the cached catalog dict's structure
does not depend on those env vars, and including them would cause spurious misses when
the build-host env differs from the runtime env (e.g. baking a pickle locally for HF).
"""
import sys as _sys
try:
with self._sqlite() as conn:
kink_count = int(conn.execute("SELECT COUNT(*) FROM kink").fetchone()[0])
edge_count = int(conn.execute("SELECT COUNT(*) FROM similarityedge").fetchone()[0])
max_kid = conn.execute("SELECT MAX(id) FROM kink").fetchone()[0] or ""
except sqlite3.OperationalError:
return ()
return (
_CATALOG_PICKLE_VERSION,
_sys.version_info[:2],
kink_count,
edge_count,
max_kid,
)
def _try_load_pickled_catalog(self, token: tuple[Any, ...]) -> dict[str, Any] | None:
path = _catalog_pickle_path(self)
if not path.is_file():
print(f"[kink_cli pickle] miss path_missing path={path}", flush=True)
return None
import pickle
try:
with path.open("rb") as f:
header = pickle.load(f)
expected = _catalog_pickle_key(self, token)
if header != expected:
print(f"[kink_cli pickle] miss key_mismatch path={path} disk={header!r} expected={expected!r}", flush=True)
return None
cache = pickle.load(f)
print(f"[kink_cli pickle] hit path={path} size={path.stat().st_size}", flush=True)
return cache
except (OSError, pickle.UnpicklingError, EOFError, ValueError, AttributeError, ImportError) as exc:
print(f"[kink_cli pickle] miss exception path={path} exc={exc!r}", flush=True)
return None
def _save_pickled_catalog(self, cache: dict[str, Any], token: tuple[Any, ...]) -> None:
path = _catalog_pickle_path(self)
tmp = path.with_suffix(path.suffix + ".tmp")
import pickle
try:
path.parent.mkdir(parents=True, exist_ok=True)
with tmp.open("wb") as f:
pickle.dump(_catalog_pickle_key(self, token), f, protocol=pickle.HIGHEST_PROTOCOL)
pickle.dump(cache, f, protocol=pickle.HIGHEST_PROTOCOL)
tmp.replace(path)
print(f"[kink_cli pickle] saved path={path} size={path.stat().st_size}", flush=True)
except OSError as exc:
print(f"[kink_cli pickle] save_failed path={path} exc={exc!r}", flush=True)
try:
tmp.unlink()
except OSError:
pass
def _catalog(self) -> dict[str, Any]:
token = _catalog_invalidation_token(self)
if self._catalog_cache is not None and getattr(self, "_catalog_token", None) == token:
return self._catalog_cache
if hasattr(self, "_catalog_ready") and not self._catalog_ready.is_set():
self._catalog_ready.wait()
error = getattr(self, "_catalog_warm_error", None)
if error is not None:
raise RuntimeError("Catalog warm-up failed") from error
token = _catalog_invalidation_token(self)
if self._catalog_cache is not None and getattr(self, "_catalog_token", None) == token:
return self._catalog_cache
lock = getattr(self, "_catalog_build_lock", None)
if lock is None:
lock = threading.Lock()
self._catalog_build_lock = lock
with lock:
token = _catalog_invalidation_token(self)
if self._catalog_cache is not None and getattr(self, "_catalog_token", None) == token:
return self._catalog_cache
cached = _try_load_pickled_catalog(self, token)
if cached is not None:
self._catalog_cache = cached
self._catalog_token = token
return self._catalog_cache
self._catalog_cache = self._build_catalog_cache()
self._catalog_token = _catalog_invalidation_token(self)
_save_pickled_catalog(self, self._catalog_cache, self._catalog_token)
return self._catalog_cache
def catalog_is_ready(self) -> bool:
return self._catalog_cache is not None
def warm_up_catalog(self) -> None:
"""Build catalog cache in background thread on startup."""
if self._catalog_cache is not None:
return
thread = getattr(self, "_catalog_warm_thread", None)
if thread is not None and thread.is_alive():
return
self._catalog_ready = threading.Event()
self._catalog_warm_error = None
def _build():
try:
lock = getattr(self, "_catalog_build_lock", None)
if lock is None:
lock = threading.Lock()
self._catalog_build_lock = lock
with lock:
token = _catalog_invalidation_token(self)
if self._catalog_cache is None or getattr(self, "_catalog_token", None) != token:
self._catalog_cache = self._build_catalog_cache()
self._catalog_token = _catalog_invalidation_token(self)
except BaseException as exc:
self._catalog_warm_error = exc
finally:
self._catalog_ready.set()
thread = threading.Thread(target=_build, daemon=True)
self._catalog_warm_thread = thread
thread.start()
def refresh_cache(self) -> None:
self._invalidate_catalog_cache(force=True)
self._catalog()