backend store_catalog updates
This commit is contained in:
217
scripts/migrate_own_products.py
Normal file
217
scripts/migrate_own_products.py
Normal file
@@ -0,0 +1,217 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Move already-stored unbranded commodities into `brand_own_products`.
|
||||
|
||||
Fixing the pipeline only helps the next upload. Rows written before the fix are
|
||||
still sitting in two kinds of wrong place:
|
||||
|
||||
* **Junk tables** - `brand_toor`, `brand_sugar`, `brand_rice`, `brand_black`,
|
||||
one per leading word, created because `infer_brand` fell back to the first
|
||||
token of the product name.
|
||||
* **Real brand catalogs** - and this is the damaging half. A single commodity
|
||||
noun matching a whole word inside a multi-word alias sent salt into
|
||||
`brand_colgate_palmolive` ("colgate active salt"), milk into `brand_cadbury`
|
||||
("cadbury dairy milk"), butter/ghee/paneer into `brand_nestle`, cheese into
|
||||
`brand_amul`, tea into `brand_hindustan_unilever` and red chilli powder into
|
||||
`brand_brooke_bond` ("brooke bond red label"). `stage_1_brand_and_fssai`
|
||||
then stamped that brand's REAL FSSAI licence number onto the row, so those
|
||||
rows are carrying another company's regulatory identifier.
|
||||
|
||||
The same `is_unbranded` predicate the pipeline now uses does both jobs, and it
|
||||
is what keeps genuine products safe: "Amul Butter 500g" contains "amul", which
|
||||
is not a commodity term, so it is never touched. Only the bare "Butter 500g" is.
|
||||
|
||||
What a moved row gets:
|
||||
* `fssai_license` cleared - it belonged to somebody else.
|
||||
* a recomputed `image_id` (`own_products_<slug>`), so a later re-ingest of the
|
||||
same sheet UPDATES the row instead of inserting a duplicate beside it.
|
||||
* `image_url` / `image_urls` cleared, so the fixed image search refetches
|
||||
them. The old S3 folder is orphaned by this; those images were mostly
|
||||
absent or belonged to the wrong product anyway.
|
||||
* `category` and `price_range` re-resolved with the staple bands.
|
||||
|
||||
Usage:
|
||||
|
||||
python -m scripts.migrate_own_products --dry-run # default
|
||||
python -m scripts.migrate_own_products --apply
|
||||
python -m scripts.migrate_own_products --apply --drop-empty
|
||||
|
||||
`--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.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import logging
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Tuple
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
|
||||
|
||||
from app.infrastructure.settings import DB_HOST, DB_NAME
|
||||
from app.services import price_estimator
|
||||
from app.services.category_registry import detect_category_from_text
|
||||
from app.services.generic_products import OWN_PRODUCTS_BRAND, is_unbranded
|
||||
from app.services.vector_store import (
|
||||
_connect,
|
||||
_sanitize_name,
|
||||
ensure_brand_schema,
|
||||
)
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s")
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
TARGET_SUFFIX = _sanitize_name(OWN_PRODUCTS_BRAND) # "own_products"
|
||||
TARGET_TABLE = f"brand_{TARGET_SUFFIX}"
|
||||
|
||||
# Columns copied across. Deliberately explicit rather than SELECT *: `id` is a
|
||||
# per-table BIGSERIAL and must not travel, and `embedding` is recomputed by the
|
||||
# next ingest rather than carried.
|
||||
_COPY_COLUMNS = (
|
||||
"product_name", "title", "description", "category", "image_id",
|
||||
"image_url", "image_urls", "price_range", "size_variants", "providers",
|
||||
"fssai_license", "product_sku", "sku_source", "hsn_code",
|
||||
"final_selling_price", "selling_price", "barcode", "barcode_type",
|
||||
"highlights", "nutrients", "search_query",
|
||||
)
|
||||
|
||||
|
||||
def _slugify(text: str) -> str:
|
||||
import re
|
||||
return re.sub(r"_+", "_", re.sub(r"[^a-z0-9]+", "_", (text or "").lower())).strip("_")
|
||||
|
||||
|
||||
def _brand_tables(cur) -> List[str]:
|
||||
cur.execute(
|
||||
"SELECT table_name FROM information_schema.tables "
|
||||
"WHERE table_schema = 'public' AND table_name LIKE 'brand_%' "
|
||||
"ORDER BY table_name"
|
||||
)
|
||||
return [r[0] for r in cur.fetchall()]
|
||||
|
||||
|
||||
def _rebuild_row(row: Dict[str, Any]) -> Dict[str, Any]:
|
||||
"""Re-resolve the fields that were wrong because the brand was wrong."""
|
||||
name = row.get("product_name") or ""
|
||||
category = detect_category_from_text(name) or "General"
|
||||
size = (row.get("size_variants") or ["Standard"])[0]
|
||||
|
||||
updated = dict(row)
|
||||
updated["category"] = category
|
||||
updated["image_id"] = f"{TARGET_SUFFIX}_{_slugify(name)}"
|
||||
# Belonged to the brand this row was wrongly filed under.
|
||||
updated["fssai_license"] = None
|
||||
# Cleared so the fixed search refetches with a query that is not poisoned
|
||||
# by a brand name the product never had.
|
||||
updated["image_url"] = None
|
||||
updated["image_urls"] = []
|
||||
lo, hi = price_estimator.estimate_price_range_for_size(
|
||||
size, name, OWN_PRODUCTS_BRAND, category
|
||||
)
|
||||
updated["price_range"] = f"₹{lo}-{hi}"
|
||||
return updated
|
||||
|
||||
|
||||
def scan(cur) -> List[Tuple[str, Dict[str, Any]]]:
|
||||
"""Every unbranded row currently sitting somewhere else."""
|
||||
found: List[Tuple[str, Dict[str, Any]]] = []
|
||||
for table in _brand_tables(cur):
|
||||
if table == TARGET_TABLE:
|
||||
continue
|
||||
try:
|
||||
cur.execute(f"SELECT {', '.join(_COPY_COLUMNS)} FROM {table}")
|
||||
except Exception as exc: # noqa: BLE001 - a legacy table may lack columns
|
||||
logger.warning("Skipping %s: %s", table, exc)
|
||||
continue
|
||||
for values in cur.fetchall():
|
||||
row = dict(zip(_COPY_COLUMNS, values))
|
||||
if is_unbranded(row.get("product_name") or ""):
|
||||
found.append((table, row))
|
||||
return found
|
||||
|
||||
|
||||
def migrate(apply: bool, drop_empty: bool) -> int:
|
||||
conn = _connect()
|
||||
if not conn:
|
||||
logger.error("Could not connect to PostgreSQL (USE_PGVECTOR off, or DB unreachable).")
|
||||
return 1
|
||||
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
candidates = scan(cur)
|
||||
|
||||
if not candidates:
|
||||
logger.info("Nothing to migrate - no unbranded rows outside %s.", TARGET_TABLE)
|
||||
return 0
|
||||
|
||||
by_table: Dict[str, int] = {}
|
||||
for table, row in candidates:
|
||||
by_table[table] = by_table.get(table, 0) + 1
|
||||
logger.info("%s: %r -> %s", table, row["product_name"], TARGET_TABLE)
|
||||
|
||||
logger.info("--- %d row(s) across %d table(s) ---", len(candidates), len(by_table))
|
||||
for table, count in sorted(by_table.items()):
|
||||
logger.info(" %-34s %d row(s)", table, count)
|
||||
|
||||
if not apply:
|
||||
logger.info("Dry run - nothing written. Re-run with --apply to commit.")
|
||||
return 0
|
||||
|
||||
ensure_brand_schema(OWN_PRODUCTS_BRAND)
|
||||
moved = 0
|
||||
with conn.cursor() as cur:
|
||||
for table, row in candidates:
|
||||
rebuilt = _rebuild_row(row)
|
||||
columns = ", ".join(_COPY_COLUMNS)
|
||||
placeholders = ", ".join(["%s"] * len(_COPY_COLUMNS))
|
||||
# ON CONFLICT: the same commodity may sit in several junk
|
||||
# tables; the first one wins and the rest are dropped rather
|
||||
# than failing the whole migration.
|
||||
cur.execute(
|
||||
f"INSERT INTO {TARGET_TABLE} ({columns}) VALUES ({placeholders}) "
|
||||
f"ON CONFLICT (image_id) DO NOTHING",
|
||||
[rebuilt[c] for c in _COPY_COLUMNS],
|
||||
)
|
||||
cur.execute(
|
||||
f"DELETE FROM {table} WHERE image_id = %s", [row["image_id"]]
|
||||
)
|
||||
moved += 1
|
||||
|
||||
if drop_empty:
|
||||
for table in sorted(by_table):
|
||||
cur.execute(f"SELECT COUNT(*) FROM {table}")
|
||||
if cur.fetchone()[0] == 0:
|
||||
logger.info("Dropping now-empty %s", table)
|
||||
cur.execute(f"DROP TABLE {table}")
|
||||
|
||||
logger.info("Moved %d row(s) into %s.", moved, TARGET_TABLE)
|
||||
logger.info(
|
||||
"Remember: %r must be in ACTIVE_BRANDS or the table stays invisible.",
|
||||
OWN_PRODUCTS_BRAND,
|
||||
)
|
||||
return 0
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def main() -> int:
|
||||
parser = argparse.ArgumentParser(
|
||||
description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter
|
||||
)
|
||||
parser.add_argument("--apply", action="store_true",
|
||||
help="write 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("--drop-empty", action="store_true",
|
||||
help="drop junk brand tables left with no rows (requires --apply)")
|
||||
args = parser.parse_args()
|
||||
|
||||
logger.info("Target database: %s/%s", DB_HOST, DB_NAME)
|
||||
logger.info("Mode: %s", "APPLY - rows will be moved" if args.apply else "DRY RUN - no writes")
|
||||
return migrate(apply=args.apply, drop_empty=args.drop_empty)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
Reference in New Issue
Block a user