412 lines
20 KiB
Python
412 lines
20 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()
|
|
counts = {} # accepted trajectories (and runs in flight) per kind
|
|
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
|
|
below = mix.below_target(counts) # only kinds below their target share run (Kral + Opus 2026-10-05)
|
|
cands = [t for t in tasks if (t, 0) not in done and (t, 0) not in self.inflight
|
|
and mix.kind_of_task_dir(t) in below]
|
|
if cands: # first attempts: the kind with the biggest deficit against the target mix goes first
|
|
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"]) \
|
|
and mix.kind_of_task_dir(t) in below:
|
|
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 acquire_controller_lock():
|
|
"""One controller only. An exclusive flock on runs/pipeline/controller.lock lives as long as the process (the
|
|
kernel drops it when the process dies, so a crash never leaves a stale lock); the pid is written for people."""
|
|
import fcntl
|
|
os.makedirs(PIPE, exist_ok=True)
|
|
path = os.path.join(PIPE, "controller.lock")
|
|
f = open(path, "a+")
|
|
try:
|
|
fcntl.flock(f, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
|
except OSError:
|
|
f.seek(0)
|
|
sys.exit("a controller is already running (pid %s, %s); not starting a second one" % (
|
|
f.read().strip() or "?", path))
|
|
f.seek(0)
|
|
f.truncate()
|
|
f.write("%d\n" % os.getpid())
|
|
f.flush()
|
|
return f
|
|
|
|
|
|
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)
|
|
guard = acquire_controller_lock()
|
|
Pipeline(a).run()
|
|
guard.close()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|