Files
abap-llm/harness/pipeline.py
2026-10-05 19:27:43 +02:00

387 lines
19 KiB
Python

"""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 mix, 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 = "<!-- stage2-summaries -->"
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()
cands = [t for t in tasks if (t, 0) not in done and (t, 0) not in self.inflight]
if cands: # first attempts: the kind with the biggest deficit against the target mix goes first
counts = {}
for r in rows:
if self.evaluate(r)["accepted"]:
k = mix.kind_of_task_dir(r["task"])
counts[k] = counts.get(k, 0) + 1
for (t, _a) in self.inflight:
k = mix.kind_of_task_dir(t)
counts[k] = counts.get(k, 0) + 1
total, tot_share = sum(counts.values()) + 1, sum(mix.TYPE_SHARE.values())
best = max(cands, key=lambda t: (total * mix.TYPE_SHARE.get(mix.kind_of_task_dir(t), 0) / tot_share
- counts.get(mix.kind_of_task_dir(t), 0), -int(t[1:])))
self.inflight.add((best, 0))
return best, 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 balanced" not in l and "--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():
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", "balanced", "--part", str(i),
"--parts", "3", "--deadline", dl], cwd=ROOT, stdout=f, stderr=f))
launched = True
log("balanced 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", "balanced", "--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()))
kinds_t, kinds_r = mix.accepted_task_counts(), {}
for r, e in zip(rows, evs):
k = mix.kind_of_task_dir(r["task"])
a = kinds_r.setdefault(k, [0, 0])
a[0] += 1
a[1] += e["accepted"]
tt, tr_ = max(sum(kinds_t.values()), 1), max(sum(v[1] for v in kinds_r.values()), 1)
share_sum = sum(mix.TYPE_SHARE.values())
kind_line = "; ".join("%s tasks %d (%.0f%%) traj %d/%d (%.0f%%) target %.0f%%" % (
k, kinds_t.get(k, 0), 100.0 * kinds_t.get(k, 0) / tt, kinds_r.get(k, [0, 0])[1], kinds_r.get(k, [0, 0])[0],
100.0 * kinds_r.get(k, [0, 0])[1] / tr_, 100.0 * v / share_sum) for k, v in mix.TYPE_SHARE.items())
ddls = {}
for r, e in zip(rows, evs):
if mix.kind_of_task_dir(r["task"]) == "DDLS" and not e["accepted"]:
for w in (e["reasons"] or ["?"]):
w = w.split("=")[0] if w.startswith(("end_reason", "score")) else w
ddls[w] = ddls.get(w, 0) + 1
ddls_n = sum(1 for r in rows if mix.kind_of_task_dir(r["task"]) == "DDLS")
ddls_acc = sum(1 for r, e in zip(rows, evs) if mix.kind_of_task_dir(r["task"]) == "DDLS" and e["accepted"])
s = spent()
used = (s - BASE_LEDGER) / LEDGER_TO_USAGE
from .ledger import _env_budget
lim, res = _env_budget()
guard = lim - res
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 {res:g} 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)}.
- Object type mix (accepted tasks; accepted/runs trajectories; target): {kind_line}.
- DDLS (CDS) runs: {ddls_acc} accepted of {ddls_n}; reject reasons {ddls or 'none'} (activation, hidden tests and ATC rejections are separate gate failures, they show as score reasons).
- 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 <noreply@anthropic.com>"],
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()