diff --git a/ingestion/opengoods/adapters/openfoodfacts.py b/ingestion/opengoods/adapters/openfoodfacts.py index 6c9a2f8..7981f30 100644 --- a/ingestion/opengoods/adapters/openfoodfacts.py +++ b/ingestion/opengoods/adapters/openfoodfacts.py @@ -27,6 +27,9 @@ _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" +# HTTP statuses worth retrying: rate limiting and transient server errors. +_RETRY_STATUS = frozenset({429, 500, 502, 503, 504}) + # Fields requested from the search API so a returned product can be transformed # without an extra per-barcode round trip. _SEARCH_FIELDS = ( @@ -46,9 +49,13 @@ class OpenFoodFactsAdapter: self, client: httpx.Client | None = None, min_interval: float = _DEFAULT_MIN_INTERVAL, + max_retries: int = 4, + backoff_base: float = 2.0, ) -> None: self._client = client or httpx.Client(headers={"User-Agent": USER_AGENT}, timeout=30.0) self._min_interval = min_interval + self._max_retries = max_retries + self._backoff_base = backoff_base self._last_call = 0.0 def _throttle(self) -> None: @@ -58,11 +65,49 @@ class OpenFoodFactsAdapter: time.sleep(wait) self._last_call = time.monotonic() + def _get(self, url: str, params: dict | None = None) -> httpx.Response: + """GET with throttling and retry/backoff on transient errors. + + Retries on connection/timeout errors and on retryable HTTP statuses + (429 and 5xx, which OFF returns intermittently when overloaded), using + exponential backoff that honours a ``Retry-After`` header when present. + """ + last_exc: Exception | None = None + for attempt in range(self._max_retries + 1): + self._throttle() + try: + resp = self._client.get(url, params=params) + except httpx.TransportError as exc: + last_exc = exc + else: + if resp.status_code < 400 or resp.status_code not in _RETRY_STATUS: + resp.raise_for_status() + return resp + last_exc = httpx.HTTPStatusError( + f"retryable status {resp.status_code}", request=resp.request, response=resp + ) + if attempt < self._max_retries: + retry_after = self._retry_after(last_exc) + time.sleep(retry_after if retry_after is not None else self._backoff_base**attempt) + assert last_exc is not None + raise last_exc + + @staticmethod + def _retry_after(exc: Exception | None) -> float | None: + resp = getattr(exc, "response", None) + if resp is None: + return None + value = resp.headers.get("Retry-After") + if not value: + return None + try: + return float(value) + except ValueError: + return None + 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() + resp = self._get(_API_URL.format(barcode=barcode)) payload = resp.json() if payload.get("status") != 1: return None @@ -91,8 +136,7 @@ class OpenFoodFactsAdapter: ``last_modified_t`` they processed as the next watermark. """ for page in range(1, max_pages + 1): - self._throttle() - resp = self._client.get( + resp = self._get( _SEARCH_URL, params={ "fields": _SEARCH_FIELDS, @@ -101,7 +145,6 @@ class OpenFoodFactsAdapter: "page_size": page_size, }, ) - resp.raise_for_status() products = resp.json().get("products") or [] if not products: return @@ -114,6 +157,60 @@ class OpenFoodFactsAdapter: if reached_old or len(products) < page_size: return + def fetch_by_country( + self, + country: str, + *, + page_size: int = 100, + max_pages: int = 10, + sort_by: str = "unique_scans_n", + ) -> Iterator[dict]: + """Yield products sold in ``country`` (an OFF ``countries_tags_en`` slug). + + Used to seed a market-specific catalogue (e.g. ``china``). Results are + sorted by ``sort_by`` (default ``unique_scans_n`` so the most-scanned, + best-known products come first) and de-duplicated across pages, since + OFF's popularity ordering is not stable between page requests. + """ + seen: set[str] = set() + for page in range(1, max_pages + 1): + resp = self._get( + _SEARCH_URL, + params={ + "fields": _SEARCH_FIELDS, + "countries_tags_en": country, + "sort_by": sort_by, + "page": page, + "page_size": page_size, + }, + ) + products = resp.json().get("products") or [] + if not products: + return + new_on_page = 0 + for prod in products: + code = str(prod.get("code") or "") + if code and code in seen: + continue + if code: + seen.add(code) + new_on_page += 1 + yield prod + if len(products) < page_size or new_on_page == 0: + return + + +def is_cn_gs1(code: str | None) -> bool: + """Return True for a GS1 China company prefix (barcodes starting 690-699). + + These identify products registered with GS1 China, i.e. genuinely domestic + items, as opposed to imported goods merely tagged as sold in China. + """ + if not code: + return False + code = code.strip() + return len(code) >= 3 and code[:2] == "69" and code[2].isdigit() + def read_dump(path: str | Path) -> Iterator[dict]: """Yield raw product records from an OFF JSONL dump file. diff --git a/ingestion/opengoods/etl/load.py b/ingestion/opengoods/etl/load.py index f052713..f90d5d0 100644 --- a/ingestion/opengoods/etl/load.py +++ b/ingestion/opengoods/etl/load.py @@ -7,6 +7,7 @@ source with field-level provenance in `product_source`. from __future__ import annotations import json +import logging import os from typing import Any @@ -18,6 +19,8 @@ from opengoods.etl.quality import update_quality OFF_HOMEPAGE = "https://world.openfoodfacts.org" +logger = logging.getLogger(__name__) + def default_dsn() -> str: return os.environ.get( @@ -196,6 +199,25 @@ def load_record(conn: psycopg.Connection, rec: dict[str, Any], source_id: str, r return product_id +def load_record_safe( + conn: psycopg.Connection, rec: dict[str, Any], source_id: str, raw: dict +) -> bool: + """Load one record inside a savepoint. + + On success the record's writes stay in the surrounding transaction. On any + error, only this record's writes are rolled back (to the savepoint) and the + batch continues, so a single malformed source record cannot abort a large + import. Returns True if loaded, False if skipped due to an error. + """ + try: + with conn.transaction(): + load_record(conn, rec, source_id, raw) + return True + except Exception as exc: # noqa: BLE001 - per-record isolation is intentional + logger.warning("skipping record gtin=%s: %s", rec.get("gtin"), exc) + return False + + def _jsonable(raw: dict) -> dict: """Drop values that are not JSON-serializable from a raw record.""" try: diff --git a/ingestion/opengoods/etl/transform.py b/ingestion/opengoods/etl/transform.py index 4201f9f..122b030 100644 --- a/ingestion/opengoods/etl/transform.py +++ b/ingestion/opengoods/etl/transform.py @@ -78,6 +78,14 @@ def map_category(raw: dict) -> str | None: return None +def _clamp(value: str | None, max_len: int) -> str | None: + """Trim a string to fit a bounded DB column; external data length varies.""" + if value is None: + return None + value = value.strip() + return value[:max_len] or None + + def _clean_tags(tags: list[str] | None, prefix: str = "") -> list[str]: out: list[str] = [] for t in tags or []: @@ -143,16 +151,16 @@ def transform(raw: dict) -> dict | None: "brand": brand, "category_path": map_category(raw), "net_content_value": net_value, - "net_content_unit": net_unit, + "net_content_unit": _clamp(net_unit, 16), "net_content_canonical": net_canonical, - "country_of_origin": (raw.get("countries") or "").split(",")[0].strip() or None, + "country_of_origin": _clamp((raw.get("countries") or "").split(",")[0].strip() or None, 64), "food": { "ingredients_text": raw.get("ingredients_text") or None, "allergens": _clean_tags(raw.get("allergens_tags")), "additives": _clean_tags(raw.get("additives_tags")), "nutriments": transform_nutriments(raw.get("nutriments") or {}), "nutrition_basis": "per_100g", - "serving_size": raw.get("serving_size") or None, + "serving_size": _clamp(raw.get("serving_size") or None, 32), "nutri_score": (raw.get("nutriscore_grade") or "").upper()[:1] or None, }, "image_url": raw.get("image_front_url") or raw.get("image_url") or None, diff --git a/ingestion/opengoods/jobs/seed_off.py b/ingestion/opengoods/jobs/seed_off.py index 03c8ab0..f7d632e 100644 --- a/ingestion/opengoods/jobs/seed_off.py +++ b/ingestion/opengoods/jobs/seed_off.py @@ -7,6 +7,11 @@ Usage: # from a downloaded OFF JSONL dump (optionally .gz), limited to N records python -m opengoods.jobs.seed_off --dump products.jsonl.gz --limit 1000 + # market-focused: the most-scanned products sold in China, restricted to + # genuine GS1-China (69x) barcodes + python -m opengoods.jobs.seed_off --country china --domestic-only \ + --max-pages 20 --limit 1000 + The OFF read API is rate-limited client-side; for large imports use a dump. """ @@ -18,25 +23,38 @@ from collections.abc import Iterator import psycopg -from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, read_dump -from opengoods.etl.load import default_dsn, ensure_source, load_record +from opengoods.adapters.openfoodfacts import ( + OpenFoodFactsAdapter, + is_cn_gs1, + read_dump, +) +from opengoods.etl.load import default_dsn, ensure_source, load_record_safe from opengoods.etl.transform import transform def _raw_records(args: argparse.Namespace) -> Iterator[dict]: if args.dump: - records = read_dump(args.dump) + records: Iterator[dict] = read_dump(args.dump) + elif args.country: + adapter = OpenFoodFactsAdapter(min_interval=args.min_interval) + records = adapter.fetch_by_country( + args.country, page_size=args.page_size, max_pages=args.max_pages + ) else: adapter = OpenFoodFactsAdapter(min_interval=args.min_interval) records = adapter.fetch(args.barcodes) - for i, rec in enumerate(records): - if args.limit and i >= args.limit: + yielded = 0 + for rec in records: + if args.domestic_only and not is_cn_gs1(str(rec.get("code") or "")): + continue + if args.limit and yielded >= args.limit: break + yielded += 1 yield rec def run(args: argparse.Namespace) -> int: - loaded = skipped = 0 + loaded = skipped = errored = 0 with psycopg.connect(args.dsn, autocommit=False) as conn: source_id = ensure_source(conn) for raw in _raw_records(args): @@ -44,10 +62,12 @@ def run(args: argparse.Namespace) -> int: if rec is None: skipped += 1 continue - load_record(conn, rec, source_id, raw) - loaded += 1 + if load_record_safe(conn, rec, source_id, raw): + loaded += 1 + else: + errored += 1 conn.commit() - print(f"loaded={loaded} skipped={skipped}") + print(f"loaded={loaded} skipped={skipped} errored={errored}") return 0 @@ -56,7 +76,18 @@ def main(argv: list[str] | None = None) -> int: src = parser.add_mutually_exclusive_group(required=True) src.add_argument("--barcodes", nargs="+", help="barcodes to fetch via the OFF API") src.add_argument("--dump", help="path to an OFF JSONL dump (.jsonl or .jsonl.gz)") + src.add_argument( + "--country", + help="OFF countries_tags_en slug to seed from, e.g. 'china' (most-scanned first)", + ) + parser.add_argument( + "--domestic-only", + action="store_true", + help="keep only genuine GS1-China (69x) barcodes; drop imported goods", + ) parser.add_argument("--limit", type=int, default=0, help="max records to load (0 = all)") + parser.add_argument("--page-size", type=int, default=100, help="search page size") + parser.add_argument("--max-pages", type=int, default=10, help="max search pages (country mode)") parser.add_argument("--min-interval", type=float, default=4.0, help="API throttle seconds") parser.add_argument("--dsn", default=default_dsn(), help="PostgreSQL DSN") return run(parser.parse_args(argv)) diff --git a/ingestion/opengoods/jobs/update_off.py b/ingestion/opengoods/jobs/update_off.py index a99dc5d..ce67495 100644 --- a/ingestion/opengoods/jobs/update_off.py +++ b/ingestion/opengoods/jobs/update_off.py @@ -17,14 +17,14 @@ import sys import psycopg from opengoods.adapters.openfoodfacts import SOURCE_NAME, OpenFoodFactsAdapter -from opengoods.etl.load import default_dsn, ensure_source, load_record +from opengoods.etl.load import default_dsn, ensure_source, load_record_safe from opengoods.etl.state import get_watermark, set_watermark from opengoods.etl.transform import transform def run(args: argparse.Namespace) -> int: adapter = OpenFoodFactsAdapter(min_interval=args.min_interval) - loaded = skipped = 0 + loaded = skipped = errored = 0 high_watermark = 0 with psycopg.connect(args.dsn, autocommit=False) as conn: source_id = ensure_source(conn) @@ -38,16 +38,21 @@ def run(args: argparse.Namespace) -> int: if rec is None: skipped += 1 continue - load_record(conn, rec, source_id, raw) - loaded += 1 + if load_record_safe(conn, rec, source_id, raw): + loaded += 1 + else: + errored += 1 set_watermark( conn, SOURCE_NAME, high_watermark, - stats={"loaded": loaded, "skipped": skipped, "since": since}, + stats={"loaded": loaded, "skipped": skipped, "errored": errored, "since": since}, ) conn.commit() - print(f"since={since} loaded={loaded} skipped={skipped} watermark={high_watermark}") + print( + f"since={since} loaded={loaded} skipped={skipped} " + f"errored={errored} watermark={high_watermark}" + ) return 0 diff --git a/ingestion/tests/test_load_integration.py b/ingestion/tests/test_load_integration.py index e0eef00..e874e72 100644 --- a/ingestion/tests/test_load_integration.py +++ b/ingestion/tests/test_load_integration.py @@ -11,7 +11,7 @@ from pathlib import Path import pytest -from opengoods.etl.load import default_dsn, ensure_source, load_record +from opengoods.etl.load import default_dsn, ensure_source, load_record, load_record_safe from opengoods.etl.transform import transform psycopg = pytest.importorskip("psycopg") @@ -59,3 +59,25 @@ def test_load_record_roundtrip(conn): assert prov[0] >= 1 conn.rollback() # keep the test DB clean + + +def test_load_record_safe_isolates_bad_record(conn): + source_id = ensure_source(conn) + + # Unique gtin so the good record is a fresh INSERT, not an upsert/update. + good = transform(FIXTURE) + good["gtin"] = "4006381333931" + assert load_record_safe(conn, good, source_id, FIXTURE) is True + after_good = conn.execute("SELECT count(*) FROM product").fetchone()[0] + + # A record whose name violates NOT NULL must not abort the batch. + bad = dict(good) + bad["gtin"] = "5000112637922" + bad["name"] = None + assert load_record_safe(conn, bad, source_id, {}) is False + + # The good record survived the bad one's rollback-to-savepoint. + after_bad = conn.execute("SELECT count(*) FROM product").fetchone()[0] + assert after_bad == after_good + + conn.rollback() # keep the test DB clean diff --git a/ingestion/tests/test_off_country.py b/ingestion/tests/test_off_country.py new file mode 100644 index 0000000..8941711 --- /dev/null +++ b/ingestion/tests/test_off_country.py @@ -0,0 +1,53 @@ +import httpx + +from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, is_cn_gs1 + + +def _adapter(handler, **kwargs): + client = httpx.Client(transport=httpx.MockTransport(handler)) + return OpenFoodFactsAdapter(client=client, min_interval=0, backoff_base=0, **kwargs) + + +def test_is_cn_gs1(): + assert is_cn_gs1("6901234567892") # GS1 China prefix + assert is_cn_gs1("690") + assert not is_cn_gs1("3017624010701") # France + assert not is_cn_gs1("5449000000996") # Belgium + assert not is_cn_gs1("") + assert not is_cn_gs1(None) + assert not is_cn_gs1("69") # too short to carry a prefix digit + + +def test_fetch_by_country_passes_filter_and_paginates(): + seen_params = [] + + def handler(request: httpx.Request) -> httpx.Response: + seen_params.append(dict(request.url.params)) + page = int(request.url.params.get("page", "1")) + if page == 1: + return httpx.Response( + 200, + json={"products": [{"code": "6901"}, {"code": "6902"}]}, + ) + return httpx.Response(200, json={"products": []}) + + adapter = _adapter(handler) + out = list(adapter.fetch_by_country("china", page_size=2, max_pages=5)) + assert [p["code"] for p in out] == ["6901", "6902"] + assert seen_params[0]["countries_tags_en"] == "china" + assert seen_params[0]["sort_by"] == "unique_scans_n" + + +def test_fetch_by_country_dedupes_across_pages(): + def handler(request: httpx.Request) -> httpx.Response: + page = int(request.url.params.get("page", "1")) + if page == 1: + return httpx.Response(200, json={"products": [{"code": "6901"}, {"code": "6902"}]}) + if page == 2: + # OFF popularity ordering is unstable; a repeat appears on page 2 + return httpx.Response(200, json={"products": [{"code": "6902"}, {"code": "6903"}]}) + return httpx.Response(200, json={"products": []}) + + adapter = _adapter(handler) + out = [p["code"] for p in adapter.fetch_by_country("china", page_size=2, max_pages=5)] + assert out == ["6901", "6902", "6903"] diff --git a/ingestion/tests/test_off_retry.py b/ingestion/tests/test_off_retry.py new file mode 100644 index 0000000..d04ff2e --- /dev/null +++ b/ingestion/tests/test_off_retry.py @@ -0,0 +1,59 @@ +import httpx +import pytest + +from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter + + +def _adapter(handler, **kwargs): + client = httpx.Client(transport=httpx.MockTransport(handler)) + return OpenFoodFactsAdapter(client=client, min_interval=0, backoff_base=0, **kwargs) + + +def test_retries_transient_5xx_then_succeeds(): + calls = {"n": 0} + + def handler(request: httpx.Request) -> httpx.Response: + calls["n"] += 1 + if calls["n"] < 3: + return httpx.Response(503) + return httpx.Response(200, json={"status": 1, "product": {"code": "x"}}) + + adapter = _adapter(handler, max_retries=4) + assert adapter.fetch_barcode("x") == {"code": "x"} + assert calls["n"] == 3 # two 503s retried, third succeeds + + +def test_retries_exhausted_raises(): + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(503) + + adapter = _adapter(handler, max_retries=2) + with pytest.raises(httpx.HTTPStatusError): + adapter.fetch_barcode("x") + + +def test_non_retryable_4xx_not_retried(): + calls = {"n": 0} + + def handler(request: httpx.Request) -> httpx.Response: + calls["n"] += 1 + return httpx.Response(404) + + adapter = _adapter(handler, max_retries=4) + with pytest.raises(httpx.HTTPStatusError): + adapter.fetch_barcode("x") + assert calls["n"] == 1 # 404 is not retried + + +def test_retries_connection_error_then_succeeds(): + calls = {"n": 0} + + def handler(request: httpx.Request) -> httpx.Response: + calls["n"] += 1 + if calls["n"] == 1: + raise httpx.ConnectError("boom") + return httpx.Response(200, json={"status": 1, "product": {"code": "y"}}) + + adapter = _adapter(handler, max_retries=4) + assert adapter.fetch_barcode("y") == {"code": "y"} + assert calls["n"] == 2 diff --git a/ingestion/tests/test_transform.py b/ingestion/tests/test_transform.py index e4bf98e..cfcbbf3 100644 --- a/ingestion/tests/test_transform.py +++ b/ingestion/tests/test_transform.py @@ -65,3 +65,14 @@ def test_transform_full_record(): def test_transform_drops_unnamed(): assert transform({"code": "0000000000000"}) is None + + +def test_transform_clamps_long_serving_size(): + raw = { + "code": "3017624010701", + "product_name": "X", + "serving_size": "1 portion (30 g) / 1 portion (30 g) / extra long descriptive text here", + } + rec = transform(raw) + assert rec is not None + assert len(rec["food"]["serving_size"]) == 32