Spaces:
Running
Running
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 Download test_parallel_ingest_e2e.py from aim4composites/AutonomousAgent: direct link, hf CLI and curl.
- Browser
- Download file 16 kB
-
https://huggingface.co/spaces/aim4composites/AutonomousAgent/resolve/main/test_parallel_ingest_e2e.py
- Command line
-
hf download hf://spaces/aim4composites/AutonomousAgent/test_parallel_ingest_e2e.py
-
curl -L -o test_parallel_ingest_e2e.py https://huggingface.co/spaces/aim4composites/AutonomousAgent/resolve/main/test_parallel_ingest_e2e.py
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) | |