Trajectory record (messages, raw tool results, metadata; reasoning apart) and trajectory runner
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
@@ -85,6 +85,7 @@ class LlmAgent:
|
|||||||
self.empty_retries = 2
|
self.empty_retries = 2
|
||||||
self.loop_guard = loop_guard # end the run after this many identical pushes in a row (None = off)
|
self.loop_guard = loop_guard # end the run after this many identical pushes in a row (None = off)
|
||||||
self.end_reason = None
|
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}
|
self.chat_template_kwargs = chat_template_kwargs # local server only, e.g. {"enable_thinking": False}
|
||||||
|
|
||||||
def _chat(self, messages, tools):
|
def _chat(self, messages, tools):
|
||||||
@@ -118,6 +119,7 @@ class LlmAgent:
|
|||||||
for t in proxy.schemas()]
|
for t in proxy.schemas()]
|
||||||
messages = [{"role": "system", "content": SYSTEM_PROMPT},
|
messages = [{"role": "system", "content": SYSTEM_PROMPT},
|
||||||
{"role": "user", "content": task.spec}]
|
{"role": "user", "content": task.spec}]
|
||||||
|
self.messages, self.tools, self.reasoning, self.turn_usage = messages, tools, [], []
|
||||||
if ":cloud" in self.model:
|
if ":cloud" in self.model:
|
||||||
check_budget()
|
check_budget()
|
||||||
final = ""
|
final = ""
|
||||||
@@ -132,6 +134,7 @@ class LlmAgent:
|
|||||||
try:
|
try:
|
||||||
msg, usage = self._chat(messages, tools)
|
msg, usage = self._chat(messages, tools)
|
||||||
add_usage(self.model, usage, kind="run", ref=proxy.prefix)
|
add_usage(self.model, usage, kind="run", ref=proxy.prefix)
|
||||||
|
self.turn_usage.append(usage)
|
||||||
except RuntimeError as e:
|
except RuntimeError as e:
|
||||||
final = f"Stopped: {e}"
|
final = f"Stopped: {e}"
|
||||||
self.end_reason = "model_error"
|
self.end_reason = "model_error"
|
||||||
@@ -139,6 +142,9 @@ class LlmAgent:
|
|||||||
proxy.note("assistant", {"content": msg.get("content"),
|
proxy.note("assistant", {"content": msg.get("content"),
|
||||||
"tool_calls": msg.get("tool_calls"), "usage": usage})
|
"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")})
|
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 []
|
calls = msg.get("tool_calls") or []
|
||||||
if not calls and not (msg.get("content") or "").strip() and empty < self.empty_retries:
|
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
|
empty += 1 # empty turn (output limit reached in reasoning): ask the same question again
|
||||||
|
|||||||
68
harness/record.py
Normal file
68
harness/record.py
Normal file
@@ -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
|
||||||
@@ -230,7 +230,8 @@ class Runner:
|
|||||||
rep["rescored"] = True
|
rep["rescored"] = True
|
||||||
else:
|
else:
|
||||||
proxy = ToolProxy(mcp, prefix, task.meta.get("budget", {}),
|
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)
|
proxy.note("spec", task.spec)
|
||||||
t1 = time.time()
|
t1 = time.time()
|
||||||
rep["final_report"] = agent.run(task, proxy)
|
rep["final_report"] = agent.run(task, proxy)
|
||||||
@@ -239,6 +240,7 @@ class Runner:
|
|||||||
rep["max_fail_streak"] = proxy.max_fail_streak
|
rep["max_fail_streak"] = proxy.max_fail_streak
|
||||||
rep["activation_failures"] = proxy.activation_failures
|
rep["activation_failures"] = proxy.activation_failures
|
||||||
rep["activation_error_messages"] = proxy.activation_errors
|
rep["activation_error_messages"] = proxy.activation_errors
|
||||||
|
rep["syntax_hints"] = proxy.syntax_hints
|
||||||
rep["end_reason"] = getattr(agent, "end_reason", None)
|
rep["end_reason"] = getattr(agent, "end_reason", None)
|
||||||
if proxy.fallbacks:
|
if proxy.fallbacks:
|
||||||
rep["adt_fallbacks"] = proxy.fallbacks
|
rep["adt_fallbacks"] = proxy.fallbacks
|
||||||
@@ -409,6 +411,9 @@ class Runner:
|
|||||||
all_objs = self._objects_with_prefix(mcp, prefix)
|
all_objs = self._objects_with_prefix(mcp, prefix)
|
||||||
rep["teardown"] = self._teardown(all_objs, run_dir) if teardown else "skipped"
|
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)
|
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
|
return rep, run_dir
|
||||||
|
|
||||||
def _score(self, task, rep):
|
def _score(self, task, rep):
|
||||||
|
|||||||
114
harness/trajectories.py
Normal file
114
harness/trajectories.py
Normal file
@@ -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()
|
||||||
Reference in New Issue
Block a user