AutonomousAgent / test_counters_e2e.py
Mathias Heider
Claude Fable 5.1
Merge the 17 Sep "more papers" lanes onto the hotfix line; parallel ingestion; throughput knobs
6bb7525 unverified
Raw History Blame Contribute Delete
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)