#!/usr/bin/env python3 """M3 验收监督(加固版):供 Hermes cron 每 5 分钟调用(no_agent)。 状态机: idle -> phase1(--batch 30 20) -> phase2(--batch 50 20) -> done - flock 单实例防重入;抢不到锁静默退出 - 冻结运行: 首次启动把 Release 输出目录递归复制到 scripts/m3_frozen/_/, v2 清单覆盖全部子目录文件 sha256;两阶段都从冻结副本运行,后续构建不会覆盖 - 真实退出码: 子进程为 sh -c 'dotnet --batch K D; echo $? > exitfile', 读 exitfile 得真实退出码(信号 128+sig;shell 被杀无 exitfile -> interrupted) - PID 身份: pid + /proc starttime + cmdline 片段全匹配(防 PID 复用) - PASS 判定收紧(全部满足才 PASS,任一不满足即 FAIL/停机): 1) exitfile 退出码 == 0 2) `batch worlds=K days=D` 与固定 PHASES[phase] 完全一致(不用 state 里可篡改值) 3) `world=` 行: 行数==K、索引恰为 0..K-1 无重复无缺失(允许流式完成序)、seed==42+i; 行尾 `elapsed_s=<秒>` 为追加观测字段,不参与判据 4) `batch_summary worlds=K passed=K failed=0` 5) 行 `M3_ACCEPTANCE=PASS` 6) 冻结清单 sha256 全部一致 + commit 绑定;v2 清单另查额外/缺失文件 - shell 死亡留下存活 dotnet(孤儿): interrupted 停机并记录孤儿 pid,绝不自动重启新 batch - FAIL/无结论/校验失败: 停机不自动重启,写 .diagnosis.txt 可诊断报告 - 状态文件原子写(tmp + os.replace);旧 v1 清单只校验所列文件(兼容正在运行的旧冻结副本, 不修改任何旧冻结目录) - --selftest: 合成断言 + 真实隔离生命周期测试(临时目录 + 短真实子进程: 退出非零/重复调用/控制器重启/子进程中断/孤儿防重启),不伪造任何 M3 结果, 不触碰项目真实 state/log/frozen """ import contextlib import fcntl import hashlib import io import json import os import re import shlex import shutil import signal 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, 20)} NEXT = {"phase1": "phase2", "phase2": None} FRAG = "LivingVillage.Headless.dll" def git_commit(root=None): root = ROOT if root is None else Path(root) try: r = subprocess.run(["git", "rev-parse", "HEAD"], cwd=str(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=None, frozen_root=None): """递归复制 Release 输出目录(含子目录)到冻结目录,生成 v2 sha256 清单。""" frozen_root = Path(frozen_root) if frozen_root else FROZEN_ROOT src = Path(src_binary) if src_binary else SRC_BINARY stamp = time.strftime("%Y%m%dT%H%M%S") base = frozen_root / f"{stamp}_{(commit or 'nocommit')[:8]}" dst = base suffix = 1 while True: try: dst.mkdir(parents=True, exist_ok=False) break except FileExistsError: dst = base.parent / f"{base.name}_{suffix}" suffix += 1 files = {} for f in sorted(src.parent.rglob("*")): if f.is_file(): rel = f.relative_to(src.parent).as_posix() target = dst / rel target.parent.mkdir(parents=True, exist_ok=True) shutil.copy2(f, target) files[rel] = sha256_file(target) manifest = {"version": 2, "commit": commit, "created": time.strftime("%Y-%m-%dT%H:%M:%S"), "source": str(src), "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): """重校验冻结目录。v2: commit 绑定 + 所列文件 sha256 + 额外/缺失文件检测 (manifest.json 自身除外)。旧 v1 清单(无 version): 只校验所列文件,兼容运行中旧副本。 返回问题列表(空=通过)。""" 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]}") if m.get("version") == 2: for f in Path(frozen_dir).rglob("*"): if not f.is_file(): continue rel = f.relative_to(Path(frozen_dir)).as_posix() if rel == "manifest.json" or rel in files: continue problems.append(f"extra file not in manifest: {rel}") return problems def load_state(path=None): # path 延迟到调用时解析(模块级常量可能被测试替换; 默认参数在 def 时绑定是已知事故源) p = Path(path) if path is not None else STATE_PATH if p.exists(): try: return json.loads(p.read_text()) except Exception: return {"phase": "corrupt_state"} return {"phase": "idle"} def save_state(state, path=None): """原子写: 同目录 tmp + os.replace。path 延迟到调用时解析。""" p = Path(path) if path is not None else STATE_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_proc_stat(pid, proc_root="/proc"): """读取 /proc//stat 的字段(字段3之后从 parts[0] 开始)。""" 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() if len(parts) < 20: return None return parts def read_starttime(pid, proc_root="/proc"): """/proc//stat 字段22 starttime(boot 后时钟滴答)。解析失败返回 None。""" parts = _read_proc_stat(pid, proc_root) if parts is None: return None # parts[i] 对应字段 i+3;字段22 -> parts[19] 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 parts = _read_proc_stat(pid, proc_root) if parts is None or parts[0] == "Z": return False if parts[19] != str(start): return False return frag in read_cmdline(pid, proc_root) def find_orphans(marker, proc_root="/proc"): """扫描 /proc,返回 cmdline 含 marker 的存活 pid(孤儿检测)。""" out = [] if not marker: return out for d in Path(proc_root).iterdir(): if not d.name.isdigit(): continue if marker in read_cmdline(int(d.name), proc_root): out.append(int(d.name)) return out def live_batch_orphans(proc_root="/proc"): """state 丢失/_idle 时兜底: 存活 cmdline 同时含冻结 dll 片段与 --batch 的进程。 命中则绝不允许再拉起新批次(防止 state 被破坏后重复 launch)。排除自身。""" out = [] for d in Path(proc_root).iterdir(): if not d.name.isdigit() or int(d.name) == os.getpid(): continue cl = read_cmdline(int(d.name), proc_root) if FRAG in cl and "--batch" in cl: out.append(int(d.name)) return out 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 world_line_verdict(line): """world 行的判据 token;忽略行尾追加的 `elapsed_s=<秒>` 观测字段。""" tokens = line.split() while tokens and re.fullmatch(r"elapsed_s=\d+(\.\d+)?", tokens[-1]): tokens.pop() return tokens[-1] if tokens else "" def verify_batch(text, worlds, days): """完整批次判据校验(固定 worlds/days 由调用方从 PHASES 取)。 返回问题列表(空=满足全部收紧条件)。""" 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}") parsed = [] malformed = [] for l in wlines: mw = re.match(r"^world=(\d+) seed=(\d+) ", l) if not mw: malformed.append((l.split() or ["?"])[0]) else: parsed.append((int(mw.group(1)), int(mw.group(2)))) if malformed: problems.append(f"malformed world lines: {malformed[:5]}") idxs = [i for i, _ in parsed] # 流式输出允许完成序非升序,判据是索引集合恰为 0..K-1(无重无缺)。 if sorted(idxs) != list(range(worlds)): dup = sorted({i for i in idxs if idxs.count(i) > 1}) missing = sorted(set(range(worlds)) - set(idxs)) extra = sorted(set(idxs) - set(range(worlds))) problems.append(f"world index set mismatch: duplicates={dup[:5]} missing={missing[:5]} extra={extra[:5]} (count {len(idxs)} vs {worlds})") badseed = [i for i, s in parsed if s != 42 + i] if badseed: problems.append(f"seed binding mismatch (expect seed=42+i): worlds={badseed[:8]}") bad = [l.split()[0] for l in wlines if world_line_verdict(l) != "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:]) orphans = state.get("orphans") or [] 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: fixed PHASES={PHASES.get(state.get('phase'))} state={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: {state.get('exitcode')}\n" f"orphan child pids (shell dead, child alive): {orphans if orphans else 'none'}\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 runner_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, runner=None): worlds, days = PHASES[phase] if runner is None: runner = runner_cmd 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") base_log = LOG_DIR / f"{phase}_{stamp}.log" log = base_log suffix = 1 # A controller restart can happen within the same second; never reuse a # previous log or its exit marker when that happens. while (log.exists() or Path(str(log) + ".exit").exists() or Path(str(log) + ".diagnosis.txt").exists()): log = LOG_DIR / f"{phase}_{stamp}_{suffix}.log" suffix += 1 exitfile = Path(str(log) + ".exit") cmd = ["sh", "-c", runner(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=str(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": [], "orphans": []}) 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 注入便于合成测试。 批次参数固定取 PHASES[phase];state 中的 worlds/days 仅作篡改检测,不作判据。""" 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 # 进程结束: 收集全部证据 problems = [] worlds, days = PHASES[phase] if state.get("worlds") != worlds or state.get("days") != days: problems.append(f"state phase params tampered: state={state.get('worlds')}x{state.get('days')}, expected fixed {worlds}x{days}") 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 "" 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["verdict"] = "FAIL" state["problems"].insert(0, "missing exit code: process ended or disappeared before exitfile was written") state["orphans"] = find_orphans(state.get("frozen") or "") state["phase"] = "interrupted" extra = f"; orphan child pids alive: {state['orphans']}" if state["orphans"] else "" return "save", state, f"{phase} interrupted without exit code (log={log}){extra}; 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 + 50x20)" 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=None): p = Path(path) if path is not None else LOCK_PATH p.parent.mkdir(parents=True, exist_ok=True) fd = os.open(str(p), 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 make_batch_log(worlds, days, ok_token="OK"): lines = [f"batch worlds={worlds} days={days} half_life_ticks=1 rel_threshold=0.50"] for i in range(worlds): ok = "OK" if ok_token == "OK" else "BAD" lines.append(f"world={i} seed={42 + i} chats_max=10 chats_avg=1.00 ratio=3.000 gini=0.100 relcnt_max=2 relcnt_min=0 relconn=5/5 nonfinite=0 check_a=PASS check_b=PASS check_c=PASS old_crit=PASS {ok}") lines.append(f"batch_summary worlds={worlds} passed={worlds if ok_token == 'OK' else 0} failed={0 if ok_token == 'OK' else worlds} old_passed={worlds}/{worlds} ratio_min=1.0 ratio_max=1.0 gini_min=0.1 gini_max=0.1 elapsed_s=1.0") lines.append(f"M3_ACCEPTANCE={'PASS' if ok_token == 'OK' else 'FAIL'}") return "\n".join(lines) + "\n" def _run_main(): buf = io.StringIO() argv_saved = sys.argv sys.argv = [argv_saved[0]] # 剥离 --selftest, 防止 main() 递归重入 selftest try: with contextlib.redirect_stdout(buf): code = main() finally: sys.argv = argv_saved return code, buf.getvalue() def lifecycle_section(check): """真实隔离生命周期测试: 临时目录 + 短真实子进程。不伪造 M3 结果。""" global ROOT, SRC_BINARY, STATE_PATH, LOCK_PATH, LOG_DIR, FROZEN_ROOT, FRAG with tempfile.TemporaryDirectory() as tds: root = Path(tds) / "proj" srcbin = root / "bin" srcbin.mkdir(parents=True) test_frag = f"M3SelftestHeadless_{os.getpid()}_{time.time_ns()}.dll" (srcbin / test_frag).write_text("fake dll payload marker") (srcbin / "runtimes").mkdir() (srcbin / "runtimes/dep.dll").write_text("dep") env = {**os.environ, "GIT_AUTHOR_NAME": "t", "GIT_AUTHOR_EMAIL": "t@t", "GIT_COMMITTER_NAME": "t", "GIT_COMMITTER_EMAIL": "t@t"} (root / "f.txt").write_text("x") subprocess.run(["git", "init", "-q"], cwd=str(root), check=True) subprocess.run(["git", "add", "-A"], cwd=str(root), check=True) subprocess.run(["git", "commit", "-qm", "init"], cwd=str(root), check=True, env=env) commit = git_commit(root) check("lifecycle tmp git commit", bool(commit)) saved = (ROOT, SRC_BINARY, STATE_PATH, LOCK_PATH, LOG_DIR, FROZEN_ROOT, FRAG, runner_cmd) try: ROOT = root FRAG = test_frag SRC_BINARY = srcbin / FRAG STATE_PATH = root / "scripts/m3_state.json" LOCK_PATH = root / "scripts/lock" LOG_DIR = root / "scripts/logs" FROZEN_ROOT = root / "frozen" frozen, mpath, man = freeze_release(commit) check("freeze recursive incl subdirs", (frozen / "runtimes/dep.dll").exists() and man["files"].get("runtimes/dep.dll") is not None) check("v2 manifest verify ok", verify_manifest(frozen, mpath, commit) == []) same_second_root = root / "same-second-frozen" old_strftime = time.strftime same_second_error = None same_second_dirs = [] try: time.strftime = lambda *args, **kwargs: "20260919T200000" same_second_dirs.append(freeze_release(commit, frozen_root=same_second_root)[0]) same_second_dirs.append(freeze_release(commit, frozen_root=same_second_root)[0]) except Exception as exc: same_second_error = exc finally: time.strftime = old_strftime check("freeze same-second paths unique", same_second_error is None and len(same_second_dirs) == 2 and same_second_dirs[0] != same_second_dirs[1]) (frozen / "extra.bin").write_text("e") check("v2 extra file detected", any("extra file" in p for p in verify_manifest(frozen, mpath, commit))) (frozen / "extra.bin").unlink() (frozen / "runtimes/dep.dll").unlink() check("v2 missing file detected", any("frozen file missing" in p for p in verify_manifest(frozen, mpath, commit))) # legacy v1 清单兼容: 所列文件通过即可,不做额外文件检测(兼容运行中的旧冻结副本) (srcbin / "runtimes/dep.dll").write_text("dep") leg = frozen / "legacy_dir" leg.mkdir() (leg / "a.dll").write_text("A") legman = leg / "manifest.json" legman.write_text(json.dumps({"commit": commit, "files": {"a.dll": hashlib.sha256(b"A").hexdigest()}})) (leg / "unlisted.txt").write_text("u") check("legacy manifest verify ok (no extra check)", verify_manifest(leg, legman, commit) == []) def long_runner(dll, batch_args, exitfile): q = shlex.quote py = q(sys.executable) code = q("import subprocess, time; subprocess.Popen(['sleep', '2']); time.sleep(2)") return f"test -f {q(str(dll))}; {py} -c {code} {q(str(dll))}; echo $? > {q(str(exitfile))}" def fail_runner(dll, batch_args, exitfile): q = shlex.quote return f"sh -c 'exit 7'; echo $? > {q(str(exitfile))}" n_logs = lambda: len(list(LOG_DIR.glob("*.log"))) # (a) 控制器启动 phase1(真实 sh + Python 子进程, ~2s 存活) save_state({"phase": "idle"}) globals()["runner_cmd"] = long_runner code, printed = _run_main() check("controller launches phase1", code == 0 and "started phase1" in printed) st = load_state() check("launch state bound to frozen+commit", st["commit"] == commit and st["worlds"] == 30 and st["days"] == 20) # (b) 运行中重复调用: 静默, 不重复拉起 code, printed = _run_main() check("repeat call silent while alive", code == 0 and printed == "" and n_logs() == 1) # (c) 子进程正常退出但日志无有效批次 -> fail-closed FAIL 停机 time.sleep(2.3) code, printed = _run_main() st = load_state() diag = Path(str(st["log"]) + ".diagnosis.txt") check("exit0 without valid batch -> FAIL stop", st["phase"] == "failed" and st["verdict"] == "FAIL" and code == 0 and diag.exists()) # (d) 真实退出码 7 传播 save_state({"phase": "idle"}) globals()["runner_cmd"] = fail_runner _run_main() time.sleep(0.4) _run_main() st = load_state() check("real child exit 7 propagated", st["exitcode"] == 7 and st["phase"] == "failed") # (e) 子进程被杀: shell 死 -> interrupted + 孤儿记录 save_state({"phase": "idle"}) globals()["runner_cmd"] = long_runner _run_main() st = load_state() time.sleep(0.2) os.kill(st["pid"], signal.SIGKILL) # 杀外层 sh, 内层 sleep sh 成为孤儿 deadline = time.monotonic() + 1.5 while proc_alive(st["pid"], st["start"], st["frag"]) and time.monotonic() < deadline: time.sleep(0.02) while not find_orphans(str(st["frozen"])) and time.monotonic() < deadline: time.sleep(0.02) _run_main() st2 = load_state() check("killed shell -> interrupted FAILED", st2["phase"] == "interrupted" and st2["verdict"] == "FAIL") check("orphan recorded", bool(st2.get("orphans"))) check("interrupted diagnosis written", Path(str(st2["log"]) + ".diagnosis.txt").exists()) # (f) 控制器重启后绝不自动重启新 batch(孤儿/中断均终态) logs_before = n_logs() code, printed = _run_main() st3 = load_state() check("no relaunch after interrupted (controller restart)", st3["phase"] == "interrupted" and n_logs() == logs_before and printed == "") # (g) state 丢失(idle) 但批次子进程存活: 拒绝重复拉起 save_state({"phase": "idle"}) guard = subprocess.Popen([ sys.executable, "-c", "import time; time.sleep(2)", str(SRC_BINARY), "--batch", "30", "20" ]) try: logs_before = n_logs() code, printed = _run_main() check("idle refuses launch while batch child alive", "refusing duplicate launch" in printed and n_logs() == logs_before) finally: if guard.poll() is None: guard.kill() guard.wait() for p in find_orphans(str(st3.get("frozen") or frozen)): try: os.kill(p, signal.SIGKILL) except OSError: pass finally: ROOT, SRC_BINARY, STATE_PATH, LOCK_PATH, LOG_DIR, FROZEN_ROOT, FRAG, rc = saved globals()["runner_cmd"] = rc def selftest(): ok = True def check(name, cond): nonlocal ok print(("PASS " if cond else "FAIL ") + name) ok = ok and cond # ---- 批次判据校验(固定参数注入) ---- check("missing everything", verify_batch("x\nM3_ACCEPTANCE=PASS\n", 1, 1) == [ "missing 'batch worlds=K days=D' header line", "world line count 0 != expected 1", "world index set mismatch: duplicates=[] missing=[0] extra=[] (count 0 vs 1)", "missing batch_summary line"]) check("verify full PASS log", verify_batch(make_batch_log(2, 20), 2, 20) == []) check("verify wrong params", any("params mismatch" in p for p in verify_batch(make_batch_log(3, 20), 2, 20))) check("1x1 log cannot certify phase2", any("params mismatch" in p for p in verify_batch(make_batch_log(1, 1), 50, 100))) check("verify BAD world", any("non-OK" in p for p in verify_batch(make_batch_log(1, 1, ok_token="BAD"), 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 seed=42 a OK\nworld=1 seed=43 a OK\nbatch_summary worlds=3 passed=3 failed=0 x\nM3_ACCEPTANCE=PASS\n", 3, 1))) check("duplicate world index", any("duplicates=[0]" in p for p in verify_batch( "batch worlds=2 days=1 x\nworld=0 seed=42 a OK\nworld=0 seed=42 a OK\nbatch_summary worlds=2 passed=2 failed=0 x\nM3_ACCEPTANCE=PASS\n", 2, 1))) check("wrong seed binding", any("seed binding mismatch" in p for p in verify_batch( "batch worlds=2 days=1 x\nworld=0 seed=42 a OK\nworld=1 seed=99 a OK\nbatch_summary worlds=2 passed=2 failed=0 x\nM3_ACCEPTANCE=PASS\n", 2, 1))) check("seed off-by-one rejected", any("seed binding mismatch" in p for p in verify_batch( "batch worlds=1 days=1 x\nworld=0 seed=43 a OK\nbatch_summary worlds=1 passed=1 failed=0 x\nM3_ACCEPTANCE=PASS\n", 1, 1))) check("malformed world line", any("malformed" in p for p in verify_batch( "batch worlds=1 days=1 x\nworld=junk\nbatch_summary worlds=1 passed=1 failed=0 x\nM3_ACCEPTANCE=PASS\n", 1, 1))) check("verify summary mismatch", any("batch_summary mismatch" in p for p in verify_batch( "batch worlds=2 days=1 x\nworld=0 seed=42 a OK\nworld=1 seed=43 a 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 seed=42 a OK\nbatch_summary worlds=1 passed=1 failed=0 x\nM3_ACCEPTANCE=FAIL\n", 1, 1))) # 流式输出:完成序非升序 + 行尾 elapsed_s 追加字段,均应通过判据。 check("verify out-of-order streaming worlds", verify_batch( "batch worlds=3 days=1 x\nworld=2 seed=44 a OK elapsed_s=3.1\nworld=0 seed=42 a OK elapsed_s=1.2\n" "world=1 seed=43 a OK elapsed_s=2.0\nbatch_summary worlds=3 passed=3 failed=0 x\nM3_ACCEPTANCE=PASS\n", 3, 1) == []) check("verify elapsed suffix tolerated with progress lines", verify_batch( "batch worlds=2 days=1 x\n[done 1/2] elapsed=0.9s world=1\nworld=1 seed=43 a OK elapsed_s=0.9\n" "[done 2/2] elapsed=1.2s world=0\nworld=0 seed=42 a OK elapsed_s=1.2\n" "batch_summary worlds=2 passed=2 failed=0 x\nM3_ACCEPTANCE=PASS\n", 2, 1) == []) check("verify elapsed suffix does not mask BAD", any("non-OK" in p for p in verify_batch( "batch worlds=1 days=1 x\nworld=0 seed=42 a BAD elapsed_s=1.0\n" "batch_summary worlds=1 passed=0 failed=1 x\nM3_ACCEPTANCE=FAIL\n", 1, 1))) check("world_line_verdict ignores elapsed", world_line_verdict("world=0 seed=42 a OK elapsed_s=1.0") == "OK" and world_line_verdict("world=0 seed=42 a BAD elapsed_s=1.0") == "BAD") # ---- 冻结清单(合成目录, v2 递归 + 额外/缺失检测) ---- with tempfile.TemporaryDirectory() as td: td = Path(td) src = td / "bin"; src.mkdir() (src / "a.dll").write_bytes(b"hello") (src / "b.json").write_bytes(b"{}") (src / "runtimes").mkdir(); (src / "runtimes/r.dll").write_bytes(b"r") dst, mpath, man = freeze_release("deadbeef" * 8, src / "a.dll", td / "frozen") check("freeze copies deps recursively", (dst / "b.json").exists() and (dst / "runtimes/r.dll").read_bytes() == b"r" and man["files"].get("runtimes/r.dll") and man["version"] == 2) 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))) (dst / "sneaky.dll").write_bytes(b"x") check("manifest detect extra", any("extra file" in p for p in verify_manifest(dst, mpath, "deadbeef" * 8))) (dst / "sneaky.dll").unlink() (dst / "runtimes/r.dll").unlink() check("manifest detect missing", any("frozen file missing" in p for p in verify_manifest(dst, mpath, "deadbeef" * 8))) # 旧 v1 清单: 只校验所列文件, 不做额外文件检测(兼容运行中的旧冻结副本) leg = td / "leg"; leg.mkdir() (leg / "a.dll").write_bytes(b"hello") (leg / "unlisted.txt").write_text("u") lm = leg / "manifest.json" lm.write_text(json.dumps({"commit": "deadbeef" * 8, "files": {"a.dll": hashlib.sha256(b"hello").hexdigest()}})) check("legacy manifest compatible", verify_manifest(leg, lm, "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 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")) (pr / "stat").write_text((pr / "stat").read_text().replace("(sh) S ", "(sh) Z ")) check("zombie rejected", not proc_alive(4242, "12345", "LivingVillage.Headless.dll", td / "proc")) (pr / "stat").write_text((pr / "stat").read_text().replace("(sh) Z ", "(sh) S ")) 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")) # state 丢失兜底: FRAG + --batch 存活进程必须被识别, 自身/无 --batch 排除 (pr / "cmdline").write_bytes(b"dotnet /x/LivingVillage.Headless.dll --batch 30 20\x00") pr2 = td / "proc" / "777"; pr2.mkdir() (pr2 / "cmdline").write_bytes(b"sh -c sleep 9\x00") pr3 = td / "proc" / str(os.getpid()); pr3.mkdir() (pr3 / "cmdline").write_bytes(b"dotnet /x/LivingVillage.Headless.dll --batch 30 20\x00") check("live batch orphan detected", live_batch_orphans(td / "proc") == [4242]) # ---- 退出码读取 ---- 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) # ---- 状态机(合成日志; worlds/days 以固定 PHASES 为准, state 篡改即 FAIL) ---- with tempfile.TemporaryDirectory() as td: td = Path(td) plog = td / "p.log" plog.write_text(make_batch_log(30, 20)) (td / "fr").mkdir() (td / "fr/k.dll").write_bytes(b"kernel") (td / "fr/manifest.json").write_text(json.dumps( {"version": 2, "commit": "c" * 40, "files": {"k.dll": hashlib.sha256(b"kernel").hexdigest()}})) base = {"phase": "phase1", "worlds": 30, "days": 20, "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": [], "orphans": []} 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", log=str(plog)), lambda pid, st, fr: False) check("phase2 with phase1 log -> FAIL (fixed PHASES)", s["phase"] == "failed" and s["verdict"] == "FAIL") plog2 = td / "p2.log" plog2.write_text(make_batch_log(50, 20)) a, s, m = decide(dict(base, phase="phase2", worlds=50, days=20, log=str(plog2)), lambda pid, st, fr: False) check("phase2 PASS -> done", s["phase"] == "done" and m and "complete" in m) (td / "e.exit").write_text("0\n") tampered = dict(base, worlds=1, days=1) a, s, _ = decide(tampered, lambda pid, st, fr: False) check("state tamper fail-closed (1x1 state cannot pass phase1)", s["phase"] == "failed" and any("tampered" in p for p in s["problems"])) (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 FAILED", s["phase"] == "interrupted" and s["verdict"] == "FAIL") (td / "e.exit").write_text("0\n") (td / "fr/manifest.json").write_text(json.dumps( {"version": 2, "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) # ---- 真实隔离生命周期测试(临时目录 + 短真实子进程) ---- lifecycle_section(check) 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": live = live_batch_orphans() if live: # state 丢失/被清但批次子进程仍存活: 拒绝重复拉起, 保留现场等人 print(f"live batch process(es) {live} found while state idle; " "refusing duplicate launch, awaiting human") return 0 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())