"""Open Food Facts (OFF) source adapter. Fetches raw product records either from the OFF read API (one product per barcode) or from a downloaded JSONL dump file. OFF data is licensed under the Open Database License (ODbL); product images are CC-BY-SA. We record OFF as the source for every field we ingest. The adapter is read-only and rate-limited to stay well within OFF's API limits (<= ~15 req/min/IP for product reads) and to be a good citizen. """ from __future__ import annotations import json import time from collections.abc import Iterator from pathlib import Path import httpx SOURCE_NAME = "openfoodfacts" OFF_LICENSE = "ODbL" USER_AGENT = "OpenGoods/0.1 (+https://github.com/baicai2026-baicai/goods) public-good product API" # Conservative client-side spacing between API calls (seconds). _DEFAULT_MIN_INTERVAL = 4.0 _API_URL = "https://world.openfoodfacts.org/api/v2/product/{barcode}.json" _SEARCH_URL = "https://world.openfoodfacts.org/api/v2/search" # Fields requested from the search API so a returned product can be transformed # without an extra per-barcode round trip. _SEARCH_FIELDS = ( "code,product_name,product_name_en,product_name_zh,brands,quantity," "categories,categories_tags,countries,ingredients_text,allergens_tags," "additives_tags,nutriments,nutriscore_grade,serving_size," "image_front_url,image_url,last_modified_t" ) class OpenFoodFactsAdapter: """Read product records from the OFF API.""" source_name = SOURCE_NAME def __init__( self, client: httpx.Client | None = None, min_interval: float = _DEFAULT_MIN_INTERVAL, ) -> None: self._client = client or httpx.Client(headers={"User-Agent": USER_AGENT}, timeout=30.0) self._min_interval = min_interval self._last_call = 0.0 def _throttle(self) -> None: elapsed = time.monotonic() - self._last_call wait = self._min_interval - elapsed if wait > 0: time.sleep(wait) self._last_call = time.monotonic() def fetch_barcode(self, barcode: str) -> dict | None: """Fetch a single product by barcode; return the raw `product` dict.""" self._throttle() resp = self._client.get(_API_URL.format(barcode=barcode)) resp.raise_for_status() payload = resp.json() if payload.get("status") != 1: return None return payload["product"] def fetch(self, barcodes: list[str]) -> Iterator[dict]: """Yield raw product records for the given barcodes.""" for code in barcodes: record = self.fetch_barcode(code) if record is not None: yield record def fetch_modified_since( self, since_t: int, *, page_size: int = 100, max_pages: int = 10, ) -> Iterator[dict]: """Yield products modified after ``since_t`` (unix ``last_modified_t``). Uses the OFF search API sorted by ``last_modified_t`` (most recent first) and paginates until it reaches products at or before the watermark, an empty/short page, or ``max_pages``. This is the incremental ingestion path: callers persist the highest ``last_modified_t`` they processed as the next watermark. """ for page in range(1, max_pages + 1): self._throttle() resp = self._client.get( _SEARCH_URL, params={ "fields": _SEARCH_FIELDS, "sort_by": "last_modified_t", "page": page, "page_size": page_size, }, ) resp.raise_for_status() products = resp.json().get("products") or [] if not products: return reached_old = False for prod in products: if int(prod.get("last_modified_t") or 0) <= since_t: reached_old = True break yield prod if reached_old or len(products) < page_size: return def read_dump(path: str | Path) -> Iterator[dict]: """Yield raw product records from an OFF JSONL dump file. Each line is one product JSON object (the format of OFF's .jsonl export). Supports plain or .gz files. """ p = Path(path) if p.suffix == ".gz": import gzip opener = lambda: gzip.open(p, "rt", encoding="utf-8") # noqa: E731 else: opener = lambda: open(p, encoding="utf-8") # noqa: E731 with opener() as fh: for line in fh: line = line.strip() if line: yield json.loads(line)