AutonomousAgent / test_parallel_ingest_e2e.py
Mathias Heider
Claude Fable 5.1
Budget hold: no downloads while Gemini refuses for budget reasons, one probe per cycle, automatic resume; adoption capped; target default
c473e91 unverified
Raw History Blame Contribute Delete
16 kB
"""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=<redacted>'")[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)