AutonomousAgent / test_keep_running.py
Mathias Heider
Claude Fable 5.1
Merge the 17 Sep "more papers" lanes onto the hotfix line; parallel ingestion; throughput knobs
6bb7525 unverified
Raw History Blame Contribute Delete
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)