Catalog feature updates on column fields
This commit is contained in:
@@ -53,6 +53,11 @@ from app.services.vector_store import (
|
||||
get_products_by_brand,
|
||||
)
|
||||
from app.services.brand_sync import upsert_products_into_catalog_file
|
||||
from app.services.enrichment.catalog_consensus import (
|
||||
consensus_rows,
|
||||
consensus_value,
|
||||
fssai_for_brand,
|
||||
)
|
||||
from app.services.embeddings_service import embed_texts
|
||||
from app.services.s3_service import s3_service
|
||||
|
||||
@@ -409,12 +414,20 @@ class _PersistOutcome:
|
||||
unavailable: Optional[str] = None
|
||||
|
||||
|
||||
# How many of a brand's rows to read when working out what they agree on.
|
||||
# Bounded because this is a nicety, not the write: the largest brand table here
|
||||
# holds ~250 rows, so this reads all of them for every real brand while still
|
||||
# refusing to degenerate into the full `SELECT *` that once ran per uploaded
|
||||
# row. One query per brand, not per row - that part is load-bearing.
|
||||
_CONSENSUS_SAMPLE_LIMIT = 300
|
||||
|
||||
|
||||
def _brand_sample(brand_parent: str) -> Dict[str, Any]:
|
||||
"""The most recently updated product for a brand, used to inherit defaults.
|
||||
|
||||
`limit=1` matters: this used to be a full `SELECT *` of the brand table,
|
||||
executed once per uploaded row. A 200-row file against a brand with a few
|
||||
thousand products meant 200 full table reads before a single insert.
|
||||
Kept for the fields where one arbitrary sibling is a defensible default
|
||||
(category, price band, size). For `fssai_license` and `providers` it is
|
||||
NOT defensible - see `_brand_defaults`.
|
||||
"""
|
||||
try:
|
||||
existing = get_products_by_brand(brand_parent, limit=1)
|
||||
@@ -424,8 +437,51 @@ def _brand_sample(brand_parent: str) -> Dict[str, Any]:
|
||||
return existing[0] if existing else {}
|
||||
|
||||
|
||||
def _brand_defaults(brand_parent: str) -> Dict[str, Any]:
|
||||
"""One arbitrary sibling row, plus what the brand's rows actually AGREE on.
|
||||
|
||||
The distinction matters for exactly two fields.
|
||||
|
||||
`fssai_license` identifies the food business legally answerable for the
|
||||
product. Inheriting it from one arbitrary sibling is already weak; the code
|
||||
this replaces was worse - it fell back to the bare constant
|
||||
"10012042000244", which is Lion Dates' real registered licence, and stamped
|
||||
it onto any brand with no sample row. `scripts/merge_haldiram.py` exists
|
||||
because that reached production.
|
||||
|
||||
`providers` is a claim about where a product can be bought. The replaced
|
||||
default asserted all six of Amazon/Flipkart/BigBasket/Jiomart/Blinkit/Zepto
|
||||
for every product nobody had checked.
|
||||
|
||||
Consensus over the brand's own rows is the honest version of both, and it
|
||||
declines to answer when the rows disagree.
|
||||
"""
|
||||
# Two reads on purpose, and the split matters on a memory-capped host.
|
||||
#
|
||||
# `_brand_sample` is SELECT * limit 1 - one row, embedding and all, because
|
||||
# the fields it seeds (category, price band, size) need the whole row.
|
||||
#
|
||||
# The consensus read is 300 rows, so it takes only the two columns it
|
||||
# actually inspects. Measured: SELECT * over 244 rows costs 3.0 MB, of
|
||||
# which 4.7 KB per row is an embedding string nothing here reads. Two named
|
||||
# columns is roughly 50 KB for the same rows.
|
||||
rows = consensus_rows(brand_parent, ["fssai_license", "providers"],
|
||||
limit=_CONSENSUS_SAMPLE_LIMIT)
|
||||
|
||||
fssai, fssai_source = fssai_for_brand(brand_parent, rows)
|
||||
providers, _why = consensus_value("providers", rows)
|
||||
|
||||
return {
|
||||
"sample": _brand_sample(brand_parent),
|
||||
"fssai_license": fssai,
|
||||
"fssai_source": fssai_source,
|
||||
"providers": providers,
|
||||
}
|
||||
|
||||
|
||||
def _build_product_dict(req: AddProductRequest, brand_parent: str,
|
||||
sample_existing: Dict[str, Any]) -> Dict[str, Any]:
|
||||
sample_existing: Dict[str, Any],
|
||||
defaults: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
||||
"""Fill in everything the catalog needs that the user did not supply.
|
||||
|
||||
Pure apart from the optional S3 image lookup - no database access and no
|
||||
@@ -438,8 +494,22 @@ def _build_product_dict(req: AddProductRequest, brand_parent: str,
|
||||
raise ValueError(f"product name {product_name!r} has no letters or digits to identify it by")
|
||||
image_id = f"{brand_slug}_{product_slug}"
|
||||
|
||||
defaults = defaults or {}
|
||||
category = req.category or sample_existing.get("category") or "Health Foods"
|
||||
fssai_license = req.fssai_license or sample_existing.get("fssai_license") or "10012042000244"
|
||||
|
||||
# NO CONSTANT FALLBACK HERE, EVER.
|
||||
#
|
||||
# This line used to end `or "10012042000244"` - Lion Dates' real registered
|
||||
# FSSAI licence - so any brand without a sample row was silently attributed
|
||||
# to a food business that had never heard of the product. An FSSAI number
|
||||
# is who is legally answerable for what is in the packet; inventing one is
|
||||
# not a cosmetic default.
|
||||
#
|
||||
# The order now is: what the uploader supplied, then the curated brand
|
||||
# registry, then what the brand's own rows unambiguously agree on, then
|
||||
# NOTHING. A blank licence is a gap someone can fill; a confidently wrong
|
||||
# one is a liability nobody knows to look for.
|
||||
fssai_license = req.fssai_license or defaults.get("fssai_license") or None
|
||||
|
||||
description = req.description or (
|
||||
f"Introducing {product_name} from the trusted {brand_parent} brand. "
|
||||
@@ -465,7 +535,11 @@ def _build_product_dict(req: AddProductRequest, brand_parent: str,
|
||||
else:
|
||||
price_range = "₹100-250"
|
||||
|
||||
providers = req.providers or list(sample_existing.get("providers") or ["Amazon", "Flipkart", "BigBasket", "Jiomart", "Blinkit", "Zepto"])
|
||||
# Likewise no invented marketplace list. Claiming a product is stocked by
|
||||
# Amazon, Flipkart, BigBasket, Jiomart, Blinkit AND Zepto because nobody
|
||||
# checked is a false availability claim on every row it touches. Consensus
|
||||
# across the brand's own rows, or empty.
|
||||
providers = req.providers or list(defaults.get("providers") or [])
|
||||
highlights = req.highlights or list(sample_existing.get("highlights") or ["100% Quality Assurance", "Authentic Brand Product"])
|
||||
nutrients = req.nutrients or list(sample_existing.get("nutrients") or ["Energy - High", "Protein - Good Source"])
|
||||
|
||||
@@ -582,11 +656,13 @@ def _persist_products(items: List[Tuple[Optional[int], AddProductRequest]]) -> _
|
||||
# row) and one embedding call for the whole upload.
|
||||
built: "OrderedDict[str, List[Tuple[Optional[int], AddProductRequest, Dict[str, Any]]]]" = OrderedDict()
|
||||
for brand_parent, rows in groups.items():
|
||||
sample = _brand_sample(brand_parent)
|
||||
defaults = _brand_defaults(brand_parent)
|
||||
sample = defaults["sample"]
|
||||
prepared = []
|
||||
for row_number, req in rows:
|
||||
try:
|
||||
prepared.append((row_number, req, _build_product_dict(req, brand_parent, sample)))
|
||||
prepared.append((row_number, req,
|
||||
_build_product_dict(req, brand_parent, sample, defaults)))
|
||||
except Exception as exc: # noqa: BLE001 - one bad row, not the file
|
||||
outcome.failures.append({
|
||||
"row": row_number,
|
||||
|
||||
@@ -91,6 +91,8 @@ from app.services.generic_products import (
|
||||
)
|
||||
from app.services.embeddings_service import embed_texts
|
||||
from app.services.enrichment.barcode.stage import BarcodeEnrichmentStage
|
||||
from app.services.enrichment.barcode.identity_stage import BarcodeIdentityStage
|
||||
from app.services.enrichment.content.stage import ContentEnrichmentStage
|
||||
from app.services.enrichment.hsn_gst.stage import HsnGstEnrichmentStage
|
||||
from app.services.enrichment.pipeline import EnrichmentPipeline
|
||||
from app.services.product_validator import validate_catalog
|
||||
@@ -658,8 +660,21 @@ async def stages_8_9_enrichment(rows: List[Dict[str, Any]], brand: str) -> List[
|
||||
stages = []
|
||||
if ENABLE_BARCODE_LOOKUP:
|
||||
stages.append(BarcodeEnrichmentStage())
|
||||
# Unconditional, and deliberately not behind ENABLE_BARCODE_LOOKUP: it
|
||||
# performs no lookup. It validates whatever barcode the row already has -
|
||||
# which for a sheet-supplied one is the first check it ever gets - and
|
||||
# derives barcode_type/gtin/ean13/upc from those digits offline. Measured
|
||||
# before it existed: upc 0.0%, ean13 6.4%, gtin 8.7%, all computable from
|
||||
# the barcode sitting in the same row.
|
||||
stages.append(BarcodeIdentityStage())
|
||||
if ENABLE_HSN_GST_ENRICHMENT:
|
||||
stages.append(HsnGstEnrichmentStage())
|
||||
# Also unconditional and also offline. `highlights` and `nutrients` land
|
||||
# empty on EVERY upload today because this pipeline has no stage that
|
||||
# generates them - only the older brand-discovery path calls the
|
||||
# generators. Fills blanks only, and refuses to put nutrients on a
|
||||
# non-consumable.
|
||||
stages.append(ContentEnrichmentStage())
|
||||
if not stages:
|
||||
return rows
|
||||
return await EnrichmentPipeline(stages).run(rows, brand)
|
||||
@@ -736,9 +751,32 @@ def _to_storage_row(row: Dict[str, Any]) -> Dict[str, Any]:
|
||||
"selling_price": row.get("selling_price"),
|
||||
"barcode": row.get("barcode"),
|
||||
"barcode_type": row.get("barcode_type"),
|
||||
# THIS PROJECTION IS THE WHOLE POINT OF THIS FUNCTION. A key the
|
||||
# enrichment stages computed but that is not named here never reaches
|
||||
# the database, however correct the stage was and however many columns
|
||||
# exist to hold it.
|
||||
#
|
||||
# That is exactly what happened to the seven fields below and to the
|
||||
# HSN/GST figures: stage 8 and stage 9 computed them on every run and
|
||||
# this dict silently dropped them. Measured before the fix - upc 0.0%,
|
||||
# ean13 6.4%, gst_percent and tax_amount 0% of upload rows.
|
||||
#
|
||||
# Anything added to an enrichment stage from here on has to be added
|
||||
# here AND to vector_store's INSERT, or it goes nowhere.
|
||||
"gtin": row.get("gtin"),
|
||||
"ean13": row.get("ean13"),
|
||||
"upc": row.get("upc"),
|
||||
"barcode_source": row.get("barcode_source"),
|
||||
"barcode_verified": row.get("barcode_verified"),
|
||||
"barcode_lookup_status": row.get("barcode_lookup_status"),
|
||||
"barcode_last_updated": row.get("barcode_last_updated"),
|
||||
"gst_percent": row.get("gst_percent"),
|
||||
"tax_amount": row.get("tax_amount"),
|
||||
"hsn_gst_needs_review": row.get("hsn_gst_needs_review"),
|
||||
"highlights": list(row.get("highlights") or []),
|
||||
"nutrients": list(row.get("nutrients") or []),
|
||||
"search_query": search_query,
|
||||
"field_sources": dict(row.get("field_sources") or {}),
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -425,6 +425,22 @@ ENABLE_SKU_WEB_LOOKUP = _bool("ENABLE_SKU_WEB_LOOKUP", "false")
|
||||
ENABLE_BARCODE_LOOKUP = _bool("ENABLE_BARCODE_LOOKUP", "false")
|
||||
ENABLE_MANUFACTURER_SITE_LOOKUP = _bool("ENABLE_MANUFACTURER_SITE_LOOKUP", "false")
|
||||
|
||||
# Barcodes AFTER the upload settles, in bulk, on the enrichment job's own
|
||||
# thread. Default TRUE where ENABLE_BARCODE_LOOKUP above is false, and the
|
||||
# difference is cost, not appetite for risk:
|
||||
#
|
||||
# ENABLE_BARCODE_LOOKUP = one search request PER PRODUCT, inline, against an
|
||||
# endpoint capped at 10 requests/minute. A 200-row
|
||||
# upload is twenty minutes of a held request.
|
||||
# ENRICH_BARCODES_ON_UPLOAD = one corpus fetch PER BRAND (~5 requests total),
|
||||
# matched offline, after the uploader has their
|
||||
# result. Cost is per brand, not per row.
|
||||
#
|
||||
# It also has to run before the nutrition phase rather than beside it: a
|
||||
# barcode makes the nutrition lookup exact (0.95) instead of fuzzy (0.32), and
|
||||
# skip_if_verified means whichever lands first wins permanently.
|
||||
ENRICH_BARCODES_ON_UPLOAD = _bool("ENRICH_BARCODES_ON_UPLOAD", "true")
|
||||
|
||||
# Offline/deterministic stages - safe to leave on.
|
||||
ENABLE_HSN_GST_ENRICHMENT = _bool("ENABLE_HSN_GST_ENRICHMENT", "true")
|
||||
ENABLE_PRODUCT_VALIDATION = _bool("ENABLE_PRODUCT_VALIDATION", "true")
|
||||
|
||||
@@ -24,6 +24,7 @@ import json
|
||||
import logging
|
||||
import os
|
||||
from collections import defaultdict
|
||||
from datetime import date, datetime
|
||||
from decimal import Decimal
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
@@ -83,13 +84,27 @@ def seed_catalog_paths(seed_dir: Path = SEED_DIR) -> List[Path]:
|
||||
# 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.
|
||||
#
|
||||
# `nutrients_per_100g` is in the same category as the two scores: mirrored from
|
||||
# nutrition_facts by nutrition_score_sync, exported for readers, never
|
||||
# round-tripped back in.
|
||||
#
|
||||
# The barcode identity/provenance columns and the HSN/GST figures below ARE
|
||||
# round-tripped. They were absent from this tuple for as long as they were
|
||||
# absent from the INSERT, which is why scripts/backfill_barcodes_from_off.py
|
||||
# refuses to call export_brand_to_seed_file() - exporting used to silently
|
||||
# strip the nine barcode keys off every product. Listing them here is what
|
||||
# makes that helper safe to use again.
|
||||
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",
|
||||
"gtin", "ean13", "upc", "barcode_source", "barcode_verified",
|
||||
"barcode_lookup_status", "barcode_last_updated",
|
||||
"gst_percent", "tax_amount", "hsn_gst_needs_review",
|
||||
"highlights", "nutrients", "search_query", "field_sources",
|
||||
"nutrition_score", "health_score", "nutrients_per_100g",
|
||||
)
|
||||
|
||||
|
||||
@@ -249,8 +264,18 @@ def _jsonable(value: Any) -> Any:
|
||||
"""Coerce a psycopg row value into something json.dumps accepts."""
|
||||
if isinstance(value, Decimal):
|
||||
return float(value)
|
||||
# `barcode_last_updated` is a TIMESTAMP column, so psycopg hands back a
|
||||
# datetime, which json.dumps refuses. Emitted as an ISO-8601 string rather
|
||||
# than an epoch float so the seed file stays human-readable; the DB write
|
||||
# path accepts either (see vector_store._epoch_to_timestamp).
|
||||
if isinstance(value, (datetime, date)):
|
||||
return value.isoformat()
|
||||
if isinstance(value, (list, tuple)):
|
||||
return [_jsonable(v) for v in value]
|
||||
# JSONB (field_sources, nutrients_per_100g) arrives as a dict; recurse so
|
||||
# a Decimal nested inside a nutrient block does not break the dump.
|
||||
if isinstance(value, dict):
|
||||
return {k: _jsonable(v) for k, v in value.items()}
|
||||
return value
|
||||
|
||||
|
||||
|
||||
117
app/services/enrichment/barcode/identity_stage.py
Normal file
117
app/services/enrichment/barcode/identity_stage.py
Normal file
@@ -0,0 +1,117 @@
|
||||
"""Derives the rest of a product's barcode identity from the barcode itself.
|
||||
|
||||
WHY THIS IS A SEPARATE STAGE FROM BarcodeEnrichmentStage
|
||||
--------------------------------------------------------
|
||||
That stage FINDS a barcode, needs the network, and is off by default
|
||||
(`ENABLE_BARCODE_LOOKUP`, see settings.py:420-423 for why). This one FINDS
|
||||
NOTHING. It takes a barcode the row already has - typed into the merchant's
|
||||
spreadsheet, seeded from a catalog, or just located by the cascade - and fills
|
||||
in the fields that are pure arithmetic on those digits:
|
||||
|
||||
barcode_type from the length (classify_barcode_type)
|
||||
gtin the validated digits (a GTIN is what a barcode encodes)
|
||||
ean13 zero-padded UPC-A (to_ean13)
|
||||
upc the digits, for UPC-A only
|
||||
|
||||
There is no lookup, no host, no rate limit and no failure mode beyond "these
|
||||
digits are not a valid GTIN", so it needs no settings flag and costs nothing.
|
||||
|
||||
THE FAILURE IT ADDRESSES
|
||||
Measured against production on 2026-09-08: `upc` was 0.0% filled, `ean13`
|
||||
6.4%, `gtin` 8.7% - against `barcode` at 18.4%. Every one of those could
|
||||
have been computed from the barcode already sitting in the same row. They
|
||||
were not, because the only code that produced them was inside the disabled
|
||||
network cascade, and the writer dropped them anyway.
|
||||
|
||||
WHY IT RUNS AFTER THE LOOKUP STAGE
|
||||
So it also normalises whatever the cascade just found. The cascade already
|
||||
validates, but a sheet-supplied barcode never passes through
|
||||
`validate_barcode` at all today - it goes straight from the spreadsheet to
|
||||
the database. This stage is the first thing that checks those digits.
|
||||
|
||||
WHAT IT WILL NOT DO
|
||||
It will not correct, reformat or delete `barcode`. If the digits fail
|
||||
checksum validation the stage returns NOTHING, leaving the merchant's value
|
||||
exactly as typed - `enrichment/base.py`'s merge guard would refuse to blank
|
||||
it anyway, and silently "fixing" a barcode a shop supplied would be worse
|
||||
than leaving it visibly wrong. The failure is recorded in `field_sources`
|
||||
so the coverage report can surface it.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
from typing import Any, Dict
|
||||
|
||||
from app.services.enrichment.base import EnrichmentStage, StageOutcome
|
||||
from app.services.enrichment.barcode.models import BarcodeType
|
||||
from app.services.enrichment.barcode.validators import (
|
||||
classify_barcode_type,
|
||||
normalize_barcode,
|
||||
to_ean13,
|
||||
validate_barcode,
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class BarcodeIdentityStage(EnrichmentStage):
|
||||
"""Offline, deterministic, additive. Never raises, never erases."""
|
||||
|
||||
name = "barcode_identity"
|
||||
|
||||
async def enrich_one(self, product: Dict[str, Any], brand: str) -> StageOutcome:
|
||||
raw = product.get("barcode")
|
||||
if not str(raw or "").strip():
|
||||
return StageOutcome(stage_name=self.name, fields={})
|
||||
|
||||
code = validate_barcode(raw)
|
||||
if not code:
|
||||
# Not a GTIN. Say so in the provenance rather than in the data, and
|
||||
# leave `barcode` untouched.
|
||||
digits = normalize_barcode(raw)
|
||||
reason = (f"{len(digits)} digits is not a GTIN-8/12/13/14 length"
|
||||
if digits else "no digits in the value")
|
||||
return StageOutcome(
|
||||
stage_name=self.name,
|
||||
fields={"field_sources": {"barcode": {
|
||||
"method": "unvalidated",
|
||||
"source": product.get("barcode_source") or "sheet",
|
||||
"note": f"failed checksum/format validation: {reason}",
|
||||
}}},
|
||||
error=f"barcode {raw!r} failed validation: {reason}",
|
||||
)
|
||||
|
||||
barcode_type = classify_barcode_type(code)
|
||||
fields: Dict[str, Any] = {
|
||||
"barcode": code, # normalised digits, same value
|
||||
"barcode_type": barcode_type.value,
|
||||
"gtin": code,
|
||||
"ean13": to_ean13(code), # None for GTIN-8, which is not a short EAN-13
|
||||
"upc": code if barcode_type is BarcodeType.UPC_A else None,
|
||||
}
|
||||
|
||||
# Only claim provenance we can stand behind. A barcode that arrived on
|
||||
# the sheet is the merchant's assertion, not ours, and is emphatically
|
||||
# not "verified" - that word is reserved for the cascade's
|
||||
# brand+size+name-matched result.
|
||||
if not str(product.get("barcode_source") or "").strip():
|
||||
fields["barcode_source"] = "sheet"
|
||||
fields["barcode_lookup_status"] = "sheet_validated"
|
||||
fields["barcode_verified"] = False
|
||||
fields["barcode_last_updated"] = time.time()
|
||||
|
||||
fields["field_sources"] = {
|
||||
"barcode": {
|
||||
"method": "sourced" if product.get("barcode_verified") else "asserted",
|
||||
"source": product.get("barcode_source") or "sheet",
|
||||
},
|
||||
# These four are arithmetic on the barcode, never a lookup. Calling
|
||||
# them "sourced" would overstate them.
|
||||
"gtin": {"method": "derived", "source": "validators.validate_barcode"},
|
||||
"ean13": {"method": "derived", "source": "validators.to_ean13"},
|
||||
"upc": {"method": "derived", "source": "validators.classify_barcode_type"},
|
||||
"barcode_type": {"method": "derived", "source": "validators.classify_barcode_type"},
|
||||
}
|
||||
|
||||
return StageOutcome(stage_name=self.name, fields=fields)
|
||||
@@ -114,9 +114,74 @@ def name_similarity(candidate_title: str, target_title: str) -> float:
|
||||
return round((overlap * 0.6) + (seq_ratio * 0.4), 3)
|
||||
|
||||
|
||||
def name_is_contained(candidate_title: str, target_title: str,
|
||||
target_brand: str = "") -> bool:
|
||||
"""True when the candidate's name is our name with only brand/size removed.
|
||||
|
||||
WHY THIS EXISTS - measured, not theoretical
|
||||
Open Food Facts stores short product names. We store long ones. Running
|
||||
`backfill_nutrition_from_barcodes` over the catalog on 2026-09-08, 149
|
||||
of 300 barcoded rows were rejected as "found, wrong product" when the
|
||||
barcode had resolved perfectly:
|
||||
|
||||
"Nestle Munch 8.9g" -> OFF "Munch" similarity 0.332
|
||||
"Coca-Cola Maaza 750ml" -> OFF "Maaza" similarity 0.304
|
||||
"Cadbury Perk 22 g" -> OFF "Perk" similarity 0.302
|
||||
|
||||
`name_similarity` divides the token overlap by the TARGET's token
|
||||
count, so a one-token candidate against a three-token target cannot
|
||||
exceed ~0.33 however right it is.
|
||||
|
||||
WHY NOT JUST LOWER THE THRESHOLD
|
||||
Because the same run also correctly rejected:
|
||||
|
||||
"Pepsico Lays 1kg" -> OFF "Spanish tomato tango" 0.133
|
||||
"Coca-Cola Fanta 750ml" -> OFF "Orange" 0.089
|
||||
"Lion Dates Powder 100g" -> OFF "PEPER NOTEN" 0.097
|
||||
|
||||
Those sit BELOW the containment cases but a threshold low enough to
|
||||
admit 0.30 also admits them. The measured yield table at
|
||||
settings.py:449-477 raised this floor to 0.78 for exactly that reason.
|
||||
Containment separates the two groups on structure rather than on a
|
||||
number: "Munch" is every token of our name minus brand and size;
|
||||
"Orange" is not a subset of "Coca-Cola Fanta 750ml" at all.
|
||||
|
||||
THE RULE
|
||||
Every token of the candidate's name must appear in the target's, once
|
||||
brand tokens and size tokens are discounted, and the candidate must
|
||||
carry at least one token that is not the brand. A bare brand name
|
||||
("Colgate", "godrej" - both real OFF titles) therefore does NOT match,
|
||||
which matters because those would otherwise attach to every product of
|
||||
that brand.
|
||||
"""
|
||||
cand_tokens = _tokens(candidate_title)
|
||||
target_tokens = _tokens(target_title)
|
||||
if not cand_tokens or not target_tokens:
|
||||
return False
|
||||
|
||||
brand_tokens = _tokens(target_brand)
|
||||
# A candidate that is only the brand identifies a brand, not a product.
|
||||
if not (cand_tokens - brand_tokens):
|
||||
return False
|
||||
|
||||
# Size tokens are not identity: our title carries the pack size, OFF's
|
||||
# usually does not, and `size_matches` has already checked the size
|
||||
# separately by the time this is consulted.
|
||||
def _meaningful(tokens):
|
||||
return {t for t in tokens if not _SIZE_TOKEN_RE.fullmatch(t)}
|
||||
|
||||
return _meaningful(cand_tokens) <= _meaningful(target_tokens | brand_tokens)
|
||||
|
||||
|
||||
# A token that is purely a quantity ("750ml", "8", "9g", "1kg"). Size is
|
||||
# compared by `size_matches`, so it must not also decide name identity.
|
||||
_SIZE_TOKEN_RE = re.compile(r"\d+(?:\.\d+)?(?:g|kg|ml|l|mg|cl|oz|gm|ltr|pcs|n)?", re.I)
|
||||
|
||||
|
||||
def is_match(candidate: BarcodeCandidate, target_brand: str, target_title: str, target_size: str,
|
||||
brand_aliases: Optional[Iterable[str]] = None,
|
||||
min_name_similarity: float = 0.45) -> tuple[bool, float]:
|
||||
min_name_similarity: float = 0.45,
|
||||
barcode_is_identity: bool = False) -> tuple[bool, float]:
|
||||
"""The combined gate a candidate must pass to be accepted:
|
||||
1. Brand matches (or overlaps a known alias).
|
||||
2. Pack size matches within a tight tolerance.
|
||||
@@ -126,16 +191,52 @@ def is_match(candidate: BarcodeCandidate, target_brand: str, target_title: str,
|
||||
product line from the same brand.
|
||||
Returns (matched, confidence) - confidence is diagnostic only, stored
|
||||
on the result for audit/QA but never used to override rule 1-3.
|
||||
|
||||
`barcode_is_identity` says the caller already knows WHICH product this is,
|
||||
because it looked the candidate up BY its GTIN rather than by searching.
|
||||
That changes what rules 2 and 4 are for: they stop being evidence of
|
||||
identity and become sanity checks against our barcode being on the wrong
|
||||
row. A sanity check cannot fail on information the source does not have, so
|
||||
under this flag:
|
||||
|
||||
* rule 2 (size) - a BLANK candidate size no longer vetoes. Open Food
|
||||
Facts leaves `quantity` null on a large share of records (57 of 146
|
||||
Amul hits), and `size_matches` returns False whenever either side is
|
||||
blank. A record with no quantity does not disagree with our pack size;
|
||||
it says nothing about it. A quantity that is PRESENT and different
|
||||
still vetoes - that is our barcode pointing at the wrong pack.
|
||||
* rule 4 (name) - see `name_is_contained`.
|
||||
|
||||
Rules 1 and 3 are unaffected: a different brand, or a "sugar free" the
|
||||
target does not have, still means a different product.
|
||||
|
||||
It defaults to False because every relaxation here is unsafe on the SEARCH
|
||||
path, where many candidates compete and name and size are the only things
|
||||
telling them apart - "Munch" with no size would match every Nestle product
|
||||
containing that word. Pass True only where a single candidate was fetched
|
||||
by barcode. Today that is `fetch_verified_nutrition_by_barcode` and
|
||||
`scripts/backfill_nutrition_from_barcodes`, and nothing else.
|
||||
|
||||
Measured on 2026-09-08: of 300 barcoded catalog rows, 149 were refused as
|
||||
"found, wrong product" with the barcode resolving perfectly. The name gate
|
||||
was the visible symptom, but the SIZE gate rejected most of them first.
|
||||
"""
|
||||
if not brand_matches(candidate.candidate_brand, target_brand, brand_aliases):
|
||||
return False, 0.0
|
||||
if not size_matches(candidate.candidate_size, target_size):
|
||||
# A blank candidate size is missing information, not a disagreement - but
|
||||
# only when the barcode already established identity. On the search path a
|
||||
# sizeless candidate is genuinely unidentifiable and must still be refused.
|
||||
size_unknown = barcode_is_identity and not str(candidate.candidate_size or "").strip()
|
||||
if not size_unknown and not size_matches(candidate.candidate_size, target_size):
|
||||
return False, 0.0
|
||||
if has_conflicting_variant_terms(candidate.candidate_title, target_title):
|
||||
return False, 0.0
|
||||
|
||||
similarity = name_similarity(candidate.candidate_title, target_title)
|
||||
if similarity < min_name_similarity:
|
||||
if barcode_is_identity and name_is_contained(
|
||||
candidate.candidate_title, target_title, target_brand):
|
||||
return True, similarity
|
||||
return False, similarity
|
||||
|
||||
return True, similarity
|
||||
|
||||
@@ -88,6 +88,25 @@ class EnrichmentStage(ABC):
|
||||
# HsnGstEnrichmentStage and BarcodeEnrichmentStage.
|
||||
if outcome.fields:
|
||||
for key, value in outcome.fields.items():
|
||||
# `field_sources` ACCUMULATES; every other key is assigned.
|
||||
#
|
||||
# It is a map keyed by column name, and each stage knows the
|
||||
# provenance of only the columns it filled. Assigning it like
|
||||
# anything else would mean the last stage to run erases what
|
||||
# every earlier stage recorded - so the barcode stage's
|
||||
# provenance would vanish the moment the HSN stage ran, and
|
||||
# the coverage report would show values with no origin.
|
||||
#
|
||||
# A shallow merge is the right depth: each key's value is one
|
||||
# flat record about one column. This mirrors the `||` in
|
||||
# vector_store's ON CONFLICT clause, so the in-memory merge
|
||||
# and the database merge agree.
|
||||
if key == "field_sources" and isinstance(value, dict):
|
||||
merged = dict(product.get("field_sources") or {})
|
||||
merged.update(value)
|
||||
product["field_sources"] = merged
|
||||
continue
|
||||
|
||||
blank_incoming = value is None or (isinstance(value, str) and not value.strip())
|
||||
existing = product.get(key)
|
||||
held = existing is not None and not (isinstance(existing, str) and not existing.strip())
|
||||
|
||||
142
app/services/enrichment/catalog_consensus.py
Normal file
142
app/services/enrichment/catalog_consensus.py
Normal file
@@ -0,0 +1,142 @@
|
||||
"""Fills a blank field from what the brand's OWN rows already agree on.
|
||||
|
||||
WHY THIS EXISTS
|
||||
`fssai_license` was 69.3% filled on 2026-09-08, sourced entirely from a
|
||||
hardcoded 34-brand map (`brand_registry.FSSAI_LICENSES`). A brand outside
|
||||
that map got nothing - except on one path, which got something far worse.
|
||||
|
||||
`user_products._build_product_dict` read:
|
||||
|
||||
fssai_license = req.fssai_license or sample_existing.get(...) or "10012042000244"
|
||||
|
||||
That constant is LION DATES' real, registered FSSAI licence. Any brand with
|
||||
no sample row was stamped with it. This is not a cosmetic default: an FSSAI
|
||||
number identifies the food business legally answerable for the product, and
|
||||
inventing one attributes a stranger's regulatory liability to a product they
|
||||
never made. `scripts/merge_haldiram.py:36` exists because this already
|
||||
reached production once, on `brand_haldirams`.
|
||||
|
||||
The honest source for a blank licence is the brand's own catalog: 400
|
||||
Britannia rows carrying one licence is good evidence for the 401st. That is
|
||||
what this module reads.
|
||||
|
||||
THE RULE IT ENFORCES
|
||||
Propagate only from UNAMBIGUOUS agreement. If a brand's rows carry two
|
||||
different licences, one of them is already wrong and this module returns
|
||||
None rather than picking. A blank field is a gap; a confidently wrong
|
||||
regulatory identifier is a liability.
|
||||
|
||||
Nothing here invents a value. Every result is a value already present on a
|
||||
row of the same brand, which is why the provenance method is
|
||||
`catalog_consensus` and never `sourced`.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from collections import Counter
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# A single dissenting row should not veto 400 agreeing ones, but a genuine
|
||||
# split must. Set so that "399 of 400 agree" propagates and "60/40" does not.
|
||||
_MIN_AGREEMENT = 0.85
|
||||
|
||||
# Below this many populated rows there is no consensus to speak of, only a
|
||||
# coincidence. Two rows agreeing proves nothing about a third.
|
||||
_MIN_ROWS = 3
|
||||
|
||||
|
||||
def _modal(values: List[Any], min_agreement: float = _MIN_AGREEMENT,
|
||||
min_rows: int = _MIN_ROWS) -> Tuple[Optional[Any], Dict[str, Any]]:
|
||||
"""The one value the population agrees on, or None with the reason why."""
|
||||
populated = [v for v in values if v not in (None, "", [], {})]
|
||||
if len(populated) < min_rows:
|
||||
return None, {"reason": "too few populated rows", "rows": len(populated)}
|
||||
|
||||
# Lists (providers) are unhashable; compare them as ordered tuples.
|
||||
keyed = [tuple(v) if isinstance(v, list) else v for v in populated]
|
||||
counts = Counter(keyed)
|
||||
winner, hits = counts.most_common(1)[0]
|
||||
agreement = hits / len(keyed)
|
||||
if agreement < min_agreement:
|
||||
return None, {"reason": "no clear majority", "agreement": round(agreement, 3),
|
||||
"distinct": len(counts)}
|
||||
|
||||
return (list(winner) if isinstance(winner, tuple) else winner), {
|
||||
"agreement": round(agreement, 3), "rows": len(keyed)}
|
||||
|
||||
|
||||
def consensus_value(column: str, rows: List[Dict[str, Any]],
|
||||
min_agreement: float = _MIN_AGREEMENT
|
||||
) -> Tuple[Optional[Any], Dict[str, Any]]:
|
||||
"""The agreed value of `column` across `rows`, plus why it was or was not
|
||||
reached. Never raises: an unreadable row set yields (None, reason)."""
|
||||
try:
|
||||
return _modal([r.get(column) for r in rows], min_agreement=min_agreement)
|
||||
except Exception as e: # pragma: no cover - defensive
|
||||
logger.debug("consensus for %s failed: %s", column, e)
|
||||
return None, {"reason": f"error: {e}"}
|
||||
|
||||
|
||||
def consensus_rows(brand: str, columns: List[str], limit: int = 300) -> List[Dict[str, Any]]:
|
||||
"""Read only the columns consensus needs, for a brand's rows.
|
||||
|
||||
Deliberately NOT `get_products_by_brand`, which is `SELECT *` and therefore
|
||||
carries the 384-dimension embedding on every row. Measured against
|
||||
production: 244 Hindustan Unilever rows cost 3.0 MB that way, 4.7 KB of the
|
||||
7.2 KB per row being an embedding string nothing here looks at.
|
||||
|
||||
Two named columns bring the same read down to roughly 50 KB. On a backend
|
||||
container capped at 2560 MB that difference is not dangerous either way -
|
||||
it is just the difference between reading what is needed and reading
|
||||
everything, once per brand per upload.
|
||||
|
||||
Probes `information_schema` first, because the column set genuinely differs
|
||||
between brand tables and a missing column would otherwise raise.
|
||||
"""
|
||||
from app.services.vector_store import _connect, _sanitize_name
|
||||
|
||||
conn = _connect()
|
||||
if conn is None:
|
||||
return []
|
||||
table = f"brand_{_sanitize_name(brand)}"
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
cur.execute(
|
||||
"SELECT column_name FROM information_schema.columns "
|
||||
"WHERE table_schema = 'public' AND table_name = %s",
|
||||
(table,),
|
||||
)
|
||||
present = {r[0] for r in cur.fetchall()}
|
||||
wanted = [c for c in columns if c in present]
|
||||
if not wanted:
|
||||
return []
|
||||
select = ", ".join(f'"{c}"' for c in wanted)
|
||||
cur.execute(f'SELECT {select} FROM "{table}" LIMIT %s', (limit,))
|
||||
return [dict(zip(wanted, row)) for row in cur.fetchall()]
|
||||
except Exception as e: # noqa: BLE001 - defaults are a nicety, not the write
|
||||
logger.debug("consensus read failed for %s: %s", brand, e)
|
||||
return []
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def fssai_for_brand(brand: str, rows: Optional[List[Dict[str, Any]]] = None) -> Tuple[Optional[str], str]:
|
||||
"""The licence to use for a new product of `brand`, and where it came from.
|
||||
|
||||
Order: the curated registry map, then the brand's own rows. Never a
|
||||
constant, never another brand's number.
|
||||
"""
|
||||
from app.services.brand_registry import get_fssai_license
|
||||
|
||||
mapped = get_fssai_license(brand)
|
||||
if mapped:
|
||||
return mapped, "brand_registry"
|
||||
|
||||
if rows:
|
||||
value, _why = consensus_value("fssai_license", rows)
|
||||
if value:
|
||||
return str(value), "catalog_consensus"
|
||||
|
||||
return None, "unknown"
|
||||
1
app/services/enrichment/content/__init__.py
Normal file
1
app/services/enrichment/content/__init__.py
Normal file
@@ -0,0 +1 @@
|
||||
"""Offline content enrichment - the display columns the store pipeline left blank."""
|
||||
112
app/services/enrichment/content/stage.py
Normal file
112
app/services/enrichment/content/stage.py
Normal file
@@ -0,0 +1,112 @@
|
||||
"""Fills `highlights` and `nutrients` for rows the store pipeline leaves empty.
|
||||
|
||||
THE FAILURE THIS ADDRESSES
|
||||
`catalog_engine.generate_product_highlights` and `generate_nutrients_info`
|
||||
have existed for a long time and `brand_discovery._build_product` calls
|
||||
both. The store-catalog pipeline never did: `_to_storage_row` simply passed
|
||||
whatever the sheet had through, so a colleague's upload - which carries
|
||||
neither column - landed `highlights=[]` and `nutrients=[]` on every row.
|
||||
|
||||
That is the whole reason those two columns look healthy in aggregate
|
||||
(95.3% / 69.7% on 2026-09-08) while being empty for exactly the rows this
|
||||
work is about.
|
||||
|
||||
WHAT IT WRITES, AND HOW HONESTLY
|
||||
`highlights` is marketing copy derived from fields we already hold - the
|
||||
category, the pack size, the brand. It is `derived`, never `sourced`.
|
||||
|
||||
`nutrients` is the display list. Where real per-100g figures exist,
|
||||
`nutrition_score_sync.sync_nutrients_to_brand_tables` renders them from
|
||||
`nutrition_facts` and overwrites whatever this stage wrote - that mirror is
|
||||
the better source and runs later. This stage only supplies the
|
||||
category-keyword fallback, flagged `estimated`, so a row is not blank while
|
||||
it waits for a nutrition lookup that may never succeed.
|
||||
|
||||
THE CONSUMABILITY GATE
|
||||
`generate_nutrients_info` works off category keywords, so a Hair Care row
|
||||
whose category or description happens to contain a matching word acquires
|
||||
entries like "Vitamin B Complex - Energy". Shampoo has no nutrients. This
|
||||
stage refuses to write the column at all for a non-consumable, which is the
|
||||
same gate `nutrition_data_service` applies on the lookup path and the same
|
||||
reason `purge_non_consumable_nutrition.py` had to exist.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import Any, Dict, List
|
||||
|
||||
from app.services.enrichment.base import EnrichmentStage, StageOutcome
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class ContentEnrichmentStage(EnrichmentStage):
|
||||
"""Offline, deterministic, fills blanks only. Never raises, never erases."""
|
||||
|
||||
name = "content"
|
||||
|
||||
async def enrich_one(self, product: Dict[str, Any], brand: str) -> StageOutcome:
|
||||
# Imported lazily: catalog_engine pulls in the image and LLM services,
|
||||
# and this stage runs inside ingestion where those are already loaded
|
||||
# but the enrichment package on its own should not require them.
|
||||
from app.core.catalog_engine import (
|
||||
generate_nutrients_info,
|
||||
generate_product_highlights,
|
||||
)
|
||||
from app.services.consumability import is_non_consumable
|
||||
|
||||
fields: Dict[str, Any] = {}
|
||||
sources: Dict[str, Any] = {}
|
||||
|
||||
title = product.get("title") or product.get("product_name") or ""
|
||||
category = product.get("category") or ""
|
||||
|
||||
if not _has_entries(product.get("highlights")):
|
||||
try:
|
||||
highlights = generate_product_highlights(product, brand)
|
||||
except Exception as e: # never abort a row
|
||||
logger.debug("highlight generation failed for %r: %s", title, e)
|
||||
highlights = []
|
||||
if highlights:
|
||||
fields["highlights"] = highlights
|
||||
sources["highlights"] = {"method": "derived",
|
||||
"source": "catalog_engine.generate_product_highlights"}
|
||||
|
||||
if not _has_entries(product.get("nutrients")):
|
||||
if is_non_consumable(category, title):
|
||||
# Not a gap - a column that cannot apply. Recording it stops
|
||||
# the coverage report counting shampoo as missing nutrition
|
||||
# forever, which is what makes someone eventually fabricate it.
|
||||
sources["nutrients"] = {"method": "not_applicable",
|
||||
"source": "non_consumable_product"}
|
||||
else:
|
||||
try:
|
||||
nutrients = generate_nutrients_info(product, brand)
|
||||
except Exception as e:
|
||||
logger.debug("nutrient generation failed for %r: %s", title, e)
|
||||
nutrients = []
|
||||
if nutrients:
|
||||
fields["nutrients"] = nutrients
|
||||
sources["nutrients"] = {
|
||||
"method": "estimated",
|
||||
"source": "catalog_engine.generate_nutrients_info",
|
||||
"note": "category keywords; replaced by real per-100g "
|
||||
"figures when a nutrition lookup succeeds",
|
||||
}
|
||||
|
||||
if sources:
|
||||
fields["field_sources"] = sources
|
||||
|
||||
return StageOutcome(stage_name=self.name, fields=fields)
|
||||
|
||||
|
||||
def _has_entries(value: Any) -> bool:
|
||||
"""True when the column already carries something worth keeping.
|
||||
|
||||
A list of empty strings counts as empty: the spreadsheet parser produces
|
||||
those from a column that exists but has no value in it, and treating one as
|
||||
"already filled" is how a row keeps `['']` forever.
|
||||
"""
|
||||
if not isinstance(value, (list, tuple)):
|
||||
return bool(value)
|
||||
return any(str(v).strip() for v in value)
|
||||
@@ -60,6 +60,17 @@ def _build_default_stages() -> List[EnrichmentStage]:
|
||||
except Exception as e:
|
||||
logger.error(f"Barcode enrichment stage unavailable: {e}")
|
||||
|
||||
# Runs AFTER the lookup so it normalises whatever that found, and runs at
|
||||
# all even when the lookup is disabled - which is the point. It derives
|
||||
# barcode_type/gtin/ean13/upc from a barcode the row already has, offline
|
||||
# and for free, so a sheet-supplied barcode finally gets validated and
|
||||
# expanded instead of going straight to the database unchecked.
|
||||
try:
|
||||
from app.services.enrichment.barcode.identity_stage import BarcodeIdentityStage
|
||||
stages.append(BarcodeIdentityStage())
|
||||
except Exception as e:
|
||||
logger.error(f"Barcode identity stage unavailable: {e}")
|
||||
|
||||
# HSN / GST & pricing enrichment (see app/services/enrichment/hsn_gst/) -
|
||||
# deterministic, offline, pure-additive. Runs AFTER the barcode stage so
|
||||
# every stored/exported row carries both sets of fields; a failure here
|
||||
@@ -70,9 +81,15 @@ def _build_default_stages() -> List[EnrichmentStage]:
|
||||
except Exception as e:
|
||||
logger.error(f"HSN/GST enrichment stage unavailable: {e}")
|
||||
|
||||
# Future stages register here, e.g.:
|
||||
# from app.services.enrichment.nutrition.stage import NutritionEnrichmentStage
|
||||
# stages.append(NutritionEnrichmentStage())
|
||||
# Offline display columns. Registered last so the nutrients fallback it
|
||||
# writes is the lowest-priority source: the real per-100g figures mirrored
|
||||
# by nutrition_score_sync overwrite it whenever a lookup succeeds.
|
||||
try:
|
||||
from app.services.enrichment.content.stage import ContentEnrichmentStage
|
||||
stages.append(ContentEnrichmentStage())
|
||||
except Exception as e:
|
||||
logger.error(f"Content enrichment stage unavailable: {e}")
|
||||
|
||||
return stages
|
||||
|
||||
|
||||
|
||||
254
app/services/enrichment/post_ingest_barcodes.py
Normal file
254
app/services/enrichment/post_ingest_barcodes.py
Normal file
@@ -0,0 +1,254 @@
|
||||
"""Finds barcodes for freshly-ingested rows, in bulk, before nutrition runs.
|
||||
|
||||
WHY THIS RUNS BEFORE THE NUTRITION JOB, NOT ALONGSIDE IT
|
||||
--------------------------------------------------------
|
||||
Ordering here is a correctness property, not a preference.
|
||||
|
||||
`fetch_verified_nutrition_by_barcode` matches on the GTIN and returns at
|
||||
confidence 0.95. The name search it falls back to accepts at a minimum of 0.32.
|
||||
The nutrition job runs with `skip_if_verified=True`, so whichever path lands
|
||||
first WINS PERMANENTLY - a 0.32 name match blocks the 0.95 barcode match from
|
||||
ever being attempted. Two jobs racing would produce exactly that, silently, and
|
||||
the catalog would end up with the worse of two available answers.
|
||||
|
||||
So this is a phase inside the same job, ahead of the nutrition phases.
|
||||
|
||||
WHY BULK, NOT THE PER-PRODUCT CASCADE
|
||||
-------------------------------------
|
||||
The per-product search endpoint Open Food Facts exposes is capped at 10
|
||||
requests per minute. A 200-row upload is twenty minutes of waiting, which is
|
||||
why `ENABLE_BARCODE_LOOKUP` defaults false (settings.py:420-423) and why the
|
||||
inline stage stays off.
|
||||
|
||||
`off_bulk.fetch_brand_corpus` fetches a brand's ENTIRE Open Food Facts
|
||||
catalogue in about five requests and matches offline against it. A brand is a
|
||||
brand whether it has 3 rows or 300, so the cost is per brand, not per product.
|
||||
That is what makes barcode enrichment affordable on the shared host at all.
|
||||
|
||||
It also uses `off_bulk.score_candidates` rather than the live matcher, because
|
||||
that scorer already fixes two measured flaws: `matching.name_similarity` is
|
||||
asymmetric ("Butter milk amul" vs "Amul Butter" scores 0.882 one way and 0.418
|
||||
the other), and `matching.size_matches` vetoes any candidate with a blank size
|
||||
when 57 of 146 Amul OFF records have `quantity: null`.
|
||||
|
||||
THE THRESHOLD IS 0.88, NOT 0.78
|
||||
-------------------------------
|
||||
`BARCODE_MIN_NAME_SIMILARITY` (0.78) is the floor for the REVERSE direction,
|
||||
where a barcode has already established identity and the name is a sanity
|
||||
check. This is the forward direction: many candidates compete and the name
|
||||
carries the whole decision. `scripts/backfill_barcodes_from_off.py` measured
|
||||
0.88 as the safe auto-apply point and 0.70-0.88 as review-only, and this reuses
|
||||
that number rather than inventing one.
|
||||
|
||||
WHAT IT WRITES
|
||||
barcode, barcode_type, gtin, ean13, barcode_source, barcode_verified=False,
|
||||
barcode_lookup_status='name_matched', barcode_last_updated, and the
|
||||
field_sources record - via targeted UPDATEs that pin the row's current
|
||||
value, never via upsert_brand_products.
|
||||
|
||||
WHY NOT THE UPSERT
|
||||
`get_products_by_brand` is `SELECT *`, so a row's `embedding` comes back as
|
||||
a pgvector string, and `upsert_brand_products` only accepts a list - it
|
||||
would write NULL and destroy the embedding. Targeted UPDATEs also cannot
|
||||
clobber a concurrent write.
|
||||
|
||||
WHAT IT WILL NOT DO
|
||||
It never overwrites a barcode a row already holds. A merchant typing one in
|
||||
is holding the pack; nothing found by name similarity outranks that.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import time
|
||||
from typing import Any, Dict, Iterable, List, Optional, Tuple
|
||||
|
||||
from psycopg.types.json import Json
|
||||
|
||||
from app.services.enrichment.barcode.validators import (
|
||||
classify_barcode_type,
|
||||
to_ean13,
|
||||
validate_barcode,
|
||||
)
|
||||
from app.services.vector_store import _connect, _sanitize_name
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Matches scripts/backfill_barcodes_from_off.py, which measured it.
|
||||
DEFAULT_MIN_SIMILARITY = 0.88
|
||||
|
||||
BARCODE_SOURCE = "openfoodfacts_bulk (search.openfoodfacts.org)"
|
||||
|
||||
|
||||
def _rows_needing_a_barcode(cur, table: str) -> List[Dict[str, Any]]:
|
||||
"""Rows with no usable barcode. Probes the column list first, because a
|
||||
table written before the schema migration may still lack the newer ones."""
|
||||
cur.execute(
|
||||
"SELECT column_name FROM information_schema.columns "
|
||||
"WHERE table_schema = 'public' AND table_name = %s",
|
||||
(table,),
|
||||
)
|
||||
present = {r[0] for r in cur.fetchall()}
|
||||
if not {"id", "product_name", "barcode"} <= present:
|
||||
return []
|
||||
|
||||
size = "size" if "size" in present else "NULL AS size"
|
||||
cur.execute(
|
||||
f'SELECT id, product_name, title, {size}, category, barcode '
|
||||
f'FROM "{table}" '
|
||||
f"WHERE barcode IS NULL OR btrim(barcode) = ''"
|
||||
)
|
||||
return [{"id": r[0], "product_name": r[1], "title": r[2], "size": r[3],
|
||||
"category": r[4], "barcode": r[5]} for r in cur.fetchall()]
|
||||
|
||||
|
||||
def enrich_brand_barcodes(brand: str, *, min_similarity: float = DEFAULT_MIN_SIMILARITY,
|
||||
dry_run: bool = False,
|
||||
progress_cb=None) -> Dict[str, int]:
|
||||
"""Fill blank barcodes for one brand from its Open Food Facts corpus.
|
||||
|
||||
Never raises: a brand whose corpus cannot be fetched reports zero and the
|
||||
caller moves to the next one. Enrichment is best-effort by contract.
|
||||
"""
|
||||
from app.services.enrichment.barcode.sources.off_bulk import (
|
||||
brand_tokens,
|
||||
fetch_brand_corpus,
|
||||
score_candidates,
|
||||
)
|
||||
|
||||
stats = {"candidates": 0, "matched": 0, "written": 0, "rejected": 0}
|
||||
|
||||
# The same refusal `store_catalog_pipeline.stages_8_9_enrichment` makes for
|
||||
# the inline stages, for the same reason: a barcode identifies a
|
||||
# manufactured article and the Own Products bucket is loose produce - an
|
||||
# apple, a bunch of coriander. There is no GTIN to find, and a name match
|
||||
# against some packaged product's corpus could only attach the wrong one.
|
||||
from app.services.generic_products import OWN_PRODUCTS_BRAND
|
||||
if brand == OWN_PRODUCTS_BRAND:
|
||||
return stats
|
||||
|
||||
table = f"brand_{_sanitize_name(brand)}"
|
||||
|
||||
conn = _connect()
|
||||
if conn is None:
|
||||
return stats
|
||||
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
rows = _rows_needing_a_barcode(cur, table)
|
||||
if not rows:
|
||||
return stats
|
||||
stats["candidates"] = len(rows)
|
||||
|
||||
try:
|
||||
# Returns the hit LIST directly (the on-disk cache file wraps it in
|
||||
# a "hits" key; the function unwraps it). About five requests for a
|
||||
# whole brand, then served from disk on later runs.
|
||||
hits = fetch_brand_corpus(brand) or []
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("OFF corpus unavailable for %s: %s", brand, e)
|
||||
return stats
|
||||
|
||||
if not hits:
|
||||
logger.info("Open Food Facts holds no India catalogue for %s", brand)
|
||||
return stats
|
||||
|
||||
# Stripped from both sides before names are compared, so "Hindustan
|
||||
# Unilever Hul Lux" reduces to "lux" on our side and matches OFF's
|
||||
# "Lux". Computed once per brand, not once per row.
|
||||
drop = brand_tokens(brand)
|
||||
|
||||
for index, row in enumerate(rows):
|
||||
if progress_cb:
|
||||
progress_cb(index, len(rows))
|
||||
|
||||
title = row.get("title") or row.get("product_name") or ""
|
||||
try:
|
||||
scored = score_candidates(
|
||||
hits, title, [row.get("size") or ""], drop,
|
||||
review_min=min_similarity,
|
||||
)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.debug("scoring failed for %r: %s", title, e)
|
||||
continue
|
||||
|
||||
# score_candidates already drops anything under review_min, so the
|
||||
# first entry is the best acceptable one. The explicit re-check is
|
||||
# kept because the ordering contract is "best first", not "all
|
||||
# above the floor" - relying on the filter alone would silently
|
||||
# break if that ever changed.
|
||||
best = scored[0] if scored else None
|
||||
if not best or best.score < min_similarity:
|
||||
stats["rejected"] += 1
|
||||
continue
|
||||
|
||||
code = validate_barcode(getattr(best, "barcode", None))
|
||||
if not code:
|
||||
stats["rejected"] += 1
|
||||
continue
|
||||
|
||||
stats["matched"] += 1
|
||||
if dry_run:
|
||||
continue
|
||||
|
||||
kind = classify_barcode_type(code)
|
||||
sources = {
|
||||
"barcode": {"method": "sourced", "source": BARCODE_SOURCE,
|
||||
"confidence": round(float(best.score), 3),
|
||||
"note": "matched on name against the brand's OFF "
|
||||
"catalogue; not verified against the pack"},
|
||||
"gtin": {"method": "derived", "source": "validators.validate_barcode"},
|
||||
"ean13": {"method": "derived", "source": "validators.to_ean13"},
|
||||
}
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
cur.execute(
|
||||
f'UPDATE "{table}" SET barcode = %s, barcode_type = %s, '
|
||||
f"gtin = %s, ean13 = %s, barcode_source = %s, "
|
||||
f"barcode_verified = FALSE, barcode_lookup_status = %s, "
|
||||
f"barcode_last_updated = NOW(), "
|
||||
f"field_sources = COALESCE(field_sources, '{{}}'::jsonb) "
|
||||
f" || %s::jsonb, "
|
||||
f"updated_at = CURRENT_TIMESTAMP "
|
||||
f"WHERE id = %s "
|
||||
f" AND (barcode IS NULL OR btrim(barcode) = '')",
|
||||
(code, kind.value, code, to_ean13(code), BARCODE_SOURCE,
|
||||
"name_matched", Json(sources), row["id"]),
|
||||
)
|
||||
stats["written"] += cur.rowcount
|
||||
conn.commit()
|
||||
except Exception as e: # noqa: BLE001
|
||||
conn.rollback()
|
||||
logger.warning("barcode write failed for %s id=%s: %s",
|
||||
table, row["id"], e)
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
return stats
|
||||
|
||||
|
||||
def enrich_barcodes_for_brands(brands: Iterable[str], *,
|
||||
min_similarity: float = DEFAULT_MIN_SIMILARITY,
|
||||
dry_run: bool = False,
|
||||
progress_cb=None) -> Dict[str, Dict[str, int]]:
|
||||
"""Run `enrich_brand_barcodes` over several brands, one at a time.
|
||||
|
||||
Deliberately sequential. `EnrichmentPipeline`'s five-way concurrency is for
|
||||
per-row work against a local corpus; firing five brand-corpus fetches at
|
||||
Open Food Facts at once is how a shared host earns a rate limit.
|
||||
"""
|
||||
out: Dict[str, Dict[str, int]] = {}
|
||||
names = [b.strip() for b in brands if b and b.strip()]
|
||||
for i, brand in enumerate(names):
|
||||
if progress_cb:
|
||||
progress_cb(i, len(names))
|
||||
try:
|
||||
stats = enrich_brand_barcodes(brand, min_similarity=min_similarity,
|
||||
dry_run=dry_run)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("barcode enrichment failed for %s: %s", brand, e)
|
||||
continue
|
||||
if stats.get("candidates"):
|
||||
out[brand] = stats
|
||||
logger.info("%s: %d without a barcode, %d matched, %d written",
|
||||
brand, stats["candidates"], stats["matched"], stats["written"])
|
||||
return out
|
||||
@@ -52,6 +52,34 @@ def run_enrich_job(job_id: str, *, label: str = "", **enrich_kwargs: Any) -> Non
|
||||
def progress_cb(done: int, total: int) -> None:
|
||||
nutrition_job_store.update(job_id, processed=done, total=total)
|
||||
|
||||
# PHASE 1 - BARCODES, BEFORE ANY NUTRITION LOOKUP.
|
||||
#
|
||||
# The ordering is a correctness property. `fetch_verified_nutrition_by_barcode`
|
||||
# matches on the GTIN and returns confidence 0.95; the name search it falls
|
||||
# back to accepts at 0.32. Because the nutrition phase runs with
|
||||
# skip_if_verified=True, whichever lands first wins permanently - so a name
|
||||
# match obtained before the barcode exists blocks the far better barcode
|
||||
# match from ever being tried.
|
||||
#
|
||||
# It is also cheap: one brand corpus (~5 requests) instead of one search per
|
||||
# product against an endpoint capped at 10 requests/minute.
|
||||
#
|
||||
# Guarded separately and never fatal: a barcode phase that fails must still
|
||||
# leave the nutrition phase to do what it always did.
|
||||
brands_for_barcodes = enrich_kwargs.get("brands") or []
|
||||
if brands_for_barcodes and settings.ENRICH_BARCODES_ON_UPLOAD:
|
||||
try:
|
||||
nutrition_job_store.update(job_id, detail="Finding barcodes")
|
||||
from app.services.enrichment.post_ingest_barcodes import (
|
||||
enrich_barcodes_for_brands,
|
||||
)
|
||||
found = enrich_barcodes_for_brands(brands_for_barcodes)
|
||||
written = sum(s.get("written", 0) for s in found.values())
|
||||
if written:
|
||||
logger.info("Barcode phase wrote %d barcode(s) before scoring", written)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("Barcode phase failed (nutrition still runs): %s", e)
|
||||
|
||||
try:
|
||||
result = nutrition_enrichment_service.enrich_all_products(
|
||||
progress_cb=progress_cb, **enrich_kwargs)
|
||||
|
||||
@@ -545,9 +545,17 @@ def fetch_verified_nutrition_by_barcode(
|
||||
candidate_size=product.get("quantity") or "",
|
||||
candidate_countries=",".join(product.get("countries_tags") or []),
|
||||
)
|
||||
# barcode_is_identity=True is safe HERE and nowhere else. The barcode has
|
||||
# already established which single product this is - there are no competing
|
||||
# candidates to disambiguate - so the name check is a guard against our
|
||||
# barcode being wrong, not the evidence for identity. Open Food Facts
|
||||
# stores short names ("Munch", "Maaza") where we store long ones, and
|
||||
# without this 149 of 300 barcoded rows were rejected as "wrong product"
|
||||
# when the barcode had resolved perfectly. See matching.name_is_contained.
|
||||
matched, similarity = is_match(
|
||||
candidate, brand, title, size,
|
||||
min_name_similarity=min_name_similarity,
|
||||
barcode_is_identity=True,
|
||||
)
|
||||
if not matched:
|
||||
logger.info(
|
||||
|
||||
@@ -49,6 +49,53 @@ class EnrichmentResult:
|
||||
duration_seconds: float = 0.0
|
||||
|
||||
|
||||
def _score_from_existing_facts(brand: str, image_id: str,
|
||||
facts: Dict[str, Any]) -> bool:
|
||||
"""Compute and store insights from facts already held. No network.
|
||||
|
||||
Exists because the two ways a `nutrition_facts` row can be written disagree
|
||||
about whose job scoring is. `enrich_one_product` fetches and scores in one
|
||||
pass; `scripts/backfill_nutrition_from_barcodes` and its siblings write
|
||||
facts only, on purpose, because they are narrow repair tools that should
|
||||
not also be recomputing narratives. Nothing then closed the gap.
|
||||
|
||||
Returns True when it wrote something. A row that already carries a score is
|
||||
left alone, so this is idempotent and safe to call on every skip.
|
||||
"""
|
||||
try:
|
||||
existing = nutrition_db.get_nutrition_insights(brand, image_id)
|
||||
if existing and existing.get("nutrition_score") is not None:
|
||||
return False
|
||||
|
||||
scores = nutrition_scoring.compute_scores(facts)
|
||||
if not scores:
|
||||
# `compute_scores` returns None rather than fabricating when there
|
||||
# is nothing scoreable. Respect that - do not write a zero.
|
||||
return False
|
||||
|
||||
positive = nutrition_scoring.generate_positive_insights(facts)
|
||||
cautions = nutrition_scoring.generate_cautions(facts)
|
||||
allergens = nutrition_scoring.normalize_allergens(facts)
|
||||
insights: Dict[str, Any] = {
|
||||
"brand": brand, "image_id": image_id,
|
||||
"positive_insights": positive, "nutritional_cautions": cautions,
|
||||
# No narrative: generating one is an LLM call, and this path exists
|
||||
# precisely to avoid doing expensive work for a row that already
|
||||
# has its data. The admin /enrich run fills narratives in later.
|
||||
"ai_summary": (existing or {}).get("ai_summary", ""),
|
||||
"diet_tags": nutrition_scoring.classify_diet_tags(facts),
|
||||
"allergens": allergens,
|
||||
"allergen_source": (facts.get("data_source") or "unavailable") if allergens else "unavailable",
|
||||
"data_status": facts.get("data_status"),
|
||||
}
|
||||
insights.update(scores)
|
||||
nutrition_db.upsert_nutrition_insights(insights)
|
||||
return True
|
||||
except Exception as e: # noqa: BLE001 - scoring is a nicety, not the write
|
||||
logger.warning("Could not score existing facts for %s/%s: %s", brand, image_id, e)
|
||||
return False
|
||||
|
||||
|
||||
def enrich_one_product(brand: str, image_id: str, product_name: str, category: str,
|
||||
skip_if_verified: bool = False, generate_narrative: bool = True) -> str:
|
||||
"""Runs the full pipeline for a single product. Returns the resulting
|
||||
@@ -66,6 +113,22 @@ def enrich_one_product(brand: str, image_id: str, product_name: str, category: s
|
||||
if skip_if_verified:
|
||||
existing = nutrition_db.get_nutrition_facts(brand, image_id)
|
||||
if existing and existing.get("data_status") == "verified":
|
||||
# SKIPPING THE FETCH MUST NOT ALSO SKIP THE SCORING.
|
||||
#
|
||||
# This used to `return "verified"` outright, which meant a product
|
||||
# whose facts arrived from a backfill script could never acquire a
|
||||
# score. Those scripts write `nutrition_facts` and deliberately do
|
||||
# not touch `nutrition_insights` - so the row looks verified here,
|
||||
# gets skipped, and its health score stays NULL forever. Measured
|
||||
# on 2026-09-08: 157 rows held usable facts with no score, which is
|
||||
# most of the gap between `nutrients_per_100g` at 65.6% and
|
||||
# `health_score` at 54.8%.
|
||||
#
|
||||
# Scoring is pure computation over facts already in hand - no
|
||||
# network, no LLM - so doing it here costs one cheap read and
|
||||
# nothing else. It is skipped when a score already exists, so a
|
||||
# re-run is a no-op.
|
||||
_score_from_existing_facts(brand, image_id, existing)
|
||||
return "verified"
|
||||
|
||||
facts = nutrition_data_service.fetch_verified_nutrition(brand, product_name, category or "")
|
||||
|
||||
@@ -57,6 +57,8 @@ from decimal import Decimal
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional, Tuple
|
||||
|
||||
from psycopg.types.json import Json
|
||||
|
||||
from app.services.vector_store import (
|
||||
_connect,
|
||||
_list_brand_table_suffixes,
|
||||
@@ -216,6 +218,175 @@ def sync_scores_to_brand_tables(brands: Optional[List[str]] = None,
|
||||
return changed
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# per-100g nutrients
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
# The flat per-100g columns on `nutrition_facts`, with the unit each is stored
|
||||
# in. Order is the order a label reads, not alphabetical, because it is also
|
||||
# the order the rendered display list comes out in.
|
||||
_NUTRIENT_COLUMNS: Tuple[Tuple[str, str, str], ...] = (
|
||||
("calories_kcal", "Energy", "kcal"),
|
||||
("protein_g", "Protein", "g"),
|
||||
("carbohydrates_g", "Carbohydrates", "g"),
|
||||
("total_sugar_g", "Total sugar", "g"),
|
||||
("added_sugar_g", "Added sugar", "g"),
|
||||
("dietary_fiber_g", "Dietary fibre", "g"),
|
||||
("total_fat_g", "Total fat", "g"),
|
||||
("saturated_fat_g", "Saturated fat", "g"),
|
||||
("trans_fat_g", "Trans fat", "g"),
|
||||
("cholesterol_mg", "Cholesterol", "mg"),
|
||||
("sodium_mg", "Sodium", "mg"),
|
||||
("potassium_mg", "Potassium", "mg"),
|
||||
("calcium_mg", "Calcium", "mg"),
|
||||
("iron_mg", "Iron", "mg"),
|
||||
("magnesium_mg", "Magnesium", "mg"),
|
||||
("zinc_mg", "Zinc", "mg"),
|
||||
("vitamin_a_mcg", "Vitamin A", "mcg"),
|
||||
("vitamin_c_mg", "Vitamin C", "mg"),
|
||||
("vitamin_d_mcg", "Vitamin D", "mcg"),
|
||||
("vitamin_e_mg", "Vitamin E", "mg"),
|
||||
("omega_3_g", "Omega-3", "g"),
|
||||
("omega_6_g", "Omega-6", "g"),
|
||||
)
|
||||
|
||||
# How many rendered lines the display column carries. The old keyword
|
||||
# generator capped at 8 and the UI is laid out for roughly that.
|
||||
_DISPLAY_LIMIT = 8
|
||||
|
||||
|
||||
def render_nutrient_lines(facts: Dict[str, Any], limit: int = _DISPLAY_LIMIT) -> List[str]:
|
||||
"""The `nutrients TEXT[]` display list, built from real measured numbers.
|
||||
|
||||
WHY THIS REPLACES WHAT WAS THERE
|
||||
`catalog_engine.generate_nutrients_info` produced this column from
|
||||
category keywords - "Energy - High", "Protein - Good Source" - with no
|
||||
connection to `nutrition_facts` at all. A product could show
|
||||
"Vitamin B Complex - Energy" because its category name matched a word.
|
||||
Where real per-100g figures exist they are strictly better, and they
|
||||
are what the merchant console is actually being asked for.
|
||||
|
||||
Values are emitted per 100 g because that is the basis `nutrition_facts`
|
||||
stores and the basis Indian FSSAI labelling uses. A missing nutrient is
|
||||
omitted rather than rendered as zero - "Trans fat 0 g" is a claim, and an
|
||||
absent measurement is not evidence of absence.
|
||||
"""
|
||||
lines: List[str] = []
|
||||
for column, label, unit in _NUTRIENT_COLUMNS:
|
||||
value = facts.get(column)
|
||||
if value is None:
|
||||
continue
|
||||
try:
|
||||
number = float(value)
|
||||
except (TypeError, ValueError):
|
||||
continue
|
||||
# Trim a trailing .0 so "Protein 12 g" does not read as "12.0".
|
||||
text = f"{number:g}"
|
||||
lines.append(f"{label} {text} {unit} per 100 g")
|
||||
if len(lines) >= limit:
|
||||
break
|
||||
return lines
|
||||
|
||||
|
||||
def sync_nutrients_to_brand_tables(brands: Optional[List[str]] = None,
|
||||
*, dry_run: bool = False) -> Dict[str, int]:
|
||||
"""Mirror `nutrition_facts` per-100g figures onto the brand tables.
|
||||
|
||||
Writes two columns, for two different readers:
|
||||
|
||||
* `nutrients_per_100g` (JSONB) - the numbers, for anything computing on
|
||||
them. This is a denormalised copy so a single product read needs no
|
||||
join, exactly as `nutrition_score`/`health_score` already are.
|
||||
* `nutrients` (TEXT[]) - the same numbers rendered for display, replacing
|
||||
the category-keyword strings.
|
||||
|
||||
`nutrition_facts` remains the source of truth; nothing here writes back to
|
||||
it. Like the score mirror this is a MIRROR, not an accumulator: a row whose
|
||||
facts have gone (the non-consumable purge removes them) has its copy
|
||||
cleared rather than left stale.
|
||||
|
||||
A row whose facts exist but are `unavailable` is left alone rather than
|
||||
blanked - `data_status='unavailable'` means "we looked and found nothing",
|
||||
which is not a reason to destroy a display list the catalog already had.
|
||||
"""
|
||||
conn = _connect()
|
||||
if conn is None:
|
||||
logger.warning("nutrient sync: no database connection")
|
||||
return {}
|
||||
|
||||
columns = ", ".join(c for c, _, _ in _NUTRIENT_COLUMNS)
|
||||
changed: Dict[str, int] = {}
|
||||
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
targets = _target_suffixes(cur, brands)
|
||||
|
||||
for suffix, display in targets:
|
||||
table = f"brand_{suffix}"
|
||||
if not dry_run:
|
||||
ensure_brand_schema(display)
|
||||
|
||||
try:
|
||||
with conn.cursor() as cur:
|
||||
cur.execute(
|
||||
f"SELECT image_id, {columns} FROM nutrition_facts "
|
||||
f"WHERE lower(brand) = lower(%s) AND data_status <> 'unavailable'",
|
||||
(display,),
|
||||
)
|
||||
rows = cur.fetchall()
|
||||
if not rows:
|
||||
continue
|
||||
|
||||
names = [c for c, _, _ in _NUTRIENT_COLUMNS]
|
||||
updates = 0
|
||||
for row in rows:
|
||||
image_id = row[0]
|
||||
facts = {n: v for n, v in zip(names, row[1:]) if v is not None}
|
||||
if not facts:
|
||||
continue
|
||||
block = {n: float(v) for n, v in facts.items()}
|
||||
lines = render_nutrient_lines(facts)
|
||||
if dry_run:
|
||||
cur.execute(
|
||||
f'SELECT 1 FROM "{table}" WHERE image_id = %s '
|
||||
f"AND nutrients_per_100g IS DISTINCT FROM %s::jsonb",
|
||||
(image_id, Json(block)),
|
||||
)
|
||||
updates += 1 if cur.fetchone() else 0
|
||||
continue
|
||||
cur.execute(
|
||||
f'UPDATE "{table}" SET nutrients_per_100g = %s::jsonb, '
|
||||
f" nutrients = %s, "
|
||||
f" field_sources = COALESCE(field_sources, '{{}}'::jsonb) || %s::jsonb, "
|
||||
f" updated_at = CURRENT_TIMESTAMP "
|
||||
f" WHERE image_id = %s "
|
||||
f" AND nutrients_per_100g IS DISTINCT FROM %s::jsonb",
|
||||
(Json(block), lines,
|
||||
Json({"nutrients": {"method": "sourced",
|
||||
"source": "nutrition_facts"},
|
||||
"nutrients_per_100g": {"method": "sourced",
|
||||
"source": "nutrition_facts"}}),
|
||||
image_id, Json(block)),
|
||||
)
|
||||
updates += cur.rowcount
|
||||
|
||||
if not dry_run:
|
||||
conn.commit()
|
||||
except Exception as e: # noqa: BLE001
|
||||
conn.rollback()
|
||||
logger.warning("nutrient sync failed for %s: %s", table, e)
|
||||
continue
|
||||
|
||||
if updates:
|
||||
changed[display] = updates
|
||||
logger.info("%s %d row(s) in %s",
|
||||
"would update" if dry_run else "updated", updates, table)
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
return changed
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# seed catalog JSON
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -356,3 +527,15 @@ def sync_after_write(brands: Optional[List[str]] = None) -> None:
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("score mirror failed (scores are still in "
|
||||
"nutrition_insights): %s", e)
|
||||
|
||||
# Separately guarded: a failure mirroring the numbers must not lose the
|
||||
# scores that were just mirrored successfully, and vice versa. Both are
|
||||
# copies of data that is safe in nutrition_facts / nutrition_insights.
|
||||
try:
|
||||
changed = sync_nutrients_to_brand_tables(brands)
|
||||
if changed:
|
||||
logger.info("mirrored nutrients onto %d brand table(s): %s",
|
||||
len(changed), ", ".join(sorted(changed)))
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("nutrient mirror failed (figures are still in "
|
||||
"nutrition_facts): %s", e)
|
||||
|
||||
@@ -6,7 +6,10 @@ import logging
|
||||
import os
|
||||
import re
|
||||
import time
|
||||
from datetime import datetime
|
||||
|
||||
import psycopg
|
||||
from psycopg.types.json import Json
|
||||
|
||||
from app.infrastructure.settings import (
|
||||
USE_PGVECTOR, DB_HOST, DB_PORT, DB_NAME, DB_USER, DB_PASSWORD,
|
||||
@@ -56,6 +59,56 @@ def _sanitize_name(name: str) -> str:
|
||||
return name.strip('_')
|
||||
|
||||
|
||||
def _epoch_to_timestamp(value: Any) -> Optional[datetime]:
|
||||
"""Convert BarcodeResult.barcode_last_updated to what the column holds.
|
||||
|
||||
`BarcodeResult` defaults this field to `time.time()` - a float epoch - but
|
||||
the seven brand tables that already carry `barcode_last_updated` were
|
||||
migrated out-of-band as TIMESTAMP. Handing psycopg a bare float for a
|
||||
timestamp column raises, so the conversion happens here rather than being
|
||||
pushed onto every caller. A datetime is passed through untouched, and
|
||||
anything unparseable degrades to None.
|
||||
"""
|
||||
if value is None:
|
||||
return None
|
||||
if isinstance(value, datetime):
|
||||
return value
|
||||
# A seed catalog round-trips this as ISO-8601 (brand_sync._jsonable), while
|
||||
# BarcodeResult and the older JSON files carry a float epoch. Both have to
|
||||
# load, or re-seeding a brand would drop the timestamp it just exported.
|
||||
if isinstance(value, str):
|
||||
text = value.strip()
|
||||
if not text:
|
||||
return None
|
||||
try:
|
||||
return datetime.fromisoformat(text.replace("Z", "+00:00"))
|
||||
except ValueError:
|
||||
pass
|
||||
try:
|
||||
return datetime.fromtimestamp(float(value))
|
||||
except (ValueError, TypeError, OSError, OverflowError):
|
||||
return None
|
||||
|
||||
|
||||
def _to_numeric_or_none(value: Any) -> Optional[float]:
|
||||
"""Coerce an enrichment stage's numeric output for a NUMERIC column.
|
||||
|
||||
Returns None for anything unparseable rather than raising, because a
|
||||
malformed tax figure must degrade to "no tax figure stored" and never
|
||||
abort a whole batch's write. Strips a leading currency symbol and commas,
|
||||
which is how these arrive when a sheet supplied them as text.
|
||||
"""
|
||||
if value is None or isinstance(value, bool):
|
||||
return None
|
||||
try:
|
||||
if isinstance(value, str):
|
||||
cleaned = re.sub(r"[^\d.\-]", "", value)
|
||||
return float(cleaned) if cleaned not in ("", "-", ".", "-.") else None
|
||||
return float(value)
|
||||
except (ValueError, TypeError):
|
||||
return None
|
||||
|
||||
|
||||
DDL_CREATE_EXTENSION = "CREATE EXTENSION IF NOT EXISTS vector;"
|
||||
def get_brand_table_ddl(brand: str) -> str:
|
||||
"""Generate DDL for brand-specific table - simplified with only essential fields"""
|
||||
@@ -88,12 +141,52 @@ def get_brand_table_ddl(brand: str) -> str:
|
||||
selling_price NUMERIC,
|
||||
barcode TEXT,
|
||||
barcode_type TEXT,
|
||||
|
||||
|
||||
-- The rest of what BarcodeResult.as_product_fields() produces. Until
|
||||
-- these existed the barcode stage returned nine fields and the INSERT
|
||||
-- named two, so seven were computed and then dropped on the floor -
|
||||
-- including the provenance needed to tell a verified GTIN from a
|
||||
-- name-matched guess.
|
||||
--
|
||||
-- TYPES ARE ADOPTED, NOT CHOSEN. Seven brand tables already carry
|
||||
-- these columns, added out-of-band before any code created them.
|
||||
-- ADD COLUMN IF NOT EXISTS does not reconcile a type difference, so
|
||||
-- picking a "better" type here would leave 7 tables disagreeing with
|
||||
-- 49 forever. Verified against information_schema on 2026-09-08:
|
||||
-- barcode_last_updated is TIMESTAMP (not the float epoch
|
||||
-- BarcodeResult carries - see _epoch_to_timestamp), and the tax
|
||||
-- columns are REAL (not NUMERIC).
|
||||
gtin TEXT,
|
||||
ean13 TEXT,
|
||||
upc TEXT,
|
||||
barcode_source TEXT,
|
||||
barcode_verified BOOLEAN,
|
||||
barcode_lookup_status TEXT,
|
||||
barcode_last_updated TIMESTAMP,
|
||||
|
||||
-- Computed by the HSN/GST stage, likewise discarded before this.
|
||||
gst_percent REAL,
|
||||
tax_amount REAL,
|
||||
hsn_gst_needs_review BOOLEAN,
|
||||
|
||||
-- Essential fields
|
||||
highlights TEXT[],
|
||||
nutrients TEXT[],
|
||||
-- Real per-100g figures mirrored from nutrition_facts by
|
||||
-- nutrition_score_sync. `nutrients` above stays the human-readable
|
||||
-- marketing list; these are the numbers.
|
||||
nutrients_per_100g JSONB,
|
||||
search_query TEXT,
|
||||
|
||||
-- Per-field provenance, keyed by column name; each value records how
|
||||
-- that column's value was arrived at. One JSONB map rather than ~60
|
||||
-- scalar columns, because every scalar column would have to be named
|
||||
-- explicitly in the INSERT below and in brand_sync.EXPORT_COLUMNS or
|
||||
-- it gets silently dropped - the exact failure the barcode columns
|
||||
-- above are here to fix.
|
||||
-- `method` is one of: sourced | estimated | derived | not_applicable.
|
||||
field_sources JSONB DEFAULT '{{}}'::jsonb,
|
||||
|
||||
-- Nutrition scores, mirrored from nutrition_insights by
|
||||
-- nutrition_score_sync. Deliberately NOT written by
|
||||
-- upsert_brand_products - see the comment on its INSERT.
|
||||
@@ -160,9 +253,30 @@ def _ensure_columns(cur, table_name: str) -> None:
|
||||
"selling_price": "NUMERIC",
|
||||
"barcode": "TEXT",
|
||||
"barcode_type": "TEXT",
|
||||
# The other seven fields BarcodeResult.as_product_fields() returns.
|
||||
# Adding them here is what puts them on all 56 existing brand tables -
|
||||
# before this only 7 tables had them, added out-of-band, and even those
|
||||
# never received a value because the INSERT did not name them.
|
||||
# Types adopted from the 7 tables that already have them, NOT chosen -
|
||||
# see get_brand_table_ddl. TIMESTAMP and REAL are what is on disk.
|
||||
"gtin": "TEXT",
|
||||
"ean13": "TEXT",
|
||||
"upc": "TEXT",
|
||||
"barcode_source": "TEXT",
|
||||
"barcode_verified": "BOOLEAN",
|
||||
"barcode_lookup_status": "TEXT",
|
||||
"barcode_last_updated": "TIMESTAMP",
|
||||
# Computed by the HSN/GST stage and discarded before this.
|
||||
"gst_percent": "REAL",
|
||||
"tax_amount": "REAL",
|
||||
"hsn_gst_needs_review": "BOOLEAN",
|
||||
"highlights": "TEXT[]",
|
||||
"nutrients": "TEXT[]",
|
||||
"nutrients_per_100g": "JSONB",
|
||||
"search_query": "TEXT",
|
||||
# Per-field provenance map - see get_brand_table_ddl for why this is
|
||||
# one JSONB column and not sixty scalar ones.
|
||||
"field_sources": "JSONB",
|
||||
# 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
|
||||
@@ -193,7 +307,11 @@ def _ensure_columns(cur, table_name: str) -> None:
|
||||
# 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"}
|
||||
#
|
||||
# `nutrients_per_100g` joins them: it is mirrored from nutrition_facts by
|
||||
# the same sync, not written by the INSERT, and is likewise a current
|
||||
# column rather than legacy debris.
|
||||
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", "gtin", "ean13", "upc", "barcode_source", "barcode_verified", "barcode_lookup_status", "barcode_last_updated", "gst_percent", "tax_amount", "hsn_gst_needs_review", "highlights", "nutrients", "nutrients_per_100g", "field_sources", "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")
|
||||
@@ -346,6 +464,30 @@ def upsert_brand_products(brand: str, products: List[Dict[str, Any]], cleanup: b
|
||||
barcode = str(p.get("barcode") or p.get("Barcode") or "").strip() or None
|
||||
barcode_type = str(p.get("barcode_type") or p.get("Barcode_Type") or "").strip() or None
|
||||
|
||||
# The remaining barcode fields. A product dict that never went through
|
||||
# the barcode stage simply has none of these, so they arrive as None
|
||||
# and the COALESCE in DO UPDATE SET below keeps whatever is already
|
||||
# stored rather than blanking it.
|
||||
gtin = str(p.get("gtin") or "").strip() or None
|
||||
ean13 = str(p.get("ean13") or "").strip() or None
|
||||
upc = str(p.get("upc") or "").strip() or None
|
||||
barcode_source = str(p.get("barcode_source") or "").strip() or None
|
||||
barcode_lookup_status = str(p.get("barcode_lookup_status") or "").strip() or None
|
||||
barcode_verified = p.get("barcode_verified")
|
||||
barcode_verified = bool(barcode_verified) if barcode_verified is not None else None
|
||||
barcode_last_updated = _epoch_to_timestamp(p.get("barcode_last_updated"))
|
||||
|
||||
# HSN/GST stage output beyond hsn_code itself.
|
||||
gst_percent = _to_numeric_or_none(p.get("gst_percent"))
|
||||
tax_amount = _to_numeric_or_none(p.get("tax_amount"))
|
||||
hsn_gst_needs_review = p.get("hsn_gst_needs_review")
|
||||
hsn_gst_needs_review = bool(hsn_gst_needs_review) if hsn_gst_needs_review is not None else None
|
||||
|
||||
# Per-field provenance. Must be a JSON object; anything else is
|
||||
# dropped rather than stored as a shape readers cannot index into.
|
||||
field_sources = p.get("field_sources")
|
||||
field_sources = Json(field_sources) if isinstance(field_sources, dict) and field_sources else None
|
||||
|
||||
# Essential fields
|
||||
highlights = p.get("highlights", [])
|
||||
if not isinstance(highlights, list):
|
||||
@@ -387,9 +529,20 @@ def upsert_brand_products(brand: str, products: List[Dict[str, Any]], cleanup: b
|
||||
selling_price,
|
||||
barcode,
|
||||
barcode_type,
|
||||
gtin,
|
||||
ean13,
|
||||
upc,
|
||||
barcode_source,
|
||||
barcode_verified,
|
||||
barcode_lookup_status,
|
||||
barcode_last_updated,
|
||||
gst_percent,
|
||||
tax_amount,
|
||||
hsn_gst_needs_review,
|
||||
highlights, # TEXT[] - psycopg will handle conversion
|
||||
nutrients, # TEXT[] - psycopg will handle conversion
|
||||
search_query,
|
||||
field_sources,
|
||||
embedding_str
|
||||
))
|
||||
|
||||
@@ -430,8 +583,11 @@ def upsert_brand_products(brand: str, products: List[Dict[str, Any]], cleanup: b
|
||||
f"""
|
||||
INSERT INTO {table_name}
|
||||
(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)
|
||||
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||
fssai_license, product_sku, sku_source, hsn_code, final_selling_price, selling_price, barcode, barcode_type,
|
||||
gtin, ean13, upc, barcode_source, barcode_verified, barcode_lookup_status, barcode_last_updated,
|
||||
gst_percent, tax_amount, hsn_gst_needs_review,
|
||||
highlights, nutrients, search_query, field_sources, embedding)
|
||||
VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
|
||||
ON CONFLICT (image_id) DO UPDATE SET
|
||||
product_name = EXCLUDED.product_name,
|
||||
title = EXCLUDED.title,
|
||||
@@ -448,12 +604,63 @@ def upsert_brand_products(brand: str, products: List[Dict[str, Any]], cleanup: b
|
||||
hsn_code = EXCLUDED.hsn_code,
|
||||
final_selling_price = EXCLUDED.final_selling_price,
|
||||
selling_price = EXCLUDED.selling_price,
|
||||
barcode = EXCLUDED.barcode,
|
||||
barcode_type = EXCLUDED.barcode_type,
|
||||
-- barcode/barcode_type are COALESCEd with the seven
|
||||
-- columns below rather than assigned like their legacy
|
||||
-- neighbours, because they are the same identity group.
|
||||
-- Proven against brand_zzsmoketest on 2026-09-08: a bare
|
||||
-- re-seed (a product dict with no barcode keys, which is
|
||||
-- what the seed loader and user_products build) set
|
||||
-- barcode to NULL while COALESCE kept gtin and ean13 -
|
||||
-- leaving a row claiming a GTIN with no barcode. That
|
||||
-- half-erased state is worse than either whole one.
|
||||
--
|
||||
-- It also aligns this statement with the invariant the
|
||||
-- rest of the barcode code already enforces: stage.py:38
|
||||
-- skips a row that has a barcode, and base.py:76-100
|
||||
-- refuses to blank a held value. The upsert was the one
|
||||
-- place that still could.
|
||||
barcode = COALESCE(EXCLUDED.barcode, {table_name}.barcode),
|
||||
barcode_type = COALESCE(EXCLUDED.barcode_type, {table_name}.barcode_type),
|
||||
highlights = EXCLUDED.highlights,
|
||||
nutrients = EXCLUDED.nutrients,
|
||||
search_query = EXCLUDED.search_query,
|
||||
embedding = EXCLUDED.embedding,
|
||||
-- COALESCE, NOT PLAIN EXCLUDED, FOR EVERY COLUMN BELOW.
|
||||
--
|
||||
-- These are enrichment outputs. Most writers that reach
|
||||
-- this statement never carry them: the seed loader,
|
||||
-- user_products._build_product_dict and brand_sync's
|
||||
-- re-seed all build a product dict from a spreadsheet,
|
||||
-- so their EXCLUDED values are NULL. A plain assignment
|
||||
-- would therefore wipe a hard-won barcode and its
|
||||
-- provenance on the next re-seed of the brand - the same
|
||||
-- failure the score-column comment above describes, which
|
||||
-- is why those two columns are omitted entirely.
|
||||
-- (Naming them here, even inside a comment, trips the
|
||||
-- substring guard in tests/test_brand_table_scores.py.
|
||||
-- That guard is crude on purpose; leave it that way.)
|
||||
--
|
||||
-- COALESCE keeps the stored value when nothing new
|
||||
-- arrives, while still letting a real incoming value
|
||||
-- correct a stored one. That matches the enrichment
|
||||
-- invariant in enrichment/base.py: a stage may fill a gap
|
||||
-- or correct a value, it may not erase one.
|
||||
gtin = COALESCE(EXCLUDED.gtin, {table_name}.gtin),
|
||||
ean13 = COALESCE(EXCLUDED.ean13, {table_name}.ean13),
|
||||
upc = COALESCE(EXCLUDED.upc, {table_name}.upc),
|
||||
barcode_source = COALESCE(EXCLUDED.barcode_source, {table_name}.barcode_source),
|
||||
barcode_verified = COALESCE(EXCLUDED.barcode_verified, {table_name}.barcode_verified),
|
||||
barcode_lookup_status = COALESCE(EXCLUDED.barcode_lookup_status, {table_name}.barcode_lookup_status),
|
||||
barcode_last_updated = COALESCE(EXCLUDED.barcode_last_updated, {table_name}.barcode_last_updated),
|
||||
gst_percent = COALESCE(EXCLUDED.gst_percent, {table_name}.gst_percent),
|
||||
tax_amount = COALESCE(EXCLUDED.tax_amount, {table_name}.tax_amount),
|
||||
hsn_gst_needs_review = COALESCE(EXCLUDED.hsn_gst_needs_review, {table_name}.hsn_gst_needs_review),
|
||||
-- Merged, not replaced: a run that learns the provenance
|
||||
-- of one field must not drop what is known about the
|
||||
-- others. `||` is a shallow merge, which is the right
|
||||
-- depth here - each key's value is one flat record.
|
||||
field_sources = COALESCE({table_name}.field_sources, '{{}}'::jsonb)
|
||||
|| COALESCE(EXCLUDED.field_sources, '{{}}'::jsonb),
|
||||
updated_at = CURRENT_TIMESTAMP
|
||||
""",
|
||||
rows,
|
||||
|
||||
Reference in New Issue
Block a user