"""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)