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)