From d5e43e1a837f78fae7eb5a3d6933a273570b43ff Mon Sep 17 00:00:00 2001 From: Kral Date: Mon, 5 Oct 2026 12:03:54 +0200 Subject: [PATCH] Trajectory record (messages, raw tool results, metadata; reasoning apart) and trajectory runner Co-Authored-By: Claude Sonnet 5.5 --- harness/agents.py | 6 +++ harness/record.py | 68 ++++++++++++++++++++++++ harness/runner.py | 7 ++- harness/trajectories.py | 114 ++++++++++++++++++++++++++++++++++++++++ 4 files changed, 194 insertions(+), 1 deletion(-) create mode 100644 harness/record.py create mode 100644 harness/trajectories.py diff --git a/harness/agents.py b/harness/agents.py index 9d50524..1019456 100644 --- a/harness/agents.py +++ b/harness/agents.py @@ -85,6 +85,7 @@ class LlmAgent: self.empty_retries = 2 self.loop_guard = loop_guard # end the run after this many identical pushes in a row (None = off) self.end_reason = None + self.messages, self.tools, self.reasoning, self.turn_usage = [], [], [], [] # for the trajectory record self.chat_template_kwargs = chat_template_kwargs # local server only, e.g. {"enable_thinking": False} def _chat(self, messages, tools): @@ -118,6 +119,7 @@ class LlmAgent: for t in proxy.schemas()] messages = [{"role": "system", "content": SYSTEM_PROMPT}, {"role": "user", "content": task.spec}] + self.messages, self.tools, self.reasoning, self.turn_usage = messages, tools, [], [] if ":cloud" in self.model: check_budget() final = "" @@ -132,6 +134,7 @@ class LlmAgent: try: msg, usage = self._chat(messages, tools) add_usage(self.model, usage, kind="run", ref=proxy.prefix) + self.turn_usage.append(usage) except RuntimeError as e: final = f"Stopped: {e}" self.end_reason = "model_error" @@ -139,6 +142,9 @@ class LlmAgent: proxy.note("assistant", {"content": msg.get("content"), "tool_calls": msg.get("tool_calls"), "usage": usage}) messages.append({k: v for k, v in msg.items() if k in ("role", "content", "tool_calls")}) + think = msg.get("reasoning") or msg.get("reasoning_content") or msg.get("thinking") + if think: # kept apart from the messages: not training data + self.reasoning.append({"message_index": len(messages) - 1, "reasoning": think}) calls = msg.get("tool_calls") or [] if not calls and not (msg.get("content") or "").strip() and empty < self.empty_retries: empty += 1 # empty turn (output limit reached in reasoning): ask the same question again diff --git a/harness/record.py b/harness/record.py new file mode 100644 index 0000000..7369dcc --- /dev/null +++ b/harness/record.py @@ -0,0 +1,68 @@ +"""Trajectory record of one run (stage 2 training data source). + +record.json : task metadata, the exact messages the teacher saw and wrote, the tool schemas, the raw tool + results (as the proxy logged them, before the 12000-character cut of the message), run metadata. +reasoning.json: the teacher reasoning per assistant message. Never training data. +""" +import json +import os + +from .ledger import LEDGER + +LEDGER_TO_USAGE = 2.7 # usage = ledger / 2.7 (measured 2026-10-03) + + +def run_cost(prefix): + """Ledger USD of the 'run' entries of this run (by prefix).""" + if not os.path.exists(LEDGER): + return 0.0 + tot = 0.0 + for line in open(LEDGER): + e = json.loads(line) + if e.get("kind") == "run" and e.get("ref") == prefix: + tot += e["usd"] + return round(tot, 5) + + +def raw_tool_results(traj_path): + out = [] + for line in open(traj_path): + e = json.loads(line) + if "tool" in e: + out.append({k: e[k] for k in ("tool", "args", "is_error", "result", "raw_result", "variant_tool") if k in e}) + return out + + +def write_record(run_dir, task, agent, rep, pool): + msgs = getattr(agent, "messages", None) + if not msgs: + return None + meta = task.meta + cost = run_cost(rep["prefix"]) + score = rep.get("score") or {} + record = { + "version": 1, + "run": rep["run"], "prefix": rep["prefix"], + "task": {"id": task.id, "pool": pool, "category": meta.get("category"), + "object_type": meta.get("object_type"), "release_target": meta.get("release_target"), + "difficulty": meta.get("difficulty"), "error_kind": meta.get("error_kind"), + "tool_schema": meta.get("tool_schema")}, + "teacher": agent.model, + "messages": msgs, + "tools": agent.tools, + "tool_results_raw": raw_tool_results(os.path.join(run_dir, "trajectory.jsonl")), + "metadata": { + "score": score.get("total"), "score_parts": score, "gates": rep.get("gates"), + "cost_ledger_usd": cost, "cost_usage_usd": round(cost / LEDGER_TO_USAGE, 5), + "end_reason": rep.get("end_reason"), "tool_calls": rep.get("tool_calls"), + "activations": rep.get("activations"), "activation_failures": rep.get("activation_failures"), + "max_fail_streak": rep.get("max_fail_streak"), "syntax_hints": rep.get("syntax_hints"), + "adt_fallbacks": bool(rep.get("adt_fallbacks")), "agent_seconds": rep.get("agent_seconds"), + "turn_usage": getattr(agent, "turn_usage", []), + "setup_failed": rep.get("setup_failed", False), + "harness_error": rep.get("harness_error"), + }, + } + json.dump(record, open(os.path.join(run_dir, "record.json"), "w"), indent=1) + json.dump(getattr(agent, "reasoning", []), open(os.path.join(run_dir, "reasoning.json"), "w"), indent=1) + return record diff --git a/harness/runner.py b/harness/runner.py index 433ffe4..b08f251 100644 --- a/harness/runner.py +++ b/harness/runner.py @@ -230,7 +230,8 @@ class Runner: rep["rescored"] = True else: proxy = ToolProxy(mcp, prefix, task.meta.get("budget", {}), - os.path.join(run_dir, "trajectory.jsonl"), task.meta.get("tool_schema")) + os.path.join(run_dir, "trajectory.jsonl"), task.meta.get("tool_schema"), + task.meta.get("release_target")) proxy.note("spec", task.spec) t1 = time.time() rep["final_report"] = agent.run(task, proxy) @@ -239,6 +240,7 @@ class Runner: rep["max_fail_streak"] = proxy.max_fail_streak rep["activation_failures"] = proxy.activation_failures rep["activation_error_messages"] = proxy.activation_errors + rep["syntax_hints"] = proxy.syntax_hints rep["end_reason"] = getattr(agent, "end_reason", None) if proxy.fallbacks: rep["adt_fallbacks"] = proxy.fallbacks @@ -409,6 +411,9 @@ class Runner: all_objs = self._objects_with_prefix(mcp, prefix) rep["teardown"] = self._teardown(all_objs, run_dir) if teardown else "skipped" json.dump(rep, open(os.path.join(run_dir, "report.json"), "w"), indent=1) + if not rescore_dir: + from .record import write_record + write_record(run_dir, task, agent, rep, os.path.basename(os.path.normpath(self.tasks_root))) return rep, run_dir def _score(self, task, rep): diff --git a/harness/trajectories.py b/harness/trajectories.py new file mode 100644 index 0000000..762fff4 --- /dev/null +++ b/harness/trajectories.py @@ -0,0 +1,114 @@ +"""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 = 41000 # 41000 + (task number - 1000) * 3 + attempt; below 36**3 * ... (prefix rule) +MODEL = "deepseek-v4.1-flash:cloud" +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 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) + 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) + 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.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()