#!/usr/bin/env python3 """M3 验收监督(加固版):供 Hermes cron 每 5 分钟调用(no_agent)。 状态机: idle -> phase1(--batch 30 20) -> phase2(--batch 50 100) -> done - flock 单实例防重入;抢不到锁静默退出 - 冻结运行: 首次启动把 Release 输出目录整体复制到 scripts/m3_frozen/_/, 记录全部文件 sha256 清单;两阶段都从冻结副本运行,后续构建不会覆盖运行中二进制 - 真实退出码: 子进程为 sh -c 'dotnet --batch K D; echo $? > exitfile', 结束后读 exitfile 得到真实退出码(含信号 128+sig;shell 被杀则无 exitfile -> interrupted) - PID 身份: 记录 /proc//stat 字段22 starttime,存活判定要求 pid+starttime+cmdline 全匹配(防 PID 复用) - PASS 判定收紧(全部满足才 PASS,任一不满足即 FAIL/停机): 1) exitfile 退出码 == 0 2) 日志含 `batch worlds=K days=D` 且 K/D 与该阶段参数完全一致 3) `world=` 行数 == K 且每行末 token 为 OK 4) `batch_summary worlds=K passed=K failed=0` 5) 行 `M3_ACCEPTANCE=PASS` 6) 冻结清单重新校验 sha256 全部一致,且清单 commit == 启动时记录 commit(commit 绑定) - FAIL/无结论/校验失败: 停机不自动重启,并在日志旁写 .diagnosis.txt 可诊断报告 - 状态文件 scripts/m3_state.json 原子写(tmp + os.replace) - --selftest: 全部合成测试(解析/判据校验/原子状态/锁互斥/PID身份解析/诊断报告/退出码读取), 不启动任何真实模拟,不伪造真实结果 """ import fcntl import hashlib import json import os import re import shlex import shutil import subprocess import sys import tempfile import time from pathlib import Path ROOT = Path(__file__).resolve().parent.parent SRC_BINARY = ROOT / "src/LivingVillage.Headless/bin/Release/net8.0/LivingVillage.Headless.dll" STATE_PATH = ROOT / "scripts/m3_state.json" LOCK_PATH = ROOT / "scripts/m3_supervise.lock" LOG_DIR = ROOT / "scripts/m3_logs" FROZEN_ROOT = ROOT / "scripts/m3_frozen" PHASES = {"phase1": (30, 20), "phase2": (50, 100)} NEXT = {"phase1": "phase2", "phase2": None} FRAG = "LivingVillage.Headless.dll" def git_commit(root=ROOT): try: r = subprocess.run(["git", "rev-parse", "HEAD"], cwd=root, capture_output=True, text=True, timeout=10) c = r.stdout.strip() return c if r.returncode == 0 and c else None except Exception: return None def sha256_file(path): h = hashlib.sha256() with open(path, "rb") as fh: for chunk in iter(lambda: fh.read(1 << 20), b""): h.update(chunk) return h.hexdigest() def freeze_release(commit, src_binary=SRC_BINARY, frozen_root=None): """复制 Release 输出目录全部文件到冻结目录并生成 sha256 清单。""" frozen_root = Path(frozen_root) if frozen_root else FROZEN_ROOT stamp = time.strftime("%Y%m%dT%H%M%S") dst = frozen_root / f"{stamp}_{(commit or 'nocommit')[:8]}" dst.mkdir(parents=True, exist_ok=False) files = {} for f in sorted(Path(src_binary).parent.iterdir()): if f.is_file(): shutil.copy2(f, dst / f.name) files[f.name] = sha256_file(dst / f.name) manifest = {"commit": commit, "created": time.strftime("%Y-%m-%dT%H:%M:%S"), "source": str(src_binary), "files": files} mpath = dst / "manifest.json" mpath.write_text(json.dumps(manifest, indent=1, ensure_ascii=False) + "\n") return dst, mpath, manifest def verify_manifest(frozen_dir, manifest_path, expected_commit): """重校验冻结目录: commit 绑定 + 全部文件 sha256。返回问题列表(空=通过)。""" problems = [] try: m = json.loads(Path(manifest_path).read_text()) except Exception as e: return [f"manifest unreadable: {e}"] if m.get("commit") != expected_commit: problems.append(f"commit mismatch manifest={m.get('commit')} expected={expected_commit}") files = m.get("files") or {} if not files: problems.append("manifest has no files") for name, want in sorted(files.items()): p = Path(frozen_dir) / name if not p.exists(): problems.append(f"frozen file missing: {name}") continue got = sha256_file(p) if got != want: problems.append(f"hash mismatch {name}: got {got[:16]} want {want[:16]}") return problems def load_state(path=STATE_PATH): if Path(path).exists(): try: return json.loads(Path(path).read_text()) except Exception: return {"phase": "corrupt_state"} return {"phase": "idle"} def save_state(state, path=STATE_PATH): """原子写: 同目录 tmp + os.replace。""" p = Path(path) p.parent.mkdir(parents=True, exist_ok=True) fd, tmp = tempfile.mkstemp(dir=str(p.parent), prefix=".m3_state.", suffix=".tmp") try: with os.fdopen(fd, "w") as fh: fh.write(json.dumps(state, indent=1, ensure_ascii=False) + "\n") os.replace(tmp, p) except BaseException: try: os.unlink(tmp) except OSError: pass raise def read_starttime(pid, proc_root="/proc"): """/proc//stat 字段22 starttime(boot 后时钟滴答)。解析失败返回 None。""" try: data = (Path(proc_root) / str(pid) / "stat").read_text() except OSError: return None rp = data.rfind(")") if rp < 0: return None parts = data[rp + 1:].split() # parts[i] 对应字段 i+3;字段22 -> parts[19] if len(parts) < 20: return None return parts[19] def read_cmdline(pid, proc_root="/proc"): try: return (Path(proc_root) / str(pid) / "cmdline").read_bytes().decode("utf-8", "replace").replace("\0", " ") except OSError: return "" def proc_alive(pid, start, frag, proc_root="/proc"): """pid + starttime + cmdline 片段三者全匹配才算同一进程(防 PID 复用)。""" if not pid or start is None: return False now = read_starttime(pid, proc_root) if now != str(start): return False return frag in read_cmdline(pid, proc_root) def read_exitcode(exitfile): p = Path(exitfile) if not p.exists(): return None try: return int(p.read_text().strip()) except ValueError: return None def verify_batch(text, worlds, days): """完整批次判据校验。返回问题列表(空=满足全部收紧条件)。""" problems = [] m = re.search(r"^batch worlds=(\d+) days=(\d+) ", text, re.M) if not m: problems.append("missing 'batch worlds=K days=D' header line") else: if int(m.group(1)) != worlds or int(m.group(2)) != days: problems.append(f"batch params mismatch: log worlds={m.group(1)} days={m.group(2)}, expected {worlds}x{days}") wlines = [l for l in text.splitlines() if l.startswith("world=")] if len(wlines) != worlds: problems.append(f"world line count {len(wlines)} != expected {worlds}") bad = [l.split()[0] for l in wlines if not l.split() or l.split()[-1] != "OK"] if bad: problems.append(f"non-OK world lines: {bad[:8]}{'...' if len(bad) > 8 else ''}") s = re.search(r"^batch_summary worlds=(\d+) passed=(\d+) failed=(\d+) ", text, re.M) if not s: problems.append("missing batch_summary line") else: if (int(s.group(1)), int(s.group(2)), int(s.group(3))) != (worlds, worlds, 0): problems.append(f"batch_summary mismatch: worlds={s.group(1)} passed={s.group(2)} failed={s.group(3)}, expected {worlds}/{worlds}/0") if not re.search(r"^M3_ACCEPTANCE=PASS$", text, re.M): problems.append("missing line M3_ACCEPTANCE=PASS") return problems def generate_diagnosis(state, problems, tag): """可诊断报告文本(纯函数,便于合成测试)。""" log = state.get("log") tail = "" if log and Path(log).exists(): lines = Path(log).read_text(errors="replace").splitlines() tail = "\n".join(lines[-60:]) exitcode = state.get("exitcode") return ( f"M3 DIAGNOSIS {tag}\n" f"time: {time.strftime('%Y-%m-%dT%H:%M:%S')}\n" f"phase: {state.get('phase')}\n" f"expected batch: {state.get('worlds')}x{state.get('days')}\n" f"commit: {state.get('commit')}\n" f"frozen dir: {state.get('frozen')}\n" f"manifest: {state.get('manifest')}\n" f"pid: {state.get('pid')} starttime: {state.get('start')}\n" f"exit code: {exitcode}\n" f"log: {log}\n" f"problems:\n" + "".join(f" - {p}\n" for p in problems) + (f"\nlog tail (last 60 lines):\n{tail}\n" if tail else "\nlog missing or empty\n") ) def write_diagnosis(state, problems, tag): log = state.get("log") if not log: return None rep = Path(str(log) + ".diagnosis.txt") rep.write_text(generate_diagnosis(state, problems, tag)) return rep def build_sh_cmd(dll, batch_args, exitfile): return f"exec 2>&1; {shlex.quote('dotnet')} {shlex.quote(str(dll))} {batch_args}; echo $? > {shlex.quote(str(exitfile))}" def launch(phase, state): worlds, days = PHASES[phase] frozen, manifest, commit = state.get("frozen"), state.get("manifest"), state.get("commit") dll = Path(frozen) / FRAG if frozen else None if not dll or not dll.exists() or not manifest or not commit: return "freeze-missing" LOG_DIR.mkdir(parents=True, exist_ok=True) stamp = time.strftime("%Y%m%dT%H%M%S") log = LOG_DIR / f"{phase}_{stamp}.log" exitfile = Path(str(log) + ".exit") cmd = ["sh", "-c", build_sh_cmd(dll, f"--batch {worlds} {days}", exitfile)] with open(log, "w") as fh: proc = subprocess.Popen(cmd, stdout=fh, stderr=subprocess.STDOUT, stdin=subprocess.DEVNULL, start_new_session=True, cwd=ROOT) start = read_starttime(proc.pid) if start is None: # 无法建立进程身份: 立即终止并按失败停机,不进入无身份运行 try: proc.kill() except OSError: pass state.update({"phase": "failed", "worlds": worlds, "days": days, "pid": proc.pid, "log": str(log), "verdict": None, "exitcode": None, "problems": [f"cannot read starttime for new pid {proc.pid}; launch aborted"]}) write_diagnosis(state, state["problems"], "LAUNCH_ABORTED") save_state(state) return f"launch-aborted {phase} pid={proc.pid} (no proc identity); stopped" state.update({ "phase": phase, "worlds": worlds, "days": days, "pid": proc.pid, "start": start, "frag": FRAG, "cmd": " ".join(cmd), "log": str(log), "exitfile": str(exitfile), "frozen": str(frozen), "manifest": str(manifest), "commit": commit, "started": time.strftime("%Y-%m-%dT%H:%M:%S"), "verdict": None, "exitcode": None, "problems": []}) save_state(state) return f"started {phase} pid={proc.pid} start={start} batch={worlds}x{days} frozen={frozen} log={log}" def decide(state, alive): """纯状态机转移。alive(pid,start,frag)->bool 注入便于合成测试。""" phase = state.get("phase", "idle") if phase == "done": return None, state, None if phase in PHASES: if state.get("pid") and alive(state["pid"], state.get("start"), state.get("frag") or FRAG): return None, state, None # 进程结束: 收集全部证据 state["exitcode"] = read_exitcode(state.get("exitfile")) log = state.get("log") p = Path(log) if log else None text = p.read_text(errors="replace") if p and p.exists() else "" worlds, days = state.get("worlds"), state.get("days") problems = [] if worlds is None or days is None: problems.append("state missing worlds/days (legacy or corrupt state)") else: problems += verify_batch(text, worlds, days) problems += verify_manifest(state.get("frozen"), state.get("manifest"), state.get("commit")) state["problems"] = problems state["pid"] = None if state["exitcode"] is None: state["phase"] = "interrupted" return "save", state, f"{phase} interrupted without exit code (log={log}); stopped, awaiting human" if state["exitcode"] == 0 and not problems: nxt = NEXT[phase] state["verdict"] = "PASS" if nxt is None: state["phase"] = "done" return "save", state, f"{phase} PASS (commit={state.get('commit')}, frozen={state.get('frozen')}) -> M3 acceptance complete (30x20 + 50x100)" return f"launch:{nxt}", state, None state["verdict"] = "FAIL" state["phase"] = "failed" why = f"exit={state['exitcode']}" + (f" problems={problems}" if problems else "") return "save", state, f"{phase} FAIL ({why}; log={log}); stopped, awaiting human" if phase == "idle": return "launch:phase1", state, None return None, state, None def acquire_lock(path=LOCK_PATH): Path(path).parent.mkdir(parents=True, exist_ok=True) fd = os.open(str(path), os.O_RDWR | os.O_CREAT, 0o644) try: fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB) except BlockingIOError: os.close(fd) return None return fd def selftest(): ok = True def check(name, cond): nonlocal ok print(("PASS " if cond else "FAIL ") + name) ok = ok and cond # ---- 解析 ---- check("parse PASS", verify_batch("x\nM3_ACCEPTANCE=PASS\n", 1, 1) == ["missing 'batch worlds=K days=D' header line", "world line count 0 != expected 1", "missing batch_summary line"]) check("parse empty", verify_batch("world=0 ...\n", 1, 1) != []) check("verify full PASS log", verify_batch( "batch worlds=2 days=20 half_life_ticks=1 rel_threshold=0.50\n" "world=0 seed=42 ... OK\nworld=1 seed=43 ... OK\n" "batch_summary worlds=2 passed=2 failed=0 old_passed=2/2 ratio_min=1.0 ratio_max=1.0 gini_min=0.1 gini_max=0.1 elapsed_s=1.0\n" "M3_ACCEPTANCE=PASS\n", 2, 20) == []) check("verify wrong params", any("params mismatch" in p for p in verify_batch( "batch worlds=3 days=20 x\nworld=0 a OK\nworld=1 b OK\nworld=2 c OK\nbatch_summary worlds=3 passed=3 failed=0 x\nM3_ACCEPTANCE=PASS\n", 2, 20))) check("verify BAD world", any("non-OK" in p for p in verify_batch( "batch worlds=1 days=1 x\nworld=0 a BAD\nbatch_summary worlds=1 passed=0 failed=1 x\nM3_ACCEPTANCE=FAIL\n", 1, 1))) check("verify truncated worlds", any("world line count" in p for p in verify_batch( "batch worlds=3 days=1 x\nworld=0 a OK\nbatch_summary worlds=3 passed=3 failed=0 x\nM3_ACCEPTANCE=PASS\n", 3, 1))) check("verify summary mismatch", any("batch_summary mismatch" in p for p in verify_batch( "batch worlds=2 days=1 x\nworld=0 a OK\nworld=1 b OK\nbatch_summary worlds=2 passed=1 failed=1 x\nM3_ACCEPTANCE=PASS\n", 2, 1))) check("verify missing verdict token", any("M3_ACCEPTANCE=PASS" in p for p in verify_batch( "batch worlds=1 days=1 x\nworld=0 a OK\nbatch_summary worlds=1 passed=1 failed=0 x\nM3_ACCEPTANCE=FAIL\n", 1, 1))) # ---- 冻结清单(合成目录) ---- with tempfile.TemporaryDirectory() as td: td = Path(td) src = td / "bin"; src.mkdir() (src / "a.dll").write_bytes(b"hello") (src / "b.json").write_text("{}") dst, mpath, man = freeze_release("deadbeef" * 8, src / "a.dll", td / "frozen") check("freeze copies deps", (dst / "b.json").exists() and (dst / "a.dll").read_bytes() == b"hello" and man["files"]["b.json"]) check("manifest verify ok", verify_manifest(dst, mpath, "deadbeef" * 8) == []) check("manifest commit bind", verify_manifest(dst, mpath, "f" * 64) != []) (dst / "a.dll").write_bytes(b"tampered") check("manifest detect tamper", any("hash mismatch" in p for p in verify_manifest(dst, mpath, "deadbeef" * 8))) # ---- 原子状态 ---- with tempfile.TemporaryDirectory() as td: sp = Path(td) / "state.json" save_state({"phase": "phase1", "n": 1}, sp) check("atomic roundtrip", load_state(sp) == {"phase": "phase1", "n": 1}) check("no tmp leftovers", not [q for q in Path(td).iterdir() if q.name.endswith(".tmp")]) # ---- flock 互斥(同进程第二个 fd 抢锁必须失败) ---- with tempfile.TemporaryDirectory() as td: l1 = acquire_lock(Path(td) / "l.lock") l2 = acquire_lock(Path(td) / "l.lock") check("flock blocks second holder", l1 is not None and l2 is None) if l1 is not None: fcntl.flock(l1, fcntl.LOCK_UN) l3 = acquire_lock(Path(td) / "l.lock") check("flock released", l3 is not None) # ---- PID 身份(合成 /proc) ---- with tempfile.TemporaryDirectory() as td: td = Path(td) pr = td / "proc" / "4242" pr.mkdir(parents=True) (pr / "stat").write_text("4242 (sh) S 1 4242 4242 0 -1 4194304 0 0 0 0 0 0 0 0 0 0 0 0 12345 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0 0\n") (pr / "cmdline").write_bytes(b"sh -c dotnet LivingVillage.Headless.dll --batch 30 20\x00") check("starttime parsed", read_starttime(4242, td / "proc") == "12345") check("alive match", proc_alive(4242, "12345", "LivingVillage.Headless.dll", td / "proc")) check("starttime mismatch rejected", not proc_alive(4242, "99999", "LivingVillage.Headless.dll", td / "proc")) check("frag mismatch rejected", not proc_alive(4242, "12345", "other.dll", td / "proc")) check("missing proc rejected", not proc_alive(9999, "12345", "LivingVillage.Headless.dll", td / "proc")) # ---- 退出码读取 ---- with tempfile.TemporaryDirectory() as td: ef = Path(td) / "x.exit" check("exit missing -> None", read_exitcode(ef) is None) ef.write_text("3\n") check("exit read", read_exitcode(ef) == 3) ef.write_text("junk") check("exit junk -> None", read_exitcode(ef) is None) # ---- sh 命令构造 ---- c = build_sh_cmd("/tmp/a b/LivingVillage.Headless.dll", "--batch 30 20", "/tmp/x y.exit") check("sh cmd quotes paths", "'/tmp/a b/LivingVillage.Headless.dll'" in c and "'/tmp/x y.exit'" in c and c.rstrip().endswith("> '/tmp/x y.exit'")) # ---- 状态机(合成日志/注入存活函数) ---- with tempfile.TemporaryDirectory() as td: td = Path(td) plog = td / "p.log" plog.write_text( "batch worlds=1 days=1 x\nworld=0 a OK\nbatch_summary worlds=1 passed=1 failed=0 x\nM3_ACCEPTANCE=PASS\n") base = {"phase": "phase1", "worlds": 1, "days": 1, "pid": 123, "start": "1", "frag": FRAG, "cmd": "c", "log": str(plog), "exitfile": str(td / "e.exit"), "frozen": str(td / "fr"), "manifest": str(td / "fr/manifest.json"), "commit": "c" * 40, "verdict": None, "exitcode": None, "started": None, "problems": []} (td / "fr").mkdir() (td / "fr/k.dll").write_bytes(b"kernel") (td / "fr/manifest.json").write_text(json.dumps( {"commit": "c" * 40, "files": {"k.dll": hashlib.sha256(b"kernel").hexdigest()}})) a, s, _ = decide(dict(base), lambda pid, st, fr: True) check("alive -> silent", a is None) (td / "e.exit").write_text("0\n") a, s, _ = decide(dict(base), lambda pid, st, fr: False) check("phase1 PASS -> launch phase2", a == "launch:phase2" and s["verdict"] == "PASS") a, s, m = decide(dict(base, phase="phase2"), lambda pid, st, fr: False) check("phase2 PASS -> done", s["phase"] == "done" and m and "complete" in m) (td / "e.exit").write_text("1\n") a, s, m = decide(dict(base), lambda pid, st, fr: False) check("exit 1 -> FAIL stop", s["phase"] == "failed" and s["verdict"] == "FAIL" and "FAIL" in m) (td / "e.exit").unlink() a, s, _ = decide(dict(base), lambda pid, st, fr: False) check("no exitcode -> interrupted", s["phase"] == "interrupted") (td / "e.exit").write_text("0\n") (td / "fr/manifest.json").write_text(json.dumps( {"commit": "d" * 40, "files": {"k.dll": hashlib.sha256(b"kernel").hexdigest()}})) a, s, m = decide(dict(base), lambda pid, st, fr: False) check("commit mismatch -> FAIL even with token", s["phase"] == "failed" and any("commit" in p for p in s["problems"])) a, s, _ = decide({"phase": "done"}, lambda pid, st, fr: False) check("done -> silent", a is None) rep = generate_diagnosis(dict(base, phase="failed", exitcode=0, problems=["boom"]), ["boom"], "FAIL") check("diagnosis contains evidence", "M3 DIAGNOSIS FAIL" in rep and "boom" in rep and "exit code: 0" in rep) print("SELFTEST=" + ("PASS" if ok else "FAIL")) return 0 if ok else 1 def main(): if "--selftest" in sys.argv: return selftest() lock = acquire_lock() if lock is None: return 0 # 已有实例在运行: 静默退出 try: state = load_state() if state.get("phase") == "corrupt_state": print("state file corrupt; stopped, awaiting human") return 1 if state.get("phase") == "idle": if not SRC_BINARY.exists(): print(f"binary missing: {SRC_BINARY}; run dotnet build src/LivingVillage.Headless -c Release") return 1 commit = git_commit() if not commit: print("cannot determine git commit; refusing to launch unbound run") return 1 frozen, manifest, _ = freeze_release(commit) state.update({"frozen": str(frozen), "manifest": str(manifest), "commit": commit}) action, state, msg = decide(state, proc_alive) if action and action.startswith("launch:"): r = launch(action.split(":", 1)[1], state) if r == "freeze-missing": print("frozen release missing at launch; stopped, awaiting human") save_state({**state, "phase": "failed", "problems": ["frozen release/manifest/commit missing at launch"]}) write_diagnosis(state, ["frozen release missing"], "LAUNCH_ABORTED") return 1 msg = r elif action == "save": save_state(state) if state.get("phase") in ("failed", "interrupted"): write_diagnosis(state, state.get("problems") or [], state.get("phase", "").upper()) if msg: print(msg) return 0 finally: try: fcntl.flock(lock, fcntl.LOCK_UN) os.close(lock) except OSError: pass if __name__ == "__main__": sys.exit(main())