#!/usr/bin/env python3 """Persona-owned execution reflex arc. The parent brain issues one intent. This runtime deterministically hands the same envelope through native eye -> execution limb -> action sense -> native eye -> world witness, then returns exactly one candidate verdict to the parent brain. Only the parent brain may close the arc as TASK_DONE or REPAIR_NEEDED. """ from __future__ import annotations import argparse import datetime import hashlib import json import os import pathlib import subprocess import sys import time HERE = pathlib.Path(__file__).resolve().parent ROOT = HERE.parents[1] LIMB = HERE / "persona_limb_agent.py" EYE = HERE / "persona_limb_eye.py" ADMISSION = ROOT / "server-tools/persona-host-write-admission/host-write-admission.mjs" FINAL_DECISIONS = {"TASK_DONE", "REPAIR_NEEDED"} def now() -> str: return datetime.datetime.now().astimezone().isoformat(timespec="milliseconds") def stable(value: object) -> bytes: return json.dumps(value, ensure_ascii=False, sort_keys=True, separators=(",", ":")).encode("utf-8") def digest(value: object) -> str: return hashlib.sha256(stable(value)).hexdigest() def file_hash(path: pathlib.Path) -> str: return hashlib.sha256(path.read_bytes()).hexdigest() def load(path: pathlib.Path) -> dict: value = json.loads(path.read_text(encoding="utf-8")) if not isinstance(value, dict): raise ValueError(f"OBJECT_REQUIRED:{path}") return value def atomic(path: pathlib.Path, value: dict) -> None: path.parent.mkdir(parents=True, exist_ok=True) temp = path.with_name(f"{path.name}.{os.getpid()}.tmp") temp.write_text(json.dumps(value, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") os.replace(temp, path) def run_process(command: list[str], *, timeout: int = 600, cwd: pathlib.Path | None = None) -> dict: started = time.monotonic_ns() try: completed = subprocess.run(command, cwd=cwd, capture_output=True, text=True, timeout=timeout, check=False) return { "exit_code": completed.returncode, "stdout": completed.stdout, "stderr": completed.stderr, "elapsed_ms": (time.monotonic_ns() - started) // 1_000_000, } except subprocess.TimeoutExpired as error: return { "exit_code": 124, "stdout": error.stdout.decode() if isinstance(error.stdout, bytes) else (error.stdout or ""), "stderr": "REFLEX_ARC_WATCHDOG_TIMEOUT", "elapsed_ms": (time.monotonic_ns() - started) // 1_000_000, } def admit(host: str, target: pathlib.Path) -> None: result = run_process(["node", str(ADMISSION), "check", "--host", host, "--path", str(target)], timeout=30) if result["exit_code"] != 0: raise PermissionError(f"CASE_ROOT_WRITE_NOT_ADMITTED:{result['stdout'][-500:]}") def probe_methods() -> dict: result = run_process([sys.executable, str(LIMB), "probe"], timeout=30) if result["exit_code"] != 0: raise RuntimeError("LIMB_PROBE_FAILED") return json.loads(result["stdout"])["runtimes"] def choose_runtime(requested: str, host: str, methods: dict) -> str: if requested != "auto": if not methods.get(requested, {}).get("available"): raise ValueError(f"REQUESTED_HOST_METHOD_UNAVAILABLE:{requested}") return requested preferred = {"codex": "codex", "zcode": "zcode", "qwen": "qwen"}.get(host) if preferred and methods.get(preferred, {}).get("available"): return preferred for candidate in ("codex", "zcode", "qwen"): if methods.get(candidate, {}).get("available"): return candidate raise ValueError("NO_HOST_METHOD_AVAILABLE") def validate_spec(spec: dict, ground: pathlib.Path) -> None: allowed = {"unchanged", "expect"} if set(spec) - allowed: raise ValueError("REFLEX_SPEC_UNKNOWN_FIELD") if not isinstance(spec.get("unchanged", []), list) or not isinstance(spec.get("expect", []), list): raise ValueError("REFLEX_SPEC_LIST_REQUIRED") for rel in spec.get("unchanged", []): if not isinstance(rel, str) or pathlib.Path(rel).is_absolute() or ".." in pathlib.Path(rel).parts: raise ValueError("REFLEX_SPEC_PATH_OUT_OF_GROUND") for item in spec.get("expect", []): if not isinstance(item, dict) or set(item) != {"file", "field", "equals"}: raise ValueError("REFLEX_EXPECT_CLOSED_FIELDS_REQUIRED") rel = pathlib.Path(str(item["file"])) if rel.is_absolute() or ".." in rel.parts or not isinstance(item["field"], str): raise ValueError("REFLEX_EXPECT_PATH_OR_FIELD_INVALID") if not ground.is_dir(): raise ValueError("REFLEX_GROUND_REQUIRED") def eye(command: list[str], *, timeout: int = 120) -> dict: return run_process([sys.executable, str(EYE), *command], timeout=timeout) def hand(persona: str, task_id: str, task: str, runtime: str, ground: pathlib.Path, access: str) -> dict: return run_process([ sys.executable, str(LIMB), "run", "--persona", persona, "--task-id", task_id, "--task", task, "--runtime", runtime, "--cwd", str(ground), "--access", access, ]) def parsed_hand_receipt(result: dict) -> dict: try: value = json.loads(result["stdout"]) if isinstance(value, dict): return value except json.JSONDecodeError: pass return { "schema": "guanghu.persona-execution-limb-runtime-receipt/v1", "outcome": "FAIL", "exit_code": result["exit_code"], "stdout": result["stdout"][-20000:], "stderr": result["stderr"][-4000:], "error": "HAND_RECEIPT_UNPARSEABLE", } def run_arc(args: argparse.Namespace) -> dict: ground = pathlib.Path(args.ground).resolve() case_root = pathlib.Path(args.case_root).resolve() case = case_root / args.task_id / f"attempt-{args.attempt:03d}" admit(args.host, case) if case.exists(): raise FileExistsError(f"REFLEX_CASE_ALREADY_EXISTS:{case}") spec_path = pathlib.Path(args.spec).resolve() spec = load(spec_path) validate_spec(spec, ground) methods = probe_methods() runtime = choose_runtime(args.runtime, args.host, methods) issued_at = now() intent = { "schema": "guanghu.persona-reflex-intent/v1", "state": "BRAIN_INTENT_ISSUED", "task_id": args.task_id, "attempt": args.attempt, "parent_persona_id": args.persona, "host": args.host, "selected_host_method": runtime, "access": args.access, "task": args.task, "spec": spec, "ground": str(ground), "issued_at": issued_at, } intent["intent_sha256"] = digest(intent) atomic(case / "intent.json", intent) stages: list[dict] = [] pre = eye(["snapshot", "--ground", str(ground), "--out", str(case / "eye-pre.json")]) stages.append({"stage": "EYE_PRE_MAP", "exit_code": pre["exit_code"], "elapsed_ms": pre["elapsed_ms"]}) if pre["exit_code"] != 0: raise RuntimeError(f"EYE_PRE_FAILED:{pre['stderr'][-300:]}") hand_result = hand(args.persona, args.task_id, args.task, runtime, ground, args.access) receipt = parsed_hand_receipt(hand_result) atomic(case / "hand-receipt.json", receipt) sense = { "schema": "guanghu.persona-action-sense/v1", "task_id": args.task_id, "attempt": args.attempt, "state": "HAND_ACTION_RETURNED", "host_method": runtime, "process_exit_code": hand_result["exit_code"], "hand_outcome": receipt.get("outcome"), "hand_receipt_sha256": file_hash(case / "hand-receipt.json"), "sensed_at": now(), "elapsed_ms": hand_result["elapsed_ms"], } atomic(case / "action-sense.json", sense) stages.append({"stage": "HAND_AND_ACTION_SENSE", "exit_code": hand_result["exit_code"], "elapsed_ms": hand_result["elapsed_ms"]}) post = eye(["observe", "--ground", str(ground), "--pre", str(case / "eye-pre.json"), "--out", str(case / "eye-post.json")]) stages.append({"stage": "EYE_POST_OBSERVE", "exit_code": post["exit_code"], "elapsed_ms": post["elapsed_ms"]}) if post["exit_code"] != 0: raise RuntimeError(f"EYE_POST_FAILED:{post['stderr'][-300:]}") witnessed = eye(["witness", "--pre", str(case / "eye-pre.json"), "--post", str(case / "eye-post.json"), "--spec", str(spec_path), "--out", str(case / "eye-witness.json")]) stages.append({"stage": "EYE_WORLD_WITNESS", "exit_code": witnessed["exit_code"], "elapsed_ms": witnessed["elapsed_ms"]}) if witnessed["exit_code"] != 0: raise RuntimeError(f"EYE_WITNESS_FAILED:{witnessed['stderr'][-300:]}") witness = load(case / "eye-witness.json") recommended = "TASK_DONE" if receipt.get("outcome") == "PASS" and witness.get("verdict") == "SEEN_MATCH" else "REPAIR_NEEDED" artifacts = { name: file_hash(case / name) for name in ("intent.json", "eye-pre.json", "hand-receipt.json", "action-sense.json", "eye-post.json", "eye-witness.json") } returned = { "schema": "guanghu.persona-reflex-brain-return/v1", "state": "AWAITING_PARENT_BRAIN_VERDICT", "task_id": args.task_id, "attempt": args.attempt, "parent_persona_id": args.persona, "selected_host_method": runtime, "recommended_decision": recommended, "hand_outcome": receipt.get("outcome"), "eye_verdict": witness.get("verdict"), "stages": stages, "total_measured_ms": sum(stage["elapsed_ms"] for stage in stages), "latency_claim": "MEASURED_NOT_ASSUMED_MILLISECOND", "artifacts": artifacts, "returned_at": now(), "case_dir": str(case), } returned["return_sha256"] = digest(returned) atomic(case / "brain-return.json", returned) atomic(case / "state.json", {"state": returned["state"], "brain_return_sha256": file_hash(case / "brain-return.json")}) return returned def validate_artifacts(case: pathlib.Path, returned: dict) -> list[str]: errors = [] for name, expected in returned.get("artifacts", {}).items(): target = case / name if not target.is_file(): errors.append(f"MISSING:{name}") elif file_hash(target) != expected: errors.append(f"HASH_MISMATCH:{name}") return errors def finalize(args: argparse.Namespace) -> dict: case = pathlib.Path(args.case).resolve() returned = load(case / "brain-return.json") if returned.get("state") != "AWAITING_PARENT_BRAIN_VERDICT": raise ValueError("BRAIN_RETURN_NOT_AWAITING_VERDICT") errors = validate_artifacts(case, returned) if errors: raise ValueError("REFLEX_ARTIFACT_INTEGRITY_FAILED:" + ",".join(errors)) if args.decision not in FINAL_DECISIONS: raise ValueError("FINAL_BRAIN_DECISION_REQUIRED") if args.decision == "TASK_DONE" and returned.get("recommended_decision") != "TASK_DONE": raise ValueError("BRAIN_CANNOT_CLOSE_FAILED_WORLD_WITNESS_AS_DONE") verdict = { "schema": "guanghu.persona-reflex-brain-verdict/v1", "state": "CLOSED_BY_PARENT_BRAIN", "task_id": returned["task_id"], "attempt": returned["attempt"], "parent_persona_id": args.persona, "decision": args.decision, "reason": args.reason, "brain_return_sha256": file_hash(case / "brain-return.json"), "artifact_chain_verified": True, "decided_at": now(), } verdict["verdict_sha256"] = digest(verdict) atomic(case / "brain-verdict.json", verdict) atomic(case / "state.json", {"state": args.decision, "brain_verdict_sha256": file_hash(case / "brain-verdict.json")}) return verdict def status(case: pathlib.Path) -> dict: returned = load(case / "brain-return.json") errors = validate_artifacts(case, returned) state = load(case / "state.json") return { "schema": "guanghu.persona-reflex-status/v1", "state": state.get("state"), "task_id": returned.get("task_id"), "attempt": returned.get("attempt"), "artifact_chain_verified": not errors, "errors": errors, "case_dir": str(case), } def parser() -> argparse.ArgumentParser: value = argparse.ArgumentParser(prog="persona_reflex_arc") sub = value.add_subparsers(dest="command", required=True) run = sub.add_parser("run") run.add_argument("--persona", default="ICE-P-ZY001") run.add_argument("--task-id", required=True) run.add_argument("--task", required=True) run.add_argument("--spec", required=True) run.add_argument("--ground", required=True) run.add_argument("--case-root", required=True) run.add_argument("--host", required=True, choices=("codex", "zcode", "qwen")) run.add_argument("--runtime", default="auto", choices=("auto", "codex", "zcode", "qwen")) run.add_argument("--access", default="read-only", choices=("read-only", "workspace-write")) run.add_argument("--attempt", type=int, default=1) finish = sub.add_parser("finalize") finish.add_argument("--persona", default="ICE-P-ZY001") finish.add_argument("--case", required=True) finish.add_argument("--decision", required=True, choices=tuple(sorted(FINAL_DECISIONS))) finish.add_argument("--reason", required=True) check = sub.add_parser("status") check.add_argument("--case", required=True) return value def main() -> int: args = parser().parse_args() try: if args.command == "run": result = run_arc(args) elif args.command == "finalize": result = finalize(args) else: result = status(pathlib.Path(args.case).resolve()) print(json.dumps(result, ensure_ascii=False, indent=2)) return 0 except Exception as error: print(json.dumps({"outcome": "FAIL", "error": str(error)}, ensure_ascii=False)) return 1 if __name__ == "__main__": raise SystemExit(main())