"""Stage 2 trajectory runs: the teacher model solves accepted training tasks on A4H; every run writes record.json (harness/record.py). Runs go to runs/traj/. python3 -m harness.trajectories run [--workers 6] [--attempts 1] [--limit N] [--tasks G1000 G1003] Stops when the budget guard (BUDGET_LIMIT_USD in .env) is reached. Resumable: a (task, attempt) with a line in runs/traj/summary.jsonl is skipped. """ import argparse import glob import json import os import threading import time import traceback from concurrent.futures import ThreadPoolExecutor from .adt_client import load_env from .agents import LlmAgent from .ledger import BudgetExceeded, check_budget, spent from .record import LEDGER_TO_USAGE from .runner import Runner ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) POOL = os.path.join(ROOT, "tasks_gen", "train") OUT = os.path.join(ROOT, "runs", "traj") RUN_BASE = 200000 # 200000 + (task number - 1000) * 3 + attempt (a digit must lead the 4-char base36 run: < 466560) MODEL = "deepseek-v4.1-flash:cloud" CDS_CALLS = 100 # tool-call budget for tasks with a CDS contract object (eval keeps 60; Kral 2026-10-05) LOCK = threading.Lock() def accepted_tasks(): out = [] for f in sorted(glob.glob(os.path.join(POOL, "G*", "generation.json"))): try: if json.load(open(f)).get("accepted"): out.append(os.path.basename(os.path.dirname(f))) except (OSError, ValueError): pass return out def done_keys(): path = os.path.join(OUT, "summary.jsonl") if not os.path.exists(path): return set() return {(r["task"], r["attempt"]) for r in map(json.loads, open(path))} def failed_first_attempts(min_score): """Tasks whose attempt 0 did not give an acceptable run (score, end reason, harness error).""" path = os.path.join(OUT, "summary.jsonl") rows = [json.loads(l) for l in open(path)] if os.path.exists(path) else [] return sorted(r["task"] for r in rows if r["attempt"] == 0 and ( r.get("harness_error") or r.get("setup_failed") or r.get("end_reason") != "report" or (r.get("score") or 0) < min_score)) def append(name, row): with LOCK, open(os.path.join(OUT, name), "a") as f: f.write(json.dumps(row) + "\n") def one(task_id, attempt, stop): if stop.is_set(): return try: check_budget() except BudgetExceeded as e: stop.set() print("BUDGET", e, flush=True) return run_no = RUN_BASE + (int("".join(c for c in task_id if c.isdigit())) - 1000) * 3 + attempt agent = LlmAgent(MODEL, loop_guard=3) runner = Runner(POOL, OUT) 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) 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"), tool_calls=rep.get("tool_calls"), record=rec, run_dir=os.path.basename(run_dir), teardown_ok=isinstance(rep.get("teardown"), dict) and all(v.get("deleted") for v in rep["teardown"].values())) except BudgetExceeded as e: stop.set() row.update(error=f"budget: {e}", harness_error=True) except Exception as e: # noqa: BLE001 a harness error is not a model result: recorded, never trained on row.update(error=repr(e)[:300], harness_error=True) append("errors.jsonl", dict(row, trace=traceback.format_exc()[-1500:])) row["seconds"] = round(time.time() - t0, 1) row["spent_ledger"] = spent() append("summary.jsonl", row) print(json.dumps(row), flush=True) def main(): load_env(os.path.join(ROOT, ".env")) ap = argparse.ArgumentParser() ap.add_argument("cmd", choices=["run", "list"]) ap.add_argument("--workers", type=int, default=6) ap.add_argument("--attempts", type=int, default=1, help="1 = one attempt per task; with --retry-failed: the second attempt only for failed tasks") ap.add_argument("--retry-failed", action="store_true") ap.add_argument("--min-score", type=float, default=80) ap.add_argument("--limit", type=int, help="maximum number of runs in this call") ap.add_argument("--tasks", nargs="*") a = ap.parse_args() os.makedirs(OUT, exist_ok=True) tasks = a.tasks or accepted_tasks() done = done_keys() todo = [(t, k) for k in range(a.attempts) for t in tasks if (t, k) not in done] if a.retry_failed: # attempt 1 only for tasks whose attempt 0 failed todo = [(t, 1) for t in failed_first_attempts(a.min_score) if t in tasks and (t, 1) not in done] if a.limit: todo = todo[:a.limit] print(len(tasks), "accepted tasks,", len(todo), "runs to do, spent", spent(), flush=True) if a.cmd == "list": return stop = threading.Event() with ThreadPoolExecutor(a.workers) as ex: for t, k in todo: ex.submit(one, t, k, stop) if __name__ == "__main__": main()