summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorSomhairle H. Marisol <[email protected]>2026-09-19 19:20:07 +0800
committerSomhairle H. Marisol <[email protected]>2026-09-19 19:20:07 +0800
commitef1f6b3d27e292a6b5e005411844eb8e14cbe4b2 (patch)
tree0b2cc55b16de379551ffcf44ff9842e7440de15b
parent557ae9acbe23a0df3d3167595aa81898e62c9eaf (diff)
downloadliving-village-ef1f6b3d27e292a6b5e005411844eb8e14cbe4b2.tar.gz
infra: P0 监督加固 — m3_supervise.py flock防重入/原子状态/冻结Release全依赖+sha256清单/真实child退出码(sh -c exitfile)/PID+starttime身份/收紧PASS判据(批次参数+world行数全OK+batch_summary+verdict token+清单hash+commit绑定)/失败停机并写diagnosis报告, selftest 33项合成断言; completion_notify.py 新completion-notice.md hash仅打印一次缺失静默, selftest 7项; Headless --batch 世界级并行(LV_BATCH_WORKERS 默认modest 4, 逐世界确定性, 输出按序, world行与串行逐字节一致); .gitignore 补冻结/锁/通知状态文件
-rw-r--r--.gitignore3
-rw-r--r--scripts/completion_notify.py132
-rwxr-xr-xscripts/m3_supervise.py508
-rw-r--r--src/LivingVillage.Headless/Program.fs97
4 files changed, 623 insertions, 117 deletions
diff --git a/.gitignore b/.gitignore
index 509472f..6f148dc 100644
--- a/.gitignore
+++ b/.gitignore
@@ -8,4 +8,7 @@ TestResults/
*.DotSettings.user
scripts/m3_logs/
scripts/m3_state.json
+scripts/m3_frozen/
+scripts/m3_supervise.lock
+scripts/completion_notify_state.json
.opencode-tmp/
diff --git a/scripts/completion_notify.py b/scripts/completion_notify.py
new file mode 100644
index 0000000..e18796b
--- /dev/null
+++ b/scripts/completion_notify.py
@@ -0,0 +1,132 @@
+#!/usr/bin/env python3
+"""completion-notice.md 变更通知:供 Hermes cron 调用。
+
+行为(严格):
+- 通知文件缺失: 静默退出 0,无任何输出
+- 文件存在: 计算 sha256;与上次已通知 hash 相同 -> 静默退出 0
+- hash 为新值: 把文件内容原样打印到 stdout 一次(供 cron 捕获后发 Slack),并原子记录该 hash
+- 每个新内容 hash 只打印一次;内容再次变化会再次打印
+- 状态文件损坏/不可读: 视为未通知过(打印一次)
+- 不依赖任何模型凭据;纯本地文件与 hash
+
+用法:
+ completion_notify.py [--notice PATH] [--state PATH] [--selftest]
+默认 notice=<cwd>/docs/completion-notice.md, state=<脚本所在目录>/completion_notify_state.json。
+脚本放项目中,由 Hermes 复制到 cron 脚本目录时建议显式传 --notice/--state。
+
+--selftest: 全部合成测试(临时文件模拟 缺失/新hash/重复/变更/损坏状态),不伪造真实结果。
+"""
+import argparse
+import hashlib
+import json
+import os
+import sys
+import tempfile
+from pathlib import Path
+
+SCRIPT_DIR = Path(__file__).resolve().parent
+DEFAULT_NOTICE = Path.cwd() / "docs/completion-notice.md"
+DEFAULT_STATE = SCRIPT_DIR / "completion_notify_state.json"
+
+
+def content_hash(data: bytes) -> str:
+ return hashlib.sha256(data).hexdigest()
+
+
+def load_state(path: Path):
+ try:
+ v = json.loads(Path(path).read_text())
+ if isinstance(v, dict) and isinstance(v.get("hash"), str):
+ return v["hash"]
+ except Exception:
+ pass
+ return None
+
+
+def save_state(path: Path, h: str):
+ p = Path(path)
+ p.parent.mkdir(parents=True, exist_ok=True)
+ fd, tmp = tempfile.mkstemp(dir=str(p.parent), prefix=".notify_state.", suffix=".tmp")
+ try:
+ with os.fdopen(fd, "w") as fh:
+ fh.write(json.dumps({"hash": h}, indent=1) + "\n")
+ os.replace(tmp, p)
+ except BaseException:
+ try:
+ os.unlink(tmp)
+ except OSError:
+ pass
+ raise
+
+
+def run_once(notice_path, state_path, out=sys.stdout):
+ """返回 'printed' | 'silent' | 'missing'。打印发生在 out 上。"""
+ p = Path(notice_path)
+ try:
+ data = p.read_bytes()
+ except OSError:
+ return "missing" # 缺失静默
+ h = content_hash(data)
+ if load_state(state_path) == h:
+ return "silent"
+ out.write(data.decode("utf-8", "replace"))
+ save_state(state_path, h)
+ return "printed"
+
+
+def run_once_str(notice_path, state_path):
+ """测试用: 返回 (action, printed_text)。"""
+ import io
+ buf = io.StringIO()
+ action = run_once(notice_path, state_path, out=buf)
+ return action, buf.getvalue()
+
+
+def selftest():
+ ok = True
+
+ def check(name, cond):
+ nonlocal ok
+ print(("PASS " if cond else "FAIL ") + name)
+ ok = ok and cond
+
+ import tempfile
+ with tempfile.TemporaryDirectory() as td:
+ td = Path(td)
+ notice = td / "completion-notice.md"
+ state = td / "state" / "notify.json"
+ action, text = run_once_str(notice, state)
+ check("missing -> silent no output", action == "missing" and text == "")
+ notice.write_text("# 完成\n- 测试内容\n")
+ action, text = run_once_str(notice, state)
+ check("new content printed once", action == "printed" and text == notice.read_text())
+ action, text = run_once_str(notice, state)
+ check("same hash silent", action == "silent" and text == "")
+ notice.write_text("# 完成 v2\n")
+ action, text = run_once_str(notice, state)
+ check("changed content printed again", action == "printed" and text == "# 完成 v2\n")
+ action, text = run_once_str(notice, state)
+ check("second view silent", action == "silent" and text == "")
+ state.write_text("not-json{{{")
+ action, text = run_once_str(notice, state)
+ check("corrupt state treated as new", action == "printed" and text == "# 完成 v2\n")
+ check("state saved after corrupt", load_state(state) == content_hash("# 完成 v2\n".encode()))
+
+ print("SELFTEST=" + ("PASS" if ok else "FAIL"))
+ return 0 if ok else 1
+
+
+def main():
+ ap = argparse.ArgumentParser(description="completion-notice.md new-hash notifier")
+ ap.add_argument("--notice", default=str(DEFAULT_NOTICE))
+ ap.add_argument("--state", default=str(DEFAULT_STATE))
+ ap.add_argument("--selftest", action="store_true")
+ args = ap.parse_args()
+ if args.selftest:
+ return selftest()
+ action = run_once(args.notice, args.state)
+ return 0
+
+
+if __name__ == "__main__":
+ sys.exit(main())
diff --git a/scripts/m3_supervise.py b/scripts/m3_supervise.py
index 3975b81..62fe94d 100755
--- a/scripts/m3_supervise.py
+++ b/scripts/m3_supervise.py
@@ -1,13 +1,32 @@
#!/usr/bin/env python3
-"""M3 验收监督:供 Hermes cron 每 10 分钟调用(no_agent)。
+"""M3 验收监督(加固版):供 Hermes cron 每 5 分钟调用(no_agent)。
状态机: idle -> phase1(--batch 30 20) -> phase2(--batch 50 100) -> done
-- 运行中(pid 存活且 cmdline 匹配): 零输出静默退出
-- 进程结束: 解析日志 M3_ACCEPTANCE -> PASS 则进入下一阶段并立即启动; FAIL/无结论则停机留给人工, 不自动重启
-- 状态文件 scripts/m3_state.json 防重复拉起; 日志写 scripts/m3_logs/<phase>_<ts>.log
-- --selftest: 仅验证解析与状态机转移, 不启动模拟
+- flock 单实例防重入;抢不到锁静默退出
+- 冻结运行: 首次启动把 Release 输出目录整体复制到 scripts/m3_frozen/<stamp>_<commit>/,
+ 记录全部文件 sha256 清单;两阶段都从冻结副本运行,后续构建不会覆盖运行中二进制
+- 真实退出码: 子进程为 sh -c 'dotnet <frozen.dll> --batch K D; echo $? > exitfile',
+ 结束后读 exitfile 得到真实退出码(含信号 128+sig;shell 被杀则无 exitfile -> interrupted)
+- PID 身份: 记录 /proc/<pid>/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/无结论/校验失败: 停机不自动重启,并在日志旁写 <log>.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
@@ -15,89 +34,297 @@ import time
from pathlib import Path
ROOT = Path(__file__).resolve().parent.parent
-BINARY = ROOT / "src/LivingVillage.Headless/bin/Release/net8.0/LivingVillage.Headless.dll"
+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"
-PHASES = {"phase1": ["--batch", "30", "20"], "phase2": ["--batch", "50", "100"]}
-NEXT = {"phase1": "phase2", "phase2": "done"}
+FROZEN_ROOT = ROOT / "scripts/m3_frozen"
+PHASES = {"phase1": (30, 20), "phase2": (50, 100)}
+NEXT = {"phase1": "phase2", "phase2": None}
+FRAG = "LivingVillage.Headless.dll"
-def load_state():
- if STATE_PATH.exists():
- return json.loads(STATE_PATH.read_text())
- return {"phase": "idle", "pid": None, "cmd": None, "log": None, "started": None, "verdict": None}
+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 save_state(state):
- STATE_PATH.write_text(json.dumps(state, indent=1, ensure_ascii=False) + "\n")
+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 parse_acceptance(text):
- verdict = None
- for line in text.splitlines():
- if line.startswith("M3_ACCEPTANCE="):
- v = line.split("=", 1)[1].strip()
- verdict = v if v in ("PASS", "FAIL") else None
- return verdict
+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 pid_matches(pid, cmd):
- if not pid:
- return False
+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/<pid>/stat 字段22 starttime(boot 后时钟滴答)。解析失败返回 None。"""
try:
- with open(f"/proc/{pid}/cmdline", "rb") as f:
- argv = f.read().decode("utf-8", "replace").replace("\0", " ")
+ 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
- return "LivingVillage.Headless.dll" in argv and (not cmd or cmd in argv)
+ 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"
- cmd = ["dotnet", str(BINARY)] + PHASES[phase]
+ 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)
+ 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, "pid": proc.pid, "cmd": str(BINARY),
- "log": str(log), "started": time.strftime("%Y-%m-%dT%H:%M:%S"),
- "verdict": None})
+ "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} log={log}"
+ 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("cmd")):
+ 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 ""
- verdict = parse_acceptance(text)
- state["verdict"] = verdict
+ 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 verdict == "PASS":
+ 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]
- if nxt == "done":
+ state["verdict"] = "PASS"
+ if nxt is None:
state["phase"] = "done"
- return "save", state, f"{phase} PASS -> M3 acceptance complete (30x20 + 50x100)"
+ 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
- if verdict == "FAIL":
- state["phase"] = "failed"
- return "save", state, f"{phase} FAIL (log={log}); stopped, awaiting human"
- state["phase"] = "interrupted"
- return "save", state, f"{phase} interrupted without verdict (log={log}); stopped, awaiting human"
+ 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
@@ -106,29 +333,119 @@ def selftest():
print(("PASS " if cond else "FAIL ") + name)
ok = ok and cond
- check("parse PASS", parse_acceptance("x\nM3_ACCEPTANCE=PASS\n") == "PASS")
- check("parse FAIL", parse_acceptance("M3_ACCEPTANCE=FAIL") == "FAIL")
- check("parse empty", parse_acceptance("world=0 ...\n") is None)
- check("parse garbage", parse_acceptance("M3_ACCEPTANCE=MAYBE") is None)
- base = {"phase": "phase1", "pid": 123, "cmd": "b.dll", "log": "nope.log", "verdict": None, "started": None}
- a, s, _ = decide(dict(base), lambda pid, cmd: True)
- check("alive -> silent", a is None)
- a, s, _ = decide(dict(base), lambda pid, cmd: False)
- check("dead without verdict -> interrupted", a == "save" and s["phase"] == "interrupted")
- with tempfile.NamedTemporaryFile("w", suffix=".log", delete=False) as fh:
- fh.write("batch_summary ...\nM3_ACCEPTANCE=PASS\n")
- plog = fh.name
- a, s, _ = decide(dict(base, log=plog), lambda pid, cmd: False)
- check("phase1 PASS -> launch phase2", a == "launch:phase2")
- a, s, m = decide(dict(base, phase="phase2", log=plog), lambda pid, cmd: False)
- check("phase2 PASS -> done", s["phase"] == "done" and m and "complete" in m)
- with tempfile.NamedTemporaryFile("w", suffix=".log", delete=False) as fh:
- fh.write("M3_ACCEPTANCE=FAIL\n")
- flog = fh.name
- a, s, _ = decide(dict(base, log=flog), lambda pid, cmd: False)
- check("FAIL -> failed", s["phase"] == "failed")
- a, s, _ = decide({"phase": "done"}, lambda pid, cmd: False)
- check("done -> silent", a is None)
+ # ---- 解析 ----
+ 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
@@ -136,18 +453,47 @@ def selftest():
def main():
if "--selftest" in sys.argv:
return selftest()
- state = load_state()
- if state.get("phase") == "idle" and not BINARY.exists():
- print(f"binary missing: {BINARY}; run dotnet build src/LivingVillage.Headless -c Release")
- return 1
- action, state, msg = decide(state, pid_matches)
- if action and action.startswith("launch:"):
- msg = launch(action.split(":", 1)[1], state)
- elif action == "save":
- save_state(state)
- if msg:
- print(msg)
- return 0
+ 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__":
diff --git a/src/LivingVillage.Headless/Program.fs b/src/LivingVillage.Headless/Program.fs
index be9c325..bcf2e41 100644
--- a/src/LivingVillage.Headless/Program.fs
+++ b/src/LivingVillage.Headless/Program.fs
@@ -175,6 +175,48 @@ let gini (xs: int64[]) : float =
num <- num + abs (float xi - float xj)
num / (2.0 * n * n * (s / n))
+type WorldVerdict =
+ { Line: string
+ Ok: bool
+ CheckOld: bool
+ Ratio: float
+ Gini: float }
+
+let evalWorld (days: int64) (k: int) : WorldVerdict =
+ let seed = 42UL + uint64 k
+ let stats = runSimulation false days seed Sim.npcCount
+ let m = Sim.relationMatrix stats.World
+ let n = Array2D.length1 m
+ let mutable hasNonFinite = stats.NonFinite > 0L
+ for i in 0 .. n - 1 do
+ for j in 0 .. n - 1 do
+ if Single.IsNaN m.[i, j] || Single.IsInfinity m.[i, j] then hasNonFinite <- true
+ let chats = stats.ChatsPerNpc
+ let chatsTotal = Array.sum chats
+ let chatsAvg = float chatsTotal / float chats.Length
+ let chatsMax = if chats.Length > 0 then Array.max chats else 0L
+ let ratio = if chatsAvg > 0.0 then float chatsMax / chatsAvg else 0.0
+ let g = gini chats
+ let counts = Sim.relationCounts m
+ let cntMax = if counts.Length > 0 then Array.max counts else 0
+ let cntMin = if counts.Length > 0 then Array.min counts else 0
+ let connected = Array.fold (fun acc c -> if c > 0 then acc + 1 else acc) 0 counts
+ let connectedNeed = int (System.Math.Ceiling(0.6 * float chats.Length))
+ let checkA = ratio > 2.0
+ let checkB = connected >= connectedNeed
+ let checkC = not hasNonFinite
+ let checkOld = cntMax > 2 * cntMin
+ let ok = checkA && checkB && checkC
+ let line =
+ sprintf "world=%d seed=%d chats_max=%d chats_avg=%.2f ratio=%.3f gini=%.3f relcnt_max=%d relcnt_min=%d relconn=%d/%d nonfinite=%d check_a=%s check_b=%s check_c=%s old_crit=%s %s"
+ k seed chatsMax chatsAvg ratio g cntMax cntMin connected connectedNeed stats.NonFinite
+ (if checkA then "PASS" else "FAIL")
+ (if checkB then "PASS" else "FAIL")
+ (if checkC then "PASS" else "FAIL")
+ (if checkOld then "PASS" else "FAIL")
+ (if ok then "OK" else "BAD")
+ { Line = line; Ok = ok; CheckOld = checkOld; Ratio = ratio; Gini = g }
+
let runBatch (worlds: int) (days: int64) : int =
printfn "batch worlds=%d days=%d half_life_ticks=%d rel_threshold=%.2f" worlds days Sim.relationHalfLifeTicks Sim.relationThreshold
// 判据语义(M3c 修正):
@@ -184,48 +226,31 @@ let runBatch (worlds: int) (days: int64) : int =
// 明星集中型分布过严,relcnt_max=0 即判死,与"非均匀分布"的验收目标不符)
// (c) 无 NaN/Inf
printfn "checks: (a) chats max/avg > 2.0 (b) npcs_with_relation >= 0.6*npc_count (c) no NaN (old) relcnt_max > 2*relcnt_min (report only)"
+ // 世界间相互独立(seed=42+k,逐世界确定性),可安全并行;输出仍按 world 序号
+ // 顺序打印,逐行内容与串行版本逐字节一致。LV_BATCH_WORKERS 可覆盖(默认 modest 4)。
+ let workers =
+ match Environment.GetEnvironmentVariable "LV_BATCH_WORKERS" with
+ | null | "" -> min 4 Environment.ProcessorCount
+ | v -> match Int32.TryParse v with true, n when n > 0 -> min n Environment.ProcessorCount | _ -> min 4 Environment.ProcessorCount
+ printfn "workers=%d" workers
+ let results = Array.zeroCreate worlds
+ System.Threading.Tasks.Parallel.For
+ (0, worlds,
+ System.Threading.Tasks.ParallelOptions(MaxDegreeOfParallelism = workers),
+ fun k -> results.[k] <- evalWorld days k)
+ |> ignore
let mutable passed = 0
let mutable failed = 0
let mutable oldPassed = 0
let ratios = ResizeArray<float> ()
let ginis = ResizeArray<float> ()
let sw = System.Diagnostics.Stopwatch.StartNew()
- for k in 0 .. worlds - 1 do
- let seed = 42UL + uint64 k
- let stats = runSimulation false days seed Sim.npcCount
- let m = Sim.relationMatrix stats.World
- let n = Array2D.length1 m
- let mutable hasNonFinite = stats.NonFinite > 0L
- for i in 0 .. n - 1 do
- for j in 0 .. n - 1 do
- if Single.IsNaN m.[i, j] || Single.IsInfinity m.[i, j] then hasNonFinite <- true
- let chats = stats.ChatsPerNpc
- let chatsTotal = Array.sum chats
- let chatsAvg = float chatsTotal / float chats.Length
- let chatsMax = if chats.Length > 0 then Array.max chats else 0L
- let ratio = if chatsAvg > 0.0 then float chatsMax / chatsAvg else 0.0
- let g = gini chats
- let counts = Sim.relationCounts m
- let cntMax = if counts.Length > 0 then Array.max counts else 0
- let cntMin = if counts.Length > 0 then Array.min counts else 0
- let connected = Array.fold (fun acc c -> if c > 0 then acc + 1 else acc) 0 counts
- let connectedNeed = int (System.Math.Ceiling(0.6 * float chats.Length))
- let checkA = ratio > 2.0
- let checkB = connected >= connectedNeed
- let checkC = not hasNonFinite
- let checkOld = cntMax > 2 * cntMin
- let ok = checkA && checkB && checkC
- if ok then passed <- passed + 1 else failed <- failed + 1
- if checkOld then oldPassed <- oldPassed + 1
- ratios.Add ratio
- ginis.Add g
- printfn "world=%d seed=%d chats_max=%d chats_avg=%.2f ratio=%.3f gini=%.3f relcnt_max=%d relcnt_min=%d relconn=%d/%d nonfinite=%d check_a=%s check_b=%s check_c=%s old_crit=%s %s"
- k seed chatsMax chatsAvg ratio g cntMax cntMin connected connectedNeed stats.NonFinite
- (if checkA then "PASS" else "FAIL")
- (if checkB then "PASS" else "FAIL")
- (if checkC then "PASS" else "FAIL")
- (if checkOld then "PASS" else "FAIL")
- (if ok then "OK" else "BAD")
+ for r in results do
+ printfn "%s" r.Line
+ if r.Ok then passed <- passed + 1 else failed <- failed + 1
+ if r.CheckOld then oldPassed <- oldPassed + 1
+ ratios.Add r.Ratio
+ ginis.Add r.Gini
sw.Stop()
let allPass = failed = 0
printfn "batch_summary worlds=%d passed=%d failed=%d old_passed=%d/%d ratio_min=%.3f ratio_max=%.3f gini_min=%.3f gini_max=%.3f elapsed_s=%.1f"