New updates on DB and JSON
This commit is contained in:
392
scripts/merge_haldiram.py
Normal file
392
scripts/merge_haldiram.py
Normal file
@@ -0,0 +1,392 @@
|
||||
#!/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())
|
||||
Reference in New Issue
Block a user