Brand image display cards

This commit is contained in:
sriram
2026-09-02 13:34:40 +05:30
parent 5399fea4cc
commit d5ec5755bf
8 changed files with 921 additions and 18 deletions

View File

@@ -57,11 +57,22 @@ def backfill_brand_images() -> None:
db_has_images = bool(image_url or (image_urls and any(image_urls)))
if not db_has_images:
# Fetch from S3 (or construct primary S3 URL)
# LIST ONLY - never construct.
#
# This used to fall back to
# s3_service.get_product_image_url(brand, image_id), which
# returns .../{image_id}/image_000.jpg whether or not the
# object exists. That is how Everest, Haldirams, MDH and
# Naga came to hold URLs that 404 on every request while
# looking perfectly populated in the database.
#
# Writing a URL we have not confirmed is worse than writing
# nothing: a row with no image renders the brand monogram,
# a row with a dead URL renders the same monogram but hides
# the fact that the image is missing. Leave it empty and let
# scripts/repair_brand_images.py find a real one.
s3_urls = s3_service.get_product_image_urls(brand, image_id)
primary_url = s3_urls[0] if s3_urls else s3_service.get_product_image_url(brand, image_id)
if not s3_urls and primary_url:
s3_urls = [primary_url]
primary_url = s3_urls[0] if s3_urls else None
if primary_url or s3_urls:
cur.execute(

View File

@@ -0,0 +1,516 @@
#!/usr/bin/env python3
"""
Find product rows whose stored image URLs no longer resolve, and repair them.
WHY THIS EXISTS
---------------
Ten of the fifty-five brand cards were rendering the initials monogram instead
of a photo, and not one of them was missing data. Every one held a non-empty
`image_url` that could not be fetched:
Balaji, Colin, Hindustan http://... - the site is https, so the
Unilever, Own Products, browser blocks it as mixed content
PepsiCo
Everest, Haldirams, https://nearledaily.s3.ap-south-1
MDH, Naga .amazonaws.com/... - every object 404s
India Gate asiancorner.pl - host is dead
The first group is already fixed in code: `_pick_sample_image` in vector_store
now prefers an https candidate, and those brands each held working https URLs
all along (Colin nine, PepsiCo fifty-five, Hindustan Unilever seventy-four).
This script is for the rest, where no usable URL exists in the row at all.
The S3 URLs were never real. `s3_service.get_product_image_url()` builds
`daily/brands/{brand}/{image_id}/image_000.jpg` and returns it whether or not
the object exists, so a row can look fully populated and render nothing. That
constructor is no longer called from the read path or from
`scripts/backfill_s3_urls.py`; if you reintroduce it, this damage comes back.
WHAT IT WILL NOT DO
-------------------
* It never rewrites a row that already has one working URL. Not re-ranked, not
re-searched, not touched.
* It never blanks a dead URL to NULL. A dead URL and no URL render identically,
so clearing gains nothing and destroys the only record of where the image was
meant to live. A row is written only when the search found a replacement.
* It never writes `image_id`. That is the UNIQUE upsert key and the S3 folder
name; changing it silently forks the product.
USAGE
-----
python -m scripts.repair_brand_images --audit
python -m scripts.repair_brand_images --audit --json
python -m scripts.repair_brand_images --brands Everest,Naga # dry run
python -m scripts.repair_brand_images --brands Everest --apply
python -m scripts.repair_brand_images --restore data/brand_image_repair_backup_X.json --apply
`--audit` is read-only by construction: it makes no image searches and opens no
write transaction. Everything else is a dry run until `--apply` is passed, and
`--apply` always writes a full backup of every targeted table first.
AFTERWARDS
----------
The running API caches the brand overview for BRAND_OVERVIEW_TTL_SECONDS in a
per-process dict, so this script cannot clear it. Either wait out the TTL or
hit `GET /api/brands/overview?refresh=true`. The script prints this reminder.
"""
from __future__ import annotations
import argparse
import json
import logging
import sys
import time
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime
from decimal import Decimal
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 # noqa: E402
from app.services.brand_registry import resolve_parent_brand # noqa: E402
from app.services.vector_store import ( # noqa: E402
_connect,
_list_brand_table_suffixes,
_sanitize_name,
invalidate_brand_overview_cache,
)
logging.basicConfig(level=logging.INFO, format="%(message)s")
logger = logging.getLogger("repair_brand_images")
DATA_DIR = Path(__file__).resolve().parents[1] / "data"
HEALTHY = "healthy" # at least one URL resolves
BROKEN = "broken" # has URLs, none resolve
EMPTY = "empty" # no URLs at all
# ---------------------------------------------------------------------------
# URL validation
# ---------------------------------------------------------------------------
# Probing the same URL once per row would be pure waste: Kaleesuwari's eighty-six
# URLs share one host and one key prefix, and the four S3 brands share another.
# Memoised across the whole run, per process.
_URL_CACHE: Dict[str, bool] = {}
def _url_is_live(url: str, timeout: int) -> bool:
"""True when the URL serves something that looks like an image.
Uses the same probe the ingestion pipeline uses, so a URL this script
accepts is one stage 6 would also have accepted.
"""
if url in _URL_CACHE:
return _URL_CACHE[url]
ok = False
try:
from app.services.image_search import validate_image_url_live
ok = bool(validate_image_url_live(url, timeout=timeout))
except Exception: # noqa: BLE001 - an unreachable URL is the thing we are measuring
ok = False
_URL_CACHE[url] = ok
return ok
def _row_urls(row: Dict[str, Any]) -> List[str]:
"""Every candidate URL on a row, deduped, order preserved."""
out: List[str] = []
seen = set()
for value in [row.get("image_url"), *(row.get("image_urls") or [])]:
if not isinstance(value, str):
continue
url = value.strip()
if url and url not in seen:
seen.add(url)
out.append(url)
return out
def _classify(rows: List[Dict[str, Any]], timeout: int, workers: int) -> Dict[str, str]:
"""image_id -> HEALTHY | BROKEN | EMPTY, validating every distinct URL once."""
distinct: List[str] = []
seen = set()
for row in rows:
for url in _row_urls(row):
if url not in seen:
seen.add(url)
distinct.append(url)
if distinct:
with ThreadPoolExecutor(max_workers=workers) as pool:
list(pool.map(lambda u: _url_is_live(u, timeout), distinct))
verdict: Dict[str, str] = {}
for row in rows:
urls = _row_urls(row)
if not urls:
verdict[row["image_id"]] = EMPTY
elif any(_URL_CACHE.get(u) for u in urls):
verdict[row["image_id"]] = HEALTHY
else:
verdict[row["image_id"]] = BROKEN
return verdict
def _host(url: str) -> str:
return url.split("/", 3)[2].lower() if url.count("/") >= 2 else "?"
# ---------------------------------------------------------------------------
# Reading
# ---------------------------------------------------------------------------
def _brand_tables(cur, brands: Optional[List[str]]) -> List[Tuple[str, str]]:
"""(suffix, table) pairs. Named brands resolve through the alias map."""
if brands:
pairs = []
for raw in brands:
suffix = _sanitize_name(resolve_parent_brand(raw.strip()))
pairs.append((suffix, f"brand_{suffix}"))
return pairs
# include_inactive: an archived brand still has a table and still renders a
# card, so it still has to be auditable.
return [(s, f"brand_{s}") for s in sorted(set(_list_brand_table_suffixes(cur, include_inactive=True)))]
def _image_columns(cur, table: str) -> List[str]:
"""Which image columns this table actually has.
Brand tables written under an older DDL carry only `image_id` - brand_anil
and brand_kaleesuwari are both like this. A hardcoded SELECT raises on them,
and swallowing that error silently drops 67 products from the audit, so the
columns are discovered the same way get_brand_overview discovers them.
"""
cur.execute(
"SELECT column_name FROM information_schema.columns "
"WHERE table_name = %s AND column_name IN ('image_url', 'image_urls')",
(table,),
)
return [r[0] for r in cur.fetchall()]
def _read_rows(cur, table: str) -> List[Dict[str, Any]]:
"""Never SELECT * - `embedding` is large and must not be read or written."""
have = _image_columns(cur, table)
cols = ["image_id", "product_name"] + have
try:
cur.execute(f"SELECT {', '.join(cols)} FROM {table} WHERE image_id IS NOT NULL")
except Exception as exc: # noqa: BLE001 - table may not exist at all
logger.warning(" skipped %s: %s", table, exc)
return []
rows = []
for record in cur.fetchall():
row = dict(zip(cols, record))
# A table with no URL columns yields no URLs, which classifies every
# row as EMPTY. That is the honest answer: those products have nothing
# stored and depend entirely on the S3 listing fallback at read time.
rows.append({
"image_id": row["image_id"],
"product_name": row.get("product_name"),
"image_url": row.get("image_url"),
"image_urls": list(row.get("image_urls") or []),
})
return rows
# ---------------------------------------------------------------------------
# Audit - read-only
# ---------------------------------------------------------------------------
def audit(cur, brands: Optional[List[str]], timeout: int, workers: int,
as_json: bool) -> Dict[str, Any]:
report: List[Dict[str, Any]] = []
dead_hosts: Dict[str, int] = {}
for suffix, table in _brand_tables(cur, brands):
rows = _read_rows(cur, table)
if not rows:
continue
verdict = _classify(rows, timeout, workers)
broken = [r for r in rows if verdict[r["image_id"]] == BROKEN]
empty = [r for r in rows if verdict[r["image_id"]] == EMPTY]
for row in broken:
for url in _row_urls(row):
if not _URL_CACHE.get(url):
dead_hosts[_host(url)] = dead_hosts.get(_host(url), 0) + 1
report.append({
"brand": suffix,
"rows": len(rows),
"healthy": len(rows) - len(broken) - len(empty),
"broken": len(broken),
"empty": len(empty),
"unusable_pct": round(100.0 * (len(broken) + len(empty)) / len(rows), 1),
})
ranked = sorted(dead_hosts.items(), key=lambda kv: -kv[1])
result = {"brands": report, "dead_hosts": ranked}
if as_json:
print(json.dumps(result, indent=2))
return result
logger.info("")
logger.info("%-24s %6s %8s %7s %6s %7s", "BRAND", "ROWS", "HEALTHY", "BROKEN", "EMPTY", "UNUSABLE")
logger.info("%s", "-" * 64)
for entry in sorted(report, key=lambda e: -e["unusable_pct"]):
flag = " <-- every image unusable" if entry["unusable_pct"] == 100.0 else ""
logger.info(
"%-24s %6d %8d %7d %6d %6.1f%%%s",
entry["brand"], entry["rows"], entry["healthy"],
entry["broken"], entry["empty"], entry["unusable_pct"], flag,
)
total_rows = sum(e["rows"] for e in report)
total_bad = sum(e["broken"] + e["empty"] for e in report)
logger.info("%s", "-" * 64)
logger.info("%d brand(s), %d product(s), %d with no usable image (%.1f%%)",
len(report), total_rows, total_bad,
(100.0 * total_bad / total_rows) if total_rows else 0.0)
if ranked:
logger.info("")
logger.info("Dead hosts, most damaging first - this is what turns a set of")
logger.info("per-brand bugs into one systemic finding:")
for host, count in ranked[:12]:
logger.info(" %-52s %5d dead URL(s)", host, count)
return result
# ---------------------------------------------------------------------------
# Repair
# ---------------------------------------------------------------------------
def _plain(value: Any) -> Any:
"""JSON-safe. merge_haldiram hit both of these traps for real."""
if isinstance(value, Decimal):
return float(value)
if isinstance(value, datetime):
return value.isoformat()
return value
def _backup(tables: Dict[str, List[Dict[str, Any]]]) -> Path:
"""Every row of every targeted table, not only the ones about to change.
That is what makes --restore meaningful: restoring a partial snapshot would
leave the table in a third state that never existed.
"""
DATA_DIR.mkdir(parents=True, exist_ok=True)
path = DATA_DIR / f"brand_image_repair_backup_{datetime.now():%Y%m%d_%H%M%S}.json"
payload = {
"created_at": datetime.now().isoformat(),
"db_host": DB_HOST,
"db_name": DB_NAME,
"tables": {
table: [{k: _plain(v) for k, v in row.items()} for row in rows]
for table, rows in tables.items()
},
}
path.write_text(json.dumps(payload, indent=2), encoding="utf-8")
return path
def _search_replacement(product_name: str, brand: str, max_results: int) -> List[str]:
"""Validated image URLs for one product, best first.
find_all_image_urls validates every candidate internally, so what comes back
is live by construction. Ranking then reuses the same _select_best_images
the ingestion pipeline uses, so a repaired row is ranked by the rules
tests/test_image_selection.py already pins.
"""
from app.services.image_search import find_all_image_urls
urls = find_all_image_urls(product_name, brand=brand, validate=True, max_results=max_results)
if not urls:
return []
try:
from app.core.catalog_engine import CatalogEngine
urls = CatalogEngine._select_best_images(
CatalogEngine, urls, product_name, brand, max_images=10
)
except Exception: # noqa: BLE001 - ranking is a nicety; a live URL is the point
pass
# https first, mirroring _pick_sample_image, so the two layers cannot
# disagree about which of a row's URLs is the good one.
return sorted(urls, key=lambda u: 0 if u.startswith("https://") else 1)
def repair(cur, conn, brands: List[str], apply: bool, timeout: int, workers: int,
limit: Optional[int], max_results: int) -> Dict[str, int]:
totals = {"scanned": 0, "healthy": 0, "repaired": 0, "unrepairable": 0}
unrepairable: List[str] = []
snapshot: Dict[str, List[Dict[str, Any]]] = {}
backup_written = False
for suffix, table in _brand_tables(cur, brands):
rows = _read_rows(cur, table)
if not rows:
logger.info("%-20s no rows", suffix)
continue
snapshot[table] = rows
verdict = _classify(rows, timeout, workers)
needs = [r for r in rows if verdict[r["image_id"]] != HEALTHY]
totals["scanned"] += len(rows)
totals["healthy"] += len(rows) - len(needs)
logger.info("%-20s %3d row(s), %d need an image", suffix, len(rows), len(needs))
if not needs:
continue
# A legacy table with no URL columns has nowhere to put the answer.
# Adding columns is a migration, not a repair, so say so and move on
# rather than raising halfway through a multi-brand run.
if len(_image_columns(cur, table)) < 2:
logger.warning(
" %s has no image_url/image_urls column - skipping. "
"It needs _ensure_columns() run against it first.", table,
)
totals["unrepairable"] += len(needs)
unrepairable.append(f"{suffix} ({len(needs)} row(s): table lacks URL columns)")
continue
if apply and not backup_written:
# Snapshot every targeted table BEFORE the first write. Done lazily
# so an audit-shaped run that finds nothing writes no file.
for other_suffix, other_table in _brand_tables(cur, brands):
snapshot.setdefault(other_table, _read_rows(cur, other_table))
logger.info("Backup written: %s", _backup(snapshot))
backup_written = True
for row in needs[: limit or len(needs)]:
name = row.get("product_name") or ""
if not name:
unrepairable.append(f"{suffix}/{row['image_id']} (no product_name to search on)")
totals["unrepairable"] += 1
continue
found = _search_replacement(name, suffix, max_results)
if not found:
logger.info(" no replacement found %s", name[:58])
unrepairable.append(f"{suffix}/{row['image_id']} ({name})")
totals["unrepairable"] += 1
time.sleep(0.5)
continue
logger.info(" %s -> %s", name[:46], found[0][:70])
if apply:
cur.execute(
f"UPDATE {table} SET image_url = %s, image_urls = %s, "
f"updated_at = CURRENT_TIMESTAMP WHERE image_id = %s",
(found[0], found, row["image_id"]),
)
totals["repaired"] += 1
# find_all_image_urls fans out across five providers per product;
# polite at twenty rows, abusive at eight thousand.
time.sleep(0.5)
if apply:
conn.commit() # per brand, so a later failure keeps earlier work
if unrepairable:
logger.info("")
logger.info("Could not repair %d row(s) - curate these by hand:", len(unrepairable))
for item in unrepairable[:40]:
logger.info(" %s", item)
return totals
def restore(cur, conn, path: Path, apply: bool) -> int:
payload = json.loads(path.read_text(encoding="utf-8"))
if payload.get("db_host") != DB_HOST:
logger.error("Backup was taken from %s but DB_HOST is %s. Refusing.",
payload.get("db_host"), DB_HOST)
return 1
count = 0
for table, rows in payload["tables"].items():
for row in rows:
if apply:
cur.execute(
f"UPDATE {table} SET image_url = %s, image_urls = %s WHERE image_id = %s",
(row.get("image_url"), row.get("image_urls") or [], row["image_id"]),
)
count += 1
if apply:
conn.commit()
logger.info("%s %d row(s) from %s", "Restored" if apply else "Would restore", count, path.name)
return 0
# ---------------------------------------------------------------------------
def main() -> int:
parser = argparse.ArgumentParser(
description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter
)
mode = parser.add_mutually_exclusive_group(required=True)
mode.add_argument("--audit", action="store_true",
help="read-only report across every brand; writes nothing")
mode.add_argument("--brands", help="comma-separated brands to repair")
mode.add_argument("--all", action="store_true",
help="repair every brand; requires --limit")
mode.add_argument("--restore", metavar="PATH", help="restore a backup file")
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("--json", action="store_true", help="audit output as JSON")
parser.add_argument("--limit", type=int, help="max rows to repair per brand")
parser.add_argument("--timeout", type=int, default=8, help="per-URL probe seconds")
parser.add_argument("--workers", type=int, default=8, help="parallel URL probes")
parser.add_argument("--max-results", type=int, default=8,
help="candidates to request per product search")
args = parser.parse_args()
apply = args.apply and not args.dry_run
if args.audit and args.apply:
parser.error("--audit is read-only; it cannot be combined with --apply")
if args.all and apply and not args.limit:
parser.error("--all --apply requires --limit: a catalogue-wide network "
"sweep must be a deliberate decision")
logger.info("Target database: %s / %s", DB_HOST, DB_NAME)
if args.audit:
logger.info("Mode: AUDIT - read-only, no searches, nothing is written")
else:
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
try:
with conn.cursor() as cur:
if args.restore:
return restore(cur, conn, Path(args.restore), apply)
if args.audit:
audit(cur, None, args.timeout, args.workers, args.json)
return 0
brands = [b for b in (args.brands or "").split(",") if b.strip()] or None
if args.all:
brands = None
totals = repair(cur, conn, brands, apply, args.timeout,
args.workers, args.limit, args.max_results)
logger.info("")
logger.info("scanned %(scanned)d | already fine %(healthy)d | "
"repaired %(repaired)d | unrepairable %(unrepairable)d", totals)
if apply and totals["repaired"]:
invalidate_brand_overview_cache() # this process only
logger.info("")
logger.info("The running API still holds its own cached overview.")
logger.info("Refresh it with: GET /api/brands/overview?refresh=true")
finally:
conn.close()
return 0
if __name__ == "__main__":
raise SystemExit(main())