Files
catalogue_backend/scripts/repair_catalog.py
2026-08-29 14:49:45 +05:30

262 lines
11 KiB
Python

#!/usr/bin/env python3
"""
Repair catalog rows already stored by a buggy ingestion run.
Fixing the pipeline only helps the NEXT upload. Rows written before the fix keep
their corrupted values, so this walks the brand tables and repairs them in
place:
* **name** - strips a trailing bare number that a mis-mapped `Quantity` /
`Case Pack` column contributed ("Coca-Cola 750ml 72").
* **category** - re-resolves a value that is a numeric id ("1", "2") or is
otherwise absent from the canonical taxonomy, using the same
`detect_category_from_text` the pipeline uses.
* **images** - clears `image_url` / `image_urls` that point at a *different*
product's image folder. That is the brand-sample bleed which
put Good Day photographs on Marie Gold and Milk Bikis. Cleared
rather than re-pointed, because there is nothing correct to
point them at; `--refetch-images` can then fill them back in.
Two things it deliberately never does:
* It never recomputes `image_id`. That column is the UNIQUE upsert key AND the
S3 folder name, so changing it would orphan the uploaded images and make the
next ingest insert a duplicate instead of updating the row. Only display
fields are rewritten.
* It never calls `upsert_brand_products(..., cleanup=True)`, which deletes
every row not in the batch it was handed.
Usage:
python -m scripts.repair_catalog --brands coca-cola,britannia # dry run
python -m scripts.repair_catalog --brands coca-cola,britannia --apply
python -m scripts.repair_catalog --all --apply --refetch-images
`--dry-run` is the default and `--apply` must be given explicitly: this rewrites
whatever database `backend/.env` points at, which may well be production. The
target host is printed on startup so it can be checked before committing.
"""
from __future__ import annotations
import argparse
import logging
import re
import sys
from pathlib import Path
from typing import Any, Dict, List, Optional, Tuple
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from app.infrastructure.settings import DB_HOST, DB_NAME
from app.services.category_registry import (
ALL_CATEGORIES,
detect_category_from_text,
)
from app.services.category_units import parse_unit
from app.services.vector_store import (
_connect,
_sanitize_name,
list_available_brands,
resolve_parent_brand,
)
logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s")
logger = logging.getLogger(__name__)
_CANONICAL = {c.strip().lower() for c in ALL_CATEGORIES}
# "General" is what the pipeline stores when detection genuinely found nothing.
# It is a real decision, not corruption, so it is left alone.
_ACCEPTED_CATEGORIES = _CANONICAL | {"general"}
# ---------------------------------------------------------------------------
# Row-level repairs (pure, and unit-testable without a database)
# ---------------------------------------------------------------------------
def repair_product_name(product_name: str) -> Optional[str]:
"""Strip a trailing bare number that is not a pack size.
Returns the corrected name, or None when nothing needed changing. Only a
trailing token is considered, and only one: the pipeline appended exactly
one size, and stripping deeper would start eating real product names
("Britannia 50-50", "Sprite 7Up 250ml").
"""
name = (product_name or "").strip()
if not name:
return None
head, _, last = name.rpartition(" ")
if not head:
return None
value, unit = parse_unit(last)
# A trailing token that parses as a number with NO unit is the bug. A real
# pack size ("750ml") has a unit; a word ("Pack") parses as (None, None).
if value is None or unit:
return None
# Guard against a name that is genuinely numeric at the end, e.g. "5 Star"
# reversed, or a variant number the store means ("Maggi 2"). Only strip when
# the remaining name still carries letters.
if not re.search(r"[a-zA-Z]", head):
return None
return head.strip()
def repair_category(category: Optional[str], product_name: str,
description: str = "") -> Optional[str]:
"""Re-resolve a category that is an id or is not a known category name."""
current = (category or "").strip()
if current.lower() in _ACCEPTED_CATEGORIES:
return None
detected = detect_category_from_text(f"{product_name} {description}".strip())
resolved = detected or "General"
return resolved if resolved != current else None
def _folder_of(url: str) -> str:
"""The `{image_id}` path segment of a stored image URL, if it has one."""
match = re.search(r"/brands/[^/]+/([^/]+)/", url or "")
return match.group(1) if match else ""
def repair_images(image_id: str, image_url: Optional[str],
image_urls: Optional[List[str]]) -> Optional[Tuple[None, list]]:
"""Detect images belonging to a different product and clear them.
Only URLs that carry an explicit `{image_id}` folder can be judged, so a
plain CDN link is left alone - it may well be correct and there is no
evidence either way.
"""
urls = list(image_urls or [])
candidates = [u for u in ([image_url] if image_url else []) + urls if u]
folders = {_folder_of(u) for u in candidates}
folders.discard("")
if not folders:
return None
if folders == {image_id}:
return None
# At least one URL is filed under another product's image_id.
return (None, [])
# ---------------------------------------------------------------------------
# Database walk
# ---------------------------------------------------------------------------
def repair_brand(conn, brand: str, *, apply: bool, refetch: bool) -> Dict[str, int]:
table = f"brand_{_sanitize_name(resolve_parent_brand(brand) or brand)}"
counts = {"scanned": 0, "name": 0, "category": 0, "images": 0, "refetched": 0}
with conn.cursor() as cur:
try:
cur.execute(
f"SELECT image_id, product_name, title, category, description, "
f"image_url, image_urls FROM {table}"
)
except Exception as exc: # noqa: BLE001 - an absent table is not fatal
logger.warning("Skipping %s: %s", table, exc)
return counts
rows = cur.fetchall()
for image_id, product_name, title, category, description, image_url, image_urls in rows:
counts["scanned"] += 1
updates: Dict[str, Any] = {}
fixed_name = repair_product_name(product_name)
if fixed_name:
updates["product_name"] = fixed_name
counts["name"] += 1
logger.info("%s: name %r -> %r", table, product_name, fixed_name)
effective_name = fixed_name or product_name
fixed_category = repair_category(category, effective_name, description or "")
if fixed_category:
updates["category"] = fixed_category
counts["category"] += 1
logger.info("%s: category %r -> %r (%s)", table, category, fixed_category, effective_name)
cleared = repair_images(image_id, image_url, image_urls)
if cleared:
updates["image_url"], updates["image_urls"] = cleared
counts["images"] += 1
logger.info("%s: clearing images filed under another product (%s)", table, image_id)
if refetch and (cleared or not (image_urls or image_url)):
found = _refetch_images(brand, effective_name)
if found:
updates["image_url"], updates["image_urls"] = found[0], found
counts["refetched"] += 1
logger.info("%s: refetched %d image(s) for %s", table, len(found), effective_name)
if updates and apply:
assignments = ", ".join(f"{col} = %s" for col in updates)
with conn.cursor() as cur:
cur.execute(
f"UPDATE {table} SET {assignments}, updated_at = CURRENT_TIMESTAMP "
f"WHERE image_id = %s",
list(updates.values()) + [image_id],
)
return counts
def _refetch_images(brand: str, product_name: str) -> List[str]:
"""Re-run the normal image path for one product. Network, hence opt-in."""
try:
from app.core.catalog_engine import catalog_engine
from app.services.image_search import find_all_image_urls
candidates = find_all_image_urls(product_name, brand=brand, max_results=24)
if not candidates:
return []
return list(catalog_engine._select_best_images(candidates, product_name, brand, max_images=10))
except Exception as exc: # noqa: BLE001 - an image is not worth failing the row
logger.warning("Image refetch failed for %r: %s", product_name, exc)
return []
def main() -> int:
parser = argparse.ArgumentParser(description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter)
group = parser.add_mutually_exclusive_group(required=True)
group.add_argument("--brands", help="comma-separated brands, e.g. coca-cola,britannia")
group.add_argument("--all", action="store_true", help="every brand table present")
parser.add_argument("--apply", action="store_true",
help="write the changes (default is a dry run)")
parser.add_argument("--refetch-images", action="store_true",
help="re-run image search for rows left without images (slow, network)")
args = parser.parse_args()
logger.info("Target database: %s/%s", DB_HOST, DB_NAME)
logger.info("Mode: %s", "APPLY - rows will be rewritten" if args.apply else "DRY RUN - no writes")
conn = _connect()
if not conn:
logger.error("Could not connect to PostgreSQL (USE_PGVECTOR off, or DB unreachable).")
return 1
brands = list_available_brands() if args.all else [
b.strip() for b in args.brands.split(",") if b.strip()
]
if not brands:
logger.error("No brands to process.")
return 1
logger.info("Brands: %s", ", ".join(brands))
totals = {"scanned": 0, "name": 0, "category": 0, "images": 0, "refetched": 0}
try:
for brand in brands:
counts = repair_brand(conn, brand, apply=args.apply, refetch=args.refetch_images)
for key, value in counts.items():
totals[key] += value
finally:
conn.close()
logger.info(
"Scanned %d row(s): %d name(s), %d category(ies), %d image set(s) to fix, %d refetched.",
totals["scanned"], totals["name"], totals["category"], totals["images"], totals["refetched"],
)
if not args.apply and any(totals[k] for k in ("name", "category", "images")):
logger.info("Dry run - nothing written. Re-run with --apply to commit these changes.")
return 0
if __name__ == "__main__":
raise SystemExit(main())