AutonomousAgent / test_candidate_queue_e2e.py
Mathias Heider
Claude Fable 5.1
boot: clear a discovery lock left by a killed container (a secret replaced mid-pass blocked every pass for 90 min)
61e0cf5 unverified
Raw History Blame Contribute Delete
16.7 kB
"""E2E tests (scratch Postgres, no network) for the candidate queue (28 Sep 2026):
* a discovery pass searches the frontier on every lane (mocked), walks the
topic pages, crawls due datasheet seeds, and queues the relevant
candidates once (dedupe across lanes; registry-seen DOIs land as 'seen');
it is an agent_runs row of kind 'discovery' and it marks its queries used;
* a cycle in queue mode pops the best candidates, downloads and ingests
them without a single search call, settles every popped row
(ingested / seen / queued-for-retry / failed) and credits the query that
found each paper without touching the frontier's LRU clock;
* a cycle that finds the queue empty searches inline and asks for a pass;
* boot repair returns rows a dead cycle left 'downloading' to the queue;
* the scheduler runs a requested pass right after a cycle.
python test_candidate_queue_e2e.py (same DB_* env as test_e2e.py)
"""
import dataclasses
import hashlib
import os
import sys
os.environ.setdefault("AGENT_WORK_DIR", "/tmp/agent_queue_test")
os.environ["GEMINI_API_KEY"] = "test-key-not-used"
os.environ["AGENT_KEEPALIVE_URL"] = ""
for _k in ("AGENT_USE_NTRS", "AGENT_USE_S2", "AGENT_USE_OPENALEX_TOPICS",
"AGENT_USE_DATASHEETS", "AGENT_SEED_QUERY_GRID", "AGENT_USE_UDSPACE", "AGENT_USE_CORE"):
os.environ[_k] = "0" # every lane is stubbed below; nothing touches the network
import fitz
import pdf_crawler
import extraction
import batch_ingest
from agent import agdb, config as C, discovery, orchestrator, scheduler
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)
QUOTE = "The tensile modulus of {abbr} was measured as 24.5 GPa at 23 C."
def make_pdf(abbr: str) -> bytes:
doc = fitz.open()
page = doc.new_page()
page.insert_text((72, 100), f"{abbr} composite datasheet")
page.insert_text((72, 130), QUOTE.format(abbr=abbr))
data = doc.tobytes()
doc.close()
return data
ABS = "thermoplastic composite tensile modulus carbon fiber PEEK laminate"
SEARCH_CALLS: list = []
def cand(key: str, source: str, query: str, doi: str = "", score_text: str = ABS,
title: str = ""):
return pdf_crawler.Candidate(title=title or f"{key} thermoplastic composite study",
pdf_url=f"https://{source}.example/{key}.pdf",
source=source, query=query, doi=doi, year="2026", abstract=score_text)
def fake_s2(query, limit):
SEARCH_CALLS.append(("semantic_scholar", query))
q = query.replace(" ", "-")
for i in range(3):
yield cand(f"s2-{q}-{i}", "semantic_scholar", query, doi=f"10.1/{q}.{i}")
def fake_ntrs(query, limit):
SEARCH_CALLS.append(("ntrs", query))
q = query.replace(" ", "-")
yield cand(f"ntrs-{q}", "ntrs", query, doi=f"ntrs:{q}")
# the same work S2 also returns (same DOI): must collapse to one row
yield cand(f"s2-{q}-0", "ntrs", query, doi=f"10.1/{q}.0")
def fake_arxiv(query, limit):
SEARCH_CALLS.append(("arxiv", query))
return iter(())
def fake_openalex(query, limit):
SEARCH_CALLS.append(("openalex", query))
q = query.replace(" ", "-")
yield cand(f"oa-{q}", "openalex", query, doi=f"10.2/{q}")
yield cand(f"oa-junk-{q}", "openalex", query, score_text="poetry anthology",
title="Selected poems of the eighteenth century") # below the gate
TOPIC_PAGES = {"*": (["t1", "t2"], "p2"), "p2": (["t3"], None)}
def fake_topic(tid, per_page, cursor, min_score, query=""):
SEARCH_CALLS.append(("topic", f"{tid}@{cursor}"))
items, nxt = TOPIC_PAGES.get(cursor, ([], None))
return [cand(f"{tid}-{k}", "openalex", query, doi=f"10.3/{tid}.{k}") for k in items], nxt
def fake_seed(seed_url):
SEARCH_CALLS.append(("datasheet", seed_url))
yield pdf_crawler.Candidate(title="Victrex PEEK 450CA30 datasheet", pdf_url=f"{seed_url}/450CA30.pdf",
source="datasheet", query="", abstract="")
DOWNLOAD_FAIL: set = set()
PDF_OF: dict = {}
def fake_download(cand, pdf_dir, state):
if cand.pdf_url in DOWNLOAD_FAIL:
return None # transient failure, nothing marked
if cand.pdf_url in state.seen_urls:
return None
state.seen_urls.add(cand.pdf_url)
key = cand.pdf_url.rsplit("/", 1)[1][:-4]
data = PDF_OF.setdefault(key, make_pdf(key.upper()[:12]))
sha = hashlib.sha256(data).hexdigest()
if sha in state.seen_hashes:
return None
state.seen_hashes.add(sha)
fname = f"{cand.source}_{key}_{sha[:8]}.pdf"
(pdf_dir / fname).write_bytes(data)
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(data)}
def fake_extract(pdf_bytes, filename, api_key):
key = filename.split("_", 1)[1].rsplit("_", 1)[0]
abbr = key.upper()[:12]
mat = extraction.Material(
material_name=f"Material {abbr}", material_abbreviation=abbr, material_class="Composite",
matrix="PEEK", fiber="Carbon", fiber_volume_fraction="30%",
properties=[extraction.Property(section="Mechanical", property_name="Tensile Modulus",
value_raw="24.5", value_num=24.5, unit="GPa", test_condition="23 C",
source_quote=QUOTE.format(abbr=abbr), page=1)])
return extraction.Extraction(materials=[mat], doc_status="ok")
pdf_crawler.search_semantic_scholar = fake_s2
pdf_crawler.search_ntrs = fake_ntrs
pdf_crawler.search_arxiv = fake_arxiv
pdf_crawler.search_openalex = fake_openalex
pdf_crawler.list_openalex_topic = fake_topic
pdf_crawler.crawl_datasheet_seed = fake_seed
pdf_crawler.download_pdf = fake_download
batch_ingest.extract_from_pdf = fake_extract
def q(sql, params=()):
conn = agdb.connect()
try:
with conn.cursor() as cur:
cur.execute(sql, params) if params else cur.execute(sql)
try:
return cur.fetchall()
except Exception:
return None
finally:
conn.commit()
conn.close()
# --- 0. boot: empty queue β†’ a pass is requested --------------------------------------
orchestrator.bootstrap()
conn = agdb.connect()
req = discovery.requested(conn)
agdb.set_config(conn, {
"queries_per_cycle": 2, "max_rounds_per_cycle": 1, "max_new_pdfs_per_cycle": 4,
"min_new_pdfs_per_cycle": 4, "figures_enabled": False, "link_figures": False,
"expand_queries": False, "ingest_workers": 2,
"use_candidate_queue": True, "queue_low_water": 5,
"discovery_queries": 3, "discovery_openalex_queries": 1, "discovery_topic_pages": 2,
"discovery_per_query": 10, "use_semantic_scholar": True, "use_ntrs": True, "use_arxiv": True,
"use_openalex": True, "use_openalex_topics": True, "use_datasheets": True,
"use_udspace": False, "use_core": False,
"openalex_topic_names": "T11664 Fiber-reinforced polymer composites",
"datasheet_seed_urls": "https://victrex.example/datasheets",
})
conn.close()
check("boot with an empty queue requests a discovery pass", bool(req) and "boot" in req.get("reason", ""))
# a DOI already in the registry: its candidate must land as 'seen'
q("INSERT INTO agent_doi_seen (doi, filename, url, ingest_status) VALUES "
"('10.3/T11664.t2', 'old.pdf', 'https://old/x.pdf', 'ingested')")
# --- 1. discovery pass -----------------------------------------------------------------
SEARCH_CALLS[:] = []
m = discovery.run_discovery(trigger="test")
run = q("SELECT kind, status, queries_used, candidates, relevance_rejected, skipped_doi, queued, duplicates, sources "
"FROM agent_runs WHERE kind = 'discovery' ORDER BY id DESC LIMIT 1")[0]
check("pass recorded as a 'discovery' run that ended ok", run[0] == "discovery" and run[1] == "ok")
check("pass searched 3 frontier queries", run[2] == 3)
lanes_hit = {l for l, _ in SEARCH_CALLS}
check("every enabled lane was asked (S2, NTRS, arXiv, OpenAlex, topic walk, datasheet seed)",
lanes_hit == {"semantic_scholar", "ntrs", "arxiv", "openalex", "topic", "datasheet"})
n_oa = sum(1 for l, _ in SEARCH_CALLS if l == "openalex")
check("OpenAlex phrase search budgeted to discovery_openalex_queries=1", n_oa == 1)
n_topic = sum(1 for l, _ in SEARCH_CALLS if l == "topic")
check("topic walk read 2 pages (then the topic was exhausted)", n_topic == 2)
check("the topic cursor is persisted as done",
q("SELECT value FROM agent_state WHERE key = 'openalex_topic_cursor:T11664'")[0][0].strip('"') == "done")
# raw relevant: per query S2 3 + NTRS 2 (one a dup of S2) = 5, Γ—3 queries = 15;
# openalex 1 (1st query only) β†’ 16; topic 3 β†’ 19; datasheet 1 β†’ 20
found = run[3]
check(f"relevant candidates found = 20 raw incl. the cross-lane duplicates ({found})", found == 20)
check("the below-gate OpenAlex hit was rejected", run[4] == 1)
by_status = dict(q("SELECT status, count(*) FROM agent_candidates GROUP BY status"))
check(f"16 queued + 1 seen (registry DOI), duplicate collapsed ({by_status})",
by_status.get("queued") == 16 and by_status.get("seen") == 1)
check("run row carries queued=16 and duplicates=0 (dup collapsed before insert)",
run[6] == 16 and run[7] == 0)
seen_row = q("SELECT status, source FROM agent_candidates WHERE doi = '10.3/T11664.t2'")
check("registry-seen DOI stored as 'seen'", seen_row == [("seen", "openalex")])
used = q("SELECT count(*) FROM agent_queries WHERE times_used = 1")[0][0]
check("the pass marked its 3 queries used (frontier rotates on passes)", used == 3)
srcs = run[8] or {}
check("per-lane counters persisted (semantic_scholar, ntrs, openalex, openalex_topic, datasheet)",
{"semantic_scholar", "ntrs", "openalex", "openalex_topic", "datasheet"} <= set(srcs))
check("the request flag was cleared by the pass",
discovery.requested(agdb.connect()) is None)
ds = q("SELECT source, status FROM agent_candidates WHERE source = 'datasheet'")
check("datasheet link queued with no DOI", ds == [("datasheet", "queued")])
ev = q("SELECT count(*) FROM agent_events WHERE run_id = (SELECT max(id) FROM agent_runs WHERE kind='discovery') "
"AND message LIKE 'discovery pass done:%'")[0][0]
check("summary event logged", ev == 1)
# --- 2. a cycle drains the queue ----------------------------------------------------------
SEARCH_CALLS[:] = []
top = q("SELECT pdf_url FROM agent_candidates WHERE status = 'queued' ORDER BY score DESC, discovered_at LIMIT 1")[0][0]
q("UPDATE agent_candidates SET score = 99 WHERE pdf_url = %s", (top,)) # always popped first
DOWNLOAD_FAIL.add(top) # the best candidate fails transiently this time
m2 = orchestrator.run_cycle(trigger="test")
row = q("SELECT status, downloaded, pdfs_ingested, rows_inserted, candidates, queries_used, kind "
"FROM agent_runs ORDER BY id DESC LIMIT 1")[0]
check("queue-mode cycle: ok, 4 downloaded, 4 ingested, 4 rows, no frontier queries",
row[:4] == ("ok", 4, 4, 4) and row[5] == 0 and row[6] == "cycle")
check("no search call was made (the queue is the plan)", SEARCH_CALLS == [])
by_status = dict(q("SELECT status, count(*) FROM agent_candidates GROUP BY status"))
check(f"4 ingested, the failed one back in the queue ({by_status})",
by_status.get("ingested") == 4 and by_status.get("queued") == 12)
fl = q("SELECT status, attempts, last_error FROM agent_candidates WHERE pdf_url = %s", (top,))[0]
check("transient failure: status queued, attempts 1, error noted",
fl[0] == "queued" and fl[1] == 1 and "download failed" in (fl[2] or ""))
credited = q("SELECT query, pdfs_found, rows_yielded, times_used FROM agent_queries WHERE pdfs_found > 0")
check("queries credited by text with pdfs_found/rows_yielded; LRU clock untouched (times_used stays 1)",
len(credited) >= 1 and all(t == 1 for _, _, _, t in credited) and sum(p for _, p, _, _ in credited) == 4)
ev = q("SELECT count(*) FROM agent_events WHERE message LIKE 'cycle start%queued candidates%'")[0][0]
check("cycle start event names the queue depth", ev == 1)
check("no discovery requested while the queue is above low water",
discovery.requested(agdb.connect()) is None)
DOWNLOAD_FAIL.clear()
# --- 3. attempts cap and settle semantics ------------------------------------------------
q("UPDATE agent_candidates SET attempts = %s WHERE pdf_url = %s", (agdb.CANDIDATE_MAX_ATTEMPTS - 1, top))
DOWNLOAD_FAIL.add(top)
conn = agdb.connect(); agdb.set_config(conn, {"max_new_pdfs_per_cycle": 1, "min_new_pdfs_per_cycle": 1}); conn.close()
orchestrator.run_cycle(trigger="test")
fl = q("SELECT status, attempts FROM agent_candidates WHERE pdf_url = %s", (top,))[0]
check("out of attempts β†’ failed", fl == ("failed", agdb.CANDIDATE_MAX_ATTEMPTS))
DOWNLOAD_FAIL.clear()
# --- 4. drain to empty β†’ inline fallback + a pass is requested -----------------------------
conn = agdb.connect(); agdb.set_config(conn, {"max_new_pdfs_per_cycle": 25, "min_new_pdfs_per_cycle": 25}); conn.close()
SEARCH_CALLS[:] = []
m4 = orchestrator.run_cycle(trigger="test")
by_status = dict(q("SELECT status, count(*) FROM agent_candidates GROUP BY status"))
check(f"queue drained ({by_status})", by_status.get("queued", 0) == 0)
check("the short cycle fell back to inline search once the queue was empty",
any(l == "semantic_scholar" for l, _ in SEARCH_CALLS))
req = discovery.requested(agdb.connect())
check("low-water mark β†’ discovery requested", bool(req) and "low-water" in req.get("reason", ""))
row = q("SELECT status, downloaded FROM agent_runs ORDER BY id DESC LIMIT 1")[0]
check("that cycle still ended ok", row[0] == "ok")
# --- 5. boot repair for rows stuck in 'downloading' --------------------------------------
q("INSERT INTO agent_candidates (doi, pdf_url, title, source, status, attempts, last_attempt_at) VALUES "
"('10.9/stuck1', 'https://x/stuck1.pdf', 'stuck one', 'ntrs', 'downloading', 1, now() - interval '3 hours'), "
"('10.9/stuck2', 'https://x/stuck2.pdf', 'stuck two', 'ntrs', 'downloading', %s, now() - interval '3 hours')",
(agdb.CANDIDATE_MAX_ATTEMPTS,))
orchestrator.bootstrap()
st = dict(q("SELECT pdf_url, status FROM agent_candidates WHERE pdf_url LIKE 'https://x/stuck%'"))
check("boot repair: stuck row with attempts left β†’ queued, exhausted one β†’ failed",
st == {"https://x/stuck1.pdf": "queued", "https://x/stuck2.pdf": "failed"})
# --- 5b. boot clears a discovery lock left by a killed container --------------------------
conn = agdb.connect()
tok = agdb.acquire_cycle_lock(conn, key=discovery.LOCK_KEY)
check("precondition: discovery lock taken (a pass in flight)", bool(tok))
check("precondition: a second pass is refused while the lock is fresh",
discovery.run_discovery(trigger="test") == {"skipped": "already_running"})
conn.close()
orchestrator.bootstrap() # a fresh container cannot have a live pass
conn = agdb.connect()
check("boot: leftover discovery lock cleared", agdb.get_state(conn, "discovery_lock") is None)
conn.close()
m_after_boot = discovery.run_discovery(trigger="test")
check("boot: a pass can run again right away", "skipped" not in m_after_boot)
# --- 6. scheduler: the cycle job runs a requested pass afterwards --------------------------
conn = agdb.connect(); discovery.request(conn, "test"); conn.close()
n_disc_before = q("SELECT count(*) FROM agent_runs WHERE kind = 'discovery'")[0][0]
n_cycles_before = q("SELECT count(*) FROM agent_runs WHERE kind = 'cycle'")[0][0]
scheduler._cycle_job()
n_disc = q("SELECT count(*) FROM agent_runs WHERE kind = 'discovery'")[0][0]
n_cycles = q("SELECT count(*) FROM agent_runs WHERE kind = 'cycle'")[0][0]
check("cycle job: one cycle, then the requested discovery pass", n_cycles == n_cycles_before + 1 and n_disc == n_disc_before + 1)
last = q("SELECT trigger, status FROM agent_runs WHERE kind = 'discovery' ORDER BY id DESC LIMIT 1")[0]
check("the pass ran with the low-water trigger and ended ok", last[0].startswith("low-water") and last[1] == "ok")
conn = agdb.connect(); agdb.set_config(conn, {"autonomy": True, "interval_hours": 0.5}); conn.close()
cfg = scheduler.sync_from_config()
check("autonomy on β†’ discovery job scheduled", scheduler.next_discovery_time() is not None)
conn = agdb.connect(); agdb.set_config(conn, {"autonomy": False}); conn.close()
scheduler.sync_from_config()
check("autonomy off β†’ discovery job removed", scheduler.next_discovery_time() is None)
print()
print("all tests passed" if not fails else f"{len(fails)} FAILED: {fails}")
sys.exit(1 if fails else 0)