feat(phase-21): T4/T5 tooling — worker Workflow + wave selector + permuter grinder

- tools/workflows/worker_wave.js: parallel toolkit-aware drafter agents (xHigh) per target,
  each writes its best matching C to the drafts dir (gate run by the orchestrator after).
  args parsed robustly (harness may serialize to a JSON string — §20 gotcha).
- tools/wave_targets.py: select a wave batch from the fuel manifest (pool/region/N, skips
  backlog walls). tools/grinder.py: token-free permuter grinder — permute backlog near-misses
  -> gate_stage -> bank; STOP/heartbeat, idles for new worker near-misses, supervisor-friendly.
This commit is contained in:
Drew T
2026-06-21 12:41:14 -06:00
parent c2e93e50cd
commit 2d7926e9f6
3 changed files with 307 additions and 0 deletions
+132
View File
@@ -0,0 +1,132 @@
#!/usr/bin/env python3
"""grinder.py — Phase 21 token-free permuter grinder (the CPU-bound worker).
Pulls the closest near-misses from the backlog, runs decomp-permuter on each (the permuter
closes regalloc/scheduling gaps — exactly the residual the worker's drafts leave), banks the
true byte-matches through the shared gate_stage (the sole arbiter, G3/P9), and re-logs any
improved-but-still-near draft. LLM-FREE — runs unattended for days under auto_supervisor.sh,
alongside the token-heavy worker waves (CPU budget vs token budget).
backlog near-miss -> decomp-permuter (output-0-* = true byte-match) -> gate_stage -> bank x reach
no winner -> re-log (closeness may improve) and move on
SAFE EXIT: touch .run/auto/STOP (tools/auto_stop.sh) — finishes the current permute+gate, exits 0.
HEARTBEAT: .run/auto/grinder_heartbeat.json — {ts, state, current, banked, fleet_pct}.
Usage: grinder.py [--permute-secs 120] [-j 14] [--batch 10] [--max-nins 220]
[--max-closeness 30] [--attempts 2] [--idle-secs 90] [--once]
"""
import argparse, glob, json, os, shutil, subprocess, sys, time
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import p16_permute, gate_stage, backlog
REPO = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
AUTODIR = ".run/auto"
STOP = f"{AUTODIR}/STOP"
HB = f"{AUTODIR}/grinder_heartbeat.json"
DRAFTS = f"{AUTODIR}/grinder_drafts"
def stop_requested():
return os.path.exists(os.path.join(REPO, STOP))
def log(m):
print(f"[{time.strftime('%H:%M:%S')}] grinder: {m}", flush=True)
def heartbeat(state, current=None, banked=0, fp=None):
os.makedirs(os.path.join(REPO, AUTODIR), exist_ok=True)
json.dump({"ts": time.strftime("%Y-%m-%d %H:%M:%S"), "state": state, "current": current,
"banked": banked, "fleet_pct": fp},
open(os.path.join(REPO, HB), "w"), indent=1)
def candidates(max_nins, max_close, tried, attempts):
"""closest still-open near-misses with a saved best draft (permuter-amenable), least-tried first."""
out = []
for r in backlog.load_best():
nm = r.get("name")
if r.get("status") != "near" or not r.get("best_draft") or not nm:
continue
if tried.get(nm, 0) >= attempts:
continue
c = r.get("closeness")
if c is None or c > max_close: # permuter closes small regalloc/sched gaps, not large rewrites
continue
if (r.get("nins") or 999) > max_nins:
continue
if not os.path.exists(os.path.join(REPO, r["best_draft"])):
continue
out.append(r)
out.sort(key=lambda r: (tried.get(r["name"], 0), r.get("closeness") or 999))
return out
def main():
ap = argparse.ArgumentParser()
ap.add_argument("--permute-secs", type=int, default=120)
ap.add_argument("-j", type=int, default=14)
ap.add_argument("--batch", type=int, default=10)
ap.add_argument("--max-nins", type=int, default=220)
ap.add_argument("--max-closeness", type=int, default=30)
ap.add_argument("--attempts", type=int, default=2)
ap.add_argument("--idle-secs", type=int, default=90)
ap.add_argument("--once", action="store_true")
a = ap.parse_args()
os.chdir(REPO)
os.makedirs(AUTODIR, exist_ok=True)
if stop_requested():
log("STOP present at startup; remove it to run."); return
tried, banked, fp = {}, 0, None
log(f"start (permute={a.permute_secs}s -j{a.j} batch={a.batch} max_close={a.max_closeness})")
while True:
if stop_requested():
log("STOP — clean exit."); heartbeat("stopped", None, banked, fp); return
cand = candidates(a.max_nins, a.max_closeness, tried, a.attempts)[:a.batch]
if not cand:
if a.once:
log("no candidates (once) — exit."); heartbeat("dry", None, banked, fp); return
heartbeat("idle", None, banked, fp)
log(f"no untried candidates; idle {a.idle_secs}s (worker may add more)")
for _ in range(a.idle_secs):
if stop_requested():
break
time.sleep(1)
if stop_requested():
continue
tried.clear() # let the stochastic permuter re-try after idle
continue
if os.path.exists(os.path.join(REPO, DRAFTS)):
shutil.rmtree(os.path.join(REPO, DRAFTS))
os.makedirs(os.path.join(REPO, DRAFTS), exist_ok=True)
won = 0
for r in cand:
if stop_requested():
break
fn = r["name"]; tried[fn] = tried.get(fn, 0) + 1
heartbeat("permuting", fn, banked, fp)
try:
draft = open(os.path.join(REPO, r["best_draft"])).read()
pd = p16_permute.setup(fn, draft)
if not pd:
continue
win = p16_permute.run_permuter(pd, a.permute_secs, a.j)
if win:
open(os.path.join(REPO, DRAFTS, fn + ".c"), "w").write(
p16_permute.winner_to_draft(open(win).read()))
won += 1; log(f"permuter WON {fn} (close was {r.get('closeness')})")
except Exception as e:
log(f"{fn}: {e}")
if won:
heartbeat("gating", None, banked, fp)
s = gate_stage.run_gate(DRAFTS, source_tag="grinder", commit=True)
banked += s.get("banked", 0); fp = s.get("fleet_pct", fp)
log(f"gate: banked {s.get('banked')} (+{s.get('propagated')} prop); total {banked}; fleet {fp}%")
heartbeat("running", None, banked, fp)
if a.once:
log(f"once done — banked {banked}."); heartbeat("done", None, banked, fp); return
if __name__ == "__main__":
main()
+93
View File
@@ -0,0 +1,93 @@
#!/usr/bin/env python3
"""wave_targets.py — select a worker-wave target batch from the fuel manifest.
Emits a JSON array of {name, addr, nins, class, asm, ghidra_c} for tools/workflows/worker_wave.js
(passed as args.targets). Filters to still-OPEN (INCLUDE_ASM) cached targets in the chosen pool,
ranked by leverage (reach*nins), skipping known walls already logged 'failed'/'stub' in the backlog.
Pools (ROI rotation): tractable (reach-134 WAVE/PINS/STRUCT <=150 ins, main region) | giants |
o0 | capped | any-reach134.
Usage: tools/wave_targets.py --pool tractable --n 24 [--region main|a|any] [--out -]
"""
import argparse, glob, json, os, re, sys
REPO = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
STUB_RE = re.compile(r"INCLUDE_ASM\([^,]+,\s*(\w+)\)")
ASM_SUBDIR = "asm/ov_SC01_077/nonmatchings/ov_SC01_077"
def live_stubs():
s = set()
for p in glob.glob(os.path.join(REPO, "src/ov_SC01_077/ov_SC01_077*.c")):
s |= set(STUB_RE.findall(open(p).read()))
return s
def backlog_walls():
"""names already logged as 'failed'/'stub' — skip (don't waste agents re-drafting known walls)."""
walls = set()
p = os.path.join(REPO, ".run/backlog.jsonl")
if os.path.exists(p):
for line in open(p):
line = line.strip()
if not line:
continue
r = json.loads(line)
if r.get("status") in ("failed", "stub") and r.get("name"):
walls.add(r["name"])
return walls
def main():
ap = argparse.ArgumentParser()
ap.add_argument("--pool", default="tractable",
choices=["tractable", "giants", "o0", "capped", "any-reach134"])
ap.add_argument("--n", type=int, default=24)
ap.add_argument("--region", default="main", choices=["main", "a", "any"])
ap.add_argument("--max-nins", type=int, default=150)
ap.add_argument("--include-walls", action="store_true", help="don't skip backlog failed/stub")
ap.add_argument("--out", default="-")
a = ap.parse_args()
m = json.load(open(os.path.join(REPO, ".run/fuel_manifest.json")))
stubs = live_stubs()
walls = set() if a.include_walls else backlog_walls()
def ok(t):
if t["name"] not in stubs or not t["cached"]:
return False
if t["name"] in walls:
return False
if a.region != "any" and t["region"] != a.region:
return False
if a.pool == "tractable":
return t["reach_134"] and t["class"] in ("WAVE", "PINS", "STRUCT") and (t["nins"] or 999) <= a.max_nins
if a.pool == "giants":
return t["class"] == "GIANT"
if a.pool == "o0":
return t["class"] == "O0"
if a.pool == "any-reach134":
return t["reach_134"]
return False
pool = [t for t in m["targets"] if ok(t)]
if a.pool == "capped": # the matched-but-local recovery set (not stubs)
pool = [{"name": n, "addr": "0x" + n[5:].lower(), "nins": None, "class": "CAPPED",
"reach": 134, "leverage": 0} for n in m.get("capped_recovery", [])]
pool.sort(key=lambda t: t.get("leverage") or 0, reverse=True)
pool = pool[:a.n]
batch = [{"name": t["name"], "addr": t["addr"], "nins": t["nins"], "class": t["class"],
"asm": f"{ASM_SUBDIR}/{t['name']}.s", "ghidra_c": f".run/ghidra_c/{t['name']}.c"}
for t in pool]
out = json.dumps(batch, indent=0)
if a.out == "-":
sys.stdout.write(out + "\n")
else:
open(os.path.join(REPO, a.out), "w").write(out)
print(f"{len(batch)} targets -> {a.out} (pool={a.pool} region={a.region})", file=sys.stderr)
if __name__ == "__main__":
main()
+82
View File
@@ -0,0 +1,82 @@
export const meta = {
name: 'worker-wave',
description: 'Phase-21 worker wave: fan out toolkit-aware drafter agents over a target batch; each writes its best matching C to the drafts dir. The orchestrator runs gate_stage on the dir afterward (byte-gate + propagate + backlog).',
phases: [
{ title: 'Draft', detail: 'one xHigh drafter agent per target' },
],
}
// args = { targets: [{name, addr, nins, class, asm, ghidra_c}], draftDir }
// Returns { draftDir, n, drafted: [{fn,status,closeness,klass,where_stuck}] }.
// NOTE: drafting only. The byte-gate / propagate / backlog (gate_stage.py) is run by the
// ORCHESTRATOR after this returns — it mutates+commits the byte-locked tree and must be
// serial (one wave at a time), and its build loop can exceed an agent's Bash timeout.
const DRAFT_SCHEMA = {
type: 'object', additionalProperties: false,
required: ['fn', 'status'],
properties: {
fn: { type: 'string' },
status: { type: 'string', enum: ['match', 'near', 'fail'] },
closeness: { type: 'integer', description: 'final match_one mismatch count; 0 if MATCH' },
klass: { type: 'string', description: 'residual class: regalloc-order|schedule|struct|loose-typing|plumbing|other' },
where_stuck: { type: 'string', description: 'one line: what is still off (for the backlog)' },
},
}
const ASM_SUBDIR = 'asm/ov_SC01_077/nonmatchings/ov_SC01_077'
function drafterPrompt(t, draftDir) {
return `Match ONE MIPS function for the Brave Fencer Musashi PS1 matching decompilation (overlay ov_SC01_077).
GOAL: write C that the pinned compiler (gcc-2.7.2-psx -O2 -G0 -mips1 -mcpu=3000 -mgas -msoft-float -fgnu-linker + maspsx --aspsx-version=2.56 --expand-div) compiles to BYTE-IDENTICAL machine code.
TARGET: ${t.name} @ ${t.addr} — ${t.nins} instructions, class hint "${t.class}".
- Target asm (the ground truth): ${t.asm}
(each line "/* off vaddr w0 w1 */ mnemonic ..." shows the exact encoded instructions.)
- Ghidra-C reference (types/locals/callee names — NOT byte-accurate, a scaffold): ${t.ghidra_c}
THE TOOLKIT (docs/matching-cookbook.md §17–§20 — read those sections for depth; the high-leverage moves):
- Register-allocation ORDER: if call-crossing locals land in the wrong saved reg vs the target,
PIN them: \`register s32 v __asm__("$16");\` ($16=$s0,$17=$s1,$18=$s2,…). Add a scheduling barrier
(a dummy volatile read or reordering) if the schedule is off. THIS is the highest-reach lever.
- Array-of-struct %lo-fold (§18): for indexed global access, declare \`extern Struct base[];\`
(sizeof(Struct)==stride) and write \`base[i].field\` — folds %lo into the load/store. Do NOT write
\`*(T*)(&sym + i*stride)\` (that materializes &sym and adds an instruction).
- For-loop vs do-while (§17): \`for(init;cond;upd)\` schedules the back-branch into the delay slot
differently than a do-while; pick the loop form the target's branch layout implies.
- Statement / for-update order: independent statements emit in source order — reorder to match.
- Declaration plumbing: don't fret callee extern types — the gate's canon_resident_calls +
cast_call_sites + sig_unify fix most extern/arity mismatches. Focus on the BODY codegen.
PROCESS (you have Bash + Read):
1. Read the target asm and the Ghidra-C.
2. Write your best C (the function definition + any externs it needs) to: ${draftDir}/${t.name}.c
3. Self-check (fast relocation-masked proxy for the byte-gate):
.venv/bin/python tools/match_one.py ${t.name} --c ${draftDir}/${t.name}.c --asm-subdir ${ASM_SUBDIR}
- "MATCH (N ins)" => byte-identical (relocation-masked). You nailed it. Stop.
- "N mismatched" => N instructions differ. Apply the toolkit, iterate to reduce N.
4. Iterate a few times; KEEP THE BEST draft in the file (always leave a file, even if imperfect —
the whole-binary gate + the permuter grinder may finish it; a close near-miss is logged for a human).
5. Return: { fn, status (match|near|fail), closeness (final mismatch count), klass, where_stuck }.
The whole-binary byte-gate (run later) is the sole arbiter — match_one is a proxy, but a MATCH there
almost always banks. Be rigorous; never fabricate a MATCH you didn't observe.`
}
// args may arrive as a JSON-encoded string (harness serialization) — parse robustly (§20 gotcha).
const A = typeof args === 'string' ? JSON.parse(args) : (args || {})
const targets = A.targets || []
const draftDir = A.draftDir || '.run/drafts-wave'
if (!targets.length) { log('worker-wave: no targets'); return { draftDir, n: 0, drafted: [] } }
log(`worker-wave: drafting ${targets.length} targets -> ${draftDir} (xHigh agents)`)
phase('Draft')
const drafted = (await parallel(targets.map(t => () =>
agent(drafterPrompt(t, draftDir), { label: `draft:${t.name}`, phase: 'Draft', schema: DRAFT_SCHEMA, effort: 'xhigh' })
.then(r => (r ? { ...r, fn: r.fn || t.name } : null))
))).filter(Boolean)
const matched = drafted.filter(d => d.status === 'match').length
const near = drafted.filter(d => d.status === 'near').length
log(`worker-wave: ${drafted.length}/${targets.length} drafted; self-assessed match=${matched} near=${near}`)
return { draftDir, n: drafted.length, matched, near, drafted }