#!/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())