Spaces:
Running
Running
File size: 15,784 Bytes
830111f 634726e 6bb7525 830111f | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 | """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)
|