diff --git a/ingestion/opengoods/adapters/openfoodfacts.py b/ingestion/opengoods/adapters/openfoodfacts.py index eb94ad8..7981f30 100644 --- a/ingestion/opengoods/adapters/openfoodfacts.py +++ b/ingestion/opengoods/adapters/openfoodfacts.py @@ -157,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/jobs/seed_off.py b/ingestion/opengoods/jobs/seed_off.py index 68ab712..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,20 +23,33 @@ from collections.abc import Iterator import psycopg -from opengoods.adapters.openfoodfacts import OpenFoodFactsAdapter, read_dump +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 @@ -58,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/tests/test_off_country.py b/ingestion/tests/test_off_country.py new file mode 100644 index 0000000..e716ba3 --- /dev/null +++ b/ingestion/tests/test_off_country.py @@ -0,0 +1,57 @@ +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"]