diff --git a/harness/dashboard.py b/harness/dashboard.py index 110d1d6..1744313 100644 --- a/harness/dashboard.py +++ b/harness/dashboard.py @@ -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 diff --git a/harness/infra.py b/harness/infra.py new file mode 100644 index 0000000..65ad2f8 --- /dev/null +++ b/harness/infra.py @@ -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) diff --git a/harness/localqwen.py b/harness/localqwen.py index 27180ef..bc8c217 100644 --- a/harness/localqwen.py +++ b/harness/localqwen.py @@ -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") diff --git a/harness/trajectories.py b/harness/trajectories.py index b2ac32f..231ecdc 100644 --- a/harness/trajectories.py +++ b/harness/trajectories.py @@ -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"), diff --git a/train/STATE.md b/train/STATE.md index f83987c..9ec1535 100644 --- a/train/STATE.md +++ b/train/STATE.md @@ -177,3 +177,5 @@ Training runs on HF Jobs with Unsloth, not on the Mac. No `mlx_lm` training. - **Incident 19:25:** restarting the controller I started a second one by mistake (wrong `pgrep` pattern) and deleted the objects of a running run. Both controllers were stopped, leftovers cleaned, one controller runs. One table `Z4AJ50UB_PO_HEAD` kept a lock from the interrupted write (SM12 needed). Interrupting a write leaks the lock: never kill a controller during a run without checking. - 2026-10-05 21:31 budget correction 4: panel 50.00 at ledger 111.22; since 45.0 (100.62): 10.6 ledger / 5.0 usage = 2.12 (generation of new types and first runs). Guard recomputed with the pessimistic ratio 1.2: panel left 60 - 50 - 3.0 reserve = 7.0 usage x 1.2 = 8.4 ledger, guard 119.6, `BUDGET_LIMIT_USD` 128 (reserve 8). Only kinds below target run (Kral + Opus 2026-10-05). + +- 2026-10-05 22:35 **pipeline stopped at 21:43 by an infrastructure outage, found 22:26 (my miss).** The MCP server (127.0.0.1:3000) refused connections for a short time; three runs got `URLError(ConnectionRefusedError)` and the pipeline counted three equal harness errors and stopped (rule). Not a model or harness bug: A4H was up (11 h), the MCP server answers again. New `harness/infra.py`: an outage of MCP or A4H is waited out (check every 30 s, two answers in a row), the run's objects are deleted, the same run starts again; it is not counted as an event or a failure (`runs/traj/outages.jsonl`). Same for series A (applies after a restart of `harness.localqwen`; the running series keeps the old code and would skip the task). The 3 interrupted runs were cleaned (11 objects) and get their second attempt. Budget: panel 50.84 at ledger 113.36 (2.54 ledger per usage in the last stretch); guard recomputed: panel left 9.16 - 3.0 reserve = 6.2 usage x 1.2 = 7.4 ledger, `BUDGET_LIMIT_USD` 129 (guard 121). Dashboard: a stopped pipeline shows no burn rate or ETA.