Spaces:
Running
Running
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 Download test_candidate_queue_e2e.py from aim4composites/AutonomousAgent: direct link, hf CLI and curl.
- Browser
- Download file 16.7 kB
-
https://huggingface.co/spaces/aim4composites/AutonomousAgent/resolve/main/test_candidate_queue_e2e.py
- Command line
-
hf download hf://spaces/aim4composites/AutonomousAgent/test_candidate_queue_e2e.py
-
curl -L -o test_candidate_queue_e2e.py https://huggingface.co/spaces/aim4composites/AutonomousAgent/resolve/main/test_candidate_queue_e2e.py
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) | |