Spaces:
Running
Running
Mathias Heider
Claude Fable 5.1
Merge the 17 Sep "more papers" lanes onto the hotfix line; parallel ingestion; throughput knobs
6bb7525 unverified Download test_keep_running.py from aim4composites/AutonomousAgent: direct link, hf CLI and curl.
- Browser
- Download file 8.27 kB
-
https://huggingface.co/spaces/aim4composites/AutonomousAgent/resolve/main/test_keep_running.py
- Command line
-
hf download hf://spaces/aim4composites/AutonomousAgent/test_keep_running.py
-
curl -L -o test_keep_running.py https://huggingface.co/spaces/aim4composites/AutonomousAgent/resolve/main/test_keep_running.py
8.27 kB
| """Tests for the keep-running changes: Retry-After cap, cycle watchdog, | |
| token-scoped lock, once-per-process boot, keep-alive. | |
| Needs the same scratch Postgres as test_e2e.py (DB_* env vars). No network: | |
| the Space's HTTP calls are never made here. | |
| """ | |
| import os | |
| import sys | |
| import threading | |
| import time | |
| os.environ.setdefault("AGENT_WORK_DIR", "/tmp/agent_keep_test") | |
| os.environ["GEMINI_API_KEY"] = "test-key-not-used" | |
| os.environ["AGENT_USE_NTRS"] = "0" # tests stub the lanes; never touch the network | |
| # No network in tests: the live sources added Sep 2026 (Semantic Scholar | |
| # bulk, OpenAlex topic walk, datasheet seed crawl) are off, and the query | |
| # grid is not seeded so the frontier is the 8 builtin intents. | |
| 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.pop("SPACE_HOST", None) # keepalive must stay off in tests | |
| os.environ["AGENT_KEEPALIVE_URL"] = "" | |
| import requests | |
| import pdf_crawler | |
| from agent import agdb, autostart, config as C, keepalive, orchestrator, scheduler | |
| FAILS = 0 | |
| def check(cond, label): | |
| global FAILS | |
| print(("PASS " if cond else "FAIL ") + label) | |
| if not cond: | |
| FAILS += 1 | |
| # --- 1. Retry-After cap ------------------------------------------------------- | |
| class _Resp: | |
| def __init__(self, headers): | |
| self.headers = headers | |
| check(pdf_crawler._retry_after(_Resp({"Retry-After": "86400"})) == 300.0, | |
| "crawler caps Retry-After: 86400 at 300 s") | |
| check(pdf_crawler._retry_after(_Resp({"Retry-After": "7"})) == 7.0, | |
| "crawler keeps a small Retry-After as is") | |
| check(pdf_crawler._retry_after(_Resp({"Retry-After": "Wed, 21 Oct 2026 07:28:00 GMT"})) is None, | |
| "crawler ignores HTTP-date Retry-After (falls back to backoff)") | |
| check(pdf_crawler._retry_after(_Resp({})) is None, "no header -> None") | |
| # --- 2. token-scoped lock ----------------------------------------------------- | |
| conn = agdb.connect() | |
| agdb.ensure_agent_schema(conn) | |
| agdb.release_cycle_lock(conn) | |
| tok1 = agdb.acquire_cycle_lock(conn) | |
| check(bool(tok1), "lock acquired with a token") | |
| check(agdb.acquire_cycle_lock(conn) is None, "second acquire refused while fresh") | |
| check(agdb.release_cycle_lock(conn, "not-the-token") is False, "wrong token cannot release") | |
| check(agdb.get_state(conn, "cycle_lock") is not None, "lock still held after wrong-token release") | |
| check(agdb.release_cycle_lock(conn, tok1) is True, "matching token releases") | |
| check(agdb.get_state(conn, "cycle_lock") is None, "lock cleared") | |
| tok2 = agdb.acquire_cycle_lock(conn) | |
| check(agdb.release_cycle_lock(conn) is True, "unconditional release (watchdog/boot) works") | |
| conn.close() | |
| # --- 3. watchdog abandons a stuck cycle ---------------------------------------- | |
| release_stuck = threading.Event() | |
| late = {} | |
| def stuck_cycle(trigger="schedule"): | |
| """Stands in for orchestrator.run_cycle: takes the lock, opens a run, | |
| then hangs (like a Retry-After sleep) until the test releases it.""" | |
| conn = agdb.connect() | |
| try: | |
| token = agdb.acquire_cycle_lock(conn) | |
| run_id = agdb.start_run(conn, trigger) | |
| 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_stuck.wait(timeout=120) | |
| # the abandoned thread eventually wakes up and tries to close normally | |
| conn = agdb.connect() | |
| try: | |
| agdb.finish_run(conn, run_id, "ok", {"rows_inserted": 5}, {}) | |
| late["released_newer_lock"] = agdb.release_cycle_lock(conn, token) | |
| finally: | |
| conn.close() | |
| orchestrator.CURRENT.pop(threading.get_ident(), None) | |
| return {"rows_inserted": 5} | |
| real_run_cycle = orchestrator.run_cycle | |
| orchestrator.run_cycle = stuck_cycle | |
| t0 = time.time() | |
| out = scheduler.run_cycle_guarded(trigger="schedule", timeout_minutes=2 / 60) # 2 s budget | |
| elapsed = time.time() - t0 | |
| check(out.get("timeout") is not None, "guarded cycle returned a timeout verdict") | |
| check(elapsed < 20, f"watchdog returned promptly ({elapsed:.1f}s)") | |
| conn = agdb.connect() | |
| row = agdb.fetch_df(conn, "SELECT status FROM agent_runs WHERE id = %s", (late["run_id"],)) | |
| check(row.iloc[0]["status"] == "timeout", "run marked 'timeout'") | |
| check(agdb.get_state(conn, "cycle_lock") is None, "watchdog released the lock") | |
| ev = agdb.fetch_df(conn, "SELECT message FROM agent_events WHERE run_id = %s AND level = 'error'", | |
| (late["run_id"],)) | |
| check(len(ev) >= 1 and "abandoned by the watchdog" in ev.iloc[0]["message"], "timeout event logged") | |
| conn.close() | |
| # a new cycle can start while the stuck one is still alive | |
| conn = agdb.connect() | |
| tok_new = agdb.acquire_cycle_lock(conn) | |
| check(bool(tok_new), "next cycle acquires the lock while the stuck thread lives") | |
| conn.close() | |
| # now the stuck thread wakes up: it must not overwrite the verdict or the new lock | |
| release_stuck.set() | |
| for _ in range(50): | |
| if "released_newer_lock" in late: | |
| break | |
| time.sleep(0.1) | |
| conn = agdb.connect() | |
| row = agdb.fetch_df(conn, "SELECT status, rows_inserted FROM agent_runs WHERE id = %s", (late["run_id"],)) | |
| check(row.iloc[0]["status"] == "timeout", "late finish_run did not overwrite 'timeout'") | |
| check(late.get("released_newer_lock") is False, "late token release refused (newer lock held)") | |
| check(agdb.get_state(conn, "cycle_lock") is not None, "newer cycle's lock intact") | |
| agdb.release_cycle_lock(conn) | |
| conn.close() | |
| # a normal fast cycle passes through the guard unchanged | |
| orchestrator.run_cycle = lambda trigger="schedule": {"rows_inserted": 1, "trigger": trigger} | |
| out = scheduler.run_cycle_guarded(trigger="manual", timeout_minutes=1) | |
| check(out == {"rows_inserted": 1, "trigger": "manual"}, "fast cycle returns its metrics through the guard") | |
| orchestrator.run_cycle = real_run_cycle | |
| # --- 4. once-per-process boot -------------------------------------------------- | |
| calls = {"bootstrap": 0, "sync": 0} | |
| real_bootstrap, real_sync = orchestrator.bootstrap, scheduler.sync_from_config | |
| orchestrator.bootstrap = lambda *a, **k: calls.__setitem__("bootstrap", calls["bootstrap"] + 1) | |
| scheduler.sync_from_config = lambda: (calls.__setitem__("sync", calls["sync"] + 1) or {}) | |
| autostart._result = None | |
| r1 = autostart.boot_once() | |
| r2 = autostart.boot_once() | |
| check(r1["db"] and r1["scheduler"], "boot_once reports db + scheduler up") | |
| check(calls == {"bootstrap": 1, "sync": 1}, "second boot_once is a no-op (no double force-sweep)") | |
| check(r1 is r2, "boot_once returns the cached result") | |
| check(r1["keepalive"] is False, "keepalive stays off without SPACE_HOST") | |
| th = autostart.start_in_background() | |
| th.join(5) | |
| check(calls["bootstrap"] == 1, "background start reuses the same boot") | |
| orchestrator.bootstrap, scheduler.sync_from_config = real_bootstrap, real_sync | |
| # --- 5. keep-alive ---------------------------------------------------------------- | |
| class _FakeResp: | |
| status_code = 200 | |
| seen = [] | |
| code = keepalive.ping("https://example.invalid/_stcore/health", | |
| get=lambda url, timeout, headers: (seen.append((url, timeout)) or _FakeResp())) | |
| check(code == 200 and seen and seen[0][0].endswith("/_stcore/health"), "ping GETs the URL and returns the status") | |
| def _boom(url, timeout, headers): | |
| raise requests.ConnectionError("down") | |
| check(keepalive.ping("https://example.invalid/", get=_boom) is None, "ping swallows transport errors") | |
| check(keepalive.start(url="", minutes=10) is False, "no URL -> keepalive not started") | |
| check(keepalive.start(url="https://example.invalid/x", minutes=0) is False, "0 minutes -> not started") | |
| # derived URL from SPACE_HOST (re-import config with the env var set) | |
| import importlib | |
| os.environ["SPACE_HOST"] = "aim4composites-autonomousagent.hf.space" | |
| os.environ.pop("AGENT_KEEPALIVE_URL", None) | |
| C2 = importlib.reload(C) | |
| check(C2.KEEPALIVE_URL == "https://aim4composites-autonomousagent.hf.space/_stcore/health", | |
| "keepalive URL derived from SPACE_HOST") | |
| check(C2.KEEPALIVE_MINUTES == 10.0 and C2.CYCLE_TIMEOUT_MINUTES == 60.0, "defaults: 10-min ping, 60-min cycle budget") | |
| os.environ.pop("SPACE_HOST", None) | |
| importlib.reload(C) | |
| print() | |
| print(f"{FAILS} checks failed") | |
| sys.exit(1 if FAILS else 0) | |