"""E2E test (scratch Postgres, no network) for parallel ingestion (28 Sep 2026): a cycle with several downloads extracts them on a thread pool, one Postgres connection per worker, and the counters / registry / events come out exactly as they did one PDF at a time. python test_parallel_ingest_e2e.py (same DB_* env as test_e2e.py) """ import dataclasses import hashlib import os import sys import threading os.environ.setdefault("AGENT_WORK_DIR", "/tmp/agent_parallel_test") os.environ["GEMINI_API_KEY"] = "test-key-not-used" os.environ["AGENT_USE_NTRS"] = "0" # lanes are stubbed; never touch the network os.environ["AGENT_USE_S2"] = "0" os.environ["AGENT_USE_OPENALEX_TOPICS"] = "0" os.environ["AGENT_USE_DATASHEETS"] = "0" os.environ["AGENT_SEED_QUERY_GRID"] = "0" import fitz import pdf_crawler import extraction import batch_ingest from agent import agdb, config as C, orchestrator C.GEMINI_API_KEY = "test-key-not-used" fails = [] def check(name, cond): print(("PASS " if cond else "FAIL ") + name) if not cond: fails.append(name) MATS = { "peek": ("Polyetheretherketone 30% Carbon Fiber", "PEEK-CF30", "PEEK", "Carbon", "Tensile Modulus", "24.5", "The tensile modulus of PEEK-CF30 was measured as 24.5 GPa at 23 C."), "pa66": ("Polyamide 66 30% Glass Fiber", "PA66-GF30", "PA66", "Glass", "Flexural Modulus", "8.9", "The flexural modulus of PA66-GF30 was measured as 8.9 GPa at 23 C."), "pps": ("Polyphenylene Sulfide 40% Glass Fiber", "PPS-GF40", "PPS", "Glass", "Tensile Strength", "190", "The tensile strength of PPS-GF40 was measured as 190 MPa at 23 C."), "pekk": ("Polyetherketoneketone 20% Carbon Fiber", "PEKK-CF20", "PEKK", "Carbon", "Tensile Modulus", "18.2", "The tensile modulus of PEKK-CF20 was measured as 18.2 GPa at 23 C."), } def make_pdf(key: str) -> bytes: name, abbr, matrix, fiber, prop, val, quote = MATS[key] doc = fitz.open() page = doc.new_page() page.insert_text((72, 100), f"{abbr} composite datasheet") page.insert_text((72, 130), quote) data = doc.tobytes() doc.close() return data PDFS = {k: make_pdf(k) for k in MATS} THREADS_SEEN: set = set() CONCURRENT = {"now": 0, "max": 0} _lock = threading.Lock() def fake_search_openalex(query, limit): for k in MATS: yield pdf_crawler.Candidate(title=f"{k} laminates study", pdf_url=f"https://example.org/{k}.pdf", source="openalex", query=query, doi=f"10.9999/aim.par.{k}", year="2026", abstract="thermoplastic composite tensile modulus carbon fiber datasheet") def fake_search_arxiv(query, limit): return iter(()) def fake_download(cand, pdf_dir, state): if cand.pdf_url in state.seen_urls: return None state.seen_urls.add(cand.pdf_url) key = cand.pdf_url.rsplit("/", 1)[1][:-4] data = PDFS[key] sha = hashlib.sha256(data).hexdigest() if sha in state.seen_hashes: return None state.seen_hashes.add(sha) fname = f"openalex_{key}_{sha[:8]}.pdf" (pdf_dir / fname).write_bytes(data) return {"filename": fname, "title": cand.title, "doi": cand.doi, "url": cand.pdf_url, "year": cand.year, "source": cand.source, "sha256": sha, "query": cand.query, "bytes": len(data)} def fake_extract(pdf_bytes, filename, api_key): import time with _lock: CONCURRENT["now"] += 1 CONCURRENT["max"] = max(CONCURRENT["max"], CONCURRENT["now"]) THREADS_SEEN.add(threading.current_thread().name) time.sleep(0.4) # a "Gemini call": long enough to overlap with _lock: CONCURRENT["now"] -= 1 key = next(k for k in MATS if f"_{k}_" in filename) name, abbr, matrix, fiber, prop, val, quote = MATS[key] unit = "MPa" if prop == "Tensile Strength" else "GPa" mat = extraction.Material( material_name=name, material_abbreviation=abbr, material_class="Composite", matrix=matrix, fiber=fiber, fiber_volume_fraction="30%", properties=[extraction.Property(section="Mechanical", property_name=prop, value_raw=val, value_num=float(val), unit=unit, test_condition="23 C", source_quote=quote, page=1)]) ext = extraction.Extraction(materials=[mat], doc_status="ok") ext.tokens_in, ext.tokens_out = 1000, 50 return ext pdf_crawler.search_openalex = fake_search_openalex pdf_crawler.search_arxiv = fake_search_arxiv pdf_crawler.download_pdf = fake_download batch_ingest.extract_from_pdf = fake_extract def q(sql, params=()): conn = agdb.connect() try: with conn.cursor() as cur: cur.execute(sql, params) if params else cur.execute(sql) try: return cur.fetchall() except Exception: return None finally: conn.commit() conn.close() orchestrator.bootstrap() conn = agdb.connect() agdb.set_config(conn, {"queries_per_cycle": 1, "max_rounds_per_cycle": 1, "max_new_pdfs_per_cycle": 5, "min_new_pdfs_per_cycle": 5, "figures_enabled": False, "link_figures": False, "expand_queries": False, "ingest_workers": 3}) conn.close() m = orchestrator.run_cycle(trigger="test") row = q("SELECT status, downloaded, pdfs_ingested, rows_inserted, errors, tokens_in, tokens_out " "FROM agent_runs ORDER BY id DESC LIMIT 1")[0] check("cycle ok: 4 downloaded, 4 ingested, 4 rows, 0 errors", row[:5] == ("ok", 4, 4, 4, 0)) check("tokens summed across workers", (row[5], row[6]) == (4000, 200)) check("extraction really ran concurrently (>= 2 PDFs in flight at once)", CONCURRENT["max"] >= 2) check("workers ran on pool threads", any(t.startswith("ingest") for t in THREADS_SEEN)) reg = q("SELECT count(*) FROM agent_doi_seen WHERE ingest_status = 'ingested'")[0][0] check("registry: all 4 rows 'ingested'", reg == 4) ev = q("SELECT count(*) FROM agent_events WHERE node = 'ingest' AND message LIKE '%rows inserted%'")[0][0] check("one ingest event per PDF", ev == 4) ev = q("SELECT count(*) FROM agent_events WHERE message LIKE 'ingesting 4 PDF(s) with 3 parallel workers'")[0][0] check("parallel-ingest event logged", ev == 1) check("PDFs removed after ingest", list(C.PDF_DIR.glob("*.pdf")) == []) rows = q('SELECT count(*) FROM "Composites_materials" WHERE source_sha1 IS NOT NULL')[0][0] check("4 material rows landed", rows == 4) credited = q("SELECT pdfs_found, rows_yielded FROM agent_queries WHERE pdfs_found > 0") check("the query is credited with 4 PDFs / 4 rows", credited == [(4, 4)]) # workers = 1 keeps the sequential path THREADS_SEEN.clear(); CONCURRENT["max"] = 0 q("DELETE FROM agent_doi_seen"); q('DELETE FROM "Composites_materials"'); q("DELETE FROM sources") conn = agdb.connect() agdb.set_config(conn, {"ingest_workers": 1}) agdb.set_state(conn, "crawler_state", {"seen_urls": [], "seen_hashes": [], "failed": {}}) conn.close() orchestrator.run_cycle(trigger="test") row = q("SELECT status, downloaded, pdfs_ingested, rows_inserted, errors FROM agent_runs ORDER BY id DESC LIMIT 1")[0] check("ingest_workers=1: same result one PDF at a time", row == ("ok", 4, 4, 4, 0) and CONCURRENT["max"] == 1) ev = q("SELECT count(*) FROM agent_events WHERE message LIKE 'ingesting % parallel workers'")[0][0] check("no parallel-ingest event when sequential", ev == 1) # --- Gemini spending cap mid-cycle: PDFs kept, next cycle adopts them; leak redacted ----- import requests QUOTA = {"on": False, "calls": 0} _real_extract = fake_extract def capped_extract(pdf_bytes, filename, api_key): QUOTA["calls"] += 1 if QUOTA["on"]: raise extraction.GeminiHTTPError( "429 RESOURCE_EXHAUSTED from Gemini: Your project has exceeded its monthly spending cap.", None, "RESOURCE_EXHAUSTED", "Your project has exceeded its monthly spending cap.") return _real_extract(pdf_bytes, filename, api_key) batch_ingest.extract_from_pdf = capped_extract q("DELETE FROM agent_doi_seen"); q('DELETE FROM "Composites_materials"'); q("DELETE FROM sources") conn = agdb.connect() agdb.set_config(conn, {"ingest_workers": 3}) agdb.set_state(conn, "crawler_state", {"seen_urls": [], "seen_hashes": [], "failed": {}}) conn.close() QUOTA["on"] = True; QUOTA["calls"] = 0 m_cap = orchestrator.run_cycle(trigger="test") row = q("SELECT status, downloaded, pdfs_ingested, rows_inserted, errors FROM agent_runs ORDER BY id DESC LIMIT 1")[0] check("spending cap: run ends 'error' (nothing ingested), 4 downloaded, 0 ingested, 0 errors", row == ("error", 4, 0, 0, 0)) check("spending cap: at most the in-flight workers called Gemini, the rest were stopped", 1 <= QUOTA["calls"] <= 3) reg = q("SELECT ingest_status, count(*) FROM agent_doi_seen GROUP BY 1") check("spending cap: every registry row stays 'downloaded'", reg == [("downloaded", 4)]) check("spending cap: the 4 PDFs stay on disk", len(list(C.PDF_DIR.glob("*.pdf"))) == 4) ev = q("SELECT message FROM agent_events WHERE run_id = (SELECT max(id) FROM agent_runs) AND level = 'warn' " "AND message LIKE 'Gemini refused for budget reasons%'") check("spending cap: one warn event naming the cause and the kept count", len(ev) == 1 and "4 PDF(s) kept" in ev[0][0] and "spending cap" in ev[0][0]) rep = q("SELECT report FROM agent_runs ORDER BY id DESC LIMIT 1")[0][0] or "" check("spending cap: the reason is in the run report (Why this run ended)", "spending cap" in rep) check("spending cap: nothing stored contains a key parameter", q("SELECT count(*) FROM agent_events WHERE message ~* '[?&]key='")[0][0] == 0) # --- budget hold: the next cycle downloads nothing and probes with ONE kept PDF ------------ hold = agdb.get_state(agdb.connect(), orchestrator.BUDGET_HOLD_KEY, None) check("spending cap: budget hold recorded with the reason", isinstance(hold, dict) and "spending cap" in hold.get("reason", "")) QUOTA["calls"] = 0 m_hold = orchestrator.run_cycle(trigger="test") row = q("SELECT status, downloaded, pdfs_ingested, candidates, queries_used FROM agent_runs ORDER BY id DESC LIMIT 1")[0] check("hold cycle: no search, no download, still refused → error", row == ("error", 0, 0, 0, 0)) check("hold cycle: exactly one Gemini probe call", QUOTA["calls"] == 1) check("hold cycle: the 4 PDFs still on disk", len(list(C.PDF_DIR.glob("*.pdf"))) == 4) ev = q("SELECT count(*) FROM agent_events WHERE node = 'plan' AND message LIKE '%Gemini budget hold since%'")[0][0] check("hold cycle: plan event says so", ev == 1) QUOTA["on"] = False; QUOTA["calls"] = 0 m_probe = orchestrator.run_cycle(trigger="test") row = q("SELECT status, downloaded, pdfs_ingested, rows_inserted FROM agent_runs ORDER BY id DESC LIMIT 1")[0] check("budget back: the probe cycle ingests its one PDF and lifts the hold (0 downloaded, 1 ingested)", row == ("ok", 0, 1, 1) and QUOTA["calls"] == 1) check("budget back: hold cleared", agdb.get_state(agdb.connect(), orchestrator.BUDGET_HOLD_KEY, None) is None) ev = q("SELECT count(*) FROM agent_events WHERE message LIKE 'Gemini budget hold lifted%'")[0][0] check("budget back: lift event logged", ev == 1) m_after = orchestrator.run_cycle(trigger="test") row = q("SELECT status, downloaded, pdfs_ingested, rows_inserted FROM agent_runs ORDER BY id DESC LIMIT 1")[0] check("budget back: the next cycle adopts the remaining 3 kept PDFs and ingests them", row == ("ok", 0, 3, 3)) check("budget back: registry rows now 'ingested'", q("SELECT count(*) FROM agent_doi_seen WHERE ingest_status = 'ingested'")[0][0] == 4) check("budget back: PDFs removed", list(C.PDF_DIR.glob("*.pdf")) == []) ev = q("SELECT count(*) FROM agent_events WHERE message LIKE 'adopted 1 PDF(s)%budget probe%'")[0][0] check("probe adoption event logged", ev >= 1) ev = q("SELECT count(*) FROM agent_events WHERE message LIKE 'adopted 3 PDF(s)%'")[0][0] check("budget back: adoption event for the rest logged", ev == 1) # --- hold with nothing on disk (container restarted): download exactly one PDF as the probe ---- conn = agdb.connect() agdb.set_state(conn, orchestrator.BUDGET_HOLD_KEY, {"at": "2026-09-29T02:45:00+00:00", "reason": "test hold"}) agdb.set_state(conn, "crawler_state", {"seen_urls": [], "seen_hashes": [], "failed": {}}) conn.close() q("DELETE FROM agent_doi_seen"); q('DELETE FROM "Composites_materials"'); q("DELETE FROM sources") QUOTA["on"] = False; QUOTA["calls"] = 0 orchestrator.run_cycle(trigger="test") row = q("SELECT status, downloaded, pdfs_ingested FROM agent_runs ORDER BY id DESC LIMIT 1")[0] check("hold with no kept PDF: exactly one download as the probe, ingested, hold lifted", row == ("ok", 1, 1) and QUOTA["calls"] == 1 and agdb.get_state(agdb.connect(), orchestrator.BUDGET_HOLD_KEY, None) is None) batch_ingest.extract_from_pdf = _real_extract # boot repair redacts a leak stored before redact() existed (planted raw) q("INSERT INTO agent_events (level, node, message) VALUES ('warn', 'ingest', " "'old.pdf: gemini_error:429 Client Error: Too Many Requests for url: https://g/v1beta/x:generateContent?key=AQ.PLANTEDSECRET')") q("UPDATE agent_runs SET report = '{\"pdf_results\":[{\"error\":\"gemini_error:429 for url: https://g/x?key=AQ.PLANTEDSECRET2\"}]}' " "WHERE id = (SELECT min(id) FROM agent_runs)") orchestrator.bootstrap() check("boot repair: planted key removed from agent_events and agent_runs.report", q("SELECT count(*) FROM agent_events WHERE message LIKE '%PLANTEDSECRET%'")[0][0] == 0 and q("SELECT count(*) FROM agent_runs WHERE report LIKE '%PLANTEDSECRET2%'")[0][0] == 0 and q("SELECT count(*) FROM agent_events WHERE message LIKE '%key='")[0][0] >= 1) ev = q("SELECT message FROM agent_events WHERE node = 'boot' AND message LIKE 'boot repair: redacted%' ORDER BY id DESC LIMIT 1") check("boot repair: event names the tables and asks for a rotation", len(ev) == 1 and "agent_events: 1" in ev[0][0] and "agent_runs: 1" in ev[0][0] and "rotate" in ev[0][0]) orchestrator.bootstrap() ev2 = q("SELECT count(*) FROM agent_events WHERE node = 'boot' AND message LIKE 'boot repair: redacted%'")[0][0] check("boot repair: idempotent (no second event)", ev2 == 1) # --- cost defaults: one-shot boot migration + the cycle applies the thinking setting ------ q("DELETE FROM agent_state WHERE key = 'cost_defaults_v1'") conn = agdb.connect() agdb.set_config(conn, {"figure_mining": True, "gemini_thinking": "dynamic", "min_new_pdfs_per_cycle": 3, "max_new_pdfs_per_cycle": 5}) # an operator's old config conn.close() orchestrator.bootstrap() conn = agdb.connect(); cfg_now = agdb.get_config(conn); conn.close() check("boot: cost defaults applied once to a stored config (mining off, thinking low, target = cap)", cfg_now.get("figure_mining") is False and cfg_now.get("gemini_thinking") == "low" and int(cfg_now.get("min_new_pdfs_per_cycle")) == int(cfg_now.get("max_new_pdfs_per_cycle"))) ev = q("SELECT count(*) FROM agent_events WHERE node = 'boot' AND message LIKE 'cost defaults applied:%'")[0][0] check("boot: cost-defaults event logged once", ev == 1) conn = agdb.connect() agdb.set_config(conn, {"gemini_thinking": "off"}) # the operator's choice conn.close() orchestrator.bootstrap() conn = agdb.connect(); cfg_now = agdb.get_config(conn); conn.close() check("boot: a later operator choice is not overridden (one-shot)", cfg_now.get("gemini_thinking") == "off") extraction.THINKING = "dynamic" q("DELETE FROM agent_doi_seen"); q('DELETE FROM "Composites_materials"'); q("DELETE FROM sources") conn = agdb.connect() agdb.set_state(conn, "crawler_state", {"seen_urls": [], "seen_hashes": [], "failed": {}}) conn.close() orchestrator.run_cycle(trigger="test") check("cycle: extraction.THINKING follows the config ('off')", extraction.THINKING == "off") print() print("all tests passed" if not fails else f"{len(fails)} FAILED: {fails}") sys.exit(1 if fails else 0)