Backend scores updates on products
This commit is contained in:
@@ -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",
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -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, "
|
||||
|
||||
327
app/services/nutrition_score_sync.py
Normal file
327
app/services/nutrition_score_sync.py
Normal file
@@ -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)
|
||||
@@ -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}
|
||||
|
||||
Reference in New Issue
Block a user