diff --git a/harness/dashboard.py b/harness/dashboard.py new file mode 100644 index 0000000..a3576be --- /dev/null +++ b/harness/dashboard.py @@ -0,0 +1,408 @@ +"""Live status page of the stage 2 pipeline (static HTML, SAP Fiori colors, rewritten every 5 minutes). + + python3 -m harness.dashboard [--interval 300] [--once] write runs/dashboard/index.html in a loop + python3 -m harness.dashboard panel 41.07 record an Ollama panel value (usage) with the ledger now + +Open runs/dashboard/index.html in a browser; the page reloads itself every 5 minutes. Nothing here calls a model or +changes A4H: it reads runs/traj, tasks_gen/train, runs/ledger.jsonl and the pipeline log. No secrets are written. +""" +import argparse +import glob +import html +import json +import os +import subprocess +import sys +import time +import urllib.request + +from .adt_client import load_env +from .ledger import _env_budget, 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 + +OUT_DIR = os.path.join(ROOT, "runs", "dashboard") +OUT = os.path.join(OUT_DIR, "index.html") +HISTORY = os.path.join(OUT_DIR, "history.jsonl") +PANEL = os.path.join(OUT_DIR, "panel.json") +TRAJ = os.path.join(ROOT, "runs", "traj") +POOL = os.path.join(ROOT, "tasks_gen", "train") +PIPE = os.path.join(ROOT, "runs", "pipeline") +LEDGER = os.path.join(ROOT, "runs", "ledger.jsonl") +BASE_LEDGER = 66.9665 +PANEL_LIMIT = 60.0 +SEED_PANEL = [ # (usage on the Ollama panel, ledger at that time, label) + [24.63, 66.9665, "2026-10-03 20:15"], [30.98, 80.1045, "2026-10-05 14:46"], [31.48, 80.9683, "2026-10-05 14:57"], + [34.79, 85.4481, "2026-10-05 15:53"], [41.07, 93.4483, "2026-10-05 17:47"]] +CACHE = {} + + +def jl(path): + if not os.path.exists(path): + return [] + out = [] + for line in open(path): + try: + out.append(json.loads(line)) + except ValueError: + pass + return out + + +def pct(v, p): + return v[min(len(v) - 1, int(p * len(v)))] if v else None + + +# ---------------------------------------------------------------------------------------------------- collect +def evaluate(row): + key = row.get("run_dir") or "%s_%s" % (row["task"], row["attempt"]) + if key in CACHE: + return CACHE[key] + ev = {"accepted": False, "repair": False, "reasons": [], "hints": 0, "ctx": None, "ledger": 0.0, "mtime": 0} + path = os.path.join(TRAJ, row.get("run_dir") or "-", "record.json") + if os.path.exists(path): + rec = json.load(open(path)) + ok, why = acc.judge(rec, row, 80) + ev["accepted"], ev["reasons"] = ok, why + md = rec["metadata"] + ev["hints"] = md.get("syntax_hints") or 0 + ev["ledger"] = md.get("cost_ledger_usd") or 0.0 + tu = md.get("turn_usage") or [] + if tu: + ev["ctx"] = (tu[-1].get("prompt_tokens") or 0) + (tu[-1].get("completion_tokens") or 0) + if ok: + ev["repair"] = acc.has_repair(acc.trim_loops(rec["messages"])[0]) + ev["mtime"] = os.path.getmtime(path) + else: + ev["reasons"] = ["harness_error" if row.get("harness_error") else "no_record"] + CACHE[key] = ev + return ev + + +def task_meta(tid, cache={}): + if tid not in cache: + try: + t = json.load(open(os.path.join(POOL, tid, "task.json"))) + except OSError: + t = {} + cache[tid] = {"category": t.get("category") or "?", "object_type": t.get("object_type") or "?", + "difficulty": t.get("difficulty") or 2} + return cache[tid] + + +def sh(cmd): + return subprocess.run(cmd, capture_output=True, text=True).stdout + + +def ping(url, timeout=4): + try: + urllib.request.urlopen(url, timeout=timeout) + return True + except Exception as e: # noqa: BLE001 an HTTP error status still means the server answers + return hasattr(e, "code") + + +def collect(): + now = time.time() + rows = jl(os.path.join(TRAJ, "summary.jsonl")) + evs = [evaluate(r) for r in rows] + accepted = [e for e in evs if e["accepted"]] + s = spent() + limit_raw, reserve = _env_budget() + guard = limit_raw - reserve + led = jl(LEDGER) + since = [e for e in led if e.get("day", "") >= "2026-10-05"] + gen_cost = sum(e["usd"] for e in since if e.get("kind") == "generate") + run_cost = sum(e["usd"] for e in since if e.get("kind") == "run") + logs = [json.load(open(f)) for f in glob.glob(os.path.join(POOL, "_logs", "G*.json"))] + d = {"now": now, "ledger": s, "guard": guard, "limit": limit_raw, "reserve": reserve, + "gen_cost": gen_cost, "run_cost": run_cost, "rows": rows, "evs": evs, "logs": logs} + # panel calibration + panel = json.load(open(PANEL)) if os.path.exists(PANEL) else SEED_PANEL + d["panel"] = panel + last = panel[-1] + ratios = [(panel[i][1] - panel[i - 1][1]) / (panel[i][0] - panel[i - 1][0]) for i in range(1, len(panel)) + if panel[i][0] > panel[i - 1][0]] + d["ratio_recent"] = ratios[-1] if ratios else None + ratio = 1.2 # the guard ratio chosen on 2026-10-05 (pessimistic, trajectory runs) + d["panel_est"] = last[0] + (s - last[1]) / ratio + d["panel_left"] = PANEL_LIMIT - d["panel_est"] + # history and burn rate + hist = jl(HISTORY) + cur = {"t": now, "ledger": s, "runs": len(rows), "acc": len(accepted), "tasks": sum(1 for l in logs if l.get("accepted")), + "slots": len(logs)} + hist.append(cur) + d["hist"] = hist + ref = next((h for h in hist if h["t"] >= now - 3600), hist[0]) + dt = max(now - ref["t"], 1) / 3600.0 + d["burn_ledger_h"] = (s - ref["ledger"]) / dt if dt > 0.2 else None + d["eta_h"] = (guard - s) / d["burn_ledger_h"] if d["burn_ledger_h"] and d["burn_ledger_h"] > 0.05 else None + d["rate_runs_h"] = (len(rows) - ref["runs"]) / dt if dt > 0.2 else None + d["rate_acc_h"] = (len(accepted) - ref["acc"]) / dt if dt > 0.2 else None + # stop rules + last30 = evs[-30:] + d["acc30"] = sum(e["accepted"] for e in last30) / len(last30) if len(last30) >= 30 else None + pl = "" + p = os.path.join(PIPE, "pipeline.log") + if os.path.exists(p): + pl = open(p, errors="replace").read() + d["log_tail"] = [l[:240] for l in pl.splitlines()[-12:]] + d["workers"] = 1 if "workers 2 -> 1" in pl else 2 + d["events"] = [l[:200] for l in pl.splitlines() if "A4H/harness event" in l][-6:] + d["stopped"] = open(os.path.join(PIPE, "STOPPED.txt")).read() if os.path.exists(os.path.join(PIPE, "STOPPED.txt")) else None + d["stop_flag"] = os.path.exists(os.path.join(PIPE, "STOP")) + # processes and systems + procs = sh(["pgrep", "-fl", r"harness\.(pipeline|trainset|trajectories|dashboard)"]).splitlines() + d["procs"] = [l.split(None, 1)[1].replace("/Library/Developer/CommandLineTools/Library/Frameworks/Python3.framework/Versions/3.9/Resources/Python.app/Contents/MacOS/Python", "python")[:90] + for l in procs if "pgrep" not in l] + d["a4h"] = ping("http://localhost:50000/sap/public/ping") + d["mcp"] = ping("http://127.0.0.1:3000/mcp") + return d + + +# ---------------------------------------------------------------------------------------------------- render +E = html.escape +CSS = """ +:root{--shell:#354A5F;--shell2:#1D2D3E;--blue:#0A6ED1;--blue2:#0854A0;--bg:#F5F6F7;--card:#fff;--line:#D9D9D9; +--text:#32363A;--muted:#6A6D70;--good:#107E3E;--warn:#E9730C;--bad:#BB0000;--goodbg:#F1FDF6;--warnbg:#FEF7F1;--badbg:#FFF4F4;--accent:#5899DA} +*{box-sizing:border-box}body{margin:0;background:var(--bg);color:var(--text);font:14px/1.45 -apple-system,"Segoe UI","72",Arial,sans-serif} +header{background:linear-gradient(90deg,var(--shell2),var(--shell));color:#fff;padding:14px 24px;display:flex;align-items:center;gap:16px;flex-wrap:wrap} +header h1{font-size:18px;font-weight:600;margin:0}header .sub{opacity:.75;font-size:12px} +.logo{background:var(--blue);padding:4px 10px;border-radius:3px;font-weight:700;letter-spacing:1px} +.chip{display:inline-block;padding:2px 10px;border-radius:12px;font-size:12px;font-weight:600;border:1px solid} +.ok{color:var(--good);background:var(--goodbg);border-color:var(--good)}.wa{color:var(--warn);background:var(--warnbg);border-color:var(--warn)} +.er{color:var(--bad);background:var(--badbg);border-color:var(--bad)}.gr{color:var(--muted);background:#fff;border-color:var(--line)} +header .chip{background:transparent;color:#fff;border-color:#ffffff88}header .chip.ok{border-color:#7fd69f;color:#b9f0cd} +main{padding:20px 24px;max-width:1500px;margin:auto} +.tiles{display:grid;grid-template-columns:repeat(auto-fit,minmax(210px,1fr));gap:14px;margin-bottom:18px} +.tile{background:var(--card);border:1px solid var(--line);border-radius:6px;padding:14px 16px;border-top:3px solid var(--blue)} +.tile.good{border-top-color:var(--good)}.tile.warn{border-top-color:var(--warn)}.tile.bad{border-top-color:var(--bad)} +.tile .l{font-size:12px;color:var(--muted);text-transform:uppercase;letter-spacing:.5px}.tile .v{font-size:28px;font-weight:300;margin:2px 0} +.tile .s{font-size:12px;color:var(--muted)} +.grid{display:grid;grid-template-columns:repeat(auto-fit,minmax(460px,1fr));gap:14px} +.card{background:var(--card);border:1px solid var(--line);border-radius:6px;padding:0 0 8px}.card h2{font-size:14px;margin:0;padding:11px 16px;border-bottom:1px solid var(--line);font-weight:600;color:var(--shell)} +.card .b{padding:8px 16px}table{border-collapse:collapse;width:100%;font-size:13px}th{text-align:left;color:var(--muted);font-weight:600;font-size:12px;padding:5px 6px;border-bottom:1px solid var(--line)} +td{padding:5px 6px;border-bottom:1px solid #eee;vertical-align:middle}td.n,th.n{text-align:right;font-variant-numeric:tabular-nums} +.bar{background:#E5EEF8;border-radius:3px;height:10px;min-width:80px;position:relative}.bar i{display:block;height:100%;background:var(--blue);border-radius:3px} +.bar.g i{background:var(--good)}.bar.w i{background:var(--warn)}.bar.r i{background:var(--bad)} +.mono{font-family:ui-monospace,Menlo,monospace;font-size:12px}.log{background:#fafafa;border:1px solid var(--line);border-radius:4px;padding:8px;max-height:230px;overflow:auto;white-space:pre-wrap;word-break:break-all} +.note{font-size:12px;color:var(--muted);padding:2px 16px}.banner{padding:10px 16px;border-radius:6px;margin-bottom:14px;font-weight:600} +svg text{font-size:10px;fill:#6A6D70}footer{color:var(--muted);font-size:12px;text-align:center;padding:16px} +""" + + +def spark(points, w=420, h=70, color="#0A6ED1", fmt="%.1f", guard=None): + pts = [(t, v) for t, v in points if v is not None] + if len(pts) < 2: + return "
not enough history yet
" + t0, t1 = pts[0][0], pts[-1][0] + lo = min(v for _, v in pts) + hi = max(max(v for _, v in pts), guard or 0) + span = max(hi - lo, 1e-9) + xs = lambda t: 4 + (w - 8) * (t - t0) / max(t1 - t0, 1) # noqa: E731 + ys = lambda v: h - 14 - (h - 24) * (v - lo) / span # noqa: E731 + path = " ".join("%s%.1f,%.1f" % ("M" if i == 0 else "L", xs(t), ys(v)) for i, (t, v) in enumerate(pts)) + g = "" + if guard is not None: + g = ('' + 'guard %s' % (w - 4, ys(guard), ys(guard), ys(guard) - 3, fmt % guard)) + return ('%s' + '%s%s' + % (w, h, h, path, color, g, h - 2, time.strftime("%H:%M", time.localtime(t0)), w - 4, h - 2, + time.strftime("%H:%M", time.localtime(t1)) + " now " + fmt % pts[-1][1])) + + +def bar(n, total, cls=""): + p = 100.0 * n / total if total else 0 + return '
' % (cls, p) + + +def tile(label, value, sub="", cls=""): + return '
%s
%s
%s
' % (cls, E(label), value, sub) + + +def agg_table(title, evs, rows, key, note=""): + agg = {} + for r, e in zip(rows, evs): + k = task_meta(r["task"])[key] + a = agg.setdefault(k, [0, 0, 0]) + a[0] += 1 + a[1] += e["accepted"] + a[2] += e["accepted"] and e["repair"] + trs = "" + for k, (n, ok, rep) in sorted(agg.items()): + rate = ok / n if n else 0 + cls = "g" if rate >= 0.8 else ("w" if rate >= 0.5 else "r") + trs += "%s%d%d%d%s%.0f%%" % ( + E(str(k)), n, ok, rep, bar(ok, n, cls), 100 * rate) + return ("

%s

" + "%s
%srunsacceptedrepairacceptance%%
%s
" + % (E(title), E(key.replace("_", " ")), trs, "
%s
" % E(note) if note else "")) + + +def gen_table(logs, key, title): + agg = {} + for l in logs: + if key == "object_type": + k = l.get("object_type") or "?" + else: + k = l.get("category") or "?" + a = agg.setdefault(k, [0, 0]) + a[0] += 1 + a[1] += bool(l.get("accepted")) + trs = "".join("%s%d%d%s%.0f%%" + % (E(str(k)), n, ok, bar(ok, n, "g" if ok / n >= 0.8 else "w"), 100.0 * ok / n) for k, (n, ok) in sorted(agg.items())) + return ("

%s

%s
%sslotsacceptedrate%%
" + % (E(title), E(key.replace("_", " ")), trs)) + + +def render(d): + rows, evs, logs = d["rows"], d["evs"], d["logs"] + acc = [e for e in evs if e["accepted"]] + n_acc, n_runs = len(acc), len(rows) + repairs = sum(e["repair"] for e in acc) + tasks_acc = sum(1 for l in logs if l.get("accepted")) + k_n = sum(1 for l in logs if l.get("category") == "K") + err_n = sum(1 for l in logs if l.get("error_kind")) + hard_n = sum(1 for l in logs if (l.get("difficulty") or 2) == 3) + tasks_with_run = {r["task"] for r in rows} + backlog = sum(1 for l in logs if l.get("accepted") and l["id"] not in tasks_with_run) + alive = any("harness.pipeline" in p for p in d["procs"]) + status, scls = ("RUNNING", "ok") if alive and not d["stopped"] else (("STOPPED", "er") if d["stopped"] else ("NOT RUNNING", "er")) + hints = sum(e["hints"] for e in evs) + ctx = sorted(e["ctx"] for e in acc if e["ctx"]) + next50 = (n_acc // 50 + 1) * 50 + est = d["panel_est"] + pcls = "good" if est < 50 else ("warn" if est < 57 else "bad") + eta = ("%.1f h (%s)" % (d["eta_h"], time.strftime("%a %H:%M", time.localtime(d["now"] + d["eta_h"] * 3600)))) if d["eta_h"] else "n/a" + a30 = d["acc30"] + a30s = "%.0f %%" % (100 * a30) if a30 is not None else "n/a (<30 runs)" + cost_per = d["run_cost"] / n_acc if n_acc else 0 + banner = "" + if d["stopped"]: + banner = "" % E(d["stopped"].replace("\n", " | ")) + elif not alive: + banner = "" + elif d["stop_flag"]: + banner = "" + tiles = "".join([ + tile("Accepted trajectories", "%d" % n_acc, "next summary at %d (%d to go)" % (next50, next50 - n_acc) + bar(n_acc % 50, 50), "good"), + tile("Acceptance (all runs)", "%.0f %%" % (100.0 * n_acc / n_runs if n_runs else 0), "%d of %d runs; last 30: %s (stop below 50 %%)" % (n_acc, n_runs, a30s), + "good" if not n_runs or n_acc / n_runs >= 0.7 else "warn"), + tile("Repair share", "%.0f %%" % (100.0 * repairs / n_acc if n_acc else 0), "%d of %d accepted contain error + fix" % (repairs, n_acc)), + tile("Training tasks", "%d" % tasks_acc, "%d slots tried, %d K, %d error-targeted, %d hard; %d wait for a first run" % (len(logs), k_n, err_n, hard_n, backlog)), + tile("Ollama usage (estimated)", "$%.1f" % est, "of $60; ~$%.1f left; last panel $%.2f" % (d["panel_left"], d["panel"][-1][0]), pcls), + tile("Ledger / guard", "%.1f" % d["ledger"], "guard %.1f (limit %.0f - reserve %.0f); %.1f left" % (d["guard"], d["limit"], d["reserve"], d["guard"] - d["ledger"]), + "good" if d["guard"] - d["ledger"] > 10 else "warn"), + tile("Time to guard", eta, "burn %s ledger/h" % ("%.1f" % d["burn_ledger_h"] if d["burn_ledger_h"] else "n/a"), "warn" if d["eta_h"] and d["eta_h"] < 3 else ""), + tile("Speed", "%s runs/h" % ("%.1f" % d["rate_runs_h"] if d["rate_runs_h"] is not None else "n/a"), + "%s accepted/h; %d trajectory workers" % ("%.1f" % d["rate_acc_h"] if d["rate_acc_h"] is not None else "n/a", d["workers"])), + tile("Syntax hints", "%d" % hints, "proxy syntaxCheck added (local abaplint)"), + tile("Cost per accepted", "%.2f" % cost_per, "ledger USD (trajectory runs only); generation %.1f, runs %.1f" % (d["gen_cost"], d["run_cost"])), + ]) + # budget calibration + pr = d["panel"] + prows = "" + for i, (u, l, lab) in enumerate(pr): + r = "" + if i and u > pr[i - 1][0]: + r = "%.2f" % ((l - pr[i - 1][1]) / (u - pr[i - 1][0])) + prows += "%s$%.2f%.1f%s" % (E(lab), u, l, r) + budget = ("

Budget and calibration

" + "%s
panel atusageledgerledger per usage
The ledger is a list-price upper bound; real usage per ledger " + "USD varies with the work (generation about 2, trajectory runs about 1.2-1.4). Guard ratio 1.2. Add a value: " + "python3 -m harness.dashboard panel 41.07
" + "
%s
") % (prows, spark([(h["t"], h["ledger"]) for h in d["hist"]], guard=d["guard"], color="#E9730C")) + prog = ("

Progress over time

accepted trajectories
%s" + "
trajectory runs
%s
accepted training tasks
%s
" + % (spark([(h["t"], h["acc"]) for h in d["hist"]], fmt="%d", color="#107E3E"), + spark([(h["t"], h["runs"]) for h in d["hist"]], fmt="%d"), + spark([(h["t"], h["tasks"]) for h in d["hist"]], fmt="%d", color="#5899DA"))) + # recent runs + rec = "" + for r, e in list(zip(rows, evs))[-14:][::-1]: + m = task_meta(r["task"]) + st = "accepted" if e["accepted"] else "%s" % E((", ".join(e["reasons"]) or "rejected")[:40]) + rec += ("%s%s%s%s%s%s%s%s%.2f" + % (E(r["task"]), E(m["category"]), E(m["object_type"]), r["attempt"], r.get("score"), E(str(r.get("end_reason"))), + r.get("tool_calls"), st + (" repair" if e["repair"] else ""), e["ledger"])) + recent = ("

Latest runs

" + "%s
taskcattypeattemptscoreendcallsresultledger
" % rec) + tok = ("

Sample size (context tokens)

" + "
p50p90p95maxn
%s%s%s%s%d
Last turn prompt + completion (DeepSeek " + "tokenizer, includes the 20 tool schemas). The Qwen token counts in train/STATE.md are about the same.
" + % (pct(ctx, .5), pct(ctx, .9), pct(ctx, .95), ctx[-1] if ctx else None, len(ctx))) + # stop rules + ev_n = len(d["events"]) + rules = [("Acceptance, last 30 runs", a30s, "stop below 50 %", "er" if a30 is not None and a30 < 0.5 else "ok"), + ("Harness error patterns", "%d seen" % ev_n, "stop at 3 equal; first one: 2 -> 1 worker", "ok" if not ev_n else "wa"), + ("Budget guard", "%.1f / %.1f" % (d["ledger"], d["guard"]), "stop at the guard", "ok" if d["ledger"] < d["guard"] - 5 else "wa"), + ("Generation deadline", "2026-10-10 18:00", "no new slot after it", "gr"), + ("Trajectory deadline", "2026-10-11 23:30", "stop and report", "gr")] + rl = "".join("%s%s%s%s" % (E(a), E(b), E(c), cl, {"ok": "ok", "wa": "watch", "er": "STOP", "gr": "set"}[cl]) for a, b, c, cl in rules) + stops = "

Stop rules

%s
rulenowlimit
%s
" % ( + rl, "
%s
" % E(" | ".join(d["events"])) if d["events"] else "") + sysc = ("

Processes and systems

" + "
A4H (HTTP 50000)%s
MCP server (3000)%s
Pipeline%s
" + "
%s
" + % ("ok" if d["a4h"] else "er", "up" if d["a4h"] else "down", "ok" if d["mcp"] else "er", "up" if d["mcp"] else "down", scls, status, + E("\n".join(d["procs"])))) + logc = "

Pipeline log

%s
" % E("\n".join(d["log_tail"])) + body = (banner + "
" + tiles + "
" + budget + prog + stops + sysc + + agg_table("Trajectories by category", evs, rows, "category") + agg_table("Trajectories by object type", evs, rows, "object_type") + + gen_table(logs, "category", "Task generation by category") + gen_table(logs, "object_type", "Task generation by object type") + + tok + recent + logc + "
") + stamp = time.strftime("%Y-%m-%d %H:%M:%S", time.localtime(d["now"])) + return ("" + "ABAP LLM Stage 2" + "

Stage 2 data pipeline

teacher DeepSeek V4.1 Flash, Qwen 3.8 27B training data, A4H
" + "%supdated %s · next in 5:00
%s
" + "" + "" + "" % (CSS, scls, status, stamp, body)) + + +# ---------------------------------------------------------------------------------------------------- main +def write_once(): + os.makedirs(OUT_DIR, exist_ok=True) + d = collect() + page = render(d) + tmp = OUT + ".tmp" + open(tmp, "w").write(page) + os.replace(tmp, OUT) + last = d["hist"][-1] + with open(HISTORY, "a") as f: + f.write(json.dumps(last) + "\n") + return d + + +def main(): + load_env(os.path.join(ROOT, ".env")) + ap = argparse.ArgumentParser() + ap.add_argument("cmd", nargs="?", default="run", choices=["run", "panel"]) + ap.add_argument("value", nargs="?", type=float) + ap.add_argument("--interval", type=int, default=300) + ap.add_argument("--once", action="store_true") + a = ap.parse_args() + if a.cmd == "panel": + os.makedirs(OUT_DIR, exist_ok=True) + pts = json.load(open(PANEL)) if os.path.exists(PANEL) else SEED_PANEL + pts.append([a.value, round(spent(), 4), time.strftime("%Y-%m-%d %H:%M")]) + json.dump(pts, open(PANEL, "w")) + print("panel value recorded:", pts[-1]) + return + while True: + try: + d = write_once() + print(time.strftime("%F %T"), "updated: runs", len(d["rows"]), "ledger", round(d["ledger"], 2), flush=True) + except Exception as e: # noqa: BLE001 keep the loop alive; the next cycle tries again + print(time.strftime("%F %T"), "update failed:", repr(e)[:300], flush=True) + if a.once: + return + time.sleep(a.interval) + + +if __name__ == "__main__": + main()