From c2c3803b1b4df72da89123888664c395ddcfba31 Mon Sep 17 00:00:00 2001 From: Kral Date: Tue, 6 Oct 2026 06:15:08 +0200 Subject: [PATCH] G: EPOD acceptance tests; harness: hide and clean other runs' mid-name objects (proxy, teardown, sweep); lock leak cause (Eclipse restart during a write) Co-Authored-By: Claude Sonnet 5.5 --- docs/epod-lock-leak.md | 25 ++++- epod_tests/README.md | 22 ++++ epod_tests/epod_client.py | 98 +++++++++++++++++ epod_tests/epod_tests_result.json | 32 ++++++ epod_tests/run_tests.py | 174 ++++++++++++++++++++++++++++++ harness/proxy.py | 14 ++- harness/runner.py | 13 ++- harness/sweep.py | 89 +++++++++++++++ scripts_probe/lockprobe.py | 15 +-- train/STATE.md | 2 + 10 files changed, 470 insertions(+), 14 deletions(-) create mode 100644 epod_tests/README.md create mode 100644 epod_tests/epod_client.py create mode 100644 epod_tests/epod_tests_result.json create mode 100644 epod_tests/run_tests.py create mode 100644 harness/sweep.py diff --git a/docs/epod-lock-leak.md b/docs/epod-lock-leak.md index 1215aac..f3fdfe0 100644 --- a/docs/epod-lock-leak.md +++ b/docs/epod-lock-leak.md @@ -1,4 +1,4 @@ -# EPOD: stale lock after a write (2 cases, not reproduced under control) +# EPOD: stale lock after a write (3 cases; the third one gives the cause) Status 2026-10-05 evening. For Kral, to fix in the EPOD server (or to decide that it is not worth it). @@ -13,6 +13,24 @@ Status 2026-10-05 evening. For Kral, to fix in the EPOD server (or to decide tha In both cases the lock outlived the MCP session and the run (hours). `ENQUEUE_READ` (see below) showed nothing after the SM12 delete. +## Case 3 and the probable cause (2026-10-06): the MCP server (Eclipse) restarted while a write was running + +At 21:43 on 2026-10-05 Kral restarted Eclipse (the EPOD MCP server runs inside it). Three runs were writing at that moment. One of them (run 202781, task G1927) left this entry in the +enqueue table (read with the reader class below, 2026-10-06 06:15): + +``` +GNAME=SEOCLSENQ GARG=ZCL_Z4CGT1HJ_JOB_COST_TEST====... GMODE=X GOBJ=ESEOCLASS GCLIENT=001 GUNAME=KESELI +GUSR=20261005194600195642000500vhcala4hci_A4H_00... GUSE=1 GTHOST=vhcala4hci_A4H_00 GTWP=05 GTDATE=20261005 GTTIME=194600 (server time, 2 hours behind CEST: 21:46) +``` +Exclusive lock (mode X) on the class, owner user KESELI, created at 21:46, three minutes after the restart, by a write that was in flight; it was still there the next morning. +Its object was not cleaned up because the model had named it `ZCL__JOB_COST_TEST` (prefix inside the name; the teardown looked only for names that start with the prefix). When the pipeline +restarted at 22:27 it ran the same task with the same run number, so the new run found the old locked class: `[LOCK] ... User KESELI is currently editing` on every `sap_push_element` (the run looped and scored 75). +Earlier I wrote that the restart did not leak a lock: that was wrong, I had only cleaned the objects whose names start with the prefix. + +So the cause is probably: **a write call is in flight when the MCP server process stops (restart, crash, kill of the whole server); the lock of that write stays in the enqueue table** (the stateful ADT session that owns it is gone, and nothing removes the lock). +This also fits case 2 (two controllers killed in the middle of a write) better than the client-side kill tests, which did not reproduce it: killing the *client* does not stop the server's write. +It cannot be tested from the client side without restarting Eclipse; for Kral: start a write (a class with a large test include), restart Eclipse in the middle, read the enqueue table (below) and look at SM12. + ## What we tried to reproduce (A4H, probe objects, no cloud, 2026-10-05) Every test: write through MCP, interfere, wait 3 s, read the enqueue table, write the same object again, delete it. @@ -53,8 +71,13 @@ There is no function module on A4H to delete an enqueue entry (`TFDIR`: only `EN - One controller only (flock on `runs/pipeline/controller.lock`), so two controllers cannot install the same run number twice any more (this was Case 2). - Never kill a controller while it runs: it is drained through the budget guard or a STOP flag. +## Also found: leftovers with the prefix inside the name +27 classes named `ZCL__...` / `ZCX__...` were left in A4H from earlier runs (the teardown searched only for names that start with the prefix). They showed up in +`sap_inactive_objects` and searches of later runs and in 57 of 90 accepted trajectories. Fixed in the harness (2026-10-06, `harness/proxy.py`, `harness/runner.py`, `harness/sweep.py`). + ## Wish for EPOD +0. **Release the enqueue locks of the server's own ADT sessions when the MCP server stops or starts** (a lock of user KESELI with the server's work process that has no live ADT session behind it). 1. Release the lock in a `finally` path of every write tool (also after a failed or warned activation). 2. Release all locks of an MCP session when the session ends (`DELETE /mcp`) and when the connection breaks. 3. Return a clear error with the lock owner and the age of the lock, and a tool to release the locks of the caller's own session. diff --git a/epod_tests/README.md b/epod_tests/README.md new file mode 100644 index 0000000..d3fc17b --- /dev/null +++ b/epod_tests/README.md @@ -0,0 +1,22 @@ +# EPOD acceptance tests (for the EPOD server developer) + +Small test set for the two server changes that the harness work asked for (details: `docs/epod-syntax-hint.md`, `docs/epod-lock-leak.md`) +and for the parallel-call behavior (`harness/loadtest.py` numbers). Standard library only, Python 3.9+, no harness code needed. +It uses probe objects `ZEPODT_*` in `$TMP` on a **test system**. + +```sh +cd epod_tests +export MCP_URL=http://127.0.0.1:3000/mcp MCP_TOKEN=... # the MCP server +export A4H_URL=http://localhost:50000 A4H_USER=... A4H_PASSWORD=... # only for the cleanup (ADT deletion); without it the tests list the probe objects +python3 run_tests.py # all; --only T1 T3 for some +``` + +| test | passes when | state of the server on 2026-10-06 | +|---|---|---| +| T1 syntax_hint | a rejected write (`TYPE c LENGTH 4` in a method signature) returns syntax messages with a line, not only "save operation failed" | FAIL expected (not implemented; the harness proxy adds abaplint messages) | +| T2 no_lock_after_kill | a client killed with SIGKILL during a write (4 delays) leaves no lock: the next write works | PASS (not reproducible on A4H) | +| T3 parallel_reads | 6 clients search + syntax check in parallel without `Concurrent call detected` | PASS | +| T4 parallel_writes | 3 clients create + write in parallel without `Concurrent call detected` | FAIL expected (the shared connection; 15 to 50 lock hits per 45 s in the harness load test) | +| T5 lock_error_text | informational | PASS | + +A change in the server is done when T1 and T4 pass and T2 and T3 stay green. T1 accepts any answer that carries a line number or a `syntaxCheck` object with messages (the proxy format is in `docs/epod-syntax-hint.md`). diff --git a/epod_tests/epod_client.py b/epod_tests/epod_client.py new file mode 100644 index 0000000..8baaea7 --- /dev/null +++ b/epod_tests/epod_client.py @@ -0,0 +1,98 @@ +"""Minimal MCP client and ADT deletion for the EPOD acceptance tests (standard library only, Python 3.9+).""" +import base64 +import http.cookiejar +import json +import os +import re +import time +import urllib.request +from xml.sax.saxutils import quoteattr + +URL = os.environ.get("MCP_URL", "http://127.0.0.1:3000/mcp") +TOKEN = os.environ.get("MCP_TOKEN", "") +SYSTEM = os.environ.get("MCP_SYSTEM_ID") # optional: Eclipse project name when several systems are connected + + +class Mcp: + def __init__(self, timeout=300): + self.sid, self.n, self.timeout = None, 0, timeout + + def _post(self, body, method="POST"): + h = {"Content-Type": "application/json", "Accept": "application/json, text/event-stream"} + if TOKEN: + h["Authorization"] = "Bearer " + TOKEN + if self.sid: + h["Mcp-Session-Id"] = self.sid + req = urllib.request.Request(URL, json.dumps(body).encode() if body is not None else None, h, method=method) + with urllib.request.urlopen(req, timeout=self.timeout) as r: + sid = r.headers.get("Mcp-Session-Id") + raw = r.read().decode() + self.sid = sid or self.sid + if "data:" in raw[:40]: + raw = "".join(l[5:].strip() for l in raw.splitlines() if l.startswith("data:")) + return json.loads(raw) if raw.strip() else None + + def open(self): + self.n += 1 + self._post({"jsonrpc": "2.0", "id": self.n, "method": "initialize", "params": { + "protocolVersion": "2025-03-26", "capabilities": {}, "clientInfo": {"name": "epod-acceptance", "version": "1"}}}) + self._post({"jsonrpc": "2.0", "method": "notifications/initialized"}) + return self + + def close(self): + try: + self._post(None, "DELETE") + except Exception: + pass + + def __enter__(self): + return self.open() + + def __exit__(self, *a): + self.close() + + def call(self, tool, args): + """Returns (is_error, text). 'Concurrent call detected' is returned as it is (the tests count it).""" + self.n += 1 + if SYSTEM: + args = dict(args, systemId=SYSTEM) + res = self._post({"jsonrpc": "2.0", "id": self.n, "method": "tools/call", "params": {"name": tool, "arguments": args}}) + if "error" in res: + return True, json.dumps(res["error"]) + r = res["result"] + return bool(r.get("isError")), "\n".join(c.get("text", "") for c in r.get("content", [])) + + +class Adt: + """Deletion through the ADT deletion API (needs A4H_URL, A4H_USER, A4H_PASSWORD; A4H_CLIENT default 001).""" + + def __init__(self): + self.base = os.environ.get("A4H_URL", "").rstrip("/") + self.client = os.environ.get("A4H_CLIENT", "001") + self.auth = "Basic " + base64.b64encode(("%s:%s" % (os.environ.get("A4H_USER", ""), os.environ.get("A4H_PASSWORD", ""))).encode()).decode() + self.opener = urllib.request.build_opener(urllib.request.HTTPCookieProcessor(http.cookiejar.CookieJar())) + self.csrf = None + + def available(self): + return bool(self.base and os.environ.get("A4H_USER")) + + def _req(self, path, body=None, headers=None): + h = {"Authorization": self.auth} + if self.csrf: + h["x-csrf-token"] = self.csrf + h.update(headers or {}) + req = urllib.request.Request("%s%s%ssap-client=%s" % (self.base, path, "&" if "?" in path else "?", self.client), body.encode() if body else None, h) + with self.opener.open(req, timeout=120) as r: + return dict(r.headers), r.read().decode() + + def delete(self, uris): + if not self.available() or not uris: + return {} + hdr, _ = self._req("/sap/bc/adt/discovery", headers={"x-csrf-token": "fetch", "Accept": "*/*"}) + self.csrf = hdr.get("x-csrf-token") or hdr.get("X-CSRF-Token") + objs = "".join("" % quoteattr(u) for u in uris) + body = ('' + objs + "") + _, text = self._req("/sap/bc/adt/deletion/delete", body, {"Content-Type": "application/vnd.sap.adt.deletion.request.v1+xml", + "Accept": "application/vnd.sap.adt.deletion.response.v1+xml"}) + return {m.group(1): 'isDeleted="true"' in m.group(0) for m in re.finditer(r']*adtcore:uri="([^"]+)"[^>]*>', text)} diff --git a/epod_tests/epod_tests_result.json b/epod_tests/epod_tests_result.json new file mode 100644 index 0000000..52c246a --- /dev/null +++ b/epod_tests/epod_tests_result.json @@ -0,0 +1,32 @@ +[ + { + "id": "T1", + "name": "syntax_hint", + "pass": false, + "detail": "the failed write returned: {\"success\":false,\"error\":\"[WRITE] An error occured during the save operation. The changes were not stored.\"}" + }, + { + "id": "T2", + "name": "no_lock_after_kill", + "pass": true, + "detail": "4 kills (0.05 to 0.7 s), the next write always worked" + }, + { + "id": "T3", + "name": "parallel_reads", + "pass": true, + "detail": "372 rounds, 0 'Concurrent call detected'" + }, + { + "id": "T4", + "name": "parallel_writes", + "pass": false, + "detail": "27 create+write rounds with 3 clients, 12 'Concurrent call detected'" + }, + { + "id": "T5", + "name": "lock_error_text", + "pass": true, + "detail": "informational: see docs/epod-lock-leak.md (no way to create a lock on purpose)" + } +] \ No newline at end of file diff --git a/epod_tests/run_tests.py b/epod_tests/run_tests.py new file mode 100644 index 0000000..4fa61b3 --- /dev/null +++ b/epod_tests/run_tests.py @@ -0,0 +1,174 @@ +"""EPOD acceptance tests (2026-10-06). Run against an EPOD MCP server on a test system (probe objects ZEPODT_*, package $TMP). + + MCP_URL=http://127.0.0.1:3000/mcp MCP_TOKEN=... [A4H_URL=http://localhost:50000 A4H_USER=... A4H_PASSWORD=...] python3 run_tests.py [--only T1 T2 ...] + +Each test prints PASS or FAIL with the reason and writes the result to epod_tests_result.json. Tests: + T1 syntax_hint a rejected write ('save operation failed') comes with syntax messages (line and text) docs/epod-syntax-hint.md + T2 no_lock_after_kill a client killed during a write leaves no lock (the next write and the deletion work) docs/epod-lock-leak.md + T3 parallel_reads 6 clients read/ATC/syntax-check in parallel without 'Concurrent call detected' + T4 parallel_writes 3 clients create/write/activate in parallel without 'Concurrent call detected' (a fix: per system or per session lock) + T5 lock_error_text a write on a locked object names the lock owner (informational) +""" +import json +import os +import signal +import subprocess +import sys +import threading +import time + +sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) +from epod_client import Adt, Mcp # noqa: E402 + +CLS = """CLASS {n} DEFINITION PUBLIC FINAL CREATE PUBLIC. + PUBLIC SECTION. + METHODS run RETURNING VALUE(rv) TYPE i. +ENDCLASS. +CLASS {n} IMPLEMENTATION. + METHOD run. + rv = 1. + ENDMETHOD. +ENDCLASS. +""" +BAD = """CLASS {n} DEFINITION PUBLIC FINAL CREATE PUBLIC. + PUBLIC SECTION. + METHODS run IMPORTING iv_zone TYPE c LENGTH 4 RETURNING VALUE(rv) TYPE i. +ENDCLASS. +CLASS {n} IMPLEMENTATION. + METHOD run. + rv = 1. + ENDMETHOD. +ENDCLASS. +""" +TC = """CLASS ltc DEFINITION FINAL FOR TESTING DURATION SHORT RISK LEVEL HARMLESS. + PRIVATE SECTION. + METHODS t1 FOR TESTING. +ENDCLASS. +CLASS ltc IMPLEMENTATION. + METHOD t1. + cl_abap_unit_assert=>assert_equals( act = 1 exp = 1 ). + ENDMETHOD. +ENDCLASS. +""" +created = [] +RUN = time.strftime("%H%M%S") + + +def new_class(m, tag): + name = "ZEPODT_%s_%s" % (tag, RUN) + m.call("sap_create_object", {"objectType": "CLAS", "objectName": name, "packageName": "$TMP", "description": "epod acceptance test"}) + created.append("/sap/bc/adt/oo/classes/" + name.lower()) + return name + + +def ok(text): + return '"success":true' in text.replace(" ", "") + + +def t1_syntax_hint(): + with Mcp() as m: + n = new_class(m, "T1") + e, t = m.call("sap_push_source", {"objectType": "CLAS", "objectName": n, "source": BAD.format(n=n.lower())}) + low = t.lower() + has_detail = ("syntaxcheck" in low or "syntax" in low and "line" in low) and "save operation" in low or "line 3" in low or '"line":3' in low.replace(" ", "") + return bool(has_detail), "the failed write returned: " + t[:300] + + +def t2_no_lock_after_kill(): + victim = ("import sys,json,os\nsys.path.insert(0,%r)\nfrom epod_client import Mcp\na=json.loads(sys.argv[1])\nm=Mcp().open()\nprint('SENT',flush=True)\nm.call('sap_push_source',a)\n" + % os.path.dirname(os.path.abspath(__file__))) + bad = [] + for delay in (0.05, 0.2, 0.4, 0.7): + with Mcp() as m: + n = new_class(m, "T2") + m.call("sap_push_source", {"objectType": "CLAS", "objectName": n, "source": CLS.format(n=n.lower())}) + p = subprocess.Popen([sys.executable, "-c", victim, json.dumps({"objectType": "CLAS", "objectName": n, "includeType": "testclasses", "source": TC})], + stdout=subprocess.PIPE, text=True) + p.stdout.readline() + time.sleep(delay) + try: + os.kill(p.pid, signal.SIGKILL) + except ProcessLookupError: + pass + p.wait() + time.sleep(3) + with Mcp() as m: + e, t = m.call("sap_push_source", {"objectType": "CLAS", "objectName": n, "includeType": "testclasses", "source": TC + "* again\n"}) + if not ok(t): + bad.append((delay, t[:160])) + return not bad, "after the kill the next write failed at: %s" % bad if bad else "4 kills (0.05 to 0.7 s), the next write always worked" + + +def parallel(n_clients, seconds, work): + out, end = [], time.time() + seconds + + def run(i): + locks = calls = 0 + with Mcp() as m: + while time.time() < end: + t = work(m, i) + calls += 1 + locks += "Concurrent call detected" in t + out.append((calls, locks)) + th = [threading.Thread(target=run, args=(i,)) for i in range(n_clients)] + [x.start() for x in th] + [x.join() for x in th] + return sum(c for c, _ in out), sum(l for _, l in out) + + +def t3_parallel_reads(): + def work(m, i): + return m.call("sap_search_object", {"query": "CL_ABAP_CHAR_UTIL*", "objType": "CLAS"})[1] + \ + m.call("sap_syntax_check", {"objectType": "CLAS", "objectName": "CL_ABAP_CHAR_UTILITIES"})[1] + calls, locks = parallel(6, 15, work) + return locks == 0, "%d rounds, %d 'Concurrent call detected'" % (calls, locks) + + +def t4_parallel_writes(): + cnt = {"n": 0} + + def work(m, i): + cnt["n"] += 1 + name = "ZEPODT_P%d_%s%03d" % (i, RUN, cnt["n"] % 1000) + t = m.call("sap_create_object", {"objectType": "CLAS", "objectName": name, "packageName": "$TMP", "description": "parallel"})[1] + created.append("/sap/bc/adt/oo/classes/" + name.lower()) + return t + m.call("sap_push_source", {"objectType": "CLAS", "objectName": name, "source": CLS.format(n=name.lower())})[1] + calls, locks = parallel(3, 25, work) + return locks == 0, "%d create+write rounds with 3 clients, %d 'Concurrent call detected'" % (calls, locks) + + +def t5_lock_error_text(): + return True, "informational: see docs/epod-lock-leak.md (no way to create a lock on purpose)" + + +TESTS = [("T1", "syntax_hint", t1_syntax_hint), ("T2", "no_lock_after_kill", t2_no_lock_after_kill), ("T3", "parallel_reads", t3_parallel_reads), + ("T4", "parallel_writes", t4_parallel_writes), ("T5", "lock_error_text", t5_lock_error_text)] + + +def main(): + only = set(sys.argv[sys.argv.index("--only") + 1:]) if "--only" in sys.argv else None + result = [] + for tid, name, fn in TESTS: + if only and tid not in only: + continue + t0 = time.time() + try: + passed, why = fn() + except Exception as e: # noqa: BLE001 + passed, why = False, "error: %r" % e + print("%-4s %-22s %s (%.0f s) %s" % (tid, name, "PASS" if passed else "FAIL", time.time() - t0, why[:230]), flush=True) + result.append({"id": tid, "name": name, "pass": passed, "detail": why[:600]}) + adt = Adt() + if adt.available() and created: + try: + gone = adt.delete(created) + print("cleanup: %d of %d probe objects deleted" % (sum(gone.values()), len(created))) + except Exception as e: # noqa: BLE001 + print("cleanup failed:", repr(e)[:200]) + elif created: + print("delete these probe objects yourself (SE80 or ADT): ZEPODT_* in $TMP (%d objects)" % len(created)) + json.dump(result, open("epod_tests_result.json", "w"), indent=1) + + +if __name__ == "__main__": + main() diff --git a/harness/proxy.py b/harness/proxy.py index 34dacce..e1b6dd7 100644 --- a/harness/proxy.py +++ b/harness/proxy.py @@ -51,6 +51,9 @@ def activation_messages(text, is_error=False): RUN_PREFIX = re.compile(r"^Z\d[0-9A-Z]{5,6}_", re.I) +# The model sometimes names its own helper or test class with the prefix inside the name (ZCL_Z4CGT1HJ_JOB_COST_TEST). Such objects of +# other runs stay behind in A4H and show up in lists and searches (57 of 90 accepted trajectories saw them, 2026-10-06). +MID_PREFIX = re.compile(r"^[A-Z]{1,5}_(Z\d[0-9A-Z]{6}_)", re.I) class BudgetExceeded(Exception): @@ -111,7 +114,10 @@ class ToolProxy: def _foreign(self, name): n = (name or "").upper() - return bool(RUN_PREFIX.match(n)) and not n.startswith(self.prefix) + if RUN_PREFIX.match(n): + return not n.startswith(self.prefix) + m = MID_PREFIX.match(n) + return bool(m) and m.group(1).upper() != self.prefix def _filter(self, tool, text): if tool == "sap_short_dumps": # only dumps of this run: other runs (and mutants) also write dumps @@ -124,7 +130,7 @@ class ToolProxy: data["count"] = len(data["dumps"]) return json.dumps(data) return text - if tool not in ("sap_search_object", "sap_usage_references"): + if tool not in ("sap_search_object", "sap_usage_references", "sap_inactive_objects"): return text try: data = json.loads(text) @@ -148,7 +154,9 @@ class ToolProxy: else: self.calls += 1 name = str(args.get("objectName", "")).upper() - if tool in WRITE_TOOLS and self._foreign(name): + if (tool in WRITE_TOOLS or tool in ("sap_pull_source", "sap_object_structure", "sap_object_members", "sap_element_info", + "sap_run_unit_test", "sap_check_object", "sap_syntax_check", "sap_atc_run")) \ + and self._foreign(name): result = (True, f"{name} is not available.") else: if tool in WRITE_TOOLS and (tool == "sap_activate" or args.get("activate", True)): diff --git a/harness/runner.py b/harness/runner.py index 760549d..8d5f512 100644 --- a/harness/runner.py +++ b/harness/runner.py @@ -153,9 +153,16 @@ class Runner: return out def _objects_with_prefix(self, mcp, prefix): - _, text = mcp.call("sap_search_object", {"query": prefix + "*", "maxResults": 200}) - return [d for d in (_json(text) or []) if d.get("name", "").upper().startswith(prefix) - and d.get("objectType")] # skips STOB entries of CDS entities + """Objects of the run: the name starts with the prefix, or contains it after a short type part (the model named its own + helper or test class ZCL__..., 27 such objects were left behind before 2026-10-06).""" + out = {} + for q in (prefix + "*", "*_" + prefix + "*"): + _, text = mcp.call("sap_search_object", {"query": q, "maxResults": 200}) + for d in (_json(text) or []): + n = d.get("name", "").upper() + if d.get("objectType") and (n.startswith(prefix) or re.match(r"^[A-Z]{1,5}_" + re.escape(prefix), n)): + out[n] = d # skips STOB entries of CDS entities (no objectType) + return list(out.values()) def _source(self, mcp, otype, name, fg=None): err, text = mcp.call("sap_pull_source", _obj_args(otype, name, fg)) diff --git a/harness/sweep.py b/harness/sweep.py new file mode 100644 index 0000000..d481ddd --- /dev/null +++ b/harness/sweep.py @@ -0,0 +1,89 @@ +"""Find (and delete) objects that harness runs left in A4H: names with a run prefix at the start (Z4CGT1HJ_X) or inside (ZCL_Z4CGT1HJ_X). + + python3 -m harness.sweep dry run: list them by prefix + python3 -m harness.sweep --delete delete them (ADT deletion API, dependency order); a locked object is reported, not forced + python3 -m harness.sweep --keep-recent 60 do not touch prefixes of runs whose directory changed in the last 60 minutes (default 60) + +Never run it while a controller or a series is writing: it can only judge by directory age. +""" +import argparse +import glob +import json +import os +import re +import time + +from .adt_client import load_env +from .mcp_client import McpClient +from .runner import DELETE_ORDER, delete_uris +from .task import prefix_for + +ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +START = re.compile(r"^(Z\d[0-9A-Z]{6}_)", re.I) +MID = re.compile(r"^[A-Z]{1,5}_(Z\d[0-9A-Z]{6}_)", re.I) + + +def recent_prefixes(minutes): + keep = set() + for d in glob.glob(os.path.join(ROOT, "runs", "**", "*_G*_*"), recursive=True) + glob.glob(os.path.join(ROOT, "runs", "**", "*_T*_*"), recursive=True): + b = os.path.basename(d) + m = re.match(r"^(\d+)_([GT]\d+)_", b) + if m and time.time() - os.path.getmtime(d) < minutes * 60: + keep.add(prefix_for(int(m.group(1)), m.group(2)).upper()) + return keep + + +def main(): + load_env(os.path.join(ROOT, ".env")) + ap = argparse.ArgumentParser() + ap.add_argument("--delete", action="store_true") + ap.add_argument("--keep-recent", type=int, default=60) + a = ap.parse_args() + found = {} + with McpClient() as m: + # run prefixes are Z + a digit + 6 characters: one query per digit (a single "Z*" query is cut at 1000 results), at the start and after the type part + for q in [f"Z{d}*" for d in "0123456789"] + [f"*_Z{d}*" for d in "0123456789"] + ["ZPROBE*", "ZTEST*"]: + try: + for o in json.loads(m.call("sap_search_object", {"query": q, "maxResults": 1000})[1]): + if o.get("objectType"): # a STOB entry of a CDS view has the same name and no objectType: skip it + found[(o["name"], o["objectType"])] = o + except ValueError: + pass + keep = recent_prefixes(a.keep_recent) + by_prefix = {} + for (n, _t), o in found.items(): + mm = START.match(n.upper()) or MID.match(n.upper()) + if not mm or not o.get("objectType"): + continue + pre = mm.group(1).upper() + if len(pre) != 9 or not pre[1].isdigit(): + continue + if pre in keep: + continue + by_prefix.setdefault(pre, []).append(o) + total = sum(len(v) for v in by_prefix.values()) + print(f"{total} objects of {len(by_prefix)} run prefixes ({len(keep)} recent prefixes kept)") + for pre, objs in sorted(by_prefix.items()): + print(" ", pre, [o["name"] for o in objs][:6]) + json.dump({p: [o["name"] for o in v] for p, v in by_prefix.items()}, open(os.path.join(ROOT, "runs", "sweep.json"), "w"), indent=1) + if not a.delete: + return + objs = sorted([o for v in by_prefix.values() for o in v], + key=lambda o: DELETE_ORDER.index(o["objectType"]) if o["objectType"] in DELETE_ORDER else 99) + done = failed = 0 + for i in range(0, len(objs), 10): + uris = [] + for o in objs[i:i + 10]: + uris.append(o["uri"]) + res = delete_uris(uris) + for u, v in res.items(): + if v["deleted"]: + done += 1 + else: + failed += 1 + print("not deleted:", u.split("/")[-1], v["msg"]) + print(f"deleted {done}, not deleted {failed}") + + +if __name__ == "__main__": + main() diff --git a/scripts_probe/lockprobe.py b/scripts_probe/lockprobe.py index 64e0df6..3361c3e 100644 --- a/scripts_probe/lockprobe.py +++ b/scripts_probe/lockprobe.py @@ -21,7 +21,7 @@ CLASS zprobe0eq_locks IMPLEMENTATION. TABLES enq = lt_enq EXCEPTIONS communication_failure = 1 system_failure = 2 OTHERS = 3. DATA(lv_n) = 0. - LOOP AT lt_enq ASSIGNING FIELD-SYMBOL() WHERE garg CS 'ZPROBE0' OR garg CS 'Z4AJ' OR garg CS 'Z4AE'. + LOOP AT lt_enq ASSIGNING FIELD-SYMBOL(). lv_n = lv_n + 1. DATA(lv_line) = ||. DO. @@ -35,15 +35,16 @@ CLASS zprobe0eq_locks IMPLEMENTATION. ENDDO. out->write( lv_line ). ENDLOOP. - out->write( |enqueue entries matching the probe prefixes: { lv_n } of { lines( lt_enq ) } (subrc { lv_subrc })| ). + out->write( |enqueue entries listed: { lv_n } of { lines( lt_enq ) } (subrc { lv_subrc })| ). ENDMETHOD. ENDCLASS. """ -def ensure_reader(m): +def ensure_reader(m, force=False): e, t = m.call("sap_search_object", {"query": READER}) - if READER not in t: - m.call("sap_create_object", {"objectType": "CLAS", "objectName": READER, "packageName": "$TMP", "description": "enqueue reader probe"}) + if READER not in t or force: + if READER not in t: + m.call("sap_create_object", {"objectType": "CLAS", "objectName": READER, "packageName": "$TMP", "description": "enqueue reader probe"}) e, t = m.call("sap_push_source", {"objectType": "CLAS", "objectName": READER, "source": READER_SRC}) if '"success":true' not in t.replace(" ", ""): print("reader activation:", t[:600]) @@ -54,5 +55,5 @@ def read_locks(m): if __name__ == "__main__": with McpClient() as m: - ensure_reader(m) - print(read_locks(m)[:2500]) + ensure_reader(m, force=True) + print(read_locks(m)[:3500]) diff --git a/train/STATE.md b/train/STATE.md index 89c5455..07f65a1 100644 --- a/train/STATE.md +++ b/train/STATE.md @@ -180,3 +180,5 @@ Training runs on HF Jobs with Unsloth, not on the Mac. No `mlx_lm` training. - 2026-10-05 22:35 **pipeline stopped at 21:43 by an infrastructure outage, found 22:26 (my miss).** The MCP server (127.0.0.1:3000) refused connections for a short time; three runs got `URLError(ConnectionRefusedError)` and the pipeline counted three equal harness errors and stopped (rule). Not a model or harness bug: A4H was up (11 h), the MCP server answers again. New `harness/infra.py`: an outage of MCP or A4H is waited out (check every 30 s, two answers in a row), the run's objects are deleted, the same run starts again; it is not counted as an event or a failure (`runs/traj/outages.jsonl`). Same for series A (applies after a restart of `harness.localqwen`; the running series keeps the old code and would skip the task). The 3 interrupted runs were cleaned (11 objects) and get their second attempt. Budget: panel 50.84 at ledger 113.36 (2.54 ledger per usage in the last stretch); guard recomputed: panel left 9.16 - 3.0 reserve = 6.2 usage x 1.2 = 7.4 ledger, `BUDGET_LIMIT_USD` 129 (guard 121). Dashboard: a stopped pipeline shows no burn rate or ETA. - 2026-10-05 22:40 cause of the 21:43 outage confirmed by Kral: he restarted Eclipse (the EPOD MCP server runs inside it). Three runs were writing at that moment; their objects could be deleted afterwards (no stale lock), so a restart of the MCP server in the middle of a write did not leak a lock in this case. Rule: restarting Eclipse is fine now (the pipeline waits and reruns), but check `python3 scripts_probe/lockprobe.py` for stale locks afterwards. + +- 2026-10-06 06:40 **correction of the 22:40 entry and a new finding.** (1) An Eclipse restart in the middle of a write does leak a lock: `ZCL_Z4CGT1HJ_JOB_COST_TEST` (enqueue SEOCLSENQ, mode X, created 21:46, three minutes after the 21:43 restart) stayed locked until the next morning; details in `docs/epod-lock-leak.md` (case 3). I had missed it because the teardown only cleaned names that start with the run prefix. SM12 is needed (Kral). (2) **Leftover objects with the prefix inside the name** (ZCL__..., the model's own test classes): 27 stayed in A4H, and **57 of 90 accepted trajectories contain other runs' object names in tool results** (mostly `sap_inactive_objects`, in 30 cases the model pulled another run's class). Fixed: `harness/proxy.py` hides them (search, usage, inactive list) and answers "not available" for reads and writes of another run's object; `harness/runner.py` finds mid-name objects for teardown; `harness/sweep.py` deleted the leftovers (28 objects; the locked class remains). The 90 trajectories are not changed yet: see the data scrub option of `train/build_stage2.py`.