"""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://: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 infra, mix from .runner import Runner from .task import prefix_for 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 if infra.is_outage(e): # MCP or A4H is down: wait, clean up, run the same task again (not a failure) self.event("infra_outage", item, repr(e)[:150]) self.state["status"] = "paused" self.save() infra.wait_until_up(max_seconds=max(60, self.window_end() - now())) try: infra.cleanup_prefix(prefix_for(run_no, item["task"])) except Exception: # noqa: BLE001 pass for d in __import__("glob").glob(os.path.join(OUT, "runs", f"{run_no}_{item['task']}_*")): self.park(d, "_paused") if now() >= self.window_end(): return {"window_over": True} restarts = min(restarts + 1, 9) continue 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://: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()