"""Stage 2 pipeline controller (rules of 2026-10-05, Opus 5.5 + Kral). Runs in the background; nobody polls it. It - starts plan 2 generation (3 workers, hard and error share raised, no fixed target) when the first plan has ended, until --gen-deadline or the budget guard; adds K (free-text) variants at about 10 % of the other tasks; - runs trajectories with 2 workers (1 after any A4H error), at most 2 attempts per task: a second attempt only when the first one failed or was accepted without a repair; - writes a summary to train/STATE.md and docs/yol-haritasi.md every 50 accepted trajectories, and commits; - stops and reports (runs/pipeline/STOPPED.txt, STOP flag for the generators) when: acceptance over the last 30 runs is below 50 %, one harness error pattern appears 3 times, the budget guard stops, or the deadline. python3 -m harness.pipeline [--wait-pid PID] [--workers 2] [--gen-deadline 2026-10-10T18:00] [--traj-deadline 2026-10-11T23:30] Budget: BUDGET_LIMIT_USD and BUDGET_RESERVE_USD in .env (the reserve stays for a second teacher test). """ import argparse import glob import json import os import re import statistics import subprocess import sys import threading import time from .adt_client import load_env from . import trainset, trajectories from .ledger import spent ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) sys.path.insert(0, os.path.join(ROOT, "train")) import accept as acc # noqa: E402 (train/accept.py) OUT = trajectories.OUT PIPE = os.path.join(ROOT, "runs", "pipeline") BASE_LEDGER = 66.9665 # ledger when Part B started (2026-10-05 morning) LEDGER_TO_USAGE = 2.7 ALLOWANCE_USAGE = 35.0 # until 11 October, without asking (Kral) MARKER = "" GIT_ENV = dict(os.environ, GIT_AUTHOR_NAME="Kral", GIT_AUTHOR_EMAIL="kral@local", GIT_COMMITTER_NAME="Kral", GIT_COMMITTER_EMAIL="kral@local") HARNESS_MATCH = acc.HARNESS_ERRORS def log(*a): print(time.strftime("%F %T"), *a, flush=True) def ts(text): return time.mktime(time.strptime(text, "%Y-%m-%dT%H:%M")) class Pipeline: def __init__(self, a): self.a = a self.lock = threading.Lock() self.stop = threading.Event() self.reason = None self.max_workers = a.workers self.inflight = set() self.cache = {} # run_dir -> evaluation self.events = {} # harness error key -> count self.task_meta = {} self.gen_deadline, self.traj_deadline = ts(a.gen_deadline), ts(a.traj_deadline) self.children = [] rows = self.rows() for r in rows: self.evaluate(r) self.milestone = self.accepted_total() // 50 # ---- data ----------------------------------------------------------------------------------------------- def rows(self): p = os.path.join(OUT, "summary.jsonl") return [json.loads(l) for l in open(p)] if os.path.exists(p) else [] def meta(self, task_id): if task_id not in self.task_meta: try: t = json.load(open(os.path.join(trajectories.POOL, task_id, "task.json"))) except OSError: t = {} self.task_meta[task_id] = {"category": t.get("category"), "object_type": t.get("object_type")} return self.task_meta[task_id] def evaluate(self, row): """{accepted, repair, reasons, event_keys, hints, tokens-free} of one summary row (cached).""" key = row.get("run_dir") or f"{row['task']}_{row['attempt']}" if key in self.cache: return self.cache[key] ev = {"task": row["task"], "attempt": row["attempt"], "accepted": False, "repair": False, "reasons": [], "events": [], "hints": 0} path = os.path.join(OUT, row.get("run_dir") or "-", "record.json") if row.get("harness_error"): ev["reasons"].append("harness_error") ev["events"].append("harness:" + str(row.get("error"))[:60]) if row.get("setup_failed"): ev["reasons"].append("setup_failed") ev["events"].append("setup_failed") if os.path.exists(path): rec = json.load(open(path)) ok, why = acc.judge(rec, row, self.a.min_score) ev["accepted"], ev["reasons"] = ok, why or ev["reasons"] ev["hints"] = rec["metadata"].get("syntax_hints") or 0 if ok: ev["repair"] = acc.has_repair(acc.trim_loops(rec["messages"])[0]) for r in rec["tool_results_raw"]: m = HARNESS_MATCH.search(r.get("result") or "") if m and not r.get("tool", "").startswith("sap_run_unit"): ev["events"].append("text:" + m.group(0)[:40]) break if row.get("teardown_ok") is False: ev["events"].append("teardown") elif not ev["reasons"]: ev["reasons"].append("no_record") self.cache[key] = ev return ev def evaluated(self): return [self.evaluate(r) for r in self.rows()] def accepted_total(self): return sum(e["accepted"] for e in self.evaluated()) # ---- jobs ----------------------------------------------------------------------------------------------- def next_job(self): with self.lock: rows = self.rows() done = {(r["task"], r["attempt"]) for r in rows} ev0 = {r["task"]: self.evaluate(r) for r in rows if r["attempt"] == 0} tasks = trajectories.accepted_tasks() for t in tasks: # first attempt for every task first if (t, 0) not in done and (t, 0) not in self.inflight: self.inflight.add((t, 0)) return t, 0 for t in tasks: # second attempt: first failed, or accepted without a repair e = ev0.get(t) if e and (t, 1) not in done and (t, 1) not in self.inflight and (not e["accepted"] or not e["repair"]): self.inflight.add((t, 1)) return t, 1 return None def worker(self, i): while not self.stop.is_set(): if time.time() > self.traj_deadline: self.halt("trajectory deadline reached") return if i >= self.max_workers: time.sleep(30) continue job = self.next_job() if not job: if self.generation_alive() or self.inflight: time.sleep(60) continue self.halt("no job left and generation ended") return trajectories.one(job[0], job[1], self.stop) with self.lock: self.inflight.discard(job) self.after_run(job) if self.stop.is_set() and not self.reason: self.halt("budget guard stopped a run") # ---- rules ---------------------------------------------------------------------------------------------- def after_run(self, job): rows = self.rows() row = next((r for r in reversed(rows) if (r["task"], r["attempt"]) == job), None) if row is None: return ev = self.evaluate(row) for k in ev["events"]: self.events[k] = self.events.get(k, 0) + 1 log("A4H/harness event:", k, "x", self.events[k], "->", row.get("run_dir")) if self.max_workers > 1: self.max_workers = 1 log("harness error seen: trajectory workers 2 -> 1") if self.events[k] >= 3: self.halt(f"harness error pattern repeated 3 times: {k}") evs = self.evaluated() if len(evs) >= 30: rate = sum(e["accepted"] for e in evs[-30:]) / 30 if rate < 0.5: self.halt(f"acceptance {rate:.0%} over the last 30 runs (below 50 %)") try: from .ledger import check_budget check_budget() except Exception as e: # noqa: BLE001 BudgetExceeded self.halt(f"budget guard: {e}") total = sum(e["accepted"] for e in evs) while total // 50 > self.milestone and not self.reason: self.milestone += 1 self.summary(self.milestone * 50) def halt(self, reason): if self.stop.is_set() and self.reason: return self.reason = reason self.stop.set() os.makedirs(PIPE, exist_ok=True) open(trainset.STOP_FLAG, "w").write(reason + "\n") # generators end after their current slot log("STOP:", reason) # ---- generation side -------------------------------------------------------------------------------------- def generation_alive(self): if self.reason: return False out = subprocess.run(["pgrep", "-f", r"harness\.trainset (run|k)"], capture_output=True, text=True).stdout return bool(out.split()) or time.time() < self.gen_deadline def old_generators(self): out = subprocess.run(["pgrep", "-fl", r"harness\.trainset run --part"], capture_output=True, text=True).stdout return [l for l in out.splitlines() if "--plan plan2" not in l] def supervisor(self): launched = False k_proc = None crashes = 0 while not self.stop.is_set(): time.sleep(60) if time.time() > self.gen_deadline or os.path.exists(trainset.STOP_FLAG): continue if not launched and not self.old_generators(): trainset.ensure_plan("plan2") dl = self.a.gen_deadline for i in range(3): f = open(os.path.join(ROOT, "runs", "gen_train", f"p2_w{i}.log"), "a") self.children.append(subprocess.Popen( [sys.executable, "-m", "harness.trainset", "run", "--plan", "plan2", "--part", str(i), "--parts", "3", "--deadline", dl], cwd=ROOT, stdout=f, stderr=f)) launched = True log("plan 2 generators started (3 workers, deadline", dl + ")") elif launched: for i, c in enumerate(self.children): if c.poll() not in (None, 0) and crashes < 3: crashes += 1 log("generator", i, "exited with", c.returncode, "- restarted") f = open(os.path.join(ROOT, "runs", "gen_train", f"p2_w{i}.log"), "a") self.children[i] = subprocess.Popen( [sys.executable, "-m", "harness.trainset", "run", "--plan", "plan2", "--part", str(i), "--parts", "3", "--deadline", self.a.gen_deadline], cwd=ROOT, stdout=f, stderr=f) # K variants: free text, about 10 % of the other accepted training tasks if k_proc is None or k_proc.poll() is not None: logs = [json.load(open(f)) for f in glob.glob(os.path.join(trainset.POOL, "_logs", "G*.json"))] base = sum(1 for l in logs if l.get("accepted") and l.get("category") != "K") have = sum(1 for l in logs if l.get("category") == "K") want = min(99, int(0.1 * base)) if want > have: f = open(os.path.join(ROOT, "runs", "gen_train", "k.log"), "a") k_proc = subprocess.Popen([sys.executable, "-m", "harness.trainset", "k", "--k-count", str(want)], cwd=ROOT, stdout=f, stderr=f) # ---- summary ---------------------------------------------------------------------------------------------- def summary(self, n): log("summary at", n, "accepted trajectories") subprocess.run([sys.executable, "train/accept.py", "--min-score", str(self.a.min_score)], cwd=ROOT, capture_output=True) vp = os.path.join(ROOT, "train", ".venv", "bin", "python") subprocess.run([vp, "train/to_qwen.py", os.path.join(OUT, "accepted.jsonl"), os.path.join(OUT, "qwen.jsonl")], cwd=ROOT, capture_output=True) toks = sorted(json.loads(l)["n_tokens"] for l in open(os.path.join(OUT, "qwen.jsonl"))) pct = lambda p: toks[min(len(toks) - 1, int(p * len(toks)))] if toks else None # noqa: E731 rows, evs = self.rows(), self.evaluated() logs = [json.load(open(f)) for f in glob.glob(os.path.join(trainset.POOL, "_logs", "G*.json"))] acc_t = sum(1 for l in logs if l.get("accepted")) accepted = [e for e in evs if e["accepted"]] repairs = sum(e["repair"] for e in accepted) def table(keyname): agg = {} for e in evs: k = self.meta(e["task"])[keyname] or "?" a = agg.setdefault(k, [0, 0, 0]) a[0] += 1 a[1] += e["accepted"] a[2] += e["accepted"] and e["repair"] return "; ".join(f"{k} {v[1]}/{v[0]} (repair {v[2]})" for k, v in sorted(agg.items())) s = spent() used = (s - BASE_LEDGER) / LEDGER_TO_USAGE guard = float(os.environ.get("BUDGET_LIMIT_USD", "0")) - float(os.environ.get("BUDGET_RESERVE_USD", "0")) text = f"""### Stage 2 summary at {n} accepted trajectories ({time.strftime('%Y-%m-%d %H:%M')}) - Tasks: {acc_t} accepted of {len(logs)} generated (K variants {sum(1 for l in logs if l.get('category') == 'K')}); \ trajectory runs {len(rows)}, accepted trajectories {len(accepted)} (acceptance {len(accepted) / max(len(rows), 1):.0%}). - Budget: used {used:.1f} usage (ledger {s - BASE_LEDGER:.1f}) since 2026-10-05; left to the guard {(guard - s) / LEDGER_TO_USAGE:.1f} usage \ (ledger {guard - s:.1f}; reserve {os.environ.get('BUDGET_RESERVE_USD', '0')} ledger kept); allowance until 11 October {ALLOWANCE_USAGE - used:.1f} usage. - Repair share (accepted trajectories with an error followed by a fix): {repairs}/{len(accepted)} = {repairs / max(len(accepted), 1):.0%}. - By category, accepted/runs: {table('category')}. - By object type, accepted/runs: {table('object_type')}. - Tokens of accepted samples (20 tool schemas kept): p50 {pct(0.5)}, p90 {pct(0.9)}, p95 {pct(0.95)}, max {toks[-1] if toks else None}, n {len(toks)}. - Syntax hints (proxy syntaxCheck added): {sum(e['hints'] for e in evs)} in {sum(1 for e in evs if e['hints'])} runs. - Harness events: {sum(self.events.values())} ({self.events or 'none'}); trajectory workers now {self.max_workers}. """ if getattr(self.a, "dry_summary", False): print(text) return with open(os.path.join(ROOT, "train", "STATE.md"), "a") as f: f.write("\n" + text) p = os.path.join(ROOT, "docs", "yol-haritasi.md") d = open(p).read() d = d.replace(MARKER, MARKER + "\n\n" + text, 1) if MARKER in d else d + "\n" + text open(p, "w").write(d) subprocess.run(["git", "add", "train/STATE.md", "docs/yol-haritasi.md", "tasks_gen/train", "harness", "train"], cwd=ROOT, env=GIT_ENV, capture_output=True) subprocess.run(["git", "commit", "-q", "-m", f"Stage 2 summary at {n} accepted trajectories\n\n" "Co-Authored-By: Claude Sonnet 5.5 "], cwd=ROOT, env=GIT_ENV, capture_output=True) # ---- main ------------------------------------------------------------------------------------------------- def run(self): os.makedirs(PIPE, exist_ok=True) if os.path.exists(trainset.STOP_FLAG): os.remove(trainset.STOP_FLAG) log("pipeline start: workers", self.max_workers, "accepted", self.accepted_total(), "spent", spent()) threads = [threading.Thread(target=self.worker, args=(i,)) for i in range(self.a.workers)] threading.Thread(target=self.supervisor, daemon=True).start() for t in threads: t.start() for t in threads: t.join() evs = self.evaluated() text = (f"reason: {self.reason}\ntime: {time.strftime('%F %T')}\nruns: {len(evs)}, accepted " f"{sum(e['accepted'] for e in evs)}\nspent ledger: {spent()}\nevents: {self.events}\n") open(os.path.join(PIPE, "STOPPED.txt"), "w").write(text) log("pipeline ended:", text.replace("\n", " | ")) def main(): load_env(os.path.join(ROOT, ".env")) ap = argparse.ArgumentParser() ap.add_argument("--wait-pid", type=int, help="start when this process has ended (the first trajectory batch)") ap.add_argument("--workers", type=int, default=2) ap.add_argument("--min-score", type=float, default=80) ap.add_argument("--gen-deadline", default="2026-10-10T18:00") ap.add_argument("--traj-deadline", default="2026-10-11T23:30") ap.add_argument("--dry-summary", action="store_true", help="print the summary text and exit (no file, no commit)") a = ap.parse_args() if a.dry_summary: Pipeline(a).summary(0) return if a.wait_pid: while True: try: os.kill(a.wait_pid, 0) except OSError: break time.sleep(30) os.makedirs(OUT, exist_ok=True) Pipeline(a).run() if __name__ == "__main__": main()