Files
BFM-decomp/tools/orchestrator.py
T
Drew T 8de18e617d feat(phase-21): T4/T5/T6 — worker validated (func_80164530 x134), grinder, orchestrator
- T4 worker Workflow validated: 3-target wave, all self-MATCH; whole-binary gate banked
  func_80164530 x134 (5 register-pins + offset-fold + scheduling). fleet 58.90->58.98%.
  §20 holds (1/3 survived the whole-binary gate); the 2 near/failed -> backlog.
- T5 grinder: token-free permuter of backlog near-misses -> gate_stage; gate correctly
  rejected a permuter over-prediction (integrity intact). auto_supervisor DRIVER-param.
- T6 orchestrator.py: ROI pool rotation (prep/finish/status); gate_stage flock serializes
  grinder+orchestrator builds. Sandbox: detached daemons build fine, foreground needs override.
2026-06-21 13:06:34 -06:00

111 lines
4.5 KiB
Python

#!/usr/bin/env python3
"""orchestrator.py — Phase 21 ROI orchestrator (the deterministic half of the worker loop).
The unattended loop is: the Claude orchestrator (/loop, Max) runs one cycle per tick —
orchestrator.py prep -> pick the current ROI pool, emit .run/auto/wave_batch.json
<orchestrator launches the worker Workflow over that batch (the ONE step only the model can do)>
orchestrator.py finish -> gate_stage the drafts (bank/propagate/log), record the close-rate,
rotate the pool when it's tapped, print a compact JSON summary
This file owns the ROI state + pool rotation; the model owns launching the Workflow. The grinder
(tools/grinder.py) runs alongside, token-free, draining the backlog the waves fill.
ROI rotation (balanced, ROI-gated — Drew): harvest a pool until its banked/drafts close-rate is
below --threshold for --patience consecutive waves, then advance: tractable -> giants -> o0 ->
capped -> (wrap to tractable). auto_stop.sh's STOP sentinel halts everything.
State: .run/auto/orch_state.json {pool, idx, low_streak, waves, banked_total, history:[...]}.
Usage:
orchestrator.py prep [--n 24] [--region main]
orchestrator.py finish --drafts .run/drafts-wave [--commit]
orchestrator.py status
"""
import argparse, json, os, subprocess, sys, time
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import gate_stage
REPO = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
PY = ".venv/bin/python"
STATE = os.path.join(REPO, ".run/auto/orch_state.json")
BATCH = ".run/auto/wave_batch.json"
POOLS = ["tractable", "giants", "o0", "capped"] # ROI rotation order
def load_state():
if os.path.exists(STATE):
return json.load(open(STATE))
return {"pool": POOLS[0], "idx": 0, "low_streak": 0, "waves": 0, "banked_total": 0, "history": []}
def save_state(s):
os.makedirs(os.path.dirname(STATE), exist_ok=True)
json.dump(s, open(STATE, "w"), indent=1)
def sh(cmd, timeout=None):
return subprocess.run(cmd, capture_output=True, text=True, cwd=REPO, timeout=timeout)
def cmd_prep(a):
s = load_state()
pool = s["pool"]
# refresh the manifest so 'cached'/stub status is current (cheap)
sh([PY, "tools/build_fuel_manifest.py"], timeout=120)
region = a.region
r = sh([PY, "tools/wave_targets.py", "--pool", pool, "--n", str(a.n),
"--region", region, "--out", BATCH], timeout=120)
n = 0
try:
n = len(json.load(open(os.path.join(REPO, BATCH))))
except Exception:
pass
print(json.dumps({"pool": pool, "n": n, "batch": BATCH, "wave": s["waves"] + 1}))
def cmd_finish(a):
s = load_state()
summary = gate_stage.run_gate(a.drafts, source_tag="worker", commit=a.commit)
banked, drafts = summary.get("banked", 0), summary.get("drafts", 0) or 1
close = banked / drafts
s["waves"] += 1
s["banked_total"] += banked
s["history"] = (s.get("history", []) + [{"pool": s["pool"], "banked": banked,
"drafts": summary.get("drafts", 0), "prop": summary.get("propagated", 0),
"fleet": summary.get("fleet_pct"), "ts": time.strftime("%Y-%m-%d %H:%M:%S")}])[-50:]
# ROI gate: rotate the pool after `patience` low-yield waves
if close < a.threshold:
s["low_streak"] = s.get("low_streak", 0) + 1
else:
s["low_streak"] = 0
rotated = False
if s["low_streak"] >= a.patience:
s["idx"] = (s.get("idx", 0) + 1) % len(POOLS)
s["pool"] = POOLS[s["idx"]]
s["low_streak"] = 0
rotated = True
save_state(s)
print(json.dumps({**summary, "close_rate": round(close, 3), "pool": s["history"][-1]["pool"],
"rotated_to": s["pool"] if rotated else None,
"waves": s["waves"], "banked_total": s["banked_total"]}))
def cmd_status(a):
s = load_state()
print(json.dumps(s, indent=1))
def main():
ap = argparse.ArgumentParser()
sub = ap.add_subparsers(dest="cmd", required=True)
p = sub.add_parser("prep"); p.add_argument("--n", type=int, default=24); p.add_argument("--region", default="main")
f = sub.add_parser("finish"); f.add_argument("--drafts", default=".run/drafts-wave")
f.add_argument("--commit", action="store_true"); f.add_argument("--threshold", type=float, default=0.15)
f.add_argument("--patience", type=int, default=2)
sub.add_parser("status")
a = ap.parse_args()
{"prep": cmd_prep, "finish": cmd_finish, "status": cmd_status}[a.cmd](a)
if __name__ == "__main__":
main()