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