perf(epa_ppls): make the monthly refresh fit the runner's 3 h budget #3

Merged
claude merged 2 commits from fix/epa-ppls-fit-in-runner-budget into main 2026-09-01 16:40:50 -04:00
4 changed files with 267 additions and 34 deletions
Showing only changes of commit 81a5eb2b76 - Show all commits
+29 -11
View File
@@ -5,11 +5,13 @@ name: Monthly corpus refresh
# reindex + image-push if the scrape produced no diff against the # reindex + image-push if the scrape produced no diff against the
# committed corpus. # committed corpus.
# #
# Bayer takes ~30 min; EPA PPLS takes ~7 h with row-crop + # Runtime budget: act_runner kills the job container at exactly 3 h
# registrant filters. The whole monthly job is ~8-9 h end-to-end. # (it runs as `/bin/sleep 10800`), which is what killed every run from
# If that's too long for the runner you can: # 2026-06-01 through 2026-09-01. The EPA step is incremental and
# - Run just one source: workflow_dispatch with sources="bayer" # parallel so a steady-state refresh is ~15-20 min: cached not-row-crop
# - Limit EPA at the scraper: edit the step to add "--limit 5000" # 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: on:
schedule: schedule:
@@ -42,6 +44,9 @@ env:
jobs: jobs:
refresh: refresh:
runs-on: docker 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: container:
image: catthehacker/ubuntu:act-latest image: catthehacker/ubuntu:act-latest
steps: steps:
@@ -80,9 +85,14 @@ jobs:
- name: Scrape EPA PPLS - name: Scrape EPA PPLS
if: ${{ inputs.sources == '' || contains(inputs.sources, 'epa_ppls') }} if: ${{ inputs.sources == '' || contains(inputs.sources, 'epa_ppls') }}
# Row-crop + registrant filters keep this to ~16K PDFs / ~7h. # Deliberately NOT --force: that re-downloaded all ~11.4K candidate
# Pass --no-row-crop-filter or --no-registrant-filter to broaden. # registrations every month (~15-20 h of work) and never finished.
run: python -m scrape.runner --source epa_ppls --force # 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 ----------------- # ---- Commit corpus changes + retry-on-race -----------------
- name: Commit corpus changes (if any) - name: Commit corpus changes (if any)
@@ -90,13 +100,21 @@ jobs:
run: | run: |
git config user.name "crop-chem-docs-refresh" git config user.name "crop-chem-docs-refresh"
git config user.email "[email protected]" git config user.email "[email protected]"
git add sources.json corpus # The filter-verdict cache is committed so next month starts warm,
if git diff --cached --quiet; then # 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 "no corpus changes — skipping reindex and image build"
echo "changed=false" >> "$GITHUB_OUTPUT" 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 exit 0
fi fi
echo "changed=true" >> "$GITHUB_OUTPUT"
ts=$(date -u +"%Y-%m-%dT%H:%MZ") ts=$(date -u +"%Y-%m-%dT%H:%MZ")
n_bayer=$(find corpus/bayer -name '*.json' 2>/dev/null | wc -l) 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) n_epa=$(find corpus/epa_ppls -name '*.json' 2>/dev/null | wc -l)
+7 -2
View File
@@ -59,9 +59,14 @@ pip install -r requirements.txt
# Sample-scrape to verify wiring: # Sample-scrape to verify wiring:
python -m scrape.runner --source bayer --limit 5 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 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: # Rebuild Chroma + BM25:
OLLAMA_URL=http://192.168.0.125:11434 PRODUCT_NAME=crop_chem \ OLLAMA_URL=http://192.168.0.125:11434 PRODUCT_NAME=crop_chem \
+226 -21
View File
@@ -51,8 +51,11 @@ import logging
import os import os
import re import re
import sys import sys
import threading
import time import time
import zipfile import zipfile
import zlib
from concurrent.futures import ThreadPoolExecutor, as_completed
from dataclasses import dataclass, field from dataclasses import dataclass, field
from datetime import UTC, datetime from datetime import UTC, datetime
from pathlib import Path 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_ROOT = Path(os.environ.get("CORPUS_ROOT") or REPO_ROOT / "corpus")
CORPUS_DIR = CORPUS_ROOT / "epa_ppls" 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) HTTP_TIMEOUT = httpx.Timeout(60.0, connect=15.0)
MAX_RETRIES = 4 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 # Row-crop scoping. Each pattern is matched case-insensitively against a
# product's "sites" array from the PPLS API. Word boundaries matter — bare # product's "sites" array from the PPLS API. Word boundaries matter — bare
# "OATS" naively matches "SHIPS, BOATS, SHIPHOLDS"; bare "RICE" matches # "OATS" naively matches "SHIPS, BOATS, SHIPHOLDS"; bare "RICE" matches
@@ -145,6 +165,103 @@ def load_registrant_allowlist() -> set[str]:
log = logging.getLogger("epa_ppls") 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 # HTTP helpers
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
@@ -165,6 +282,7 @@ def _get_with_retries(
last_exc: Exception | None = None last_exc: Exception | None = None
for attempt in range(1, MAX_RETRIES + 1): for attempt in range(1, MAX_RETRIES + 1):
try: try:
_LIMITER.acquire()
resp = client.get(url) resp = client.get(url)
if resp.status_code in (429, 500, 502, 503, 504): if resp.status_code in (429, 500, 502, 503, 504):
wait = min(2 ** attempt, 30) 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]: def download_pdf(
"""Download a label PDF; return (bytes, Last-Modified header or None).""" client: httpx.Client, url: str, *, deadline: float = PDF_DEADLINE_SECONDS
resp = _get_with_retries(client, url) ) -> tuple[bytes, str | None]:
last_modified = resp.headers.get("last-modified") """Download a label PDF; return (bytes, Last-Modified header or None).
return resp.content, last_modified
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]: 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" 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( def process_one(
client: httpx.Client, client: httpx.Client,
regno: str, regno: str,
*, *,
force: bool = False, force: bool = False,
row_crop_filter: bool = True, row_crop_filter: bool = True,
filter_cache: "FilterCache | None" = None,
) -> str: ) -> str:
"""Fetch + extract one product. Returns """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) md_path = _md_path(regno)
json_path = _json_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) if not force and row_crop_filter and filter_cache and filter_cache.is_fresh(regno):
return "skipped" log.debug("[%s] filtered (cached verdict)", regno)
return "filtered-cached"
try: try:
record = fetch_product_record(client, regno) record = fetch_product_record(client, regno)
except Exception as exc: except Exception as exc:
log.error("[%s] API fetch failed: %s", regno, exc) log.error("[%s] API fetch failed: %s", regno, exc)
return "error" return "error"
time.sleep(REQUEST_DELAY_SECONDS)
if row_crop_filter and not matches_row_crop(record): if row_crop_filter and not matches_row_crop(record):
log.info("[%s] filtered (not row-crop)", regno) log.info("[%s] filtered (not row-crop)", regno)
if filter_cache:
filter_cache.remember(regno)
return "filtered" 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( def _build_sidecar(
*, *,
label_url: str | None, label_url: str | None,
@@ -527,7 +701,6 @@ def process_one(
except Exception as exc: except Exception as exc:
log.error("[%s] PDF download failed: %s", regno, exc) log.error("[%s] PDF download failed: %s", regno, exc)
return "error" return "error"
time.sleep(REQUEST_DELAY_SECONDS)
text, has_text = extract_pdf_text(pdf_bytes) text, has_text = extract_pdf_text(pdf_bytes)
last_modified_iso = _http_date_to_iso(last_modified_raw) 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 " "of the PPIS universe is non-ag. --no-registrant-filter to enumerate "
"everything.", "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( parser.add_argument(
"--log-level", default="INFO", "--log-level", default="INFO",
choices=["DEBUG", "INFO", "WARNING", "ERROR"], 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) 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: with _client() as client:
for regno in _iter_regnos(args, client): with ThreadPoolExecutor(max_workers=workers) as pool:
result = process_one( futures = [
client, regno, pool.submit(
force=args.force, process_one, client, regno,
row_crop_filter=args.row_crop_filter, force=args.force,
) row_crop_filter=args.row_crop_filter,
summary[result] = summary.get(result, 0) + 1 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) print(json.dumps(summary), file=sys.stderr)
return 0 return 0
+5
View File
@@ -0,0 +1,5 @@
{
"version": 1,
"ttl_days": 180,
"entries": {}
}