Brand valid image generation
This commit is contained in:
@@ -10,34 +10,50 @@ resulting cascade "misses far more often than it hits". That module is wired
|
||||
into the live enrichment pipeline and is deliberately NOT touched by this file.
|
||||
|
||||
This module inverts the problem: fetch a brand's **entire** OFF catalogue in one
|
||||
or two requests from the Search-a-licious endpoint, cache it on disk, then match
|
||||
every one of our products against that corpus offline. A whole-catalogue
|
||||
backfill costs ~5 HTTP requests instead of ~600, and re-tuning the similarity
|
||||
threshold costs zero network because the corpus is cached.
|
||||
or two requests from the v2 search API, cache it on disk, then match every one
|
||||
of our products against that corpus offline. A whole-catalogue backfill costs
|
||||
~5 HTTP requests instead of ~600, and re-tuning the similarity threshold costs
|
||||
zero network because the corpus is cached.
|
||||
|
||||
Nothing here is imported by the running app - the only consumer is
|
||||
`scripts/backfill_barcodes_from_off.py`. Everything except `fetch_brand_corpus`
|
||||
is a pure function so it can be unit-tested without network or database.
|
||||
`fetch_brand_corpus` is the one network function here and is shared by the
|
||||
backfill scripts, `product_grounding`, `post_ingest_barcodes` and
|
||||
`brand_discovery`. Everything else is a pure function so it can be unit-tested
|
||||
without network or database.
|
||||
|
||||
ENDPOINT NOTES (verified empirically, 2026-09)
|
||||
GET https://search.openfoodfacts.org/search
|
||||
?q=brands:amul AND countries_tags:"en:india"
|
||||
ENDPOINT NOTES (verified empirically, 2026-09-11)
|
||||
GET https://world.openfoodfacts.org/api/v2/search
|
||||
?brands_tags=amul
|
||||
&countries_tags=en:india
|
||||
&fields=code,product_name,quantity,...
|
||||
&page_size=250
|
||||
* The double quotes around "en:india" are REQUIRED. Without them the query
|
||||
parses as a bare term and silently returns count=0 rather than erroring.
|
||||
* page_size up to 1000 is accepted; 250 keeps responses small.
|
||||
&page_size=100&page=1
|
||||
* `brands_tags` is an EXACT match on the brand's tag slug ("naga", not
|
||||
"Naga"), which is what a brand catalogue needs. The earlier Search-a-licious
|
||||
query `q=brands:Naga` was a free-text match: it returned "Mr Naga" and
|
||||
"Bombay Naga Jhal" - other companies - and MISSED the real Naga rows,
|
||||
because Search-a-licious reads a separate Elasticsearch index that was
|
||||
stale for them (barcode 8906011830068 was still filed brandless under
|
||||
Kuwait). The v2 API reads the live product database. Measured on the same
|
||||
day: Amul 216 products here vs 140 there; Naga 2 vs 0.
|
||||
* `page_count` in this API is the number of products ON THIS PAGE, not the
|
||||
number of pages, and `page_size` in the RESPONSE is what the server
|
||||
actually applied (a larger request is capped to 100 without complaint).
|
||||
Paginate from `count` / that served `page_size`.
|
||||
* OFF rate-limits all search endpoints to 10 requests/minute per IP.
|
||||
PAUSE_SECONDS keeps a multi-page brand under that.
|
||||
* `world.openfoodfacts.org` intermittently serves an HTML "Page temporarily
|
||||
unavailable" page with a 200 status, so every response is content-type
|
||||
checked before parsing.
|
||||
unavailable" page with a 200 status, and plain 503s under load, so every
|
||||
response is status- and content-type-checked before parsing, and a failed
|
||||
fetch is never written to the cache as an empty corpus.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import math
|
||||
import os
|
||||
import re
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, Iterable, List, Optional, Sequence, Tuple
|
||||
|
||||
@@ -63,21 +79,50 @@ from app.services.quantity_utils import quantities_match
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
SEARCH_URL = "https://search.openfoodfacts.org/search"
|
||||
SEARCH_URL = "https://world.openfoodfacts.org/api/v2/search"
|
||||
|
||||
# OFF's usage policy requires a contactable custom User-Agent; requests sent
|
||||
# with the default python-requests agent are treated as anonymous crawling.
|
||||
USER_AGENT = "BrandCatalogRAG/1.0 (suriya@tenext.in)"
|
||||
|
||||
FIELDS = "code,product_name,product_name_en,brands,quantity,countries_tags"
|
||||
PAGE_SIZE = 250
|
||||
MAX_PAGES = 10 # 2500 products per brand is far beyond any real brand
|
||||
PAUSE_SECONDS = 2.0 # search endpoints are the strictly-limited ones
|
||||
PAGE_SIZE = 100 # the v2 API caps this at 100 and silently serves that
|
||||
MAX_PAGES = 25 # 2500 products per brand is far beyond any real brand
|
||||
PAUSE_SECONDS = 6.0 # 10 search requests/minute is OFF's published limit
|
||||
|
||||
# Bumped when the cache file's meaning changes. Files written before this key
|
||||
# existed came from the Search-a-licious endpoint; an EMPTY one of those is
|
||||
# more likely to be that endpoint's stale index (or a swallowed 503) than a
|
||||
# real absence, so it is re-fetched. A non-empty one is real data and is kept.
|
||||
CACHE_SCHEMA = 2
|
||||
|
||||
# backend/app/services/enrichment/barcode/sources/off_bulk.py -> backend/
|
||||
_BACKEND_DIR = Path(__file__).resolve().parents[5]
|
||||
CACHE_DIR = _BACKEND_DIR / "data" / "cache" / "off_brand_corpus"
|
||||
|
||||
|
||||
class OffUnavailable(requests.exceptions.ConnectionError):
|
||||
"""OFF answered but could not serve (429 / 5xx).
|
||||
|
||||
Subclasses ConnectionError so `retry.with_retry` treats it exactly like a
|
||||
dropped socket - a few seconds later the same request usually succeeds -
|
||||
without widening the retry set for every other barcode source.
|
||||
"""
|
||||
|
||||
|
||||
@dataclass
|
||||
class CorpusFetch:
|
||||
"""What `fetch_brand_corpus_result` learned.
|
||||
|
||||
`error` is None when OFF answered every page, even if it answered with
|
||||
nothing - that is a real "this brand is not on Open Food Facts". When
|
||||
`error` is set the fetch did not complete and NOTHING was cached, so a
|
||||
caller can say "could not be reached" instead of "has nothing".
|
||||
"""
|
||||
hits: List[Dict[str, Any]] = field(default_factory=list)
|
||||
error: Optional[str] = None
|
||||
from_cache: bool = False
|
||||
|
||||
# Pack sizes embedded in a product name ("Marie Gold 250g", "Butter 1L").
|
||||
# Mirrors the unit list in scripts/backfill_nutrition_from_barcodes.py, widened
|
||||
# with the count-based units this catalog also uses.
|
||||
@@ -122,84 +167,149 @@ def _cache_path(slug: str, cache_dir: Optional[Path] = None) -> Path:
|
||||
return (cache_dir or CACHE_DIR) / f"{slug}.json"
|
||||
|
||||
|
||||
def brand_tag_slug(brand: str) -> str:
|
||||
"""The brand as Open Food Facts tags it: lowercase, runs of anything that
|
||||
is not a letter or digit collapsed to one hyphen. "Hindustan Unilever" ->
|
||||
"hindustan-unilever", "P&G" -> "p-g", "Naga" -> "naga". OFF matches
|
||||
`brands_tags` case-insensitively, but sending the slug form is what its
|
||||
own site does and avoids depending on that."""
|
||||
return re.sub(r"[^a-z0-9]+", "-", (brand or "").lower()).strip("-")
|
||||
|
||||
|
||||
@with_retry(max_attempts=3, min_wait=2.0, max_wait=10.0)
|
||||
def _get_page(brand: str, country: Optional[str], page: int) -> Dict[str, Any]:
|
||||
"""One page of OFF search results. Raises on transport errors (retried by
|
||||
the decorator); returns an empty result dict for any response that is not
|
||||
parseable JSON, which is how the HTML "temporarily unavailable" page and
|
||||
any future error page are absorbed without killing the run."""
|
||||
query = f"brands:{brand}"
|
||||
def _get_page(brand: str, country: Optional[str], page: int) -> Optional[Dict[str, Any]]:
|
||||
"""One page of OFF v2 search results.
|
||||
|
||||
Raises on transport errors and on 429 / 5xx (both retried by the
|
||||
decorator); returns None for any other response that is not parseable
|
||||
JSON - an HTML "temporarily unavailable" page, a 4xx, a truncated body.
|
||||
None means "this page FAILED", which the caller must keep distinct from a
|
||||
page that parsed fine and simply held no products.
|
||||
"""
|
||||
params: Dict[str, Any] = {
|
||||
"brands_tags": brand_tag_slug(brand),
|
||||
"fields": FIELDS,
|
||||
"page_size": PAGE_SIZE,
|
||||
"page": page,
|
||||
}
|
||||
if country:
|
||||
# The quotes are load-bearing - see the module docstring.
|
||||
query += f' AND countries_tags:"en:{country}"'
|
||||
params["countries_tags"] = f"en:{country}"
|
||||
|
||||
resp = requests.get(
|
||||
SEARCH_URL,
|
||||
params={"q": query, "fields": FIELDS, "page_size": PAGE_SIZE, "page": page},
|
||||
params=params,
|
||||
headers={"User-Agent": USER_AGENT, "Accept": "application/json"},
|
||||
timeout=BARCODE_LOOKUP_TIMEOUT_SECONDS,
|
||||
)
|
||||
|
||||
if resp.status_code == 429 or resp.status_code >= 500:
|
||||
raise OffUnavailable(f"HTTP {resp.status_code} from Open Food Facts")
|
||||
if resp.status_code != 200:
|
||||
logger.warning("OFF search returned HTTP %s for brand %r page %s",
|
||||
resp.status_code, brand, page)
|
||||
return {}
|
||||
return None
|
||||
if "json" not in (resp.headers.get("content-type") or "").lower():
|
||||
logger.warning("OFF search returned non-JSON (%s) for brand %r - "
|
||||
"the service is probably serving an error page",
|
||||
resp.headers.get("content-type"), brand)
|
||||
return {}
|
||||
return None
|
||||
try:
|
||||
payload = resp.json()
|
||||
except ValueError as e:
|
||||
logger.warning("OFF search returned unparseable JSON for brand %r: %s", brand, e)
|
||||
return {}
|
||||
return payload if isinstance(payload, dict) else {}
|
||||
return None
|
||||
return payload if isinstance(payload, dict) else None
|
||||
|
||||
|
||||
def fetch_brand_corpus(brand: str,
|
||||
country: Optional[str] = None,
|
||||
refresh: bool = False,
|
||||
cache_dir: Optional[Path] = None) -> List[Dict[str, Any]]:
|
||||
"""Every OFF product for `brand`, from disk cache unless `refresh`.
|
||||
def _read_cache(path: Path, brand: str) -> Optional[List[Dict[str, Any]]]:
|
||||
"""Cached hits, or None when the cache must not be trusted: unreadable, or
|
||||
an empty file from before CACHE_SCHEMA existed (see that constant)."""
|
||||
try:
|
||||
cached = json.loads(path.read_text(encoding="utf-8"))
|
||||
except Exception as e: # noqa: BLE001 - a corrupt cache must not be fatal
|
||||
logger.warning(" Ignoring unreadable OFF cache %s: %s", path.name, e)
|
||||
return None
|
||||
hits = cached.get("hits") or []
|
||||
if not hits and cached.get("schema") != CACHE_SCHEMA:
|
||||
logger.info(" Ignoring empty pre-v%s OFF cache for %r - re-fetching",
|
||||
CACHE_SCHEMA, brand)
|
||||
return None
|
||||
logger.info(" OFF corpus for %r: %d product(s) (cached %s)",
|
||||
brand, len(hits), cached.get("fetched_at_human", "?"))
|
||||
return hits
|
||||
|
||||
|
||||
def fetch_brand_corpus_result(brand: str,
|
||||
country: Optional[str] = None,
|
||||
refresh: bool = False,
|
||||
cache_dir: Optional[Path] = None) -> CorpusFetch:
|
||||
"""Every OFF product for `brand`, from disk cache unless `refresh`, plus
|
||||
whether the fetch actually completed.
|
||||
|
||||
Hits with no usable product name are dropped here rather than at match time
|
||||
(6 of 146 Amul hits, 1 of 233 Britannia hits) - a nameless hit can never
|
||||
clear a name-similarity threshold, so carrying it forward only inflates the
|
||||
corpus. Returns [] rather than raising when OFF is unreachable, so one bad
|
||||
brand does not abort a multi-brand backfill.
|
||||
corpus.
|
||||
|
||||
A fetch that fails part-way returns what it got with `error` set and
|
||||
writes NO cache file. Caching a failure as `hits: []` is how "Open Food
|
||||
Facts has nothing for this brand" was being asserted for brands OFF had
|
||||
simply been too busy to answer about.
|
||||
"""
|
||||
country = BARCODE_COUNTRY_TAG if country is None else (country or None)
|
||||
slug = re.sub(r"[^a-z0-9]+", "_", brand.lower()).strip("_")
|
||||
path = _cache_path(slug, cache_dir)
|
||||
|
||||
if not refresh and path.exists():
|
||||
try:
|
||||
cached = json.loads(path.read_text(encoding="utf-8"))
|
||||
hits = cached.get("hits") or []
|
||||
logger.info(" OFF corpus for %r: %d product(s) (cached %s)",
|
||||
brand, len(hits), cached.get("fetched_at_human", "?"))
|
||||
return hits
|
||||
except Exception as e: # noqa: BLE001 - a corrupt cache must not be fatal
|
||||
logger.warning(" Ignoring unreadable OFF cache %s: %s", path.name, e)
|
||||
cached = _read_cache(path, brand)
|
||||
if cached is not None:
|
||||
return CorpusFetch(hits=cached, from_cache=True)
|
||||
|
||||
hits: List[Dict[str, Any]] = []
|
||||
error: Optional[str] = None
|
||||
page = 1
|
||||
while page <= MAX_PAGES:
|
||||
total_pages = 1
|
||||
while page <= min(total_pages, MAX_PAGES):
|
||||
if page > 1:
|
||||
time.sleep(PAUSE_SECONDS)
|
||||
payload = _get_page(brand, country, page)
|
||||
batch = payload.get("hits") or []
|
||||
try:
|
||||
payload = _get_page(brand, country, page)
|
||||
except requests.exceptions.RequestException as e:
|
||||
# Retries exhausted (OffUnavailable is a ConnectionError too).
|
||||
error = str(e) or e.__class__.__name__
|
||||
break
|
||||
if payload is None:
|
||||
error = "Open Food Facts returned an unusable response"
|
||||
break
|
||||
batch = payload.get("products") or []
|
||||
hits.extend(h for h in batch if isinstance(h, dict) and (h.get("product_name") or "").strip())
|
||||
page_count = payload.get("page_count") or 0
|
||||
if page >= page_count or not batch:
|
||||
if page == 1:
|
||||
# `page_count` is products-on-this-page in the v2 API; the real
|
||||
# page total is `count` over the page size the SERVER applied -
|
||||
# it silently caps the requested size (250 asked, 100 served for
|
||||
# Amul), so dividing by PAGE_SIZE under-pages.
|
||||
count = payload.get("count") or 0
|
||||
served = payload.get("page_size") or len(batch) or PAGE_SIZE
|
||||
try:
|
||||
total_pages = max(1, math.ceil(int(count) / int(served)))
|
||||
except (TypeError, ValueError, ZeroDivisionError):
|
||||
total_pages = 1
|
||||
if not batch:
|
||||
break
|
||||
page += 1
|
||||
|
||||
if error:
|
||||
logger.warning(" OFF corpus for %r: fetch failed after %d usable product(s) - %s "
|
||||
"(nothing cached)", brand, len(hits), error)
|
||||
return CorpusFetch(hits=hits, error=error)
|
||||
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
tmp = path.with_suffix(".json.tmp")
|
||||
tmp.write_text(json.dumps({
|
||||
"schema": CACHE_SCHEMA,
|
||||
"endpoint": SEARCH_URL,
|
||||
"brand": brand,
|
||||
"brand_tag": brand_tag_slug(brand),
|
||||
"country": country,
|
||||
"fetched_at": time.time(),
|
||||
"fetched_at_human": time.strftime("%Y-%m-%d %H:%M:%S"),
|
||||
@@ -208,7 +318,20 @@ def fetch_brand_corpus(brand: str,
|
||||
os.replace(tmp, path)
|
||||
|
||||
logger.info(" OFF corpus for %r: %d usable product(s) fetched", brand, len(hits))
|
||||
return hits
|
||||
return CorpusFetch(hits=hits)
|
||||
|
||||
|
||||
def fetch_brand_corpus(brand: str,
|
||||
country: Optional[str] = None,
|
||||
refresh: bool = False,
|
||||
cache_dir: Optional[Path] = None) -> List[Dict[str, Any]]:
|
||||
"""`fetch_brand_corpus_result(...).hits` - the list-only form every
|
||||
backfill and ingestion caller uses. Returns [] rather than raising when
|
||||
OFF is unreachable, so one bad brand does not abort a multi-brand run;
|
||||
callers that need to tell "unreachable" from "empty" use the result form.
|
||||
"""
|
||||
return fetch_brand_corpus_result(brand, country=country, refresh=refresh,
|
||||
cache_dir=cache_dir).hits
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@@ -243,6 +243,16 @@ def _round2(value: Optional[float]) -> Optional[float]:
|
||||
return round(value, 2)
|
||||
|
||||
|
||||
def _held_number(value: object) -> Optional[float]:
|
||||
"""A price the row already carries, or None for blank / unparseable."""
|
||||
if value is None or (isinstance(value, str) and not value.strip()):
|
||||
return None
|
||||
try:
|
||||
return float(value)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
|
||||
def enrich_pricing_fields(product: dict) -> Dict[str, object]:
|
||||
"""Compute the pricing fields added by this stage for one catalog row.
|
||||
|
||||
@@ -253,8 +263,24 @@ def enrich_pricing_fields(product: dict) -> Dict[str, object]:
|
||||
`cost_price`, `profit_before_tax` and `profit_after_tax` are not
|
||||
derivable from the data this pipeline generates, so they stay None
|
||||
(exactly as the existing enriched exports store them).
|
||||
|
||||
A PRICE THE ROW ALREADY HOLDS WINS. A store sheet that says "Retail Price
|
||||
155" has stated a fact; the band derived from it is Rs143-167, and this
|
||||
stage used to read the band's ceiling back as "selling_price = 167" and
|
||||
hand that to `EnrichmentStage.apply`, which overwrites a held value with
|
||||
a different one (it only refuses to BLANK one). Held prices are therefore
|
||||
returned as None here so `apply` keeps them, the base for the tax is the
|
||||
held selling price, then the held final price, and only then the band.
|
||||
"""
|
||||
selling_price = _extract_selling_price(product.get("price_range"))
|
||||
held_selling = _held_number(product.get("selling_price"))
|
||||
held_final = _held_number(product.get("final_selling_price"))
|
||||
if held_selling is not None:
|
||||
selling_price = held_selling
|
||||
elif held_final is not None:
|
||||
selling_price = held_final
|
||||
else:
|
||||
selling_price = _extract_selling_price(product.get("price_range"))
|
||||
|
||||
gst = product.get("gst_percent")
|
||||
tax_amount = None
|
||||
final_selling_price = None
|
||||
@@ -263,10 +289,10 @@ def enrich_pricing_fields(product: dict) -> Dict[str, object]:
|
||||
final_selling_price = _round2(selling_price + tax_amount)
|
||||
|
||||
return {
|
||||
"selling_price": _round2(selling_price),
|
||||
"selling_price": None if held_selling is not None else _round2(selling_price),
|
||||
"cost_price": None,
|
||||
"tax_amount": tax_amount,
|
||||
"final_selling_price": final_selling_price,
|
||||
"final_selling_price": None if held_final is not None else final_selling_price,
|
||||
"profit_before_tax": None,
|
||||
"profit_after_tax": None,
|
||||
}
|
||||
|
||||
@@ -76,7 +76,7 @@ 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)"
|
||||
BARCODE_SOURCE = "openfoodfacts_bulk (world.openfoodfacts.org/api/v2)"
|
||||
|
||||
|
||||
def _rows_needing_a_barcode(cur, table: str) -> List[Dict[str, Any]]:
|
||||
|
||||
Reference in New Issue
Block a user