feat(persona): add reflex arc and embedded time system
This commit is contained in:
parent
4423b4bc2f
commit
28a60b4e80
19 changed files with 1382 additions and 81 deletions
342
server-tools/persona-execution-limb-agent/persona_reflex_arc.py
Normal file
342
server-tools/persona-execution-limb-agent/persona_reflex_arc.py
Normal file
|
|
@ -0,0 +1,342 @@
|
|||
#!/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())
|
||||
Loading…
Reference in a new issue