Valid Barcode Generation
This commit is contained in:
516
scripts/backfill_barcodes_from_off.py
Normal file
516
scripts/backfill_barcodes_from_off.py
Normal file
@@ -0,0 +1,516 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Backfill product barcodes from Open Food Facts into the brand tables and the
|
||||
seed catalog JSON files.
|
||||
|
||||
WHY
|
||||
611 catalog products, 588 of them with no barcode. OFF has an India
|
||||
catalogue for the food brands (amul 146, cadbury 123, tata 50, britannia
|
||||
233 products), keyed by a name that is usually recognisable but rarely
|
||||
identical to ours - OFF says "Britannia Marie Gold", we say "Britannia
|
||||
Marie Gold 250g". This script closes that gap by name similarity.
|
||||
|
||||
It fetches each brand's ENTIRE OFF catalogue in one or two requests and
|
||||
matches offline. The pre-existing per-product path
|
||||
(`sources/open_food_facts.py`, driven by the enrichment pipeline) searches
|
||||
OFF once per product against an endpoint capped at 10 requests/minute;
|
||||
this run costs about five requests in total and is fully re-runnable from
|
||||
the on-disk corpus cache with no further network traffic.
|
||||
|
||||
WHAT IT TOUCHES
|
||||
* brand_<slug>.barcode and .barcode_type - only where the current value is
|
||||
absent OR fails checksum validation. A genuine barcode (Tata has 23
|
||||
manually sourced ones) is never overwritten; a value that is not a
|
||||
barcode at all ('8900000000000.0', 38 rows) is repairable rather than
|
||||
permanently locked in. The UPDATE pins the row's own current value in its
|
||||
WHERE clause, so a concurrent real write is never lost.
|
||||
* backend/data/seed_catalogs/brand_catalog_<slug>.json - the nine barcode
|
||||
keys on matched products only, written in place. Every other key on
|
||||
every product, including `embedding`, `size`, `variant_key` and
|
||||
`gst_percent`, is preserved byte-for-byte.
|
||||
* backend/data/cache/off_brand_corpus/<slug>.json - the fetched OFF corpus.
|
||||
* The CSV review report.
|
||||
|
||||
WHAT IT DOES NOT TOUCH
|
||||
Nothing in the running app. `stage.py`, `service.py`, `pipeline.py`,
|
||||
`sources/open_food_facts.py`, `matching.py`, `validators.py` and
|
||||
`vector_store.py` are all consumed read-only; ENABLE_BARCODE_LOOKUP stays
|
||||
off and the ingestion pipeline behaves exactly as before.
|
||||
|
||||
In particular this does NOT call `brand_sync.export_brand_to_seed_file()`.
|
||||
That helper rebuilds each product from EXPORT_COLUMNS and
|
||||
`upsert_products_into_catalog_file` then REPLACES the whole product dict
|
||||
(brand_sync.py:232) after stripping the embedding (:221) - running it here
|
||||
would silently delete `size`, `variant_key`, `gst_percent`, `embedding` and
|
||||
the seven JSON-only barcode keys from every product in the file.
|
||||
|
||||
MATCHING RULES
|
||||
Products are grouped by their size-less `title`, matched once per group,
|
||||
and the winning barcode is written to every size variant in that group.
|
||||
That is a deliberate choice: a GTIN identifies exactly one pack, so
|
||||
"Amul Butter 90ml/200ml/500ml/1L" all carrying one real pack's barcode is
|
||||
knowingly approximate. Every such write is tagged
|
||||
barcode_source="openfoodfacts_bulk" and barcode_verified=false so these
|
||||
rows stay selectable for later correction or a clean bulk revert.
|
||||
|
||||
score >= --min-similarity (0.88) -> written
|
||||
score >= --review-min (0.70) -> CSV only, nothing written
|
||||
below -> dropped
|
||||
|
||||
Usage:
|
||||
|
||||
python -m scripts.backfill_barcodes_from_off # dry run, all brands
|
||||
python -m scripts.backfill_barcodes_from_off --brand amul # repeatable flag
|
||||
python -m scripts.backfill_barcodes_from_off --brand amul --apply
|
||||
python -m scripts.backfill_barcodes_from_off --refresh-cache --report out.csv
|
||||
python -m scripts.backfill_barcodes_from_off --apply --no-db # JSON only
|
||||
|
||||
`--dry-run` is the default and `--apply` must be explicit: backend/.env points
|
||||
at the PRODUCTION database, so an accidental run must not be able to write.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import csv
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
from collections import OrderedDict, defaultdict
|
||||
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.brand_sync import SEED_DIR, brand_slug, seed_catalog_paths
|
||||
from app.services.enrichment.barcode.models import BarcodeResult, LookupStatus
|
||||
from app.services.enrichment.barcode.sources.off_bulk import (
|
||||
brand_tokens,
|
||||
fetch_brand_corpus,
|
||||
normalize_for_match,
|
||||
score_candidates,
|
||||
)
|
||||
from app.services.enrichment.barcode.validators import validate_barcode
|
||||
from app.services.vector_store import _connect
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="%(message)s")
|
||||
logger = logging.getLogger("backfill_barcodes_off")
|
||||
|
||||
BARCODE_SOURCE = "openfoodfacts_bulk (search.openfoodfacts.org)"
|
||||
|
||||
# The nine keys BarcodeResult.as_product_fields() emits. Named here so the JSON
|
||||
# writer can prove it touches nothing else.
|
||||
BARCODE_KEYS = (
|
||||
"barcode", "barcode_type", "gtin", "ean13", "upc",
|
||||
"barcode_source", "barcode_verified", "barcode_lookup_status",
|
||||
"barcode_last_updated",
|
||||
)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Catalog reading
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def load_catalogs(only: Optional[List[str]] = None,
|
||||
include_archive: bool = False) -> "OrderedDict[Path, Dict[str, Any]]":
|
||||
"""Every active seed catalog, keyed by path. The JSON is the read source
|
||||
rather than the database so that a dry run needs no DB connection at all -
|
||||
which matters when the configured database is production.
|
||||
|
||||
`seed_catalog_paths` deliberately returns the archive as well (28 more
|
||||
brands, plus a stray-copy file), but those brands are not live and their
|
||||
tables may not exist, so they are excluded unless asked for.
|
||||
"""
|
||||
wanted = {b.lower().strip() for b in (only or [])}
|
||||
out: "OrderedDict[Path, Dict[str, Any]]" = OrderedDict()
|
||||
for path in sorted(seed_catalog_paths(SEED_DIR)):
|
||||
if not include_archive and path.parent != SEED_DIR:
|
||||
continue
|
||||
try:
|
||||
data = json.loads(path.read_text(encoding="utf-8-sig"))
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("Skipping unreadable catalog %s: %s", path.name, e)
|
||||
continue
|
||||
if not isinstance(data, dict):
|
||||
continue
|
||||
products = data.get("products")
|
||||
brand = data.get("brand")
|
||||
if not brand or not isinstance(products, list) or not products:
|
||||
continue
|
||||
if wanted and brand.lower().strip() not in wanted and brand_slug(brand) not in wanted:
|
||||
continue
|
||||
out[path] = data
|
||||
return out
|
||||
|
||||
|
||||
def product_size(product: Dict[str, Any]) -> str:
|
||||
"""The product's pack size, from whichever field carries it."""
|
||||
size = str(product.get("size") or "").strip()
|
||||
if size and size.lower() != "standard":
|
||||
return size
|
||||
variants = product.get("size_variants") or []
|
||||
if isinstance(variants, list) and variants:
|
||||
return str(variants[0]).strip()
|
||||
return ""
|
||||
|
||||
|
||||
def group_by_title(products: List[Dict[str, Any]]) -> "OrderedDict[str, List[Dict[str, Any]]]":
|
||||
"""Bucket barcode-less products by their size-less title.
|
||||
|
||||
Products that already carry a barcode are excluded entirely - they are
|
||||
neither re-looked-up nor overwritten.
|
||||
"""
|
||||
groups: "OrderedDict[str, List[Dict[str, Any]]]" = OrderedDict()
|
||||
for p in products:
|
||||
if str(p.get("barcode") or "").strip():
|
||||
continue
|
||||
if not str(p.get("image_id") or "").strip():
|
||||
continue
|
||||
title = str(p.get("title") or p.get("product_name") or "").strip()
|
||||
if not title:
|
||||
continue
|
||||
groups.setdefault(title, []).append(p)
|
||||
return groups
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Writers
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def build_fields(candidate) -> Dict[str, Any]:
|
||||
"""The nine catalog barcode fields for one accepted candidate, built
|
||||
through BarcodeResult so this script cannot drift from the shape the rest
|
||||
of the codebase writes."""
|
||||
result = BarcodeResult(
|
||||
barcode=candidate.barcode,
|
||||
barcode_type=candidate.barcode_type,
|
||||
gtin=candidate.barcode,
|
||||
ean13=candidate.barcode if len(candidate.barcode) == 13 else None,
|
||||
upc=candidate.source_upc,
|
||||
barcode_source=BARCODE_SOURCE,
|
||||
barcode_verified=False,
|
||||
barcode_lookup_status=LookupStatus.NOT_FOUND.value,
|
||||
match_confidence=candidate.score,
|
||||
)
|
||||
fields = result.as_product_fields()
|
||||
# Not VERIFIED: nothing here was confirmed against the physical pack, and
|
||||
# for a multi-variant group the code is right for at most one member.
|
||||
fields["barcode_lookup_status"] = "name_matched"
|
||||
return fields
|
||||
|
||||
|
||||
def write_json(path: Path, data: Dict[str, Any],
|
||||
updates: Dict[str, Dict[str, Any]]) -> int:
|
||||
"""Apply barcode fields to products in one catalog file, in place.
|
||||
|
||||
Only the nine barcode keys of matched products are assigned; the product
|
||||
dicts are mutated rather than rebuilt, so `embedding`, `size`,
|
||||
`variant_key`, `gst_percent` and anything else survive untouched. The write
|
||||
is atomic so a crash cannot leave a half-written catalog that fails to
|
||||
parse on the next boot.
|
||||
"""
|
||||
changed = 0
|
||||
for product in data.get("products") or []:
|
||||
fields = updates.get(str(product.get("image_id") or ""))
|
||||
if not fields:
|
||||
continue
|
||||
product.update(fields)
|
||||
changed += 1
|
||||
if not changed:
|
||||
return 0
|
||||
|
||||
payload = json.dumps(data, indent=2, ensure_ascii=False)
|
||||
tmp = path.with_suffix(".json.tmp")
|
||||
tmp.write_text(payload, encoding="utf-8")
|
||||
os.replace(tmp, path)
|
||||
return changed
|
||||
|
||||
|
||||
def write_db(brand: str, updates: Dict[str, Dict[str, Any]]) -> Tuple[int, str]:
|
||||
"""UPDATE barcode/barcode_type on a brand table, keyed on image_id.
|
||||
|
||||
Returns (rows_updated, note). A row is written only when its current
|
||||
barcode is absent OR fails checksum validation, so this is idempotent and
|
||||
can never clobber a genuine manually-sourced barcode.
|
||||
|
||||
Treating an INVALID barcode as writable matters: 41 rows across the brand
|
||||
tables hold values that are not barcodes at all - 38 of them the Excel
|
||||
float '8900000000000.0'. Guarding on "IS NULL OR = ''" alone would declare
|
||||
those rows already-done and leave them unrepairable forever, and because
|
||||
the seed JSON has no such value the two stores would silently diverge on
|
||||
exactly those rows. The UPDATE still pins the row's own current value in
|
||||
its WHERE clause, so a real barcode written concurrently is never lost.
|
||||
"""
|
||||
table = f"brand_{brand_slug(brand)}"
|
||||
conn = _connect()
|
||||
if conn is None:
|
||||
return 0, "no database connection (USE_PGVECTOR off or connect failed)"
|
||||
|
||||
updated = 0
|
||||
kept = 0
|
||||
try:
|
||||
with conn, conn.cursor() as cur:
|
||||
cur.execute(
|
||||
"SELECT column_name FROM information_schema.columns WHERE table_name = %s",
|
||||
(table,),
|
||||
)
|
||||
cols = {r[0] for r in cur.fetchall()}
|
||||
if not cols:
|
||||
return 0, f"table {table} does not exist"
|
||||
missing = {"barcode", "barcode_type"} - cols
|
||||
if missing:
|
||||
return 0, f"table {table} is missing column(s) {sorted(missing)}"
|
||||
|
||||
image_ids = list(updates)
|
||||
cur.execute(
|
||||
f"SELECT image_id, barcode FROM {table} WHERE image_id = ANY(%s)",
|
||||
(image_ids,),
|
||||
)
|
||||
current = dict(cur.fetchall())
|
||||
|
||||
for image_id, fields in updates.items():
|
||||
existing = current.get(image_id)
|
||||
if existing and validate_barcode(existing) is not None:
|
||||
kept += 1
|
||||
continue
|
||||
cur.execute(
|
||||
f"UPDATE {table} SET barcode = %s, barcode_type = %s, "
|
||||
f"updated_at = NOW() "
|
||||
f"WHERE image_id = %s AND (barcode IS NULL OR barcode = '' "
|
||||
f"OR barcode = %s)",
|
||||
(fields["barcode"], fields["barcode_type"], image_id, existing),
|
||||
)
|
||||
updated += cur.rowcount
|
||||
except Exception as e: # noqa: BLE001 - one brand failing must not abort the run
|
||||
return updated, f"error: {e}"
|
||||
finally:
|
||||
conn.close()
|
||||
note = f"table {table}"
|
||||
if kept:
|
||||
note += f", {kept} row(s) left alone (valid barcode already present)"
|
||||
return updated, note
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Main
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def sync_json_to_db(path: Path, data: Dict[str, Any]) -> Tuple[int, str]:
|
||||
"""Push barcodes the seed JSON already has into the database.
|
||||
|
||||
The matching path writes JSON and DB together, but it skips any product
|
||||
whose JSON barcode is already set - so if the two stores ever disagree, a
|
||||
rerun cannot heal them. They can disagree because a row may hold a value
|
||||
that is not a barcode ('8900000000000.0' appears 38 times), which the
|
||||
matcher happily replaced in JSON while the DB guard correctly refused to
|
||||
overwrite what looked like existing data.
|
||||
|
||||
This reconciles in the safe direction only: a DB barcode that passes
|
||||
checksum validation always wins and is left alone. Only absent or invalid
|
||||
DB values are replaced, and only from a JSON value that itself validates.
|
||||
"""
|
||||
updates: Dict[str, Dict[str, Any]] = {}
|
||||
for product in data.get("products") or []:
|
||||
image_id = str(product.get("image_id") or "").strip()
|
||||
code = validate_barcode(str(product.get("barcode") or "").strip())
|
||||
if not image_id or not code:
|
||||
continue
|
||||
updates[image_id] = {
|
||||
"barcode": code,
|
||||
"barcode_type": product.get("barcode_type") or "EAN-13",
|
||||
}
|
||||
if not updates:
|
||||
return 0, "no valid barcodes in JSON"
|
||||
return write_db(data["brand"], updates)
|
||||
|
||||
|
||||
def main() -> int:
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Backfill barcodes from Open Food Facts into brand tables and seed JSON.")
|
||||
parser.add_argument("--brand", action="append", default=None,
|
||||
help="Brand to process (repeatable). Default: every seed catalog.")
|
||||
parser.add_argument("--apply", action="store_true",
|
||||
help="Actually write. Without this nothing is modified.")
|
||||
parser.add_argument("--dry-run", action="store_true",
|
||||
help="Force a dry run even if --apply is passed.")
|
||||
parser.add_argument("--min-similarity", type=float, default=0.88,
|
||||
help="Auto-apply threshold (default 0.88).")
|
||||
parser.add_argument("--review-min", type=float, default=0.70,
|
||||
help="Report-only threshold (default 0.70).")
|
||||
parser.add_argument("--limit", type=int, default=None,
|
||||
help="Stop after this many title groups per brand.")
|
||||
parser.add_argument("--refresh-cache", action="store_true",
|
||||
help="Re-fetch each brand's OFF corpus instead of using the disk cache.")
|
||||
parser.add_argument("--report", type=str, default=None,
|
||||
help="CSV report path (default backend/data/barcode_review_<ts>.csv).")
|
||||
parser.add_argument("--no-db", action="store_true", help="Do not write to the database.")
|
||||
parser.add_argument("--no-json", action="store_true", help="Do not write to the seed JSON.")
|
||||
parser.add_argument("--sync-json-to-db", action="store_true",
|
||||
help="Skip matching; push barcodes the seed JSON already has "
|
||||
"into the DB wherever the DB value is absent or fails "
|
||||
"checksum validation. Reconciles the two stores.")
|
||||
parser.add_argument("--allow-foreign-prefix", action="store_true",
|
||||
help="Accept barcodes not issued in India (GS1 prefix 890). "
|
||||
"Off by default: the non-890 hits observed were the same "
|
||||
"product line sold in another market, with a different pack.")
|
||||
parser.add_argument("--include-archive", action="store_true",
|
||||
help="Also process the archived catalogs under seed_catalogs/archive/.")
|
||||
args = parser.parse_args()
|
||||
|
||||
apply = args.apply and not args.dry_run
|
||||
|
||||
logger.info("Database: %s / %s", DB_HOST, DB_NAME)
|
||||
logger.info("Mode: %s", "APPLY (writes)" if apply else "DRY RUN (no writes)")
|
||||
logger.info("Thresholds: apply >= %.2f, review >= %.2f",
|
||||
args.min_similarity, args.review_min)
|
||||
logger.info("")
|
||||
|
||||
catalogs = load_catalogs(args.brand, include_archive=args.include_archive)
|
||||
if not catalogs:
|
||||
logger.error("No matching seed catalogs found in %s", SEED_DIR)
|
||||
return 1
|
||||
|
||||
if args.sync_json_to_db:
|
||||
total = 0
|
||||
for path, data in catalogs.items():
|
||||
if not apply:
|
||||
logger.info("=== %s - dry run, would reconcile from %s", data["brand"], path.name)
|
||||
continue
|
||||
n, note = sync_json_to_db(path, data)
|
||||
total += n
|
||||
logger.info("=== %s - %d row(s) reconciled (%s)", data["brand"], n, note)
|
||||
logger.info("---------------------------------------------------------")
|
||||
logger.info("Rows reconciled: %d", total)
|
||||
if not apply:
|
||||
logger.info("DRY RUN - re-run with --apply to write.")
|
||||
return 0
|
||||
|
||||
report_rows: List[Dict[str, Any]] = []
|
||||
totals = defaultdict(int)
|
||||
|
||||
for path, data in catalogs.items():
|
||||
brand = data["brand"]
|
||||
products = data.get("products") or []
|
||||
groups = group_by_title(products)
|
||||
already = sum(1 for p in products if str(p.get("barcode") or "").strip())
|
||||
logger.info("=== %s (%s) - %d product(s), %d with a barcode already, "
|
||||
"%d title group(s) to match",
|
||||
brand, path.name, len(products), already, len(groups))
|
||||
|
||||
if not groups:
|
||||
logger.info(" nothing to do")
|
||||
logger.info("")
|
||||
continue
|
||||
|
||||
corpus = fetch_brand_corpus(brand, refresh=args.refresh_cache)
|
||||
if not corpus:
|
||||
logger.info(" OFF has no products for this brand - skipping")
|
||||
logger.info("")
|
||||
totals["brands_without_corpus"] += 1
|
||||
continue
|
||||
|
||||
tokens = brand_tokens(brand)
|
||||
used: Dict[str, Tuple[float, str]] = {} # barcode -> (score, title)
|
||||
updates: Dict[str, Dict[str, Any]] = {}
|
||||
applied_groups = 0
|
||||
|
||||
for i, (title, rows) in enumerate(groups.items()):
|
||||
if args.limit is not None and i >= args.limit:
|
||||
break
|
||||
sizes = [product_size(r) for r in rows]
|
||||
candidates = score_candidates(
|
||||
corpus, title, sizes, tokens, args.review_min,
|
||||
require_india_prefix=not args.allow_foreign_prefix)
|
||||
if not candidates:
|
||||
totals["no_candidate"] += 1
|
||||
continue
|
||||
|
||||
best = candidates[0]
|
||||
decision = "apply" if best.score >= args.min_similarity else "review"
|
||||
|
||||
# One OFF product cannot be the barcode for two different title
|
||||
# groups. The higher-scoring group keeps it; this one is demoted to
|
||||
# the report rather than silently duplicating a code.
|
||||
claim = used.get(best.barcode)
|
||||
if decision == "apply" and claim is not None:
|
||||
decision = "review"
|
||||
note = f"barcode already claimed by {claim[1]!r} (score {claim[0]:.3f})"
|
||||
else:
|
||||
note = ""
|
||||
if decision == "apply":
|
||||
used[best.barcode] = (best.score, title)
|
||||
|
||||
if decision == "apply":
|
||||
fields = build_fields(best)
|
||||
for row in rows:
|
||||
updates[str(row["image_id"])] = fields
|
||||
applied_groups += 1
|
||||
totals["rows_matched"] += len(rows)
|
||||
else:
|
||||
totals["needs_review"] += 1
|
||||
|
||||
for row in rows:
|
||||
report_rows.append({
|
||||
"brand": brand,
|
||||
"image_id": row.get("image_id"),
|
||||
"our_title": title,
|
||||
"our_product_name": row.get("product_name"),
|
||||
"our_size": product_size(row),
|
||||
"off_name": best.off_name,
|
||||
"off_qty": best.off_qty or "",
|
||||
"barcode": best.barcode,
|
||||
"barcode_type": best.barcode_type,
|
||||
"score": f"{best.score:.3f}",
|
||||
"size_bonus": best.size_bonus,
|
||||
"variants_in_group": len(rows),
|
||||
"decision": decision,
|
||||
"note": note,
|
||||
"runner_up": candidates[1].off_name if len(candidates) > 1 else "",
|
||||
})
|
||||
|
||||
logger.info(" %d group(s) matched at >= %.2f -> %d product row(s)",
|
||||
applied_groups, args.min_similarity, len(updates))
|
||||
|
||||
if updates and apply:
|
||||
if not args.no_json:
|
||||
n = write_json(path, data, updates)
|
||||
logger.info(" JSON: updated %d product(s) in %s", n, path.name)
|
||||
totals["json_rows"] += n
|
||||
if not args.no_db:
|
||||
n, note = write_db(brand, updates)
|
||||
logger.info(" DB: updated %d row(s) (%s)", n, note)
|
||||
totals["db_rows"] += n
|
||||
elif updates:
|
||||
logger.info(" (dry run - nothing written)")
|
||||
logger.info("")
|
||||
|
||||
report_path = Path(args.report) if args.report else (
|
||||
_BACKEND_DATA / f"barcode_review_{time.strftime('%Y%m%d_%H%M%S')}.csv")
|
||||
if report_rows:
|
||||
report_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
with report_path.open("w", newline="", encoding="utf-8") as fh:
|
||||
writer = csv.DictWriter(fh, fieldnames=list(report_rows[0].keys()))
|
||||
writer.writeheader()
|
||||
writer.writerows(report_rows)
|
||||
|
||||
logger.info("---------------------------------------------------------")
|
||||
logger.info("Product rows matched at >= %.2f : %d", args.min_similarity, totals["rows_matched"])
|
||||
logger.info("Title groups needing review : %d", totals["needs_review"])
|
||||
logger.info("Title groups with no candidate : %d", totals["no_candidate"])
|
||||
if apply:
|
||||
logger.info("JSON products written : %d", totals["json_rows"])
|
||||
logger.info("DB rows written : %d", totals["db_rows"])
|
||||
else:
|
||||
logger.info("DRY RUN - re-run with --apply to write.")
|
||||
if report_rows:
|
||||
logger.info("Report: %s (%d row(s))", report_path, len(report_rows))
|
||||
return 0
|
||||
|
||||
|
||||
_BACKEND_DATA = Path(__file__).resolve().parents[1] / "data"
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -1,103 +1,223 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Fix duplicate image_ids in seed catalog JSON files.
|
||||
Give each pack-size variant its own image_id in the seed catalogs.
|
||||
|
||||
Products sharing the same `title` (e.g. "Nestle Kitkat" with sizes 30g, 50g, 70g)
|
||||
previously all got the same `image_id` (generated from `title` only). Since the
|
||||
database uses `ON CONFLICT (image_id) DO UPDATE`, only one variant per title
|
||||
survived the upsert.
|
||||
WHY
|
||||
Six catalogs share one `image_id` across several genuinely different
|
||||
products. In brand_catalog_dabur.json, "Red Chyawanprash" 20g, 50g and
|
||||
100g all carry `red_chyawanprash_f122a959` - they differ in
|
||||
product_name, size, variant_key, product_sku, price_range and highlights,
|
||||
so they are three real products, not three copies of one.
|
||||
|
||||
This script re-generates `image_id` from `product_name` (or `title` + `size`)
|
||||
so every product variant gets a unique, deterministic image_id.
|
||||
`image_id` is `TEXT UNIQUE NOT NULL` and is the seeder's dedup key, so
|
||||
only ONE variant per group ever reaches the database; the rest are
|
||||
silently discarded on every seed. That is why brand_dabur holds 22 rows
|
||||
for a 54-product catalog. The ids look scraper-generated
|
||||
(`<slug>_<uuid4-fragment>`), which is the path that does not derive a
|
||||
per-variant id.
|
||||
|
||||
WHAT IT CHANGES
|
||||
Only the `image_id` field, only on the extra members of a duplicate group,
|
||||
only in backend/data/seed_catalogs/ (including archive/). No product is
|
||||
added, removed or reordered, and no other field is touched.
|
||||
|
||||
The variant that already owns the database row KEEPS its original id. That
|
||||
is deliberate: `image_id` is also the S3 path component
|
||||
(`daily/brands/{brand}/{image_id}/image_000.jpg` - s3_service.py:246), so
|
||||
renaming the resident variant would orphan its row and point its image URL
|
||||
at a prefix that holds no objects. The other variants have no row and no S3
|
||||
objects today, so giving them a new id cannot break anything that works.
|
||||
|
||||
New ids come from `variant_key`, which is already unique within every
|
||||
affected file and never collides with an existing image_id (both verified
|
||||
before writing).
|
||||
|
||||
WHAT IT DOES NOT DO
|
||||
It does not touch the database. The renamed variants gain rows only when
|
||||
someone deliberately reseeds; archived catalogs are not loaded at boot.
|
||||
|
||||
Usage:
|
||||
python scripts/fix_duplicate_image_ids.py
|
||||
|
||||
python -m scripts.fix_duplicate_image_ids # dry run
|
||||
python -m scripts.fix_duplicate_image_ids --apply
|
||||
python -m scripts.fix_duplicate_image_ids --file brand_catalog_cadbury.json
|
||||
|
||||
`--dry-run` is the default and `--apply` must be explicit.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
import uuid
|
||||
from collections import Counter, defaultdict
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s")
|
||||
logger = logging.getLogger(__name__)
|
||||
from app.infrastructure.settings import DB_HOST, DB_NAME
|
||||
from app.services.brand_sync import SEED_DIR, brand_slug, seed_catalog_paths
|
||||
from app.services.vector_store import _connect
|
||||
|
||||
SEED_DIR = Path(__file__).resolve().parents[1] / "data" / "seed_catalogs"
|
||||
logging.basicConfig(level=logging.INFO, format="%(message)s")
|
||||
logger = logging.getLogger("fix_duplicate_image_ids")
|
||||
|
||||
|
||||
def make_image_id(source: str) -> str:
|
||||
"""Deterministic image_id: sanitize source name + short hash suffix."""
|
||||
sanitized = source.lower()
|
||||
sanitized = ''.join(c if c.isalnum() else '_' for c in sanitized)
|
||||
sanitized = '_'.join(filter(None, sanitized.split('_')))
|
||||
if len(sanitized) > 50:
|
||||
sanitized = sanitized[:50]
|
||||
# Use a hash of the source string so the same source always produces the same ID
|
||||
short_hash = str(uuid.uuid5(uuid.NAMESPACE_DNS, source))[:8]
|
||||
return f"{sanitized}_{short_hash}"
|
||||
|
||||
|
||||
def fix_seed_file(path: Path) -> int:
|
||||
with open(path, encoding="utf-8-sig") as f:
|
||||
data = json.load(f)
|
||||
|
||||
products = data.get("products", [])
|
||||
fixed = 0
|
||||
seen_ids: set[str] = set()
|
||||
|
||||
def duplicate_groups(products: List[Dict[str, Any]]) -> Dict[str, List[Dict[str, Any]]]:
|
||||
by_id: Dict[str, List[Dict[str, Any]]] = defaultdict(list)
|
||||
for product in products:
|
||||
# Determine the source for the image_id: prefer product_name, fall back to title+size, then title
|
||||
source = product.get("product_name") or ""
|
||||
if not source:
|
||||
title = product.get("title", "")
|
||||
size = product.get("size", "")
|
||||
source = f"{title}_{size}" if size else title
|
||||
image_id = str(product.get("image_id") or "").strip()
|
||||
if image_id:
|
||||
by_id[image_id].append(product)
|
||||
return {k: v for k, v in by_id.items() if len(v) > 1}
|
||||
|
||||
if not source:
|
||||
|
||||
def db_resident_names(brand: str, image_ids: List[str]) -> Dict[str, str]:
|
||||
"""product_name of the row each duplicated image_id currently points at.
|
||||
|
||||
Empty when the database is unreachable - the caller then falls back to
|
||||
keeping the first entry, which is still safe because every candidate in the
|
||||
group is equally unbacked in that case.
|
||||
"""
|
||||
if not image_ids:
|
||||
return {}
|
||||
conn = _connect()
|
||||
if conn is None:
|
||||
logger.warning(" no database connection - keeping the first entry of each group")
|
||||
return {}
|
||||
table = f"brand_{brand_slug(brand)}"
|
||||
out: Dict[str, str] = {}
|
||||
try:
|
||||
with conn, conn.cursor() as cur:
|
||||
cur.execute("SELECT to_regclass(%s)", (table,))
|
||||
if cur.fetchone()[0] is None:
|
||||
return {}
|
||||
cur.execute(
|
||||
f"SELECT image_id, product_name FROM {table} WHERE image_id = ANY(%s)",
|
||||
(image_ids,),
|
||||
)
|
||||
out = {i: (n or "") for i, n in cur.fetchall()}
|
||||
except Exception as e: # noqa: BLE001 - one file failing must not abort the sweep
|
||||
logger.warning(" could not read %s: %s", table, e)
|
||||
finally:
|
||||
conn.close()
|
||||
return out
|
||||
|
||||
|
||||
def plan_renames(products: List[Dict[str, Any]],
|
||||
resident: Dict[str, str]) -> List[Tuple[Dict[str, Any], str, str]]:
|
||||
"""(product, old_id, new_id) for every entry that must be renamed.
|
||||
|
||||
The keeper is the entry whose product_name matches the database row; with
|
||||
no row, the first entry keeps the id. Everything else takes its
|
||||
`variant_key`.
|
||||
"""
|
||||
taken = {str(p.get("image_id") or "").strip() for p in products}
|
||||
renames: List[Tuple[Dict[str, Any], str, str]] = []
|
||||
|
||||
for image_id, group in duplicate_groups(products).items():
|
||||
db_name = resident.get(image_id)
|
||||
keeper = next(
|
||||
(p for p in group if db_name and str(p.get("product_name") or "") == db_name),
|
||||
group[0],
|
||||
)
|
||||
for product in group:
|
||||
if product is keeper:
|
||||
continue
|
||||
new_id = str(product.get("variant_key") or "").strip()
|
||||
if not new_id or new_id in taken:
|
||||
logger.warning(" cannot rename %r (%s): variant_key %r is missing or taken",
|
||||
product.get("product_name"), image_id, new_id)
|
||||
continue
|
||||
taken.add(new_id)
|
||||
renames.append((product, image_id, new_id))
|
||||
return renames
|
||||
|
||||
|
||||
def main() -> int:
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Give each pack-size variant a unique image_id in the seed catalogs.")
|
||||
parser.add_argument("--apply", action="store_true", help="Actually write.")
|
||||
parser.add_argument("--dry-run", action="store_true", help="Force a dry run.")
|
||||
parser.add_argument("--file", action="append", default=None,
|
||||
help="Catalog filename to process (repeatable). Default: all.")
|
||||
args = parser.parse_args()
|
||||
|
||||
apply = args.apply and not args.dry_run
|
||||
wanted = set(args.file or [])
|
||||
|
||||
logger.info("Database: %s / %s (read-only - used to find the resident row)", DB_HOST, DB_NAME)
|
||||
logger.info("Mode: %s", "APPLY (writes JSON)" if apply else "DRY RUN (no writes)")
|
||||
logger.info("")
|
||||
|
||||
total_renamed = 0
|
||||
total_files = 0
|
||||
|
||||
for path in sorted(seed_catalog_paths(SEED_DIR)):
|
||||
if wanted and path.name not in wanted:
|
||||
continue
|
||||
try:
|
||||
data = json.loads(path.read_text(encoding="utf-8-sig"))
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("Skipping unreadable catalog %s: %s", path.name, e)
|
||||
continue
|
||||
if not isinstance(data, dict):
|
||||
continue
|
||||
products = data.get("products")
|
||||
brand = data.get("brand")
|
||||
if not brand or not isinstance(products, list) or not products:
|
||||
continue
|
||||
|
||||
new_id = make_image_id(source)
|
||||
groups = duplicate_groups(products)
|
||||
if not groups:
|
||||
continue
|
||||
|
||||
# Ensure no collisions across all products
|
||||
while new_id in seen_ids:
|
||||
new_id = make_image_id(source + "_" + str(uuid.uuid4())[:4])
|
||||
total_files += 1
|
||||
logger.info("=== %s (%s) - %d entries, %d duplicated id(s)",
|
||||
path.name, brand, len(products), len(groups))
|
||||
|
||||
old_id = product.get("image_id", "")
|
||||
if old_id != new_id:
|
||||
logger.debug(f" {old_id} -> {new_id} (source={source!r})")
|
||||
product["image_id"] = new_id
|
||||
seen_ids.add(new_id)
|
||||
fixed += 1
|
||||
resident = db_resident_names(brand, list(groups))
|
||||
renames = plan_renames(products, resident)
|
||||
if not renames:
|
||||
logger.info(" nothing renameable")
|
||||
continue
|
||||
|
||||
data["total_products"] = len(products)
|
||||
with open(path, "w", encoding="utf-8") as f:
|
||||
json.dump(data, f, indent=2, ensure_ascii=False)
|
||||
for product, old_id, new_id in renames[:4]:
|
||||
logger.info(" %-42s %s -> %s", str(product.get("product_name"))[:42], old_id, new_id)
|
||||
if len(renames) > 4:
|
||||
logger.info(" ... and %d more", len(renames) - 4)
|
||||
|
||||
logger.info(f"Fixed {fixed} products in {path.name}")
|
||||
return fixed
|
||||
if apply:
|
||||
before = len(products)
|
||||
for product, _old, new_id in renames:
|
||||
product["image_id"] = new_id
|
||||
ids = [str(p.get("image_id") or "") for p in products]
|
||||
still_dup = [i for i, n in Counter(ids).items() if n > 1 and i]
|
||||
assert len(products) == before, "product count changed"
|
||||
if still_dup:
|
||||
logger.error(" ABORTING %s - %d id(s) still duplicated", path.name, len(still_dup))
|
||||
continue
|
||||
|
||||
payload = json.dumps(data, indent=2, ensure_ascii=False)
|
||||
tmp = path.with_suffix(".json.tmp")
|
||||
tmp.write_text(payload, encoding="utf-8")
|
||||
os.replace(tmp, path)
|
||||
logger.info(" wrote %s (%d renamed)", path.name, len(renames))
|
||||
else:
|
||||
logger.info(" (dry run - nothing written)")
|
||||
|
||||
def main() -> None:
|
||||
if not SEED_DIR.exists():
|
||||
logger.error("Seed directory not found: %s", SEED_DIR)
|
||||
sys.exit(1)
|
||||
total_renamed += len(renames)
|
||||
logger.info("")
|
||||
|
||||
files = sorted(SEED_DIR.glob("*.json"))
|
||||
if not files:
|
||||
logger.warning("No JSON files found in %s", SEED_DIR)
|
||||
return
|
||||
|
||||
total = 0
|
||||
for f in files:
|
||||
total += fix_seed_file(f)
|
||||
|
||||
print(f"\nDone. Fixed {total} total products across {len(files)} file(s).")
|
||||
print("Now re-run: python scripts/seed_sample_data.py")
|
||||
logger.info("---------------------------------------------------------")
|
||||
logger.info("Catalogs with duplicates : %d", total_files)
|
||||
logger.info("Entries renamed : %d", total_renamed)
|
||||
if not apply:
|
||||
logger.info("DRY RUN - re-run with --apply to write.")
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
raise SystemExit(main())
|
||||
|
||||
Reference in New Issue
Block a user