342 lines
14 KiB
Python
342 lines
14 KiB
Python
#!/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())
|