diff --git a/docs/remote-model.md b/docs/remote-model.md new file mode 100644 index 0000000..e2538ca --- /dev/null +++ b/docs/remote-model.md @@ -0,0 +1,87 @@ +# Local Qwen on the MacBook, harness on the Mac mini (series A) + +Roles: the **MacBook only serves the model** (MLX, 4-bit). The harness, A4H, the MCP server, the records and the dashboard +stay on the **Mac mini**. The series is `python3 -m harness.localqwen` (see below). Written 2026-10-05. + +## 1. Start the server on the MacBook (step by step) + +1. **Plug the MacBook in and keep the lid open.** Close all heavy apps (the model needs about 16 GB, the prompt cache up to 6 GB). +2. **Check that this is the official baseline build** (same 4-bit weights as the baseline of 2026-10-04): + ```sh + cd ~/models/Qwen3.8-27B-4bit + shasum -a 256 config.json model.safetensors.index.json tokenizer.json chat_template.jinja + ``` + The values must be: + ``` + 14b65a0ee06517060a6bbd979bb1a8ff54e7b304b1a1f01d54344b88b8285e85 config.json + 13b840162b4cb35c66fef7df072f7dbb4717908204364f5e5d9f9655a2758fa8 model.safetensors.index.json + 06b9509352d2af50381ab2247e083b80d32d5c0aba91c272ca9ff729b6a0e523 tokenizer.json + c3cf9e34abf4f9e36c2d72165aa9c132d3e2a725b6c2586aaa3a8af9d7a81041 chat_template.jinja + ``` + Full check (about 2 minutes, the weights): `shasum -a 256 model-0000*-of-00003.safetensors` + ``` + 6cc1508e96fb5d0865dfd5753a79f4ec60651bf3e2a82844a7e8ae9c60528c0d model-00001-of-00003.safetensors + 83f2a20ca8058f486a3634a27faf99587f4cd3c156a83dee34fb99e6ac178670 model-00002-of-00003.safetensors + 31b8c91ef899f79efaaa69e3d2c096f6e2ebeb2ff20e29222abbd9ebc79e560a model-00003-of-00003.safetensors + ``` + (These are the hashes of `mlx-community/Qwen3.8-27B-4bit` on the Mac mini, `~/models/Qwen3.8-27B-4bit`.) If one differs, stop. +3. **Get the start script** `train/serve_remote.sh` onto the MacBook. Either, if the repo is on the MacBook, `cd ~/projects/abap-llm/harness && git pull` + (path `train/serve_remote.sh`), or copy it: `scp erhankeseli@192.168.178.40:~/projects/abap-llm/harness/train/serve_remote.sh ~/serve_remote.sh` + (192.168.178.40 is the Mac mini; use `.29` if that is its other address, check with `ifconfig` on the mini), then `chmod +x ~/serve_remote.sh`. + The script uses the same flags as the baseline (`train/serve.sh`): temp 0.2, top-p 0.95, top-k 20, min-p 0, max tokens 32768, prompt cache 4 / 6 GB; + only the host is `0.0.0.0` (port 8080), and `caffeinate -dimsu -w ` keeps the MacBook awake as long as the server lives. +4. **Start it** (a terminal window on the MacBook, leave it open; or in the background): + ```sh + nohup ~/projects/abap-llm/harness/train/serve_remote.sh > ~/qwen_server.log 2>&1 & + ``` + (use `~/serve_remote.sh` if you copied it). The first start loads the model for about 1 minute. macOS may ask + "Do you want the application Python to accept incoming network connections?": **Allow**. +5. **Check on the MacBook:** `tail -f ~/qwen_server.log` shows `Starting httpd at 0.0.0.0 on port 8080`, and + `curl -s http://127.0.0.1:8080/v1/models` returns the model path. +6. **Find the MacBook IP** (the harness needs it): `ipconfig getifaddr en0` (Wi-Fi; if empty try `en1`, or look at System Settings > Wi-Fi > Details). + Both Macs must be in the same network (the mini is 192.168.178.x). +7. **Check from the Mac mini** (replace the IP): `curl -s http://:8080/v1/models`. It must show + `/Users/I301710/models/Qwen3.8-27B-4bit`. +8. **Start the series on the Mac mini:** + ```sh + cd ~/projects/abap-llm/harness + python3 -m harness.localqwen --base-url http://:8080/v1 --wait 600 + ``` + (`--wait 600`: waits up to 10 minutes for the server at the start. Run it with `nohup ... > runs/local_qwen.log 2>&1 &`.) + The 12 hour window starts with the first run, not with the command. + +## 2. Stop + +- **The series** (Mac mini): `pkill -f harness.localqwen` stops it, but the run in progress is cut and its objects stay in A4H. Better: wait for the + window to end (hard stop, clean) or ask Claude to stop it cleanly. +- **The server** (MacBook): `kill $(cat ~/qwen_server.pid)` (`caffeinate -w` ends when the server ends), or `pkill -f mlx_lm.server`. + Check: `curl -s -m 3 http://127.0.0.1:8080/v1/models` gives no answer. +- Closing the lid or a sleeping MacBook stops the server: the series then pauses by itself and resumes when the server answers (below). + +## 3. What the harness does (series A) + +- **Before the first run:** the server must answer, the served model path must end in `Qwen3.8-27B-4bit`, and a 16-token probe must work. + The answer is saved in `runs/local_qwen/preflight.json` (model path, speed of the probe). The weights themselves cannot be + checked remotely: the shasum comparison in step 2 is the proof of the same build. +- **Tasks and order** (`--plan-only` prints them): 3 x INTF, 3 x TABL, 3 x STRU, 3 x MSAG, 3 x exception, then 5 x DDLS (accepted + training tasks, no K variants, lowest ids). One task at a time. +- **Settings as the official baseline:** thinking off (`enable_thinking=false` per request), max_tokens 16384, loop guard 3, + tool budget 60 (CDS too), temperature 0.2, same system prompt, at most 80 turns. +- **Window:** 12 hours from the first run. At the end there is a hard stop, also in the middle of a run: the model request is dropped, + the run is not scored (no failure), its objects are deleted on A4H (the teardown needs about 20 s). +- **Server does not answer:** during a request the server is pinged every 30 s (3 misses in a row); between requests a failed + request is checked with a ping. Then the run ends cleanly (objects deleted, no scoring, folder `runs/local_qwen/_paused/`), the series + waits (a check every 30 s), and when the server answers again the same task starts again. The time of the pause counts in the 12 hours. + This is not counted as a failure. +- **Output** (`runs/local_qwen/`): `state.json` (the dashboard card), `results.jsonl` (one line per task), `summary.json` (at the end: pass rate and + failure types per kind), `accepted.jsonl` (accepted trajectories, **not for the first SFT**: `metadata.use_for_sft = false`, `series = local_qwen_A`), + `runs/` (one folder per run with `record.json`, trajectory, report). Resumable: start it again with the same arguments. +- **Acceptance** is the same filter as for the DeepSeek trajectories: score at least 80, end reason `report`, no harness error text. + Failure types per run: `loop`, `tool_budget`, `empty_response`, `model_error`, `not_active`, `contract`, `hidden_tests_not_run`, + `hidden_tests_failed`, `out_of_scope_changed`, `atc_priority1`, `release_syntax`, `low_score`. +- **Dashboard:** the card "Series A: local Qwen on the MacBook" in `runs/dashboard/index.html` (5-minute refresh). + +## 4. Speed and cost + +About 12 tokens per second with the 4-bit model (`train/README.md`); a task needs 20 to 40 minutes, so the 20 tasks take 7 to 13 hours: +at the slow end the window ends before the last DDLS tasks. No cloud cost: the ledger is not touched (the model is not a `:cloud` model). diff --git a/harness/agents.py b/harness/agents.py index eaa2dc5..e92dbc6 100644 --- a/harness/agents.py +++ b/harness/agents.py @@ -4,9 +4,20 @@ import os import time import urllib.request +import threading +import urllib.error + from .ledger import add_usage, check_budget from .proxy import BudgetExceeded + +class ServerDown(Exception): + """The model server does not answer (remote local model). The run ends cleanly and is not a result.""" + + +class WindowEnd(Exception): + """The time window of the series is over (hard stop, also in the middle of a request).""" + SYSTEM_PROMPT = """You are an ABAP developer. You implement a plan on an SAP system with the tools. Rules: @@ -72,7 +83,8 @@ class LlmAgent: """ def __init__(self, model, base_url=None, api_key=None, max_turns=80, temperature=0.2, - max_seconds=None, max_tokens=None, chat_template_kwargs=None, loop_guard=None): + max_seconds=None, max_tokens=None, chat_template_kwargs=None, loop_guard=None, + deadline=None, watch=False): self.model = model self.name = f"llm:{model}" self.base_url = (base_url or os.environ.get("LLM_BASE_URL", "http://127.0.0.1:11434/v1")).rstrip("/") @@ -86,6 +98,8 @@ class LlmAgent: self.max_tokens = max_tokens or (32000 if ":cloud" in (model or "") else None) self.empty_retries = 2 self.loop_guard = loop_guard # end the run after this many identical pushes in a row (None = off) + self.deadline = deadline # absolute time (time.time()) of the hard stop, or None + self.watch = watch # remote model: ping the server during a request, end the run when it is gone self.end_reason = None self.messages, self.tools, self.reasoning, self.turn_usage = [], [], [], [] # for the trajectory record self.chat_template_kwargs = chat_template_kwargs # local server only, e.g. {"enable_thinking": False} @@ -103,17 +117,60 @@ class LlmAgent: last = None for attempt in range(4): # model server errors (HTTP 5xx, timeouts): retry with backoff try: - with urllib.request.urlopen(req, timeout=self.request_timeout) as r: - data = json.loads(r.read().decode()) - return data["choices"][0]["message"], data.get("usage", {}) + if self.watch or self.deadline: + data = self._post_watched(req) + else: + with urllib.request.urlopen(req, timeout=self.request_timeout) as r: + data = json.loads(r.read().decode()) + return data["choices"][0]["message"], data.get("usage", {}) + except (ServerDown, WindowEnd): + raise except Exception as e: # noqa: BLE001 last = e code = getattr(e, "code", None) if code is not None and code < 500 and code != 429: raise + if self.deadline and time.time() >= self.deadline: + raise WindowEnd() + if self.watch and not self._ping(): + raise ServerDown(str(e)[:200]) time.sleep(10 * (attempt + 1)) raise RuntimeError(f"model request failed after retries: {last}") + def _ping(self, timeout=10): + try: + urllib.request.urlopen(urllib.request.Request(f"{self.base_url}/models"), timeout=timeout).read() + return True + except Exception: # noqa: BLE001 + return False + + def _post_watched(self, req): + """The request runs in a thread; this thread watches the deadline and (remote model) the server: a hard stop + or a dead server ends the wait at once, also in the middle of a long generation.""" + box = {} + + def work(): + try: + with urllib.request.urlopen(req, timeout=self.request_timeout) as r: + box["data"] = json.loads(r.read().decode()) + except BaseException as e: # noqa: BLE001 + box["err"] = e + th = threading.Thread(target=work, daemon=True) + th.start() + last_ping, fails = time.time(), 0 + while th.is_alive(): + th.join(5) + if self.deadline and time.time() >= self.deadline: + raise WindowEnd() + if self.watch and th.is_alive() and time.time() - last_ping >= 30: + last_ping = time.time() + fails = 0 if self._ping() else fails + 1 + if fails >= 3: + raise ServerDown("no answer to 3 pings in a row during a request") + if "err" in box: + raise box["err"] + return box["data"] + def run(self, task, proxy): tools = [{"type": "function", "function": {"name": t["name"], "description": t.get("description", ""), @@ -129,6 +186,10 @@ class LlmAgent: self.end_reason = "max_turns" start = time.time() for _ in range(self.max_turns): + if self.deadline and time.time() >= self.deadline: + final = "Stopped: time window over." + self.end_reason = "window_end" + break if self.max_seconds and time.time() - start > self.max_seconds: final = f"Stopped: time budget exceeded ({self.max_seconds} s)." self.end_reason = "time_budget" @@ -137,6 +198,14 @@ class LlmAgent: msg, usage = self._chat(messages, tools) add_usage(self.model, usage, kind="run", ref=proxy.prefix) self.turn_usage.append(usage) + except WindowEnd: + final = "Stopped: time window over." + self.end_reason = "window_end" + break + except ServerDown as e: + final = f"Stopped: model server not reachable ({e})." + self.end_reason = "server_down" + break except RuntimeError as e: final = f"Stopped: {e}" self.end_reason = "model_error" diff --git a/harness/dashboard.py b/harness/dashboard.py index 04d1351..110d1d6 100644 --- a/harness/dashboard.py +++ b/harness/dashboard.py @@ -243,6 +243,65 @@ def agg_table(title, evs, rows, key, note=""): % (E(title), E(key.replace("_", " ")), trs, "
%s
" % E(note) if note else "")) +def local_card(): + """Card for series A (local Qwen on the MacBook): state.json of harness.localqwen. Empty string without a series.""" + path = os.path.join(ROOT, "runs", "local_qwen", "state.json") + if not os.path.exists(path): + return "" + try: + st = json.load(open(path)) + except ValueError: + return "" + now = time.time() + url = st.get("base_url") or "" + up = ping(url + "/models", timeout=4) if url else False # the dashboard asks the MacBook itself + status = st.get("status", "?") + left = st.get("window_start", now) + st.get("window_hours", 12) * 3600 - now + cur = st.get("current") + plan, res = st.get("plan", []), st.get("results", []) + chips = "server %s %s" % ( + "ok" if up else "er", "answers" if up else "does not answer", + {"running": "ok", "paused": "wa", "finished": "gr", "window_over": "gr"}.get(status, "gr"), E(status)) + rows = "" + per = {} + for r in res: + k = per.setdefault(r["kind"], {"runs": 0, "ok": 0, "scores": [], "fail": {}}) + k["runs"] += 1 + k["ok"] += r["accepted"] + if r.get("score") is not None: + k["scores"].append(r["score"]) + for f in r.get("failure_types", []): + k["fail"][f] = k["fail"].get(f, 0) + 1 + planned = {} + for p_ in plan: + planned[p_["kind"]] = planned.get(p_["kind"], 0) + 1 + for kind in [x for x in ("INTF", "TABL", "STRU", "MSAG", "EXC", "DDLS") if x in planned or x in per]: + k = per.get(kind, {"runs": 0, "ok": 0, "scores": [], "fail": {}}) + mean = "%.0f" % (sum(k["scores"]) / len(k["scores"])) if k["scores"] else "-" + rate = "%.0f%%" % (100.0 * k["ok"] / k["runs"]) if k["runs"] else "-" + fails = ", ".join("%s %d" % (f, n) for f, n in sorted(k["fail"].items(), key=lambda x: -x[1])) or "-" + rows += ("%s%d/%d%d%s%s%s" + % (kind, k["runs"], planned.get(kind, 0), k["ok"], rate, mean, E(fails))) + if cur: + elapsed = (now - cur["start"]) / 60 + curtxt = "%s (%s, %s), run %d min%s" % (cur["task"], cur["kind"], cur["category"], elapsed, + ", restart %d" % cur["restarts"] if cur.get("restarts") else "") + else: + curtxt = "none (%s)" % status + evs = "".join("
%s %s %s: %s
" % (time.strftime("%H:%M", time.localtime(e["t"])), E(e["kind"]), + E(e["task"]), E(e["text"])) for e in st.get("events", [])[-3:]) + hdr = ("

Series A: local Qwen 3.8 27B on the MacBook (thinking off, 16384 tokens, loop guard 3, budget 60)

" + "
%s   model %s
" + "
" + "
Current task%sDone / planned%d / %dTime left in the window%sPaused%.0f min
" + % (chips, E(os.path.basename(str(st.get("model_id", "")))), E(curtxt), len(res), len(plan), + ("%dh %02dm" % (left // 3600, (left % 3600) // 60)) if left > 0 else "over", st.get("paused_seconds", 0) / 60)) + tbl = ("
" + "%s
kindruns/plannedacceptedpass ratemean scorefailure types
%s" + "
Accepted trajectories of this series are kept apart (runs/local_qwen/accepted.jsonl), not for the first SFT.
" % (rows, evs)) + return hdr + tbl + + def mix_table(rows, evs): kt = mix.accepted_task_counts() kr = {} @@ -376,7 +435,7 @@ def render(d): % ("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 + body = (banner + "
" + tiles + "
" + local_card() + budget + prog + stops + sysc + mix_table(rows, evs) + 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 + "
") diff --git a/harness/localqwen.py b/harness/localqwen.py new file mode 100644 index 0000000..27180ef --- /dev/null +++ b/harness/localqwen.py @@ -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://: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://: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() diff --git a/harness/runner.py b/harness/runner.py index 08ea75d..760549d 100644 --- a/harness/runner.py +++ b/harness/runner.py @@ -209,11 +209,13 @@ class Runner: "line": i.get("start", {}).get("row"), "msg": i.get("description")} for i in issues] # ---------- main ---------- - def run(self, task_id, agent, run_no, teardown=True, rescore_dir=None, cds_calls=None): + def run(self, task_id, agent, run_no, teardown=True, rescore_dir=None, cds_calls=None, tool_budget=None): """cds_calls: raise max_tool_calls to this value for tasks with a CDS (DDLS) contract object (trajectory runs, 2026-10-05: the own CDS test class needs more than 60 calls; eval stays at the task value).""" prefix = prefix_for(run_no, task_id) task = Task(os.path.join(self.tasks_root, task_id), prefix) + if tool_budget: # fixed budget for every task (local model series: 60) + task.meta.setdefault("budget", {})["max_tool_calls"] = tool_budget if cds_calls and any(c.get("type") == "DDLS" for c in task.meta.get("contract", [])): b = task.meta.setdefault("budget", {}) b["max_tool_calls"] = max(b.get("max_tool_calls", 60), cds_calls) @@ -261,6 +263,15 @@ class Runner: if proxy.fallbacks: rep["adt_fallbacks"] = proxy.fallbacks + if rep.get("end_reason") in ("window_end", "server_down"): # not a result: no scoring, only the teardown + rep["score"] = {"total": None, "note": f"not scored: {rep['end_reason']}"} + rep["not_a_result"] = True + rep["seconds"] = round(time.time() - t0, 1) + all_objs = self._objects_with_prefix(mcp, prefix) + rep["teardown"] = self._teardown(all_objs, run_dir) if teardown else "skipped" + json.dump(rep, open(os.path.join(run_dir, "report.json"), "w"), indent=1) + return rep, run_dir + # 3 collect hidden_names = {o["name"].upper() for o in task.objects("hidden_tests")} seed_names = set(seed_src) diff --git a/train/serve_remote.sh b/train/serve_remote.sh new file mode 100755 index 0000000..a7170a2 --- /dev/null +++ b/train/serve_remote.sh @@ -0,0 +1,24 @@ +#!/bin/sh +# Serve Qwen 3.8 27B (4-bit, the official baseline build) for the harness on the Mac mini. Run this on the MacBook. +# Same server flags as train/serve.sh (the official baseline); only the host differs (0.0.0.0) and caffeinate keeps the +# MacBook awake as long as the server process lives (caffeinate -w ). Usage: ~/serve_remote.sh Stop: kill $(cat ~/qwen_server.pid) +MODEL="$HOME/models/Qwen3.8-27B-4bit" +[ -f "$MODEL/config.json" ] || { echo "model not found: $MODEL" >&2; exit 1; } +# the mlx_lm of the harness venv if the repo is on this Mac, else the one on PATH +if [ -x "$HOME/projects/abap-llm/harness/train/.venv/bin/mlx_lm.server" ]; then + SERVER="$HOME/projects/abap-llm/harness/train/.venv/bin/mlx_lm.server" +elif command -v mlx_lm.server >/dev/null 2>&1; then + SERVER="$(command -v mlx_lm.server)" +else + echo "mlx_lm.server not found (pip install mlx-lm==0.32.0)" >&2; exit 1 +fi +"$SERVER" \ + --model "$MODEL" \ + --host 0.0.0.0 --port 8080 \ + --temp 0.2 --top-p 0.95 --top-k 20 --min-p 0 \ + --max-tokens 32768 \ + --prompt-cache-size 4 --prompt-cache-bytes 6000000000 \ + --chat-template-args '{"enable_thinking": true, "reasoning_effort": "medium"}' & +SP=$! +echo "$SP" > "$HOME/qwen_server.pid" # the pid of the server itself +exec caffeinate -dimsu -w "$SP" # keeps the Mac awake until the server process ends