diff --git a/app/api/routers/brands.py b/app/api/routers/brands.py index 741ba88..0732216 100644 --- a/app/api/routers/brands.py +++ b/app/api/routers/brands.py @@ -86,6 +86,18 @@ def _row_to_product_out(row: dict, fallback_brand: str) -> ProductOut: if fssai is not None: fssai = str(fssai).strip() or None + # NUMERIC arrives as decimal.Decimal, which Pydantic would coerce but JSON + # would not. Converted here so 0.0 survives - `or None` would turn a real + # score of zero into "not scored", and zero is a valid health score. + def _score(key: str) -> Optional[float]: + raw = row.get(key) + if raw is None: + return None + try: + return float(raw) + except (TypeError, ValueError): + return None + return ProductOut( image_id=image_id, image_url=primary_url, @@ -108,6 +120,8 @@ def _row_to_product_out(row: dict, fallback_brand: str) -> ProductOut: selling_price=sp, barcode=bcd, barcode_type=bcd_type, + nutrition_score=_score("nutrition_score"), + health_score=_score("health_score"), ) diff --git a/app/api/routers/upload.py b/app/api/routers/upload.py index edb9a2e..419fa4b 100644 --- a/app/api/routers/upload.py +++ b/app/api/routers/upload.py @@ -569,6 +569,8 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]: skipped = 0 skipped_non_consumable = 0 status_counts = {"verified": 0, "partial": 0, "unavailable": 0} + # Which brands to mirror onto the brand tables once the import commits. + touched_brands: set = set() try: with conn, conn.cursor() as cur: @@ -752,6 +754,7 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]: ) imported_count += 1 + touched_brands.add(brand) except HTTPException: raise @@ -761,6 +764,14 @@ async def upload_nutrition_file(file: UploadFile = File(...)) -> Dict[str, Any]: finally: conn.close() + # Mirror onto the brand tables, so a consumer reading product rows directly + # sees an uploaded score too. After the commit and outside the try, because + # this must never turn a successful import into a 500 - sync_after_write + # swallows its own failures for the same reason. + if touched_brands: + from app.services.nutrition_score_sync import sync_after_write + sync_after_write(sorted(touched_brands)) + return _finish( file.filename, len(df), imported_count, errors, skipped, "nutritional intelligence items", diff --git a/app/api/schemas.py b/app/api/schemas.py index f2a78c4..355dd37 100644 --- a/app/api/schemas.py +++ b/app/api/schemas.py @@ -32,6 +32,10 @@ class ProductOut(BaseModel): selling_price: Optional[float] = None barcode: Optional[str] = None barcode_type: Optional[str] = None + # Mirrored from nutrition_insights onto the brand table by + # nutrition_score_sync. None means "not scored yet", never "scored zero". + nutrition_score: Optional[float] = None + health_score: Optional[float] = None class SourceProductOut(BaseModel): @@ -56,6 +60,10 @@ class SourceProductOut(BaseModel): selling_price: Optional[float] = None barcode: Optional[str] = None barcode_type: Optional[str] = None + # Mirrored from nutrition_insights onto the brand table by + # nutrition_score_sync. None means "not scored yet", never "scored zero". + nutrition_score: Optional[float] = None + health_score: Optional[float] = None similarity: float diff --git a/app/services/brand_sync.py b/app/services/brand_sync.py index 3ca9d32..50abe2e 100644 --- a/app/services/brand_sync.py +++ b/app/services/brand_sync.py @@ -76,12 +76,20 @@ def seed_catalog_paths(seed_dir: Path = SEED_DIR) -> List[Path]: # separately because it is huge and optional). Keeping these in sync is what # makes an exported file round-trip back through the seeder without losing # hsn/price/barcode/sku data - the failing mode of scripts/export_seed_data.py. +# +# `nutrition_score` and `health_score` are the exception to the first sentence: +# they are brand-table columns that upsert_brand_products deliberately does NOT +# write (see the comment on its INSERT). They are exported so a catalog file +# carries them, but they do NOT round-trip back into the database - re-seeding +# ignores them, and app/services/nutrition_score_sync.py is what restores them +# from nutrition_insights, which is their source of truth. EXPORT_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", + "nutrition_score", "health_score", ) diff --git a/app/services/nutrition_enrichment_service.py b/app/services/nutrition_enrichment_service.py index 17def6e..bdb5e55 100644 --- a/app/services/nutrition_enrichment_service.py +++ b/app/services/nutrition_enrichment_service.py @@ -213,6 +213,16 @@ def enrich_all_products( if progress_cb: progress_cb(i + 1, result.total_products) + # Mirror the new scores onto the brand tables, which is where consumers + # reading product rows directly will look for them. Scoped to the brands + # this run actually touched, so a --brands run does not sweep the whole + # catalogue. Imported here rather than at module scope to keep the + # nutrition service free of a load-time dependency on the brand-table + # layer, matching how _brands_to_enrich reaches vector_store. + if all_products: + from app.services.nutrition_score_sync import sync_after_write + sync_after_write(sorted({p["brand"] for p in all_products})) + result.duration_seconds = round(time.time() - start, 1) logger.info( f"Nutrition enrichment complete: {result.verified} verified, {result.partial} partial, " diff --git a/app/services/nutrition_score_sync.py b/app/services/nutrition_score_sync.py new file mode 100644 index 0000000..b391913 --- /dev/null +++ b/app/services/nutrition_score_sync.py @@ -0,0 +1,327 @@ +"""Mirror `nutrition_insights` scores onto the brand tables and the seed JSON. + +WHY THIS EXISTS +--------------- +`nutrition_score` and `health_score` live on `nutrition_insights`, keyed by +`(brand, image_id)`. Nothing joins that table to `brand_*`: the product card +fetches the two sides in two separate HTTP requests and merges them in the +browser. A consumer reading product rows straight out of Postgres therefore has +no way to see a score at all - which is exactly the report this module answers. + +So the two scores are denormalised onto every brand table. That buys one thing +and costs one thing: a plain `SELECT * FROM brand_cadbury` now carries the +score, and the copy can go stale. This module is what stops it going stale. + +THE SOURCE OF TRUTH IS `nutrition_insights`, ALWAYS +--------------------------------------------------- +Everything here is a one-way copy out of that table. Nothing writes a score +into `nutrition_insights` from a brand table, and a failed sync is never fatal +to the caller - a missing mirror is a display gap, a corrupted source table +would be data loss. + +WHY IT READS THE VALUE BACK INSTEAD OF BEING PASSED IT +------------------------------------------------------ +There are two score-write paths and they disagree about what "the value" is. +`nutrition_enrichment_service` passes a computed score straight through +`upsert_nutrition_insights`. The Excel path in `api/routers/upload.py` writes +`COALESCE(EXCLUDED.health_score, nutrition_insights.health_score)`, so a sheet +that omits a score leaves the previous one standing - the value stored is not +the value supplied. Mirroring what the caller passed would put a NULL on the +brand table while the real score sat in `nutrition_insights`. So the sync +always re-reads. + +WHY THE BRAND TABLE IS NOT THE JOIN KEY YOU EXPECT +--------------------------------------------------- +Brand tables have no `brand` column - the brand is the table name. The key +`nutrition_insights.brand` holds is `display_name_for_suffix(suffix)`, the +string `enrich_all_products` iterated, so that is what every statement below +binds. Matching on anything else silently returns zero rows. + +AND WHY THAT MATCH IS CASE-INSENSITIVE +--------------------------------------- +The two write paths disagree about capitalisation as well as value. Enrichment +stores the display name, `Cadbury`. The Excel path lowercases it first +(`api/routers/upload.py`, `brand = brand.lower()`), storing `cadbury`. An `=` +match on the display name would therefore mirror everything enrichment wrote +and silently skip everything an operator uploaded - the harder failure to +notice, because the table would look populated. Distinct display names stay +distinct when folded, so folding costs nothing. +""" +from __future__ import annotations + +import json +import logging +import os +import re +from decimal import Decimal +from pathlib import Path +from typing import Any, Dict, List, Optional, Tuple + +from app.services.vector_store import ( + _connect, + _list_brand_table_suffixes, + display_name_for_suffix, + ensure_brand_schema, +) + +logger = logging.getLogger(__name__) + +SCORE_COLUMNS: Tuple[str, str] = ("nutrition_score", "health_score") + +# Suffixes come from information_schema, so they are already real table names. +# Validated anyway because they are interpolated into SQL that cannot take them +# as a parameter. +_SAFE_TABLE = re.compile(r"^[A-Za-z0-9_]+$") + + +# --------------------------------------------------------------------------- +# brand tables +# --------------------------------------------------------------------------- + +def _target_suffixes(cur, brands: Optional[List[str]]) -> List[Tuple[str, str]]: + """[(table_suffix, display_brand)] for the brands this run covers.""" + suffixes = sorted(set(_list_brand_table_suffixes(cur, include_inactive=True))) + pairs = [(s, display_name_for_suffix(s)) for s in suffixes if _SAFE_TABLE.match(s)] + if not brands: + return pairs + wanted = {b.strip().casefold() for b in brands if b and b.strip()} + return [(s, d) for s, d in pairs + if d.casefold() in wanted or s.casefold() in wanted] + + +def sync_scores_to_brand_tables(brands: Optional[List[str]] = None, + *, dry_run: bool = False) -> Dict[str, int]: + """Copy scores from `nutrition_insights` onto each brand table. + + Returns {display_brand: rows_changed}, omitting brands with nothing to do. + + Two statements per table. The first copies current scores in; the second + clears scores whose insight row has gone, so this is a mirror rather than + an accumulator - without it a product removed by the non-consumable purge + keeps a health score on its brand row forever. + + `IS DISTINCT FROM` (not `!=`, which is never true against NULL) both makes + the run idempotent and makes `rowcount` mean "rows actually changed", so a + second run reporting zero is a real assertion and not an artefact. + """ + conn = _connect() + if conn is None: + logger.warning("nutrition score sync: no database connection") + return {} + + changed: Dict[str, int] = {} + try: + with conn.cursor() as cur: + targets = _target_suffixes(cur, brands) + + for suffix, display in targets: + table = f"brand_{suffix}" + # The columns may not exist yet on a table that has not been + # written since the migration landed; this is the migration. + if not dry_run: + ensure_brand_schema(display) + + try: + with conn.cursor() as cur: + if dry_run: + cur.execute( + f""" + SELECT count(*) FROM {table} b + JOIN nutrition_insights i + ON i.image_id = b.image_id AND lower(i.brand) = lower(%s) + WHERE b.nutrition_score IS DISTINCT FROM i.nutrition_score + OR b.health_score IS DISTINCT FROM i.health_score + """, + (display,), + ) + copied = cur.fetchone()[0] + cur.execute( + f""" + SELECT count(*) FROM {table} b + WHERE (b.nutrition_score IS NOT NULL + OR b.health_score IS NOT NULL) + AND NOT EXISTS ( + SELECT 1 FROM nutrition_insights i + WHERE lower(i.brand) = lower(%s) AND i.image_id = b.image_id) + """, + (display,), + ) + cleared = cur.fetchone()[0] + else: + cur.execute( + f""" + UPDATE {table} b + SET nutrition_score = i.nutrition_score, + health_score = i.health_score + FROM nutrition_insights i + WHERE i.image_id = b.image_id + AND lower(i.brand) = lower(%s) + AND (b.nutrition_score IS DISTINCT FROM i.nutrition_score + OR b.health_score IS DISTINCT FROM i.health_score) + """, + (display,), + ) + copied = cur.rowcount + cur.execute( + f""" + UPDATE {table} b + SET nutrition_score = NULL, health_score = NULL + WHERE (b.nutrition_score IS NOT NULL + OR b.health_score IS NOT NULL) + AND NOT EXISTS ( + SELECT 1 FROM nutrition_insights i + WHERE lower(i.brand) = lower(%s) AND i.image_id = b.image_id) + """, + (display,), + ) + cleared = cur.rowcount + conn.commit() + except Exception as e: # noqa: BLE001 + # One unusable table (a legacy shape, a lock) must not abort + # the rest of the catalogue. + conn.rollback() + logger.warning("score sync failed for %s: %s", table, e) + continue + + total = (copied or 0) + (cleared or 0) + if total: + changed[display] = total + logger.info("%s %d row(s) in %s (%d copied, %d cleared)", + "would update" if dry_run else "updated", + total, table, copied or 0, cleared or 0) + finally: + conn.close() + + return changed + + +# --------------------------------------------------------------------------- +# seed catalog JSON +# --------------------------------------------------------------------------- + +def _score_index() -> Dict[Tuple[str, str], Tuple[Optional[float], Optional[float]]]: + """{(casefolded_brand, image_id): (nutrition_score, health_score)} + + Brand is folded for the same reason the SQL folds it: enrichment stores + `Cadbury` and the Excel path stores `cadbury`, and a file must match both. + """ + conn = _connect() + if conn is None: + return {} + try: + with conn.cursor() as cur: + cur.execute( + "SELECT brand, image_id, nutrition_score, health_score " + "FROM nutrition_insights" + ) + out: Dict[Tuple[str, str], Tuple[Optional[float], Optional[float]]] = {} + for brand, image_id, ns, hs in cur.fetchall(): + out[((brand or "").casefold(), image_id)] = ( + float(ns) if isinstance(ns, Decimal) else ns, + float(hs) if isinstance(hs, Decimal) else hs, + ) + return out + except Exception as e: # noqa: BLE001 + logger.warning("could not read nutrition_insights: %s", e) + return {} + finally: + conn.close() + + +def _canonical_brand(raw: str) -> str: + """A brand string as `nutrition_insights` spells it.""" + from app.services.brand_sync import brand_slug + return display_name_for_suffix(brand_slug(raw)) + + +def sync_scores_to_seed_files(brands: Optional[List[str]] = None, + *, dry_run: bool = False) -> Dict[str, int]: + """Set the two score keys on every product in every seed catalog. + + Returns {filename: products_changed}. + + Deliberately NOT `brand_sync.export_brand_to_seed_file`. That rebuilds each + product from EXPORT_COLUMNS and `upsert_products_into_catalog_file` then + does `products_list[idx] = clean` - a REPLACE, not a merge. Running it over + brand_catalog_cadbury.json, whose products carry 33 keys, would drop + validation_status, confidence_score, size, variant_key, price_analysis and + the rest on the floor. This patcher touches two keys and nothing else. + + A file's own `brand` decides its table, not its filename - and several + files can feed one table - so each product is resolved individually rather + than assuming the filename. + """ + from app.services.brand_sync import _read_catalog, seed_catalog_paths + + index = _score_index() + if not index: + logger.info("no scores to write into seed files") + return {} + + wanted = {b.strip().casefold() for b in (brands or []) if b and b.strip()} + changed: Dict[str, int] = {} + + for path in seed_catalog_paths(): + data = _read_catalog(path) + if not data: + continue + file_brand = data.get("brand") or "" + products: List[Dict[str, Any]] = data.get("products") or [] + + touched = 0 + for product in products: + image_id = product.get("image_id") + if not image_id: + continue + raw = product.get("brand") or product.get("brand_name") or file_brand + if not raw: + continue + display = _canonical_brand(raw) + if wanted and display.casefold() not in wanted: + continue + ns, hs = index.get((display.casefold(), image_id), (None, None)) + # Written even when None, so the key is always present and a + # consumer can tell "not scored" from "field missing". + if (product.get("nutrition_score") != ns + or product.get("health_score") != hs + or "nutrition_score" not in product + or "health_score" not in product): + product["nutrition_score"] = ns + product["health_score"] = hs + touched += 1 + + if not touched: + continue + changed[path.name] = touched + if dry_run: + logger.info("would update %d product(s) in %s", touched, path.name) + continue + + # Atomic, matching upsert_products_into_catalog_file: a crash must not + # leave a half-written catalog that fails to parse on the next boot. + payload = json.dumps(data, indent=2, ensure_ascii=False) + tmp_path = path.with_suffix(".json.tmp") + tmp_path.write_text(payload, encoding="utf-8") + os.replace(tmp_path, path) + logger.info("updated %d product(s) in %s", touched, path.name) + + return changed + + +# --------------------------------------------------------------------------- + +def sync_after_write(brands: Optional[List[str]] = None) -> None: + """Fire-and-forget mirror, for callers that have just written scores. + + Never raises. The mirror is a convenience; the caller's own write to + `nutrition_insights` has already succeeded and must not be reported as + failed because a copy did not land. + """ + try: + changed = sync_scores_to_brand_tables(brands) + if changed: + logger.info("mirrored scores onto %d brand table(s): %s", + len(changed), ", ".join(sorted(changed))) + except Exception as e: # noqa: BLE001 + logger.warning("score mirror failed (scores are still in " + "nutrition_insights): %s", e) diff --git a/app/services/vector_store.py b/app/services/vector_store.py index c6a9805..4a70cc9 100644 --- a/app/services/vector_store.py +++ b/app/services/vector_store.py @@ -93,7 +93,13 @@ def get_brand_table_ddl(brand: str) -> str: highlights TEXT[], nutrients TEXT[], search_query TEXT, - + + -- Nutrition scores, mirrored from nutrition_insights by + -- nutrition_score_sync. Deliberately NOT written by + -- upsert_brand_products - see the comment on its INSERT. + nutrition_score NUMERIC, + health_score NUMERIC, + -- Vector embedding for search embedding vector(384), @@ -157,6 +163,12 @@ def _ensure_columns(cur, table_name: str) -> None: "highlights": "TEXT[]", "nutrients": "TEXT[]", "search_query": "TEXT", + # Mirrored from nutrition_insights by nutrition_score_sync, never by + # the INSERT below. This dict is the only migration mechanism there + # is (no Alembic), so listing them here is what creates them on every + # existing brand table at the next write. + "nutrition_score": "NUMERIC", + "health_score": "NUMERIC", "embedding": "vector(384)", } cur.execute( @@ -175,7 +187,13 @@ def _ensure_columns(cur, table_name: str) -> None: logger.info(f"Added missing column '{col}' to {table_name}") # 2. Relax legacy NOT NULL constraints on columns not present in standard insert - inserted_cols = {"id", "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", "embedding", "created_at", "updated_at"} + # + # `nutrition_score` and `health_score` are listed even though the INSERT + # does not write them: they are current columns owned by + # nutrition_score_sync, not legacy debris from an older DDL, and this sweep + # must never claim them. (Both are nullable, so it would be a no-op today - + # the entry is here so it stays a no-op if that ever changes.) + inserted_cols = {"id", "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", "nutrition_score", "health_score", "embedding", "created_at", "updated_at"} for col, is_nullable, col_def in col_info: if col not in inserted_cols and is_nullable == 'NO' and col_def is None: cur.execute(f"ALTER TABLE {table_name} ALTER COLUMN {col} DROP NOT NULL") @@ -393,6 +411,21 @@ def upsert_brand_products(brand: str, products: List[Dict[str, Any]], cleanup: b try: with conn, conn.cursor() as cur: + # `nutrition_score` and `health_score` are OMITTED FROM THIS + # STATEMENT ON PURPOSE - do not "fix" it by adding them. + # + # DO UPDATE SET overwrites every column it names. The rows arriving + # here are built by _to_storage_row (store_catalog_pipeline.py), + # _build_product_dict (user_products.py) and the seed loader, none + # of which carry a score - so naming the score columns would set + # them to NULL on every re-seed, pipeline re-run and product edit, + # silently wiping the data. A column this statement never names is + # a column it cannot damage. + # + # They are written only by app/services/nutrition_score_sync.py, + # which mirrors nutrition_insights. Extra score keys present in an + # incoming product dict (get_products_by_brand is SELECT *, so + # stage_11's merge carries them back in) are simply ignored here. cur.executemany( f""" INSERT INTO {table_name} diff --git a/scripts/sync_nutrition_scores.py b/scripts/sync_nutrition_scores.py new file mode 100644 index 0000000..36b3900 --- /dev/null +++ b/scripts/sync_nutrition_scores.py @@ -0,0 +1,116 @@ +#!/usr/bin/env python3 +""" +Copy `nutrition_insights` scores onto the brand tables and the seed catalogs. + +WHY THIS EXISTS +--------------- +`nutrition_score` and `health_score` are stored on `nutrition_insights`, keyed +by (brand, image_id). Nothing joined that table to `brand_*`, so anyone reading +product rows straight out of Postgres - which is how the merchant console does +it - could not see a score at all. The two columns now exist on every brand +table; this is what fills them the first time. + +Afterwards it maintains itself: `enrich_all_products` and the `/upload/nutrition` +endpoint both call `nutrition_score_sync.sync_after_write` when they finish, so +the mirror stays within one run of the source. Run this by hand after anything +that edits `nutrition_insights` outside those two paths - notably +`purge_non_consumable_nutrition`, which deletes rows and would otherwise leave +stale scores behind on the brand tables. + +SAFE TO RE-RUN +-------------- +Every statement is guarded with `IS DISTINCT FROM`, so a second run changes +nothing and reports 0. That makes "0 rows changed" a real assertion that the +mirror is current, not an artefact of the script doing nothing. + +USAGE +----- + python -m scripts.sync_nutrition_scores # dry run + python -m scripts.sync_nutrition_scores --apply + python -m scripts.sync_nutrition_scores --apply --brands "Cadbury,Amul" + python -m scripts.sync_nutrition_scores --apply --db-only + python -m scripts.sync_nutrition_scores --apply --json-only + +Dry run by default. `--apply` writes. +""" +from __future__ import annotations + +import argparse +import logging +import sys +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parents[1])) + +from app.infrastructure.settings import DB_HOST, DB_NAME # noqa: E402 +from app.services.nutrition_score_sync import ( # noqa: E402 + sync_scores_to_brand_tables, + sync_scores_to_seed_files, +) + +logging.basicConfig(level=logging.INFO, format="%(message)s") +logger = logging.getLogger(__name__) + + +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("--brands", default=None, + help="comma-separated brand names to limit the run to") + parser.add_argument("--db-only", action="store_true", + help="brand tables only, leave the seed JSON alone") + parser.add_argument("--json-only", action="store_true", + help="seed JSON only, leave the database alone") + args = parser.parse_args() + + if args.db_only and args.json_only: + parser.error("--db-only and --json-only are mutually exclusive") + + apply = args.apply and not args.dry_run + brands = [b.strip() for b in (args.brands or "").split(",") if b.strip()] or None + + 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") + if brands: + logger.info("Limited to: %s", ", ".join(brands)) + logger.info("") + + total = 0 + + if not args.json_only: + logger.info("Brand tables") + logger.info("-" * 52) + changed = sync_scores_to_brand_tables(brands, dry_run=not apply) + for brand, count in sorted(changed.items(), key=lambda kv: -kv[1]): + logger.info(" %-38s %5d", brand, count) + if not changed: + logger.info(" nothing to change - the mirror is current") + total += sum(changed.values()) + logger.info("") + + if not args.db_only: + logger.info("Seed catalogs") + logger.info("-" * 52) + changed = sync_scores_to_seed_files(brands, dry_run=not apply) + for name, count in sorted(changed.items(), key=lambda kv: -kv[1]): + logger.info(" %-38s %5d", name, count) + if not changed: + logger.info(" nothing to change - the files are current") + total += sum(changed.values()) + logger.info("") + + if not apply and total: + logger.info("Re-run with --apply to write %d change(s).", total) + elif apply: + logger.info("Done. %d change(s) written.", total) + + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_brand_table_scores.py b/tests/test_brand_table_scores.py new file mode 100644 index 0000000..d08cb84 --- /dev/null +++ b/tests/test_brand_table_scores.py @@ -0,0 +1,437 @@ +"""The score columns on the brand tables, and the mirror that fills them. + +THE FAILURE THIS FILE EXISTS FOR +-------------------------------- +`nutrition_score` and `health_score` live on `nutrition_insights`. A consumer +reading product rows straight out of Postgres - the merchant console does +exactly this - could not see them at all, because nothing joined the two. + +The fix denormalises both onto every `brand_*` table. That creates a second +copy of a number, and the two ways a second copy goes wrong are what most of +this file pins: + +1. `upsert_brand_products` overwrites every column it names. If the score + columns were added to that INSERT, every re-seed and pipeline re-run would + set them to NULL, because the rows arriving there are built by + `_to_storage_row` and carry no score. So the statement must NEVER name + them, and `test_the_insert_does_not_name_the_score_columns` is the guard. + +2. The two score-write paths spell the brand differently - enrichment stores + `Cadbury`, the Excel upload stores `cadbury`. An `=` match would mirror one + and silently skip the other. + +No database is involved: `_connect` is a recorder, so the assertions are about +the exact SQL and parameters sent. +""" +from __future__ import annotations + +import json + +import pytest + +from app.services import brand_sync, nutrition_score_sync, vector_store + + +# --------------------------------------------------------------------------- +# The columns exist +# --------------------------------------------------------------------------- + +class MigrationCursor: + """A brand table that exists but has none of the current columns.""" + + def __init__(self): + self.statements = [] + + def execute(self, sql, params=None): + self.statements.append(" ".join(str(sql).split())) + + def fetchall(self): + return [] # no existing columns -> every column is missing + + +def test_the_migration_adds_both_columns_to_an_existing_table(): + """`_ensure_columns` is the only migration mechanism there is - no Alembic, + no migrations directory. A column it does not add is a column that never + appears on a brand table that already exists, which is all of them.""" + cur = MigrationCursor() + + vector_store._ensure_columns(cur, "brand_cadbury") + + assert ("ALTER TABLE brand_cadbury ADD COLUMN IF NOT EXISTS " + "nutrition_score NUMERIC") in cur.statements + assert ("ALTER TABLE brand_cadbury ADD COLUMN IF NOT EXISTS " + "health_score NUMERIC") in cur.statements + + +def test_the_migration_still_adds_the_columns_it_always_did(): + cur = MigrationCursor() + + vector_store._ensure_columns(cur, "brand_cadbury") + + joined = " | ".join(cur.statements) + for col in ("product_name", "image_id", "hsn_code", "barcode", + "selling_price", "search_query", "embedding"): + assert f"ADD COLUMN IF NOT EXISTS {col} " in joined + + +def test_a_brand_new_table_is_born_with_both_columns(): + ddl = vector_store.get_brand_table_ddl("Cadbury") + + assert "nutrition_score NUMERIC" in ddl + assert "health_score NUMERIC" in ddl + + +def test_the_scores_are_numeric_not_text(): + """A TEXT score sorts lexically: '9' > '80'. Any ORDER BY or range filter + the consumer writes would be quietly wrong.""" + ddl = vector_store.get_brand_table_ddl("Amul") + + assert "health_score TEXT" not in ddl + assert "nutrition_score TEXT" not in ddl + + +# --------------------------------------------------------------------------- +# The guard that the whole design rests on +# --------------------------------------------------------------------------- + +def _upsert_source() -> str: + import inspect + return " ".join(inspect.getsource(vector_store.upsert_brand_products).split()) + + +def test_the_insert_does_not_name_the_score_columns(): + """DO NOT "FIX" THIS BY ADDING THEM. + + `upsert_brand_products` writes rows built by `_to_storage_row`, + `_build_product_dict` and the seed loader, none of which carry a score. + Naming the score columns in the INSERT would therefore set them to NULL on + every re-seed, every pipeline re-run and every product edit - wiping the + data this feature exists to provide. A column the statement never names is + a column it cannot damage. + """ + src = _upsert_source() + + insert_start = src.index("INSERT INTO {table_name}") + statement = src[insert_start:src.index('""",', insert_start)] + + assert "nutrition_score" not in statement, ( + "upsert_brand_products must not write nutrition_score - it would NULL " + "it on every catalogue re-write. nutrition_score_sync owns this column." + ) + assert "health_score" not in statement, ( + "upsert_brand_products must not write health_score - it would NULL it " + "on every catalogue re-write. nutrition_score_sync owns this column." + ) + + +def test_the_insert_still_writes_the_columns_it_always_did(): + """The negative above must not have been achieved by breaking the write.""" + src = _upsert_source() + + for col in ("product_name", "image_id", "hsn_code", "barcode", + "selling_price", "search_query", "embedding"): + assert f"{col} = EXCLUDED.{col}" in src or f"({col}," in src or f" {col}," in src + + +def test_image_id_is_still_the_fifth_element_of_the_row_tuple(): + """`image_ids = [r[4] for r in rows]` is a hard-coded positional index. + Anything inserted before image_id in the tuple breaks the conflict target + silently, and products start overwriting each other.""" + src = _upsert_source() + + assert "image_ids = [r[4] for r in rows if r[4]]" in src + columns = src[src.index("INSERT INTO {table_name}"):] + columns = columns[columns.index("(") + 1:columns.index(")")] + names = [c.strip() for c in columns.split(",")] + assert names[4] == "image_id" + + +# --------------------------------------------------------------------------- +# The mirror: brand tables +# --------------------------------------------------------------------------- + +class FakeCursor: + def __init__(self, conn): + self.conn = conn + + def execute(self, sql, params=None): + self.conn.calls.append((" ".join(str(sql).split()), params)) + + def fetchone(self): + return (self.conn.next_count,) + + @property + def rowcount(self): + return self.conn.next_count + + def __enter__(self): + return self + + def __exit__(self, *exc): + return False + + +class FakeConn: + def __init__(self, next_count=1): + self.calls = [] + self.next_count = next_count + self.commits = 0 + self.rollbacks = 0 + + def cursor(self): + return FakeCursor(self) + + def commit(self): + self.commits += 1 + + def rollback(self): + self.rollbacks += 1 + + def close(self): + pass + + +@pytest.fixture +def mirror(monkeypatch): + """One brand table, `brand_cadbury`, displayed as 'Cadbury'.""" + conn = FakeConn() + monkeypatch.setattr(nutrition_score_sync, "_connect", lambda: conn) + monkeypatch.setattr(nutrition_score_sync, "ensure_brand_schema", lambda b: "") + monkeypatch.setattr(nutrition_score_sync, "_list_brand_table_suffixes", + lambda cur, include_inactive=False: ["cadbury"]) + monkeypatch.setattr(nutrition_score_sync, "display_name_for_suffix", + lambda s, brand_map=None: s.title()) + return conn + + +def test_the_mirror_copies_scores_onto_the_brand_table(mirror): + changed = nutrition_score_sync.sync_scores_to_brand_tables() + + assert changed == {"Cadbury": 2} # one copied + one cleared + copy_sql = mirror.calls[0][0] + assert "UPDATE brand_cadbury b" in copy_sql + assert "SET nutrition_score = i.nutrition_score" in copy_sql + assert "health_score = i.health_score" in copy_sql + assert mirror.commits == 1 + + +def test_the_brand_match_is_case_insensitive(mirror): + """Enrichment writes 'Cadbury', the Excel upload writes 'cadbury'. An `=` + match would mirror the first and silently skip the second - and the table + would still look populated, which is why this needs a test rather than a + comment.""" + nutrition_score_sync.sync_scores_to_brand_tables() + + for sql, params in mirror.calls: + assert "lower(i.brand) = lower(%s)" in sql + assert params == ("Cadbury",) + assert "i.brand = %s" not in sql + + +def test_the_mirror_uses_is_distinct_from_so_a_second_run_is_a_no_op(mirror): + """`!=` is never true against NULL, so an unscored row would be rewritten + on every run and `rowcount` would stop meaning 'rows actually changed'.""" + nutrition_score_sync.sync_scores_to_brand_tables() + + copy_sql = mirror.calls[0][0] + assert "IS DISTINCT FROM" in copy_sql + assert "!=" not in copy_sql + + +def test_a_score_whose_insight_row_has_gone_is_cleared(mirror): + """Otherwise a product deleted by the non-consumable purge keeps a health + score on its brand row forever - the mirror would be an accumulator.""" + nutrition_score_sync.sync_scores_to_brand_tables() + + clear_sql = mirror.calls[1][0] + assert "SET nutrition_score = NULL, health_score = NULL" in clear_sql + assert "NOT EXISTS" in clear_sql + + +def test_nothing_is_written_on_a_dry_run(mirror): + nutrition_score_sync.sync_scores_to_brand_tables(dry_run=True) + + assert mirror.commits == 0 + for sql, _ in mirror.calls: + assert sql.startswith("SELECT count(*)") + + +def test_a_brand_filter_narrows_the_run(mirror): + assert nutrition_score_sync.sync_scores_to_brand_tables(["Amul"]) == {} + assert mirror.calls == [] + + assert nutrition_score_sync.sync_scores_to_brand_tables(["Cadbury"]) + + +def test_the_table_suffix_is_also_accepted_as_a_brand(mirror): + """`sync_after_write` is called from the Excel path, which lowercases the + brand, so 'cadbury' must reach brand_cadbury.""" + assert nutrition_score_sync.sync_scores_to_brand_tables(["cadbury"]) + + +def test_one_broken_table_does_not_abort_the_catalogue(monkeypatch): + class Exploding(FakeConn): + """Hands out one working cursor (for the suffix listing), then fails on + the first table and recovers for the second.""" + + def __init__(self): + super().__init__() + self.handed_out = 0 + + def cursor(self): + self.handed_out += 1 + if self.handed_out == 2: + raise RuntimeError("table is locked") + return FakeCursor(self) + + conn = Exploding() + monkeypatch.setattr(nutrition_score_sync, "_connect", lambda: conn) + monkeypatch.setattr(nutrition_score_sync, "ensure_brand_schema", lambda b: "") + monkeypatch.setattr(nutrition_score_sync, "_list_brand_table_suffixes", + lambda cur, include_inactive=False: ["a", "b"]) + monkeypatch.setattr(nutrition_score_sync, "display_name_for_suffix", + lambda s, brand_map=None: s.title()) + + nutrition_score_sync.sync_scores_to_brand_tables() # must not raise + + +def test_sync_after_write_never_raises(monkeypatch): + """The mirror is a convenience. The caller's write to nutrition_insights + has already succeeded and must not be reported as failed.""" + def boom(*a, **k): + raise RuntimeError("database on fire") + + monkeypatch.setattr(nutrition_score_sync, "sync_scores_to_brand_tables", boom) + + nutrition_score_sync.sync_after_write(["Cadbury"]) + + +def test_no_connection_is_not_an_error(monkeypatch): + monkeypatch.setattr(nutrition_score_sync, "_connect", lambda: None) + + assert nutrition_score_sync.sync_scores_to_brand_tables() == {} + + +# --------------------------------------------------------------------------- +# The mirror: seed catalog JSON +# --------------------------------------------------------------------------- + +# A product carrying the keys the Cadbury catalog really has. The point of the +# test below is that all of them survive. +FULL_PRODUCT = { + "product_name": "Cadbury Oreo", + "title": "Cadbury Oreo 30g", + "image_id": "cadbury_cadbury_oreo_30g", + "brand_name": "Cadbury", + "category": "Biscuits", + "size": "30g", + "variant_key": "oreo-30g", + "validation_status": "passed", + "confidence_score": 0.91, + "validation_issues": [], + "price_analysis": {"median": 10}, + "best_deals": ["x"], + "primary_image": "http://example/i.jpg", + "image_source": "wikimedia", + "total_variants": 5, +} + + +@pytest.fixture +def seed(tmp_path, monkeypatch): + """One seed catalog on disk, with one fully-populated product.""" + path = tmp_path / "brand_catalog_cadbury.json" + path.write_text(json.dumps({ + "brand": "Cadbury", + "total_products": 1, + "products": [dict(FULL_PRODUCT)], + }, indent=2), encoding="utf-8") + + monkeypatch.setattr(brand_sync, "seed_catalog_paths", lambda *a, **k: [path]) + monkeypatch.setattr(nutrition_score_sync, "_canonical_brand", lambda raw: "Cadbury") + monkeypatch.setattr(nutrition_score_sync, "_score_index", + lambda: {("cadbury", "cadbury_cadbury_oreo_30g"): (13.2, 11.2)}) + return path + + +def _products(path): + return json.loads(path.read_text(encoding="utf-8"))["products"] + + +def test_the_scores_are_written_into_the_seed_file(seed): + changed = nutrition_score_sync.sync_scores_to_seed_files() + + assert changed == {"brand_catalog_cadbury.json": 1} + product = _products(seed)[0] + assert product["nutrition_score"] == 13.2 + assert product["health_score"] == 11.2 + + +def test_every_other_key_survives(seed): + """THE REGRESSION GUARD FOR THE OBVIOUS SHORTCUT. + + `export_brand_to_seed_file` would have been the easy way to do this, but + `upsert_products_into_catalog_file` does `products_list[idx] = clean` - a + REPLACE - and rebuilds each product from EXPORT_COLUMNS. Running it over + this file would drop validation_status, confidence_score, size, + variant_key, price_analysis and the rest on the floor. + """ + nutrition_score_sync.sync_scores_to_seed_files() + + product = _products(seed)[0] + for key, value in FULL_PRODUCT.items(): + assert product[key] == value, f"{key} was lost or changed" + + assert set(product) == set(FULL_PRODUCT) | {"nutrition_score", "health_score"} + + +def test_an_unscored_product_gets_an_explicit_null(seed, monkeypatch): + """Present-and-null, not absent. A consumer must be able to tell 'we have + not scored this' from 'this build of the file predates the field'.""" + monkeypatch.setattr(nutrition_score_sync, "_score_index", + lambda: {("cadbury", "something_else"): (1.0, 2.0)}) + + nutrition_score_sync.sync_scores_to_seed_files() + + product = _products(seed)[0] + assert product["nutrition_score"] is None + assert product["health_score"] is None + + +def test_a_dry_run_leaves_the_file_untouched(seed): + before = seed.read_text(encoding="utf-8") + + changed = nutrition_score_sync.sync_scores_to_seed_files(dry_run=True) + + assert changed == {"brand_catalog_cadbury.json": 1} + assert seed.read_text(encoding="utf-8") == before + + +def test_a_second_run_reports_nothing_to_do(seed): + nutrition_score_sync.sync_scores_to_seed_files() + + assert nutrition_score_sync.sync_scores_to_seed_files() == {} + + +def test_no_temp_file_is_left_behind(seed): + nutrition_score_sync.sync_scores_to_seed_files() + + assert list(seed.parent.glob("*.tmp")) == [] + + +# --------------------------------------------------------------------------- +# The export column list +# --------------------------------------------------------------------------- + +def test_the_export_list_carries_the_scores(): + assert "nutrition_score" in brand_sync.EXPORT_COLUMNS + assert "health_score" in brand_sync.EXPORT_COLUMNS + + +def test_the_export_list_did_not_lose_anything(): + """It is the round-trip contract - a name dropped here silently strips that + field out of every exported catalog.""" + for col in ("product_name", "image_id", "hsn_code", "final_selling_price", + "selling_price", "barcode", "barcode_type", "fssai_license", + "product_sku", "search_query"): + assert col in brand_sync.EXPORT_COLUMNS diff --git a/tests/test_own_products_fields.py b/tests/test_own_products_fields.py index e5a4a47..b3ea3db 100644 --- a/tests/test_own_products_fields.py +++ b/tests/test_own_products_fields.py @@ -168,6 +168,24 @@ def test_the_sheets_values_are_kept_and_nothing_else_is_invented(store): assert row[field] is None, f"{field} was invented: {row[field]!r}" +def test_the_pipeline_never_sends_a_nutrition_score(store): + """`brand_*` carries nutrition_score/health_score, but the pipeline is not + what fills them - `nutrition_score_sync` mirrors them out of + nutrition_insights after enrichment has computed one. + + The keys must be ABSENT from the row handed to upsert_brand_products, not + present-and-None. `upsert_brand_products` deliberately omits both columns + from its INSERT so a catalogue re-write cannot NULL a real score; a row + that started carrying them would be the first step towards someone + "fixing" that omission. + """ + _run(["Product Name", "Weight", "Selling Price"], [["Apple", "500g", 155]]) + + row = store[0][1] + assert "nutrition_score" not in row + assert "health_score" not in row + + def test_the_price_band_is_eight_percent_of_the_sheet_price(store): _run(["Product Name", "Weight", "Selling Price"], [["Tomato", "1kg", 100]])