Files
catalogue_backend/scripts/merge_haldiram.py
2026-09-01 13:55:15 +05:30

393 lines
16 KiB
Python

#!/usr/bin/env python3
"""
Fold `brand_haldiram` into `brand_haldirams` and drop the singular table.
WHY THERE WERE TWO
------------------
Nothing told the system they were one brand. `BRAND_ALIASES` had no haldiram
entry, so `resolve_parent_brand` was the identity for both spellings and each
upload built whichever table its sheet happened to name. That has been fixed in
`app/services/brand_registry.py`; this script cleans up the rows the split left
behind. **Deploy the alias first** - with it in place, an upload arriving while
this runs lands in `brand_haldirams` instead of recreating the singular behind
us.
WHAT IS ACTUALLY BEING MERGED
-----------------------------
Not a missing product. `brand_haldiram` holds ONE row which is a corrupted
duplicate of a product already in `brand_haldirams`:
brand_haldiram 'Haldiram Aloo Bhujia 200g 45' category '3' size ['45']
brand_haldirams "Haldiram's Aloo Bhujia 200g" Snacks ['200g']
The stray "45", the raw category id and the bare-number size all predate guards
that now exist. Moving that row across would put a second, worse Aloo Bhujia
inside the good table - duplication made worse rather than fixed. So the row is
dropped, and only what it holds that the target does NOT is carried over:
* the PRICE. The corrupted row is 12 days newer (29 Aug vs 17 Aug) and says
58 / Rs52-58 against the target's 48 / Rs45-50. A more recent upload is the
better evidence of what the shop charges.
* the FSSAI LICENCE, under --fix-fssai (off by default; see below).
THE LICENCE PROBLEM
-------------------
Both `brand_haldirams` rows carry 10012042000244. That is LION DATES' licence -
the same number on all 21 Lion Dates products - not Haldiram's. Haldiram's own,
per FSSAI_LICENSES, is 10012011000140, which is what the corrupted row has. So
the junk row holds the correct regulatory identifier and the clean rows hold
another company's.
Shipping one manufacturer's licence on another's product is the same class of
fault that `generic_products.py` exists to prevent, but repairing it is a
separate decision from merging two tables - so it is behind `--fix-fssai` and
off unless asked for. The seed file
`data/seed_catalogs/archive/brand_catalog_haldirams.json` carries the same wrong
number and is corrected in the same pass.
THE TABLE COMES BACK IF YOU ONLY DROP IT
----------------------------------------
`/app/data` is a named Docker volume and `reconcile_brand_catalogs()` exports
any table that has rows but no seed file, on boot and every 300s. Production has
very likely already written `brand_catalog_haldiram.json` into that volume. Left
there it is at best a stale file claiming to be a brand catalogue, so step 5
deletes it and reports whether it existed.
Usage:
python -m scripts.merge_haldiram # dry run, the default
python -m scripts.merge_haldiram --apply
python -m scripts.merge_haldiram --apply --fix-fssai
`--dry-run` is the default and `--apply` must be explicit: this rewrites
whatever database `backend/.env` points at, which is production. The target host
is printed on startup so it can be checked before committing, and both tables
are backed up to disk before the first write.
"""
from __future__ import annotations
import argparse
import json
import logging
import sys
from datetime import datetime
from decimal import Decimal
from pathlib import Path
from typing import Any, Dict, List, Optional
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from app.infrastructure.settings import DB_HOST, DB_NAME
from app.services.vector_store import _connect
logging.basicConfig(level=logging.INFO, format="%(message)s")
logger = logging.getLogger("merge_haldiram")
SOURCE_TABLE = "brand_haldiram"
TARGET_TABLE = "brand_haldirams"
TARGET_IMAGE_ID = "haldirams_haldiram_s_aloo_bhujia_200g"
SOURCE_IMAGE_ID = "haldiram_haldiram_aloo_bhujia_200g_45"
# Haldiram's own, from brand_registry.FSSAI_LICENSES.
HALDIRAM_FSSAI = "10012011000140"
# What the target rows wrongly carry today: Lion Dates'.
WRONG_FSSAI = "10012042000244"
ROOT = Path(__file__).resolve().parents[1]
SEED_DIR = ROOT / "data" / "seed_catalogs"
STALE_SEED = SEED_DIR / "brand_catalog_haldiram.json"
TARGET_SEED = SEED_DIR / "archive" / "brand_catalog_haldirams.json"
# ---------------------------------------------------------------------------
# Reading
# ---------------------------------------------------------------------------
def _table_exists(cur, table: str) -> bool:
cur.execute(
"SELECT 1 FROM information_schema.tables "
"WHERE table_schema='public' AND table_name=%s",
(table,),
)
return cur.fetchone() is not None
def _rows(cur, table: str) -> List[Dict[str, Any]]:
cur.execute(f"SELECT * FROM {table} ORDER BY id")
cols = [d[0] for d in cur.description]
return [dict(zip(cols, r)) for r in cur.fetchall()]
def _jsonable(row: Dict[str, Any]) -> Dict[str, Any]:
"""Make a DB row writable as JSON - datetimes and vectors are not."""
out = {}
for key, value in row.items():
if isinstance(value, datetime):
out[key] = value.isoformat()
elif key == "embedding" and value is not None:
out[key] = list(value) if not isinstance(value, str) else value
else:
out[key] = value
return out
def _backup(source: List[Dict], target: List[Dict]) -> Path:
"""Write both tables to disk BEFORE anything is changed. This is the undo."""
stamp = datetime.now().strftime("%Y%m%d_%H%M%S")
path = ROOT / "data" / f"haldiram_merge_backup_{stamp}.json"
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(
json.dumps(
{
"taken_at": stamp,
"db_host": DB_HOST,
"db_name": DB_NAME,
SOURCE_TABLE: [_jsonable(r) for r in source],
TARGET_TABLE: [_jsonable(r) for r in target],
},
ensure_ascii=False,
indent=1,
default=str,
),
encoding="utf-8",
)
return path
# ---------------------------------------------------------------------------
# Preconditions
# ---------------------------------------------------------------------------
def _check(source: List[Dict], target: List[Dict]) -> Optional[str]:
"""Refuse to run against data that is not what this script was written for.
The merge is hand-derived from one specific pair of rows. If either side has
moved since, the reasoning above may no longer hold and a blind UPDATE could
overwrite something real - so stop and let a human look.
"""
if len(source) != 1:
return f"{SOURCE_TABLE} has {len(source)} rows, expected exactly 1"
if source[0].get("image_id") != SOURCE_IMAGE_ID:
return (f"{SOURCE_TABLE} row is {source[0].get('image_id')!r}, "
f"expected {SOURCE_IMAGE_ID!r}")
if not any(r.get("image_id") == TARGET_IMAGE_ID for r in target):
return f"{TARGET_TABLE} has no row {TARGET_IMAGE_ID!r} to merge into"
return None
# ---------------------------------------------------------------------------
# Reporting
# ---------------------------------------------------------------------------
def _describe(source: Dict, target: Dict, fix_fssai: bool) -> None:
logger.info("")
logger.info(" %s -> %s", SOURCE_TABLE, TARGET_TABLE)
logger.info("")
logger.info(" row being DROPPED (corrupted duplicate):")
logger.info(" product_name %r", source.get("product_name"))
logger.info(" category %r size %s",
source.get("category"), source.get("size_variants"))
logger.info(" price %s (%s)",
source.get("selling_price"), source.get("price_range"))
logger.info("")
logger.info(" row being KEPT and updated:")
logger.info(" product_name %r", target.get("product_name"))
logger.info(" category %r size %s",
target.get("category"), target.get("size_variants"))
logger.info(" price %s (%s) -> %s (%s)",
target.get("selling_price"), target.get("price_range"),
source.get("selling_price"), source.get("price_range"))
if fix_fssai:
logger.info(" fssai %s -> %s",
target.get("fssai_license"), HALDIRAM_FSSAI)
else:
logger.info(" fssai %s (unchanged - pass --fix-fssai to "
"correct it; see the docstring)", target.get("fssai_license"))
logger.info(" name / category / size / image / description / embedding"
" all kept as they are")
logger.info("")
# ---------------------------------------------------------------------------
# The seed file
# ---------------------------------------------------------------------------
def _plain(value: Any) -> Any:
"""psycopg returns NUMERIC as Decimal, which json.dumps refuses.
Caught the hard way: the database work committed and then the seed write
blew up on `final_selling_price`, leaving the two halves out of step. The
conversion happens here rather than at the call site so every field copied
into a JSON file goes through it.
"""
return float(value) if isinstance(value, Decimal) else value
def _update_seed(live_rows: List[Dict], fix_fssai: bool, apply: bool) -> None:
"""Bring the archived catalogue into step with the TABLE.
Mirrored from the live rows rather than from the row being merged in, for
two reasons. It is idempotent - re-running compares the file against the
database and changes only what differs - and it still works once the source
table has been dropped, which matters because the database write and this
file write are not in one transaction. They came apart once already: the
merge committed and then the JSON write failed on a Decimal, leaving the two
halves disagreeing until this ran again.
"""
if not TARGET_SEED.exists():
logger.info(" seed file %s not found - skipping", TARGET_SEED.name)
return
by_id = {r["image_id"]: r for r in live_rows}
doc = json.loads(TARGET_SEED.read_text(encoding="utf-8-sig"))
changes: List[str] = []
for product in doc.get("products", []):
row = by_id.get(product.get("image_id"))
if not row:
continue
for field in ("price_range", "selling_price", "final_selling_price"):
new = _plain(row.get(field))
if product.get(field) != new:
changes.append(f"{product['image_id']}.{field}: "
f"{product.get(field)!r} -> {new!r}")
product[field] = new
if fix_fssai and product.get("fssai_license") != row.get("fssai_license"):
changes.append(f"{product['image_id']}.fssai_license: "
f"{product.get('fssai_license')!r} -> "
f"{row.get('fssai_license')!r}")
product["fssai_license"] = row.get("fssai_license")
if not changes:
logger.info(" %s: already in step with the table", TARGET_SEED.name)
return
logger.info(" %s: %d field(s) differ from the table", TARGET_SEED.name, len(changes))
for line in changes:
logger.info(" %s", line)
if apply:
TARGET_SEED.write_text(
json.dumps(doc, ensure_ascii=False, indent=1), encoding="utf-8")
logger.info(" written")
# ---------------------------------------------------------------------------
def main() -> int:
parser = argparse.ArgumentParser(
description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter
)
parser.add_argument("--apply", action="store_true",
help="commit the changes (default is a dry run)")
parser.add_argument("--dry-run", action="store_true",
help="explicit no-op; this is already the default")
parser.add_argument("--fix-fssai", action="store_true",
help="also replace Lion Dates' licence on the Haldirams "
"rows with Haldiram's own")
args = parser.parse_args()
apply = args.apply and not args.dry_run
logger.info("Target database: %s / %s", DB_HOST, DB_NAME)
logger.info("Mode: %s", "APPLY - this writes" if apply else "DRY RUN - nothing is written")
conn = _connect()
if conn is None:
logger.error("No database connection.")
return 1
with conn:
with conn.cursor() as cur:
merged_already = not _table_exists(cur, SOURCE_TABLE)
if merged_already:
# The table is gone, but the seed files may still be behind -
# they are written outside the database transaction, so a run
# that committed the merge and then failed on the file leaves
# exactly this state. Fall through to the file work.
logger.info("")
logger.info("%s is already gone; checking the seed files.",
SOURCE_TABLE)
target_rows = _rows(cur, TARGET_TABLE)
source_rows = []
else:
source_rows = _rows(cur, SOURCE_TABLE)
target_rows = _rows(cur, TARGET_TABLE)
problem = None if merged_already else _check(source_rows, target_rows)
if problem:
logger.error("")
logger.error("ABORTING - the data is not what this script expects:")
logger.error(" %s", problem)
logger.error("")
logger.error("Re-read the rows and update the script rather than "
"forcing it; a blind UPDATE here could overwrite a "
"real product.")
return 1
if not merged_already:
source = source_rows[0]
target = next(r for r in target_rows
if r["image_id"] == TARGET_IMAGE_ID)
_describe(source, target, args.fix_fssai)
if apply and not merged_already:
path = _backup(source_rows, target_rows)
logger.info(" backup written to %s", path)
cur.execute(
f"UPDATE {TARGET_TABLE} SET selling_price=%s, "
f"final_selling_price=%s, price_range=%s, updated_at=now() "
f"WHERE image_id=%s",
(source.get("selling_price"), source.get("final_selling_price"),
source.get("price_range"), TARGET_IMAGE_ID),
)
logger.info(" updated %s (%d row)", TARGET_TABLE, cur.rowcount)
if args.fix_fssai:
cur.execute(
f"UPDATE {TARGET_TABLE} SET fssai_license=%s, "
f"updated_at=now() WHERE fssai_license=%s",
(HALDIRAM_FSSAI, WRONG_FSSAI),
)
logger.info(" corrected fssai on %d row(s)", cur.rowcount)
cur.execute(f"DROP TABLE {SOURCE_TABLE}")
logger.info(" dropped %s", SOURCE_TABLE)
conn.commit()
# Re-read so the seed mirrors what was actually committed.
target_rows = _rows(cur, TARGET_TABLE)
elif not merged_already:
logger.info(" would UPDATE %s then DROP TABLE %s",
TARGET_TABLE, SOURCE_TABLE)
# ---- the seed files ----------------------------------------------------
logger.info("")
_update_seed(target_rows, args.fix_fssai, apply)
if STALE_SEED.exists():
logger.info(" %s EXISTS (auto-exported by reconcile) - deleting",
STALE_SEED.name)
if apply:
STALE_SEED.unlink()
logger.info(" deleted")
else:
logger.info(" %s not present locally", STALE_SEED.name)
logger.info(" NOTE: production keeps data/ on a named volume, so check "
"there too - reconcile exports any table with rows and no file.")
# ---- caches ------------------------------------------------------------
if apply:
from app.services.vector_store import invalidate_brand_overview_cache
invalidate_brand_overview_cache()
try:
from app.services.query_intent import invalidate_brand_mention_cache
invalidate_brand_mention_cache()
except Exception: # noqa: BLE001 - a cold cache is not a failure
pass
logger.info("")
logger.info("Done. Brand caches invalidated.")
else:
logger.info("")
logger.info("Dry run - nothing written. Re-run with --apply to commit.")
return 0
if __name__ == "__main__":
raise SystemExit(main())