diff --git a/.gitea/workflows/refresh.yml b/.gitea/workflows/refresh.yml index 47d1eb0..7599a08 100644 --- a/.gitea/workflows/refresh.yml +++ b/.gitea/workflows/refresh.yml @@ -5,11 +5,13 @@ name: Monthly corpus refresh # reindex + image-push if the scrape produced no diff against the # committed corpus. # -# Bayer takes ~30 min; EPA PPLS takes ~7 h with row-crop + -# registrant filters. The whole monthly job is ~8-9 h end-to-end. -# If that's too long for the runner you can: -# - Run just one source: workflow_dispatch with sources="bayer" -# - Limit EPA at the scraper: edit the step to add "--limit 5000" +# Runtime budget: act_runner kills the job container at exactly 3 h +# (it runs as `/bin/sleep 10800`), which is what killed every run from +# 2026-06-01 through 2026-09-01. The EPA step is incremental and +# parallel so a steady-state refresh is ~15-20 min: cached not-row-crop +# verdicts cost no request, and a product already on disk is only +# re-downloaded when EPA reports a new label date. `timeout-minutes` +# below fails the job legibly before the container disappears. on: schedule: @@ -42,6 +44,9 @@ env: jobs: refresh: runs-on: docker + # Below act_runner's own 3 h container lifetime, so an overrun fails + # as a timeout instead of "container ... does not exist". + timeout-minutes: 170 container: image: catthehacker/ubuntu:act-latest steps: @@ -80,9 +85,14 @@ jobs: - name: Scrape EPA PPLS if: ${{ inputs.sources == '' || contains(inputs.sources, 'epa_ppls') }} - # Row-crop + registrant filters keep this to ~16K PDFs / ~7h. - # Pass --no-row-crop-filter or --no-registrant-filter to broaden. - run: python -m scrape.runner --source epa_ppls --force + # Deliberately NOT --force: that re-downloaded all ~11.4K candidate + # registrations every month (~15-20 h of work) and never finished. + # Without it the run is incremental — one cheap API call per product, + # a PDF only when the label date actually moved — and the committed + # filter cache skips the ~7.3K non-row-crop products outright. + # Workers share one 5 req/sec ceiling, so this is faster without + # being ruder. Local full re-fetch: add --force --no-filter-cache. + run: python -m scrape.runner --source epa_ppls --workers 6 # ---- Commit corpus changes + retry-on-race ----------------- - name: Commit corpus changes (if any) @@ -90,13 +100,21 @@ jobs: run: | git config user.name "crop-chem-docs-refresh" git config user.email "actions@jpaul.io" - git add sources.json corpus - if git diff --cached --quiet; then + # The filter-verdict cache is committed so next month starts warm, + # but it changes on every run — only a real corpus diff may trigger + # the reindex + image build, so `changed` is decided on the corpus + # paths alone. + git add sources.json corpus scrape/state + if git diff --cached --quiet -- sources.json corpus; then echo "no corpus changes — skipping reindex and image build" echo "changed=false" >> "$GITHUB_OUTPUT" + else + echo "changed=true" >> "$GITHUB_OUTPUT" + fi + if git diff --cached --quiet; then + echo "nothing staged — no commit" exit 0 fi - echo "changed=true" >> "$GITHUB_OUTPUT" ts=$(date -u +"%Y-%m-%dT%H:%MZ") n_bayer=$(find corpus/bayer -name '*.json' 2>/dev/null | wc -l) n_epa=$(find corpus/epa_ppls -name '*.json' 2>/dev/null | wc -l) diff --git a/README.md b/README.md index b94059d..4f865de 100644 --- a/README.md +++ b/README.md @@ -59,9 +59,14 @@ pip install -r requirements.txt # Sample-scrape to verify wiring: python -m scrape.runner --source bayer --limit 5 -# Full refresh (be polite — bayer is small, epa_ppls is hours): +# Incremental refresh — what CI runs monthly (~15-20 min for epa_ppls): +# one cheap API call per product, a PDF only when the label date moved, +# and cached "not row-crop" verdicts skipped outright. python -m scrape.runner --source bayer --force -python -m scrape.runner --source epa_ppls --force +python -m scrape.runner --source epa_ppls --workers 6 + +# Full re-fetch of every label (hours — does NOT fit CI's 3 h runner cap): +python -m scrape.runner --source epa_ppls --force --no-filter-cache # Rebuild Chroma + BM25: OLLAMA_URL=http://192.168.0.125:11434 PRODUCT_NAME=crop_chem \ diff --git a/scrape/sources/epa_ppls.py b/scrape/sources/epa_ppls.py index 7e1edfa..a6e7c32 100644 --- a/scrape/sources/epa_ppls.py +++ b/scrape/sources/epa_ppls.py @@ -51,8 +51,11 @@ import logging import os import re import sys +import threading import time import zipfile +import zlib +from concurrent.futures import ThreadPoolExecutor, as_completed from dataclasses import dataclass, field from datetime import UTC, datetime from pathlib import Path @@ -79,10 +82,27 @@ REPO_ROOT = Path(__file__).resolve().parents[2] CORPUS_ROOT = Path(os.environ.get("CORPUS_ROOT") or REPO_ROOT / "corpus") CORPUS_DIR = CORPUS_ROOT / "epa_ppls" -REQUEST_DELAY_SECONDS = 1.1 # polite: ~1 req/sec +# Politeness ceiling, shared by every worker thread: raising --workers makes +# the run finish sooner without hitting EPA any harder than this. +RATE_LIMIT_RPS = 5.0 +DEFAULT_WORKERS = 6 + HTTP_TIMEOUT = httpx.Timeout(60.0, connect=15.0) MAX_RETRIES = 4 +# httpx timeouts are per-read, so a label body that trickles in at a few KB/s +# never trips them. One www3.epa.gov stall burned 11 min of a single CI run +# (100-1262) and another 10 min on a 33 MB label (241-441). Cap the whole +# download instead, and don't retry it four times. +PDF_DEADLINE_SECONDS = 180.0 +PDF_MAX_ATTEMPTS = 2 + +# "Not row-crop" verdicts are cached so the ~7,300 products we discard every +# month cost zero API calls. Verdicts expire so a product whose registered +# sites later change still gets re-examined. +FILTER_CACHE_TTL_DAYS = 180 +FILTER_CACHE_PATH = REPO_ROOT / "scrape" / "state" / "epa_ppls_filtered.json" + # Row-crop scoping. Each pattern is matched case-insensitively against a # product's "sites" array from the PPLS API. Word boundaries matter — bare # "OATS" naively matches "SHIPS, BOATS, SHIPHOLDS"; bare "RICE" matches @@ -145,6 +165,103 @@ def load_registrant_allowlist() -> set[str]: log = logging.getLogger("epa_ppls") +# --------------------------------------------------------------------------- +# Politeness + verdict cache +# --------------------------------------------------------------------------- + + +class _RateLimiter: + """Aggregate request pacing across worker threads.""" + + def __init__(self, rps: float) -> None: + self._interval = 1.0 / rps if rps > 0 else 0.0 + self._lock = threading.Lock() + self._next = 0.0 + + def acquire(self) -> None: + if not self._interval: + return + with self._lock: + now = time.monotonic() + wait = max(0.0, self._next - now) + self._next = max(now, self._next) + self._interval + if wait: + time.sleep(wait) + + +_LIMITER = _RateLimiter(RATE_LIMIT_RPS) + + +class FilterCache: + """Remembers which reg nos were judged not-row-crop, and when. + + Committed to git so CI benefits: without it every run re-fetches ~7,300 + products from the ORDS API purely to discard them again. + """ + + def __init__(self, path: Path, ttl_days: int = FILTER_CACHE_TTL_DAYS) -> None: + self.path = path + self.ttl_days = ttl_days + self._lock = threading.Lock() + self._entries: dict[str, str] = {} + self._dirty = False + + def load(self) -> "FilterCache": + try: + data = json.loads(self.path.read_text(encoding="utf-8")) + self._entries = dict(data.get("entries") or {}) + except FileNotFoundError: + pass + except (OSError, ValueError) as exc: + log.warning("filter cache at %s unreadable (%s) — starting empty", + self.path, exc) + return self + + def _ttl_for(self, regno: str) -> int: + """TTL with a deterministic per-product jitter. + + A cold run stamps every verdict on the same day, so a flat TTL would + expire all ~7,300 of them in the same month and resurrect the 3-hour + run we are fixing. Spreading expiry over a +/- 30 day window keeps the + re-examination cost roughly level month to month. + """ + jitter = zlib.crc32(regno.encode()) % 61 - 30 + return max(1, self.ttl_days + jitter) + + def is_fresh(self, regno: str) -> bool: + with self._lock: + stamp = self._entries.get(regno) + if not stamp: + return False + try: + decided = datetime.fromisoformat(stamp) + except ValueError: + return False + return (datetime.now(UTC) - decided).days < self._ttl_for(regno) + + def remember(self, regno: str) -> None: + with self._lock: + self._entries[regno] = datetime.now(UTC).isoformat() + self._dirty = True + + def save(self) -> None: + if not self._dirty: + return + self.path.parent.mkdir(parents=True, exist_ok=True) + with self._lock: + payload = { + "version": 1, + "ttl_days": self.ttl_days, + "entries": dict(sorted(self._entries.items())), + } + self.path.write_text(json.dumps(payload, indent=1) + "\n", encoding="utf-8") + log.info("filter cache: %d verdicts → %s", len(payload["entries"]), self.path) + + +class PdfDownloadTimeout(Exception): + """A label body was still arriving when its overall deadline expired.""" + + # --------------------------------------------------------------------------- # HTTP helpers # --------------------------------------------------------------------------- @@ -165,6 +282,7 @@ def _get_with_retries( last_exc: Exception | None = None for attempt in range(1, MAX_RETRIES + 1): try: + _LIMITER.acquire() resp = client.get(url) if resp.status_code in (429, 500, 502, 503, 504): wait = min(2 ** attempt, 30) @@ -395,11 +513,40 @@ def fetch_product_record(client: httpx.Client, regno: str) -> ProductRecord: # --------------------------------------------------------------------------- -def download_pdf(client: httpx.Client, url: str) -> tuple[bytes, str | None]: - """Download a label PDF; return (bytes, Last-Modified header or None).""" - resp = _get_with_retries(client, url) - last_modified = resp.headers.get("last-modified") - return resp.content, last_modified +def download_pdf( + client: httpx.Client, url: str, *, deadline: float = PDF_DEADLINE_SECONDS +) -> tuple[bytes, str | None]: + """Download a label PDF; return (bytes, Last-Modified header or None). + + Streams the body and enforces an overall wall-clock deadline. A plain + ``client.get(...).content`` cannot do this: httpx's read timeout is + per-read, so a body arriving a few KB at a time keeps resetting it and the + call can hang for as long as EPA keeps dribbling bytes. + """ + last_exc: Exception | None = None + for attempt in range(1, PDF_MAX_ATTEMPTS + 1): + _LIMITER.acquire() + started = time.monotonic() + try: + with client.stream("GET", url) as resp: + resp.raise_for_status() + last_modified = resp.headers.get("last-modified") + buf = bytearray() + for chunk in resp.iter_bytes(): + buf += chunk + if time.monotonic() - started > deadline: + raise PdfDownloadTimeout( + f"still downloading after {deadline:.0f}s " + f"({len(buf)} bytes received)" + ) + return bytes(buf), last_modified + except (PdfDownloadTimeout, httpx.TransportError, httpx.HTTPError) as exc: + last_exc = exc + log.warning( + "PDF download failed on %s (attempt %d/%d, %.0fs): %s", + url, attempt, PDF_MAX_ATTEMPTS, time.monotonic() - started, exc, + ) + raise RuntimeError(f"GET {url} failed after {PDF_MAX_ATTEMPTS} attempts: {last_exc}") def extract_pdf_text(pdf_bytes: bytes) -> tuple[str, bool]: @@ -441,32 +588,59 @@ def _json_path(regno: str) -> Path: return CORPUS_DIR / f"{regno}.json" +def _label_is_current(json_path: Path, record: "ProductRecord") -> bool: + """True when the committed sidecar already describes EPA's current label.""" + try: + sidecar = json.loads(json_path.read_text(encoding="utf-8")) + except (OSError, ValueError): + return False + label = sidecar.get("label") or {} + return ( + label.get("accepted_date") == record.label_accepted_date + and label.get("url") == record.label_pdf_url + ) + + def process_one( client: httpx.Client, regno: str, *, force: bool = False, row_crop_filter: bool = True, + filter_cache: "FilterCache | None" = None, ) -> str: """Fetch + extract one product. Returns - 'skipped'|'wrote'|'no-pdf'|'error'|'filtered'.""" + 'wrote'|'unchanged'|'no-pdf'|'error'|'filtered'|'filtered-cached'. + + Without ``force`` this is an incremental refresh: a cached not-row-crop + verdict costs nothing, and a product already on disk is re-downloaded only + when EPA reports a different label acceptance date or PDF URL. + """ md_path = _md_path(regno) json_path = _json_path(regno) - if not force and md_path.exists() and json_path.exists(): - log.info("[%s] skip (already on disk)", regno) - return "skipped" + + if not force and row_crop_filter and filter_cache and filter_cache.is_fresh(regno): + log.debug("[%s] filtered (cached verdict)", regno) + return "filtered-cached" try: record = fetch_product_record(client, regno) except Exception as exc: log.error("[%s] API fetch failed: %s", regno, exc) return "error" - time.sleep(REQUEST_DELAY_SECONDS) if row_crop_filter and not matches_row_crop(record): log.info("[%s] filtered (not row-crop)", regno) + if filter_cache: + filter_cache.remember(regno) return "filtered" + if not force and md_path.exists() and json_path.exists(): + if _label_is_current(json_path, record): + log.debug("[%s] unchanged (label %s)", regno, record.label_accepted_date) + return "unchanged" + log.info("[%s] label changed — re-fetching", regno) + def _build_sidecar( *, label_url: str | None, @@ -527,7 +701,6 @@ def process_one( except Exception as exc: log.error("[%s] PDF download failed: %s", regno, exc) return "error" - time.sleep(REQUEST_DELAY_SECONDS) text, has_text = extract_pdf_text(pdf_bytes) last_modified_iso = _http_date_to_iso(last_modified_raw) @@ -670,6 +843,18 @@ def main(argv: list[str] | None = None) -> int: "of the PPIS universe is non-ag. --no-registrant-filter to enumerate " "everything.", ) + parser.add_argument( + "--workers", type=int, default=DEFAULT_WORKERS, metavar="N", + help=f"Parallel product workers (default {DEFAULT_WORKERS}). The " + f"aggregate request rate stays capped at {RATE_LIMIT_RPS:g}/sec " + "regardless, so this shortens the run without hitting EPA harder.", + ) + parser.add_argument( + "--filter-cache", action=argparse.BooleanOptionalAction, default=True, + help="Reuse recorded 'not row-crop' verdicts instead of re-fetching " + f"those products (verdicts expire after {FILTER_CACHE_TTL_DAYS} " + "days). Default on; --no-filter-cache re-examines everything.", + ) parser.add_argument( "--log-level", default="INFO", choices=["DEBUG", "INFO", "WARNING", "ERROR"], @@ -684,17 +869,37 @@ def main(argv: list[str] | None = None) -> int: CORPUS_DIR.mkdir(parents=True, exist_ok=True) - summary = {"wrote": 0, "skipped": 0, "no-pdf": 0, "filtered": 0, "error": 0} + cache = FilterCache(FILTER_CACHE_PATH).load() if args.filter_cache else None + workers = max(1, args.workers) + summary = { + "wrote": 0, "unchanged": 0, "no-pdf": 0, + "filtered": 0, "filtered-cached": 0, "error": 0, + } + started = time.monotonic() with _client() as client: - for regno in _iter_regnos(args, client): - result = process_one( - client, regno, - force=args.force, - row_crop_filter=args.row_crop_filter, - ) - summary[result] = summary.get(result, 0) + 1 + with ThreadPoolExecutor(max_workers=workers) as pool: + futures = [ + pool.submit( + process_one, client, regno, + force=args.force, + row_crop_filter=args.row_crop_filter, + filter_cache=cache, + ) + for regno in _iter_regnos(args, client) + ] + log.info("processing %d products with %d workers", len(futures), workers) + for future in as_completed(futures): + try: + result = future.result() + except Exception as exc: # a worker should never escape, but + log.error("worker raised: %s", exc) + result = "error" + summary[result] = summary.get(result, 0) + 1 - log.info("done: %s", summary) + if cache: + cache.save() + + log.info("done in %.1f min: %s", (time.monotonic() - started) / 60, summary) print(json.dumps(summary), file=sys.stderr) return 0 diff --git a/scrape/state/epa_ppls_filtered.json b/scrape/state/epa_ppls_filtered.json new file mode 100644 index 0000000..cf0ec63 --- /dev/null +++ b/scrape/state/epa_ppls_filtered.json @@ -0,0 +1,5 @@ +{ + "version": 1, + "ttl_days": 180, + "entries": {} +}