Pipeline controller: plan 2 (hard +, error +50 %), backlog throttle, 2 trajectory workers, stop rules, summaries every 50; budget reserve
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
353
harness/pipeline.py
Normal file
353
harness/pipeline.py
Normal file
@@ -0,0 +1,353 @@
|
||||
"""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 = "<!-- 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()
|
||||
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 <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()
|
||||
Reference in New Issue
Block a user