393 lines
16 KiB
Python
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())
|