Spaces:
Running
Running
Mathias Heider
Claude Fable 5.1
Merge the 17 Sep "more papers" lanes onto the hotfix line; parallel ingestion; throughput knobs
6bb7525 unverified Download test_counters_e2e.py from aim4composites/AutonomousAgent: direct link, hf CLI and curl.
- Browser
- Download file 15.8 kB
-
https://huggingface.co/spaces/aim4composites/AutonomousAgent/resolve/main/test_counters_e2e.py
- Command line
-
hf download hf://spaces/aim4composites/AutonomousAgent/test_counters_e2e.py
-
curl -L -o test_counters_e2e.py https://huggingface.co/spaces/aim4composites/AutonomousAgent/resolve/main/test_counters_e2e.py
15.8 kB
| """E2E for the 22 Sep instrumentation, against a local Postgres. | |
| Covers, with network discovery + Gemini stubbed and everything else real: | |
| * agent_runs counters that start the funnel before the download gate | |
| (relevance_rejected, download_failed, url_seen_skipped), Gemini tokens, | |
| rows_linked, the per-source jsonb, heartbeat_at | |
| * figure citation linking wired into node_ingest (a text row whose quote | |
| cites "Figure 1" gets figure_id/figure_ref + the embedded PNG) and the | |
| figures_ready_problems gate degrading both figure stages with a warn | |
| * the stall detector: a cycle that stops emitting events is abandoned | |
| long before the hard budget; one that keeps reporting progress is not | |
| * sweep_stale_runs measuring age from the last heartbeat | |
| """ | |
| import os | |
| import sys | |
| import threading | |
| import time | |
| os.environ.setdefault("AGENT_WORK_DIR", "/tmp/agent_test") | |
| os.environ["GEMINI_API_KEY"] = "test-key-not-used" | |
| os.environ["AGENT_USE_NTRS"] = "0" # tests stub the lanes; 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" | |
| os.environ["AGENT_KEEPALIVE_URL"] = "" | |
| import fitz # noqa: E402 | |
| import batch_ingest # noqa: E402 | |
| import extraction # noqa: E402 | |
| import pdf_crawler # noqa: E402 | |
| import pg_mirror # noqa: E402 | |
| from agent import agdb, config as C, orchestrator, scheduler # noqa: E402 | |
| 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) | |
| # --- a PDF whose page-1 quote cites a raster Figure 1 on page 2 -------------- | |
| QUOTE = "Average tensile strength (608 MPa) of the PPS composite, see Figure 1." | |
| def make_pdf() -> bytes: | |
| doc = fitz.open() | |
| pm = fitz.Pixmap(fitz.csRGB, fitz.IRect(0, 0, 600, 400), False) | |
| pm.clear_with(255) | |
| for i in range(400): | |
| pm.set_pixel(min(599, int(i * 1.4)), 399 - i, (200, 40, 40)) | |
| plot_png = pm.tobytes("png") | |
| p1 = doc.new_page() | |
| p1.insert_text((72, 72), "PPS carbon fibre composite datasheet", fontsize=11) | |
| p1.insert_text((72, 100), QUOTE, fontsize=9) | |
| p2 = doc.new_page() | |
| p2.insert_image(fitz.Rect(72, 80, 472, 346), stream=plot_png) | |
| p2.insert_text((72, 370), "Figure 1. Tensile strength of the PPS composite", fontsize=10) | |
| data = doc.tobytes() | |
| doc.close() | |
| return data | |
| PDF = make_pdf() | |
| SEEN_URL = "https://example.org/already_final.pdf" | |
| def fake_search_openalex(query, limit): | |
| # 1 relevant, 2 below the relevance gate | |
| yield pdf_crawler.Candidate( | |
| title="Tensile behavior of carbon fiber PPS thermoplastic composite laminates", | |
| pdf_url="https://example.org/pps_cf.pdf", source="openalex", query=query, | |
| doi="10.9999/aim.counters.0001", year="2026", | |
| abstract="thermoplastic composite tensile strength carbon fiber PPS") | |
| yield pdf_crawler.Candidate(title="Annual report of the bird society", | |
| pdf_url="https://example.org/birds.pdf", source="openalex", | |
| query=query, doi="10.9999/birds", abstract="birds nests eggs") | |
| yield pdf_crawler.Candidate(title="Municipal water pricing 2025", | |
| pdf_url="https://example.org/water.pdf", source="openalex", | |
| query=query, doi="10.9999/water", abstract="tariffs households") | |
| def fake_search_arxiv(query, limit): | |
| # relevant, but its download will fail (403 from the fake) | |
| yield pdf_crawler.Candidate( | |
| title="Fatigue of glass fiber PA66 thermoplastic composite", | |
| pdf_url="https://example.org/forbidden.pdf", source="arxiv", query=query, | |
| doi="10.9999/aim.counters.0002", year="2026", | |
| abstract="thermoplastic composite fatigue glass fiber polyamide") | |
| # DOI-less candidate whose only URL is already final in the crawler state | |
| yield pdf_crawler.Candidate( | |
| title="Thermoplastic composite tensile modulus datasheet PEEK carbon", | |
| pdf_url=SEEN_URL, source="arxiv", query=query, doi="", year="2026", | |
| abstract="thermoplastic composite tensile modulus carbon fiber PEEK") | |
| def fake_download(cand, pdf_dir, state): | |
| import hashlib | |
| if cand.pdf_url in state.seen_urls: | |
| return None | |
| if "forbidden" in cand.pdf_url: | |
| raise RuntimeError("HTTP 403 (scripted client refused)") | |
| state.seen_urls.add(cand.pdf_url) | |
| sha = hashlib.sha256(PDF).hexdigest() | |
| if sha in state.seen_hashes: | |
| return None | |
| state.seen_hashes.add(sha) | |
| fname = f"{cand.source}_{pdf_crawler.slugify(cand.title)}_{sha[:8]}.pdf" | |
| (pdf_dir / fname).write_bytes(PDF) | |
| 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(PDF)} | |
| def fake_extract(pdf_bytes, filename, api_key): | |
| mat = extraction.Material( | |
| material_name="Polyphenylene sulfide carbon fibre composite", | |
| material_abbreviation="PPS-CF", material_class="Composite", | |
| matrix="PPS", fiber="Carbon", fiber_volume_fraction="", | |
| properties=[extraction.Property( | |
| section="Mechanical", property_name="Tensile Strength", | |
| value_raw="608", value_num=608.0, unit="MPa", test_condition="", | |
| source_quote=QUOTE, page=1)]) | |
| return extraction.Extraction(materials=[mat], doc_status="ok", | |
| tokens_in=12345, tokens_out=678) | |
| 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 | |
| orchestrator.bootstrap() | |
| conn = agdb.connect() | |
| try: | |
| # figure mining would need vision calls; keep the linking stage only | |
| # one query, one round: 3 openalex + 2 arxiv candidates, deterministic counts | |
| agdb.set_config(conn, {"figures_enabled": False, "link_figures": True, | |
| "use_openalex": True, "use_arxiv": True, | |
| "queries_per_cycle": 1, "max_rounds_per_cycle": 1, | |
| "min_new_pdfs_per_cycle": 1}) | |
| st = pdf_crawler.CrawlerState(C.WORK_DIR / "state.json") | |
| agdb.load_crawler_state(conn, st) | |
| st.seen_urls.add(SEEN_URL) | |
| agdb.save_crawler_state(conn, st) | |
| finally: | |
| conn.close() | |
| m1 = orchestrator.run_cycle(trigger="test") | |
| print("cycle 1 metrics:", {k: v for k, v in m1.items() if k != "sources"}, m1.get("sources")) | |
| conn = agdb.connect() | |
| try: | |
| run = agdb.fetch_df(conn, "SELECT * FROM agent_runs ORDER BY id DESC LIMIT 1").iloc[0] | |
| check("cycle downloaded 1, ingested 1", run["downloaded"] == 1 and run["pdfs_ingested"] == 1) | |
| check("relevance_rejected = 2 (the two off-topic OpenAlex hits)", run["relevance_rejected"] == 2) | |
| check("download_failed = 1 (the 403)", run["download_failed"] == 1) | |
| check("url_seen_skipped = 1 (DOI-less candidate whose URL was already final)", | |
| run["url_seen_skipped"] == 1) | |
| check("tokens_in / tokens_out persisted from the extraction", | |
| run["tokens_in"] == 12345 and run["tokens_out"] == 678) | |
| check("rows_linked = 1", run["rows_linked"] == 1) | |
| src = run["sources"] if isinstance(run["sources"], dict) else {} | |
| check("sources jsonb has per-source candidate counts", | |
| src.get("openalex", {}).get("candidates") == 3 | |
| and src.get("openalex", {}).get("relevant") == 1 | |
| and src.get("arxiv", {}).get("candidates") == 2) | |
| check("heartbeat_at set and not before started_at", | |
| run["heartbeat_at"] is not None and run["heartbeat_at"] >= run["started_at"]) | |
| # `candidates` keeps its campaign meaning (relevant candidates); the raw | |
| # count is candidates + relevance_rejected = 5 here. | |
| check("finish_run still records the legacy counters (candidates = relevant = 3)", | |
| run["candidates"] == 3 and run["candidates"] + run["relevance_rejected"] == 5 | |
| and run["rows_inserted"] == 1) | |
| with conn.cursor() as cur: | |
| cur.execute('SELECT figure_id, figure_ref, figure_link_score, figure_link_signals, ' | |
| 'octet_length(image), status FROM "Composites_materials"') | |
| rows = cur.fetchall() | |
| cur.execute("SELECT count(*) FROM figures WHERE image_bytes IS NOT NULL") | |
| n_figs = cur.fetchone()[0] | |
| check("one text row inserted", len(rows) == 1) | |
| r = rows[0] if rows else (None,) * 6 | |
| check("row links to the cited figure (figure_id + 'Figure 1' ref)", | |
| bool(r[0]) and str(r[1]).lower().startswith("fig")) | |
| check("citation-tier score (>= 0.9) and citation signal", | |
| (r[2] or 0) >= 0.9 and "citation" in str(r[3])) | |
| check("PNG embedded in the row's image column", (r[4] or 0) > 500) | |
| check("linked figure stored in figures with its bytes", n_figs >= 1) | |
| check("status untouched by linking (ok)", r[5] == "ok") | |
| ev = agdb.fetch_df(conn, "SELECT message FROM agent_events WHERE run_id = %s", (int(run["id"]),)) | |
| msgs = "\n".join(ev["message"]) | |
| check("events mention links, tokens and the relevance gate", | |
| "rows linked" in msgs and "tokens 12345 in" in msgs and "below the relevance gate" in msgs) | |
| check("cost derived from tokens is positive", C.cost_usd(run["tokens_in"], run["tokens_out"]) > 0) | |
| finally: | |
| conn.close() | |
| # --- gate: figures not ready -> both figure stages disabled, warn event ----- | |
| real_problems = pg_mirror.figures_ready_problems | |
| pg_mirror.figures_ready_problems = lambda conn: ["figures table lacks the image_bytes column"] | |
| seen_link_opts = {} | |
| real_process = batch_ingest.process_pdf | |
| def spy_process(pdf_path, conn, api_key, db=None, figure_opts=None, link_opts=None): | |
| seen_link_opts["link_opts"] = link_opts | |
| seen_link_opts["figure_opts"] = figure_opts | |
| return real_process(pdf_path, conn, api_key, db=db, figure_opts=figure_opts, link_opts=link_opts) | |
| batch_ingest.process_pdf = spy_process | |
| conn = agdb.connect() | |
| try: | |
| # a fresh DOI so the cycle downloads again | |
| with conn.cursor() as cur: | |
| cur.execute("DELETE FROM agent_doi_seen") | |
| cur.execute("DELETE FROM sources") | |
| conn.commit() | |
| st = pdf_crawler.CrawlerState(C.WORK_DIR / "state.json") | |
| agdb.load_crawler_state(conn, st) | |
| st.seen_urls.discard("https://example.org/pps_cf.pdf") | |
| st.seen_hashes.clear() | |
| agdb.save_crawler_state(conn, st) | |
| finally: | |
| conn.close() | |
| m2 = orchestrator.run_cycle(trigger="test") | |
| pg_mirror.figures_ready_problems = real_problems | |
| batch_ingest.process_pdf = real_process | |
| conn = agdb.connect() | |
| try: | |
| run2 = agdb.fetch_df(conn, "SELECT id, downloaded, pdfs_ingested, rows_linked FROM agent_runs " | |
| "ORDER BY id DESC LIMIT 1").iloc[0] | |
| ev = agdb.fetch_df(conn, "SELECT level, message FROM agent_events WHERE run_id = %s", | |
| (int(run2["id"]),)) | |
| gate = ev[ev["message"].str.contains("figure stages")] | |
| check("gate: cycle still ingested the PDF", run2["downloaded"] == 1 and run2["pdfs_ingested"] == 1) | |
| check("gate: link_opts and figure_opts were withheld from process_pdf", | |
| seen_link_opts.get("link_opts") is None and seen_link_opts.get("figure_opts") is None) | |
| check("gate: warn event names the problem and the remedy", | |
| len(gate) == 1 and gate.iloc[0]["level"] == "warn" | |
| and "pg_migrate.py --apply" in gate.iloc[0]["message"]) | |
| check("gate: rows_linked = 0 for that cycle", run2["rows_linked"] == 0) | |
| finally: | |
| conn.close() | |
| # --- stall detector ----------------------------------------------------------- | |
| real_run_cycle = orchestrator.run_cycle | |
| release = threading.Event() | |
| late = {} | |
| def silent_cycle(trigger="schedule"): | |
| """Starts a run, reports once, then goes quiet (a hung network call).""" | |
| conn = agdb.connect() | |
| try: | |
| token = agdb.acquire_cycle_lock(conn) | |
| run_id = agdb.start_run(conn, trigger) | |
| agdb.log_event(conn, "cycle start", node="plan", run_id=run_id) | |
| finally: | |
| conn.close() | |
| orchestrator.CURRENT[threading.get_ident()] = {"run_id": run_id, "token": token, | |
| "trigger": trigger, "started": time.time()} | |
| late["run_id"] = run_id | |
| release.wait(timeout=120) | |
| orchestrator.CURRENT.pop(threading.get_ident(), None) | |
| return {} | |
| def chatty_cycle(trigger="schedule"): | |
| """Slow but alive: an event every 0.5 s for 4 s, then finishes normally.""" | |
| conn = agdb.connect() | |
| try: | |
| token = agdb.acquire_cycle_lock(conn) | |
| run_id = agdb.start_run(conn, trigger) | |
| orchestrator.CURRENT[threading.get_ident()] = {"run_id": run_id, "token": token, | |
| "trigger": trigger, "started": time.time()} | |
| for i in range(8): | |
| agdb.log_event(conn, f"ingesting pdf {i}", node="ingest", run_id=run_id) | |
| time.sleep(0.5) | |
| agdb.finish_run(conn, run_id, "ok", {"rows_inserted": 8}, {}) | |
| agdb.release_cycle_lock(conn, token) | |
| finally: | |
| conn.close() | |
| orchestrator.CURRENT.pop(threading.get_ident(), None) | |
| return {"rows_inserted": 8} | |
| try: | |
| orchestrator.run_cycle = silent_cycle | |
| t0 = time.time() | |
| out = scheduler.run_cycle_guarded(trigger="schedule", timeout_minutes=5, | |
| stall_minutes=2 / 60, poll_seconds=0.5) | |
| elapsed = time.time() - t0 | |
| release.set() | |
| check("stall: abandoned on silence, far inside the 5-min hard budget", | |
| out.get("timeout") is not None and out.get("stalled") is not None and elapsed < 30) | |
| conn = agdb.connect() | |
| try: | |
| row = agdb.fetch_df(conn, "SELECT status, report FROM agent_runs WHERE id = %s", | |
| (late["run_id"],)).iloc[0] | |
| check("stall: run marked 'timeout' with the stall reason", | |
| row["status"] == "timeout" and "no progress event" in str(row["report"])) | |
| check("stall: lock released", agdb.get_state(conn, "cycle_lock") is None) | |
| finally: | |
| conn.close() | |
| orchestrator.run_cycle = chatty_cycle | |
| out = scheduler.run_cycle_guarded(trigger="schedule", timeout_minutes=5, | |
| stall_minutes=2 / 60, poll_seconds=0.25) | |
| check("heartbeat: a slow but reporting cycle (4 s > 2 s stall limit) is NOT abandoned", | |
| out.get("rows_inserted") == 8 and "timeout" not in out) | |
| finally: | |
| orchestrator.run_cycle = real_run_cycle | |
| # --- stale sweep measures from the heartbeat ---------------------------------- | |
| conn = agdb.connect() | |
| try: | |
| with conn.cursor() as cur: | |
| cur.execute("INSERT INTO agent_runs (trigger, status, started_at, heartbeat_at) " | |
| "VALUES ('test', 'running', now() - interval '3 hours', now()) RETURNING id") | |
| alive_id = cur.fetchone()[0] | |
| cur.execute("INSERT INTO agent_runs (trigger, status, started_at, heartbeat_at) " | |
| "VALUES ('test', 'running', now() - interval '3 hours', " | |
| "now() - interval '3 hours') RETURNING id") | |
| dead_id = cur.fetchone()[0] | |
| conn.commit() | |
| n = agdb.sweep_stale_runs(conn) | |
| st = agdb.fetch_df(conn, "SELECT id, status FROM agent_runs WHERE id = ANY(%s)", | |
| ([alive_id, dead_id],)).set_index("id")["status"] | |
| check("sweep: old start but fresh heartbeat stays 'running'", st[alive_id] == "running") | |
| check("sweep: old heartbeat -> 'stale'", st[dead_id] == "stale" and n == 1) | |
| with conn.cursor() as cur: | |
| cur.execute("UPDATE agent_runs SET status = 'ok', finished_at = now() WHERE id = %s", (alive_id,)) | |
| conn.commit() | |
| finally: | |
| conn.close() | |
| print("\n%d checks failed" % len(fails)) | |
| sys.exit(1 if fails else 0) | |