Series A: local Qwen on the MacBook (remote server), 12 h window with hard stop, clean pause on server loss, dashboard card, docs/remote-model.md
Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
320
harness/localqwen.py
Normal file
320
harness/localqwen.py
Normal file
@@ -0,0 +1,320 @@
|
||||
"""Series A: the local Qwen 3.8 27B (4-bit, served on the MacBook) on training tasks of the new kinds (2026-10-05).
|
||||
|
||||
Harness, A4H, MCP server and records stay on the Mac mini; the MacBook only serves the model (docs/remote-model.md).
|
||||
|
||||
python3 -m harness.localqwen --base-url http://<MacBook LAN IP>:8080/v1 [--window-hours 12] [--plan-only]
|
||||
|
||||
Order: 3 tasks each of INTF, TABL, STRU, MSAG, exception, then 5 DDLS; one task at a time. Settings as the official
|
||||
baseline (docs/stage1-baseline.md): thinking off, max_tokens 16384, loop guard 3, tool budget 60, temperature 0.2, same
|
||||
system prompt. Time window: 12 hours from the first run, hard stop at the end also in the middle of a run (the run is
|
||||
not scored, its objects are deleted on A4H). If the server does not answer, the run ends cleanly, the series pauses and
|
||||
resumes when it answers again; such a run is not a result and not a failure.
|
||||
Output (runs/local_qwen/): state.json (dashboard), runs/, accepted.jsonl (accepted trajectories, NOT for the first SFT),
|
||||
results.jsonl, summary.json. Resumable: call again with the same arguments.
|
||||
"""
|
||||
import argparse
|
||||
import json
|
||||
import os
|
||||
import shutil
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
import urllib.request
|
||||
|
||||
from .adt_client import load_env
|
||||
from .agents import LlmAgent
|
||||
from . import mix
|
||||
from .runner import Runner
|
||||
|
||||
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
|
||||
|
||||
POOL = os.path.join(ROOT, "tasks_gen", "train")
|
||||
OUT = os.path.join(ROOT, "runs", "local_qwen")
|
||||
STATE = os.path.join(OUT, "state.json")
|
||||
RUN_BASE = 410000 # + task index * 10 + restart number (a digit must lead the 4-char base36 run: < 466560)
|
||||
EXPECTED_MODEL = "Qwen3.8-27B-4bit" # basename of the served model path (same build as the official baseline)
|
||||
PLAN_ORDER = [("INTF", 3), ("TABL", 3), ("STRU", 3), ("MSAG", 3), ("EXC", 3), ("DDLS", 5)]
|
||||
SETTINGS = {"max_tokens": 16384, "enable_thinking": False, "temperature": 0.2, "tool_call_budget": 60, "loop_guard": 3,
|
||||
"max_turns": 80, "one_task_at_a_time": True}
|
||||
MIN_SCORE = 80
|
||||
|
||||
|
||||
def now():
|
||||
return time.time()
|
||||
|
||||
|
||||
def http_json(url, timeout=10, body=None):
|
||||
req = urllib.request.Request(url, json.dumps(body).encode() if body else None,
|
||||
{"Content-Type": "application/json", "Authorization": "Bearer none"})
|
||||
with urllib.request.urlopen(req, timeout=timeout) as r:
|
||||
return json.loads(r.read().decode())
|
||||
|
||||
|
||||
def preflight(base_url, wait=0):
|
||||
"""Server answers, serves the expected build, can generate. Returns a dict (ok, model id, speed, models json)."""
|
||||
deadline = now() + wait
|
||||
while True:
|
||||
out = {"ok": False, "base_url": base_url, "time": time.strftime("%F %T")}
|
||||
try:
|
||||
models = http_json(base_url + "/models", timeout=10)
|
||||
ids = [m.get("id", "") for m in models.get("data", [])]
|
||||
out["models"] = ids
|
||||
out["model_id"] = ids[0] if ids else ""
|
||||
out["build_ok"] = bool(ids) and os.path.basename(ids[0].rstrip("/")) == EXPECTED_MODEL
|
||||
if out["build_ok"]:
|
||||
t0 = now()
|
||||
r = http_json(base_url + "/chat/completions", timeout=180,
|
||||
body={"model": out["model_id"], "messages": [{"role": "user", "content": "Say OK."}],
|
||||
"max_tokens": 16, "temperature": 0, "chat_template_kwargs": {"enable_thinking": False}})
|
||||
out["probe_seconds"] = round(now() - t0, 1)
|
||||
out["probe_reply"] = (r["choices"][0]["message"].get("content") or "")[:40]
|
||||
out["ok"] = True
|
||||
else:
|
||||
out["error"] = f"served model {ids} is not {EXPECTED_MODEL} (the official baseline build)"
|
||||
except Exception as e: # noqa: BLE001
|
||||
out["error"] = repr(e)[:200]
|
||||
if out["ok"] or now() >= deadline:
|
||||
return out
|
||||
time.sleep(15)
|
||||
|
||||
|
||||
def build_plan():
|
||||
"""3 x INTF, TABL, STRU, MSAG, EXC, then 5 DDLS: accepted training tasks, no K variants, lowest ids first."""
|
||||
import glob
|
||||
by_kind = {}
|
||||
for f in sorted(glob.glob(os.path.join(POOL, "G*", "generation.json"))):
|
||||
tid = os.path.basename(os.path.dirname(f))
|
||||
try:
|
||||
if not json.load(open(f)).get("accepted"):
|
||||
continue
|
||||
meta = json.load(open(os.path.join(POOL, tid, "task.json")))
|
||||
except (OSError, ValueError):
|
||||
continue
|
||||
if meta.get("category") == "K":
|
||||
continue
|
||||
k = mix.kind_of_task_dir(tid)
|
||||
by_kind.setdefault(k, []).append((tid, meta.get("category")))
|
||||
plan = []
|
||||
for kind, n in PLAN_ORDER:
|
||||
for tid, cat in by_kind.get(kind, [])[:n]:
|
||||
plan.append({"task": tid, "kind": kind, "category": cat})
|
||||
return plan
|
||||
|
||||
|
||||
def failure_types(rep, why):
|
||||
"""Why a run is not accepted (list of short labels)."""
|
||||
out = []
|
||||
er = rep.get("end_reason")
|
||||
if er and er != "report":
|
||||
out.append(er) # loop, tool_budget, empty_response, model_error, max_turns, ...
|
||||
g = rep.get("gates") or {}
|
||||
names = {"G1_active": "not_active", "G2_contract": "contract", "G3_hidden_runs": "hidden_tests_not_run",
|
||||
"G4_out_of_scope": "out_of_scope_changed", "G5_no_p1": "atc_priority1", "G6_release": "release_syntax"}
|
||||
out += [names[k] for k, v in g.items() if not v and k in names]
|
||||
h = rep.get("hidden_tests") or {}
|
||||
if h.get("total") and h.get("passed", 0) < h["total"]:
|
||||
out.append("hidden_tests_failed")
|
||||
if not out and any(w.startswith("score=") for w in why):
|
||||
out.append("low_score")
|
||||
if "no_final_report" in why and "no_report" not in out and er == "report":
|
||||
out.append("no_report")
|
||||
return out or ["other"]
|
||||
|
||||
|
||||
class Series:
|
||||
def __init__(self, a):
|
||||
self.a = a
|
||||
os.makedirs(OUT, exist_ok=True)
|
||||
self.state = json.load(open(STATE)) if os.path.exists(STATE) else {}
|
||||
self.stop_hb = threading.Event()
|
||||
self.lock = threading.Lock()
|
||||
|
||||
# ---- state -------------------------------------------------------------------------------------------------
|
||||
def save(self):
|
||||
with self.lock:
|
||||
self.state["updated"] = now()
|
||||
tmp = STATE + ".tmp"
|
||||
json.dump(self.state, open(tmp, "w"), indent=1, default=str)
|
||||
os.replace(tmp, STATE)
|
||||
|
||||
def heartbeat(self):
|
||||
while not self.stop_hb.wait(30):
|
||||
try:
|
||||
ok = True
|
||||
http_json(self.a.base_url + "/models", timeout=8)
|
||||
except Exception: # noqa: BLE001
|
||||
ok = False
|
||||
self.state["server_ok"] = ok
|
||||
self.state["server_checked"] = now()
|
||||
self.save()
|
||||
|
||||
def window_end(self):
|
||||
return self.state["window_start"] + self.a.window_hours * 3600
|
||||
|
||||
# ---- one run ---------------------------------------------------------------------------------------------
|
||||
def one(self, item, idx, model_id):
|
||||
restarts = 0
|
||||
while True:
|
||||
run_no = RUN_BASE + idx * 10 + restarts
|
||||
self.state["current"] = {"task": item["task"], "kind": item["kind"], "category": item["category"], "start": now(),
|
||||
"run": run_no, "restarts": restarts}
|
||||
self.state["status"] = "running"
|
||||
self.save()
|
||||
agent = LlmAgent(model_id, self.a.base_url, max_tokens=SETTINGS["max_tokens"],
|
||||
chat_template_kwargs={"enable_thinking": False}, loop_guard=SETTINGS["loop_guard"],
|
||||
deadline=self.window_end(), watch=True)
|
||||
runner = Runner(POOL, os.path.join(OUT, "runs"))
|
||||
try:
|
||||
rep, run_dir = runner.run(item["task"], agent, run_no, teardown=True, tool_budget=SETTINGS["tool_call_budget"])
|
||||
except Exception as e: # noqa: BLE001 a harness exception is not a model result
|
||||
self.event("harness_exception", item, repr(e)[:200])
|
||||
return {"task": item["task"], "kind": item["kind"], "category": item["category"], "harness_error": repr(e)[:200]}
|
||||
er = rep.get("end_reason")
|
||||
if er == "server_down":
|
||||
self.pause(item, run_dir)
|
||||
if now() >= self.window_end():
|
||||
return {"window_over": True}
|
||||
restarts = min(restarts + 1, 9)
|
||||
continue
|
||||
if er == "window_end":
|
||||
self.event("window_end_mid_run", item, "run stopped and cleaned up; not scored")
|
||||
self.park(run_dir, "_window_end")
|
||||
return {"window_over": True}
|
||||
return self.finish(item, rep, run_dir, agent)
|
||||
|
||||
def park(self, run_dir, folder):
|
||||
dst = os.path.join(OUT, folder)
|
||||
os.makedirs(dst, exist_ok=True)
|
||||
shutil.move(run_dir, os.path.join(dst, os.path.basename(run_dir) + "_" + str(int(now()))))
|
||||
|
||||
def event(self, kind, item, text):
|
||||
self.state.setdefault("events", []).append({"t": now(), "kind": kind, "task": item["task"], "text": text})
|
||||
self.save()
|
||||
print(time.strftime("%F %T"), kind, item["task"], text, flush=True)
|
||||
|
||||
def pause(self, item, run_dir):
|
||||
self.event("server_down", item, "run ended cleanly (objects deleted); series paused")
|
||||
self.park(run_dir, "_paused")
|
||||
self.state["status"] = "paused"
|
||||
self.state["pause_start"] = now()
|
||||
self.save()
|
||||
while now() < self.window_end():
|
||||
pf = preflight(self.a.base_url)
|
||||
if pf["ok"]:
|
||||
break
|
||||
time.sleep(30)
|
||||
paused = now() - self.state.pop("pause_start", now())
|
||||
self.state["paused_seconds"] = self.state.get("paused_seconds", 0) + paused
|
||||
self.event("server_back", item, f"paused {paused / 60:.1f} min; the task is run again (no failure counted)")
|
||||
|
||||
def finish(self, item, rep, run_dir, agent):
|
||||
rec_path = os.path.join(run_dir, "record.json")
|
||||
rec = json.load(open(rec_path)) if os.path.exists(rec_path) else None
|
||||
row = {"harness_error": False, "setup_failed": bool(rep.get("setup_failed"))}
|
||||
ok, why = acc.judge(rec, row, MIN_SCORE) if rec else (False, ["no_record"])
|
||||
if rep.get("setup_failed"):
|
||||
ok, why = False, ["setup_failed"]
|
||||
res = {"task": item["task"], "kind": item["kind"], "category": item["category"], "accepted": ok,
|
||||
"score": (rep.get("score") or {}).get("total"), "end_reason": rep.get("end_reason"),
|
||||
"tool_calls": rep.get("tool_calls"), "seconds": rep.get("seconds"), "agent_seconds": rep.get("agent_seconds"),
|
||||
"hidden": "%s/%s" % ((rep.get("hidden_tests") or {}).get("passed"), (rep.get("hidden_tests") or {}).get("total")),
|
||||
"activation_failures": rep.get("activation_failures"), "syntax_hints": rep.get("syntax_hints"),
|
||||
"failure_types": [] if ok else failure_types(rep, why), "reasons": why,
|
||||
"run_dir": os.path.relpath(run_dir, ROOT), "teardown_ok": isinstance(rep.get("teardown"), dict)
|
||||
and all(v.get("deleted") for v in rep["teardown"].values())}
|
||||
if ok:
|
||||
rec["metadata"]["use_for_sft"] = False # not for the first SFT (Kral 2026-10-05)
|
||||
rec["metadata"]["series"] = "local_qwen_A"
|
||||
rec["metadata"]["served_model"] = agent.model
|
||||
with open(os.path.join(OUT, "accepted.jsonl"), "a") as f:
|
||||
f.write(json.dumps({"id": f"{item['task']}_r{rep.get('run')}", "task": rec["task"], "teacher": agent.model,
|
||||
"metadata": rec["metadata"], "messages": rec["messages"], "tools": rec["tools"]}) + "\n")
|
||||
with open(os.path.join(OUT, "results.jsonl"), "a") as f:
|
||||
f.write(json.dumps(res) + "\n")
|
||||
return res
|
||||
|
||||
# ---- series ------------------------------------------------------------------------------------------------
|
||||
def run(self):
|
||||
if self.a.plan_only:
|
||||
for p in self.state.get("plan") or build_plan():
|
||||
print(p)
|
||||
return
|
||||
pf = preflight(self.a.base_url, wait=self.a.wait)
|
||||
json.dump(pf, open(os.path.join(OUT, "preflight.json"), "w"), indent=1)
|
||||
print("preflight:", json.dumps(pf)[:400], flush=True)
|
||||
if not pf["ok"]:
|
||||
sys.exit("the model server is not usable: %s" % pf.get("error"))
|
||||
plan = self.state.get("plan") or build_plan()
|
||||
if self.a.only:
|
||||
plan = [p for p in plan if p["task"] in self.a.only]
|
||||
self.state.update(plan=plan, settings=SETTINGS, base_url=self.a.base_url, model_id=pf["model_id"],
|
||||
window_hours=self.a.window_hours, preflight=pf, server_ok=True, results=self.state.get("results", []))
|
||||
self.state.setdefault("window_start", now()) # 12 hours from the first run
|
||||
self.save()
|
||||
threading.Thread(target=self.heartbeat, daemon=True).start()
|
||||
done = {r["task"] for r in self.state["results"]}
|
||||
for idx, item in enumerate(plan):
|
||||
if item["task"] in done:
|
||||
continue
|
||||
if now() >= self.window_end():
|
||||
break
|
||||
res = self.one(item, idx, pf["model_id"])
|
||||
if res.get("window_over"):
|
||||
break
|
||||
if res.get("harness_error"):
|
||||
continue
|
||||
self.state["results"].append(res)
|
||||
self.state["current"] = None
|
||||
self.save()
|
||||
print(json.dumps({k: res[k] for k in ("task", "kind", "accepted", "score", "end_reason", "failure_types")}), flush=True)
|
||||
left = self.window_end() - now()
|
||||
self.state["status"] = "window_over" if left <= 0 else "finished"
|
||||
self.state["current"] = None
|
||||
self.save()
|
||||
self.stop_hb.set()
|
||||
self.summary()
|
||||
|
||||
def summary(self):
|
||||
per = {}
|
||||
for r in self.state["results"]:
|
||||
k = per.setdefault(r["kind"], {"runs": 0, "accepted": 0, "scores": [], "failures": {}})
|
||||
k["runs"] += 1
|
||||
k["accepted"] += r["accepted"]
|
||||
if r["score"] is not None:
|
||||
k["scores"].append(r["score"])
|
||||
for ft in r["failure_types"]:
|
||||
k["failures"][ft] = k["failures"].get(ft, 0) + 1
|
||||
for k in per.values():
|
||||
k["mean_score"] = round(sum(k["scores"]) / len(k["scores"]), 1) if k["scores"] else None
|
||||
k["pass_rate"] = round(k["accepted"] / k["runs"], 2) if k["runs"] else None
|
||||
out = {"per_kind": per, "runs": len(self.state["results"]), "accepted": sum(r["accepted"] for r in self.state["results"]),
|
||||
"window_hours": self.a.window_hours, "paused_seconds": self.state.get("paused_seconds", 0),
|
||||
"events": self.state.get("events", []), "settings": SETTINGS, "model_id": self.state.get("model_id")}
|
||||
json.dump(out, open(os.path.join(OUT, "summary.json"), "w"), indent=1)
|
||||
print("summary:", json.dumps({k: (v["accepted"], v["runs"]) for k, v in per.items()}), flush=True)
|
||||
|
||||
|
||||
def main():
|
||||
load_env(os.path.join(ROOT, ".env"))
|
||||
ap = argparse.ArgumentParser()
|
||||
ap.add_argument("--base-url", default=os.environ.get("LOCAL_MODEL_URL"), help="http://<MacBook LAN IP>:8080/v1")
|
||||
ap.add_argument("--window-hours", type=float, default=12.0)
|
||||
ap.add_argument("--wait", type=int, default=0, help="seconds to wait for the server at the start")
|
||||
ap.add_argument("--plan-only", action="store_true")
|
||||
ap.add_argument("--only", nargs="*", help="test: only these tasks")
|
||||
ap.add_argument("--out-dir", help="test: another output directory")
|
||||
a = ap.parse_args()
|
||||
if a.out_dir:
|
||||
global OUT, STATE
|
||||
OUT = os.path.abspath(a.out_dir)
|
||||
STATE = os.path.join(OUT, "state.json")
|
||||
if not a.base_url:
|
||||
sys.exit("--base-url (or LOCAL_MODEL_URL) is required, for example http://192.168.1.20:8080/v1")
|
||||
a.base_url = a.base_url.rstrip("/")
|
||||
Series(a).run()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user