Infrastructure outage handling (MCP/A4H down: wait, clean up, rerun), dashboard fix, budget limit 129
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -361,7 +361,9 @@ def render(d):
|
||||
next50 = (n_acc // 50 + 1) * 50
|
||||
est = d["panel_est"]
|
||||
pcls = "good" if est < 50 else ("warn" if est < 57 else "bad")
|
||||
eta = ("%.1f h (%s)" % (d["eta_h"], time.strftime("%a %H:%M", time.localtime(d["now"] + d["eta_h"] * 3600)))) if d["eta_h"] else "n/a"
|
||||
eta = ("%.1f h (%s)" % (d["eta_h"], time.strftime("%a %H:%M", time.localtime(d["now"] + d["eta_h"] * 3600)))) if d["eta_h"] and alive else "n/a"
|
||||
if not alive: # a stopped pipeline has no burn rate and no speed
|
||||
d["burn_ledger_h"] = d["rate_runs_h"] = d["rate_acc_h"] = None
|
||||
a30 = d["acc30"]
|
||||
a30s = "%.0f %%" % (100 * a30) if a30 is not None else "n/a (<30 runs)"
|
||||
cost_per = d["run_cost"] / n_acc if n_acc else 0
|
||||
|
||||
62
harness/infra.py
Normal file
62
harness/infra.py
Normal file
@@ -0,0 +1,62 @@
|
||||
"""Infrastructure outages: the MCP server (ADT in Eclipse, 127.0.0.1:3000) or A4H (localhost:50000) is not reachable.
|
||||
|
||||
On 2026-10-05 21:43 the MCP server refused connections for a short time; three runs got `Connection refused` and the
|
||||
pipeline stopped because it counted three equal harness errors. An outage is not a model result and not a harness bug:
|
||||
the run is cleaned up, the work waits until the services answer, and the same run is started again.
|
||||
"""
|
||||
import errno
|
||||
import os
|
||||
import socket
|
||||
import time
|
||||
import urllib.error
|
||||
import urllib.request
|
||||
|
||||
MCP_URL = os.environ.get("MCP_URL", "http://127.0.0.1:3000/mcp")
|
||||
A4H_URL = os.environ.get("A4H_URL", "http://localhost:50000").rstrip("/") + "/sap/public/ping"
|
||||
|
||||
|
||||
def is_outage(exc):
|
||||
"""True for a refused or lost connection to a local service (not for an HTTP error status)."""
|
||||
if isinstance(exc, urllib.error.HTTPError):
|
||||
return False
|
||||
if isinstance(exc, urllib.error.URLError):
|
||||
return is_outage(exc.reason) if isinstance(exc.reason, BaseException) else True
|
||||
return isinstance(exc, (ConnectionError, socket.timeout, TimeoutError)) or \
|
||||
(isinstance(exc, OSError) and exc.errno in (errno.ECONNREFUSED, errno.ECONNRESET, errno.EHOSTUNREACH, errno.EPIPE))
|
||||
|
||||
|
||||
def _answers(url):
|
||||
try:
|
||||
urllib.request.urlopen(url, timeout=6)
|
||||
return True
|
||||
except urllib.error.HTTPError:
|
||||
return True # 401 / 403 / 404: the server answers
|
||||
except Exception: # noqa: BLE001
|
||||
return False
|
||||
|
||||
|
||||
def up():
|
||||
return _answers(MCP_URL) and _answers(A4H_URL)
|
||||
|
||||
|
||||
def wait_until_up(max_seconds=3600, poll=30, log=print):
|
||||
"""Wait until MCP and A4H answer twice in a row (a restart shows a short flicker). False when the time is over."""
|
||||
end, ok_in_row = time.time() + max_seconds, 0
|
||||
while time.time() < end:
|
||||
ok_in_row = ok_in_row + 1 if up() else 0
|
||||
if ok_in_row >= 2:
|
||||
return True
|
||||
time.sleep(poll)
|
||||
return False
|
||||
|
||||
|
||||
def cleanup_prefix(prefix):
|
||||
"""Delete the A4H objects of an interrupted run (best effort; returns the number of objects found)."""
|
||||
from .mcp_client import McpClient
|
||||
from .runner import DELETE_ORDER, Runner, delete_uris
|
||||
with McpClient() as m:
|
||||
objs = Runner("", "")._objects_with_prefix(m, prefix.upper())
|
||||
objs.sort(key=lambda o: DELETE_ORDER.index(o["objectType"]) if o["objectType"] in DELETE_ORDER else 99)
|
||||
if objs:
|
||||
delete_uris([o["uri"] for o in objs])
|
||||
return len(objs)
|
||||
@@ -23,8 +23,9 @@ import urllib.request
|
||||
|
||||
from .adt_client import load_env
|
||||
from .agents import LlmAgent
|
||||
from . import mix
|
||||
from . import infra, mix
|
||||
from .runner import Runner
|
||||
from .task import prefix_for
|
||||
|
||||
ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||||
sys.path.insert(0, os.path.join(ROOT, "train"))
|
||||
@@ -169,6 +170,21 @@ class Series:
|
||||
try:
|
||||
rep, run_dir = runner.run(item["task"], agent, run_no, teardown=True, tool_budget=SETTINGS["tool_call_budget"])
|
||||
except Exception as e: # noqa: BLE001 a harness exception is not a model result
|
||||
if infra.is_outage(e): # MCP or A4H is down: wait, clean up, run the same task again (not a failure)
|
||||
self.event("infra_outage", item, repr(e)[:150])
|
||||
self.state["status"] = "paused"
|
||||
self.save()
|
||||
infra.wait_until_up(max_seconds=max(60, self.window_end() - now()))
|
||||
try:
|
||||
infra.cleanup_prefix(prefix_for(run_no, item["task"]))
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
for d in __import__("glob").glob(os.path.join(OUT, "runs", f"{run_no}_{item['task']}_*")):
|
||||
self.park(d, "_paused")
|
||||
if now() >= self.window_end():
|
||||
return {"window_over": True}
|
||||
restarts = min(restarts + 1, 9)
|
||||
continue
|
||||
self.event("harness_exception", item, repr(e)[:200])
|
||||
return {"task": item["task"], "kind": item["kind"], "category": item["category"], "harness_error": repr(e)[:200]}
|
||||
er = rep.get("end_reason")
|
||||
|
||||
@@ -17,7 +17,9 @@ from concurrent.futures import ThreadPoolExecutor
|
||||
|
||||
from .adt_client import load_env
|
||||
from .agents import LlmAgent
|
||||
from . import infra
|
||||
from .ledger import BudgetExceeded, check_budget, spent
|
||||
from .task import prefix_for
|
||||
from .record import LEDGER_TO_USAGE
|
||||
from .runner import Runner
|
||||
|
||||
@@ -77,7 +79,25 @@ def one(task_id, attempt, stop):
|
||||
t0 = time.time()
|
||||
row = {"task": task_id, "attempt": attempt, "run": run_no, "model": MODEL}
|
||||
try:
|
||||
rep, run_dir = runner.run(task_id, agent, run_no, teardown=True, cds_calls=CDS_CALLS)
|
||||
for _try in range(4): # an outage of MCP or A4H is waited out and the same run starts again
|
||||
try:
|
||||
rep, run_dir = runner.run(task_id, agent, run_no, teardown=True, cds_calls=CDS_CALLS)
|
||||
break
|
||||
except Exception as e: # noqa: BLE001
|
||||
if not infra.is_outage(e) or _try == 3:
|
||||
raise
|
||||
print(time.strftime("%F %T"), "infra outage at", task_id, "->", repr(e)[:120], "; waiting", flush=True)
|
||||
append("outages.jsonl", {"t": time.time(), "task": task_id, "attempt": attempt, "error": repr(e)[:200]})
|
||||
if not infra.wait_until_up():
|
||||
raise
|
||||
try:
|
||||
infra.cleanup_prefix(prefix_for(run_no, task_id))
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
for d in glob.glob(os.path.join(OUT, f"{run_no}_{task_id}_*")):
|
||||
os.makedirs(os.path.join(OUT, "_aborted"), exist_ok=True)
|
||||
os.rename(d, os.path.join(OUT, "_aborted", os.path.basename(d) + "_" + str(int(time.time()))))
|
||||
agent = LlmAgent(MODEL, loop_guard=3)
|
||||
score = (rep.get("score") or {}).get("total")
|
||||
rec = os.path.exists(os.path.join(run_dir, "record.json"))
|
||||
row.update(score=score, setup_failed=bool(rep.get("setup_failed")), end_reason=rep.get("end_reason"),
|
||||
|
||||
Reference in New Issue
Block a user