test(827): capture Gate0 runtime evidence
This commit is contained in:
450
docs/compaction-refresh/probes/pi_gate0_run.py
Normal file
450
docs/compaction-refresh/probes/pi_gate0_run.py
Normal file
@@ -0,0 +1,450 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Drive real Pi 0.80.x RPC for P2/P3/P5/P6 runtime evidence."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import queue
|
||||
import shutil
|
||||
import signal
|
||||
import socket
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Any, Callable
|
||||
|
||||
HERE = Path(__file__).resolve().parent
|
||||
BLOCK = "\n".join(
|
||||
[
|
||||
"GATE0_PI_ATOMIC_BEGIN",
|
||||
"segment-01=alpha-7e31",
|
||||
"segment-02=middle-9c42",
|
||||
"segment-03=omega-5b83",
|
||||
"GATE0_PI_ATOMIC_END",
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
def wait_path(path: Path, timeout: float = 20) -> None:
|
||||
deadline = time.monotonic() + timeout
|
||||
while time.monotonic() < deadline:
|
||||
if path.exists():
|
||||
return
|
||||
time.sleep(0.05)
|
||||
raise TimeoutError(f"timed out waiting for {path}")
|
||||
|
||||
|
||||
def socket_request(path: Path, payload: dict[str, object]) -> None:
|
||||
conn = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
||||
conn.connect(str(path))
|
||||
conn.sendall((json.dumps(payload) + "\n").encode())
|
||||
conn.makefile("r", encoding="utf-8").readline()
|
||||
conn.close()
|
||||
|
||||
|
||||
def jsonl(path: Path) -> list[dict[str, Any]]:
|
||||
if not path.exists():
|
||||
return []
|
||||
return [json.loads(line) for line in path.read_text().splitlines() if line]
|
||||
|
||||
|
||||
class PiRpc:
|
||||
def __init__(self, command: list[str], cwd: Path, env: dict[str, str]):
|
||||
self.process = subprocess.Popen(
|
||||
command,
|
||||
cwd=cwd,
|
||||
env=env,
|
||||
stdin=subprocess.PIPE,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.PIPE,
|
||||
text=True,
|
||||
bufsize=1,
|
||||
start_new_session=True,
|
||||
)
|
||||
self.events: queue.Queue[dict[str, Any]] = queue.Queue()
|
||||
self.raw_lines: list[str] = []
|
||||
self.stderr_lines: list[str] = []
|
||||
threading.Thread(target=self._read_stdout, daemon=True).start()
|
||||
threading.Thread(target=self._read_stderr, daemon=True).start()
|
||||
|
||||
def _read_stdout(self) -> None:
|
||||
assert self.process.stdout is not None
|
||||
for line in self.process.stdout:
|
||||
stripped = line.rstrip("\n")
|
||||
self.raw_lines.append(stripped)
|
||||
try:
|
||||
event = json.loads(stripped)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
self.events.put(event)
|
||||
|
||||
def _read_stderr(self) -> None:
|
||||
assert self.process.stderr is not None
|
||||
for line in self.process.stderr:
|
||||
self.stderr_lines.append(line.rstrip("\n"))
|
||||
|
||||
def send(self, payload: dict[str, object]) -> None:
|
||||
assert self.process.stdin is not None
|
||||
self.process.stdin.write(json.dumps(payload) + "\n")
|
||||
self.process.stdin.flush()
|
||||
|
||||
def wait(self, predicate: Callable[[dict[str, Any]], bool], description: str, timeout: float = 180) -> dict[str, Any]:
|
||||
deadline = time.monotonic() + timeout
|
||||
while time.monotonic() < deadline:
|
||||
if self.process.poll() is not None and self.events.empty():
|
||||
raise RuntimeError(
|
||||
f"Pi exited {self.process.returncode} while waiting for {description}: "
|
||||
+ " | ".join(self.stderr_lines[-5:])
|
||||
)
|
||||
try:
|
||||
event = self.events.get(timeout=0.2)
|
||||
except queue.Empty:
|
||||
continue
|
||||
if predicate(event):
|
||||
return event
|
||||
raise TimeoutError(f"timed out waiting for {description}")
|
||||
|
||||
def response(self, request_id: str, timeout: float = 180) -> dict[str, Any]:
|
||||
return self.wait(
|
||||
lambda event: event.get("type") == "response" and event.get("id") == request_id,
|
||||
f"response {request_id}",
|
||||
timeout,
|
||||
)
|
||||
|
||||
def prompt_and_settle(self, request_id: str, message: str) -> None:
|
||||
self.send({"id": request_id, "type": "prompt", "message": message})
|
||||
response = self.response(request_id)
|
||||
if not response.get("success"):
|
||||
raise RuntimeError(f"prompt rejected: {response}")
|
||||
self.wait(lambda event: event.get("type") == "agent_settled", f"agent_settled {request_id}")
|
||||
|
||||
def close(self) -> None:
|
||||
if self.process.poll() is None:
|
||||
try:
|
||||
os.killpg(self.process.pid, signal.SIGTERM)
|
||||
except ProcessLookupError:
|
||||
pass
|
||||
try:
|
||||
self.process.wait(timeout=8)
|
||||
except subprocess.TimeoutExpired:
|
||||
os.killpg(self.process.pid, signal.SIGKILL)
|
||||
self.process.wait(timeout=5)
|
||||
|
||||
|
||||
def manifest(path: Path, fragment: Path, expected_hash: str, max_bytes: int = 64) -> None:
|
||||
path.write_text(
|
||||
json.dumps(
|
||||
{
|
||||
"maxBytes": max_bytes,
|
||||
"fragments": [{"path": str(fragment), "sha256": expected_hash}],
|
||||
},
|
||||
sort_keys=True,
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
def run_open(root: Path) -> tuple[list[dict[str, Any]], list[dict[str, Any]], list[str], list[str]]:
|
||||
workspace = root / "workspace"
|
||||
workspace.mkdir()
|
||||
session_dir = root / "sessions"
|
||||
session_dir.mkdir()
|
||||
pi_log = root / "pi-hooks.jsonl"
|
||||
generation_log = root / "generation.jsonl"
|
||||
generation_socket = root / "generation.sock"
|
||||
source_manifest = root / "manifest.json"
|
||||
valid_fragment = root / "fragment.md"
|
||||
valid_fragment.write_text("NORMATIVE-FRAGMENT-v1\n")
|
||||
expected = hashlib.sha256(valid_fragment.read_bytes()).hexdigest()
|
||||
manifest(source_manifest, valid_fragment, expected)
|
||||
|
||||
broker = subprocess.Popen(
|
||||
[
|
||||
sys.executable,
|
||||
str(HERE / "p3_generation_broker.py"),
|
||||
"--socket",
|
||||
str(generation_socket),
|
||||
"--log",
|
||||
str(generation_log),
|
||||
],
|
||||
text=True,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.STDOUT,
|
||||
)
|
||||
wait_path(generation_socket)
|
||||
env = os.environ.copy()
|
||||
env.update(
|
||||
{
|
||||
"GATE0_PI_LOG": str(pi_log),
|
||||
"GATE0_GENERATION_SOCKET": str(generation_socket),
|
||||
"GATE0_SOURCE_MANIFEST": str(source_manifest),
|
||||
"GATE0_PI_CONTEXT_BLOCK": BLOCK,
|
||||
"MOSAIC_PI_FORCE_SKILLS": "",
|
||||
"PI_SKIP_VERSION_CHECK": "1",
|
||||
}
|
||||
)
|
||||
command = [
|
||||
"mosaic",
|
||||
"yolo",
|
||||
"pi",
|
||||
"--mode",
|
||||
"rpc",
|
||||
"--session-dir",
|
||||
str(session_dir),
|
||||
"--no-extensions",
|
||||
"--no-context-files",
|
||||
"--no-prompt-templates",
|
||||
"--model",
|
||||
"openai-codex/gpt-5.6-sol",
|
||||
"--thinking",
|
||||
"medium",
|
||||
"--extension",
|
||||
str(HERE / "pi_gate0_extension.ts"),
|
||||
]
|
||||
pi = PiRpc(command, workspace, env)
|
||||
try:
|
||||
pi.send({"id": "state-0", "type": "get_state"})
|
||||
state0 = pi.response("state-0")
|
||||
original_session = state0["data"]["sessionFile"]
|
||||
|
||||
pi.prompt_and_settle(
|
||||
"p2",
|
||||
"Call gate0_nonce_probe exactly once with label p2. After the tool finishes, copy the exact full GATE0_PI_ATOMIC_BEGIN through GATE0_PI_ATOMIC_END block from context, with no commentary.",
|
||||
)
|
||||
|
||||
# P3 immediately follows the valid P2 promotion so reload must revoke a
|
||||
# genuinely VERIFIED prior generation, not an already-invalid source run.
|
||||
pi.send({"id": "reload", "type": "prompt", "message": "/gate0-reload"})
|
||||
reload_response = pi.response("reload")
|
||||
if not reload_response.get("success"):
|
||||
raise RuntimeError(f"reload command failed: {reload_response}")
|
||||
|
||||
pi.send({"id": "clone", "type": "clone"})
|
||||
clone_response = pi.response("clone")
|
||||
if not clone_response.get("success") or clone_response.get("data", {}).get("cancelled"):
|
||||
raise RuntimeError(f"clone failed: {clone_response}")
|
||||
|
||||
pi.send({"id": "new", "type": "new_session"})
|
||||
new_response = pi.response("new")
|
||||
if not new_response.get("success") or new_response.get("data", {}).get("cancelled"):
|
||||
raise RuntimeError(f"new session failed: {new_response}")
|
||||
|
||||
pi.send(
|
||||
{
|
||||
"id": "resume",
|
||||
"type": "switch_session",
|
||||
"sessionPath": original_session,
|
||||
}
|
||||
)
|
||||
resume_response = pi.response("resume")
|
||||
if not resume_response.get("success") or resume_response.get("data", {}).get("cancelled"):
|
||||
raise RuntimeError(f"resume failed: {resume_response}")
|
||||
|
||||
# P5 missing fragment: action-time source validation must revoke/refuse.
|
||||
manifest(source_manifest, root / "absent-fragment.md", expected)
|
||||
pi.prompt_and_settle(
|
||||
"p5-missing",
|
||||
"Call gate0_nonce_probe exactly once with label p5-missing, then stop.",
|
||||
)
|
||||
|
||||
# P5 oversize fragment: expected hash is correct, size limit is not.
|
||||
oversize = root / "oversize.md"
|
||||
oversize.write_text("X" * 65)
|
||||
manifest(source_manifest, oversize, hashlib.sha256(oversize.read_bytes()).hexdigest(), 64)
|
||||
pi.prompt_and_settle(
|
||||
"p5-oversize",
|
||||
"Call gate0_nonce_probe exactly once with label p5-oversize, then stop.",
|
||||
)
|
||||
|
||||
# P5 hash mismatch: size is valid but bytes differ from expected.
|
||||
mismatch = root / "mismatch.md"
|
||||
mismatch.write_text("tampered\n")
|
||||
manifest(source_manifest, mismatch, expected, 64)
|
||||
pi.prompt_and_settle(
|
||||
"p5-hash",
|
||||
"Call gate0_nonce_probe exactly once with label p5-hash-mismatch, then stop.",
|
||||
)
|
||||
|
||||
time.sleep(1)
|
||||
return jsonl(pi_log), jsonl(generation_log), list(pi.raw_lines), list(pi.stderr_lines)
|
||||
finally:
|
||||
pi.close()
|
||||
try:
|
||||
socket_request(generation_socket, {"action": "shutdown-broker"})
|
||||
except OSError:
|
||||
pass
|
||||
try:
|
||||
broker.wait(timeout=5)
|
||||
except subprocess.TimeoutExpired:
|
||||
broker.kill()
|
||||
broker.wait()
|
||||
|
||||
|
||||
def run_closed(root: Path) -> list[dict[str, Any]]:
|
||||
workspace = root / "closed-workspace"
|
||||
workspace.mkdir()
|
||||
pi_log = root / "closed-hooks.jsonl"
|
||||
env = os.environ.copy()
|
||||
env.update(
|
||||
{
|
||||
"GATE0_PI_LOG": str(pi_log),
|
||||
"MOSAIC_PI_FORCE_SKILLS": "",
|
||||
"PI_SKIP_VERSION_CHECK": "1",
|
||||
}
|
||||
)
|
||||
command = [
|
||||
"mosaic",
|
||||
"yolo",
|
||||
"pi",
|
||||
"--mode",
|
||||
"rpc",
|
||||
"--no-session",
|
||||
"--no-extensions",
|
||||
"--no-context-files",
|
||||
"--no-prompt-templates",
|
||||
"--extension",
|
||||
str(HERE / "pi_gate0_extension.ts"),
|
||||
"--extension",
|
||||
str(HERE / "pi_later_extension.ts"),
|
||||
]
|
||||
pi = PiRpc(command, workspace, env)
|
||||
try:
|
||||
pi.send({"id": "closed-state", "type": "get_state"})
|
||||
pi.response("closed-state")
|
||||
time.sleep(0.5)
|
||||
return jsonl(pi_log)
|
||||
finally:
|
||||
pi.close()
|
||||
|
||||
|
||||
def main() -> None:
|
||||
with tempfile.TemporaryDirectory(prefix="gate0-pi-") as temp:
|
||||
root = Path(temp)
|
||||
records, generations, rpc_lines, stderr_lines = run_open(root)
|
||||
closed = run_closed(root)
|
||||
|
||||
p2_message = next(
|
||||
r for r in records if r["event"] == "message_end" and r.get("nonceMappings")
|
||||
)
|
||||
p2_tool = next(r for r in records if r["event"] == "tool_call" and r.get("allowed"))
|
||||
mapped = p2_message["nonceMappings"][0]
|
||||
assert mapped["toolCallId"] == p2_tool["toolCallId"]
|
||||
assert mapped["requestNonce"] == p2_tool["mapping"]["nonce"]
|
||||
assert next(r for r in records if r["event"] == "session_start")["lastPosition"] is True
|
||||
assert next(r for r in closed if r["event"] == "session_start")["gateState"] == "CLOSED_NOT_LAST"
|
||||
reload_revoke = next(
|
||||
r
|
||||
for r in generations
|
||||
if r["event"] == "runtime_generation_bump"
|
||||
and r.get("reason") == "reload"
|
||||
and r.get("phase") == "shutdown"
|
||||
)
|
||||
assert reload_revoke["prior_lease"] == "VERIFIED"
|
||||
assert reload_revoke["prior_lease_revoked"] is True
|
||||
for reason in {"missing", "oversize", "hash-mismatch"}:
|
||||
assert any(
|
||||
r["event"] == "context_return"
|
||||
and r.get("sourceValidation", {}).get("reason") == reason
|
||||
and r.get("injectionDecision") == "REFUSED"
|
||||
and r.get("promotion") is False
|
||||
for r in records
|
||||
)
|
||||
assert any(
|
||||
r["event"] == "tool_call"
|
||||
and r.get("mapping", {}).get("sourceReason") == reason
|
||||
and r.get("allowed") is False
|
||||
for r in records
|
||||
)
|
||||
assert any(
|
||||
r["event"] == "message_end" and r.get("exactContextBlockCopied") is True
|
||||
for r in records
|
||||
)
|
||||
|
||||
print("$ python3 docs/compaction-refresh/probes/pi_gate0_run.py")
|
||||
print("machine_assertions=PASS")
|
||||
print("runtime_versions:")
|
||||
print(" " + subprocess.check_output(["pi", "--version"], text=True).strip())
|
||||
print(" " + subprocess.check_output(["mosaic", "--version"], text=True).strip())
|
||||
|
||||
print("\nP2_EVENT_ORDER_AND_NONCE_MAP:")
|
||||
for record in records:
|
||||
if record["seq"] <= 12 and record["event"] in {
|
||||
"after_provider_response",
|
||||
"message_end",
|
||||
"tool_call",
|
||||
"tool_execute",
|
||||
} and (
|
||||
record["event"] != "message_end"
|
||||
or record.get("role") == "assistant"
|
||||
):
|
||||
print(json.dumps(record, sort_keys=True))
|
||||
|
||||
print("\nP2_LAST_OR_CLOSED:")
|
||||
print(json.dumps(next(r for r in records if r["event"] == "session_start"), sort_keys=True))
|
||||
print(json.dumps(next(r for r in closed if r["event"] == "session_start"), sort_keys=True))
|
||||
|
||||
print("\nP3_GENERATION_BROKER:")
|
||||
for record in generations:
|
||||
if record["event"] in {"probe_lease_promoted", "runtime_generation_bump"}:
|
||||
print(json.dumps(record, sort_keys=True))
|
||||
|
||||
print("\nP5_SOURCE_INVALIDATION:")
|
||||
fault_reasons = {"missing", "oversize", "hash-mismatch"}
|
||||
emitted_context: set[str] = set()
|
||||
emitted_tool: set[str] = set()
|
||||
for record in records:
|
||||
source_reason = record.get("sourceValidation", {}).get("reason")
|
||||
if (
|
||||
record["event"] == "context_return"
|
||||
and source_reason in fault_reasons
|
||||
and source_reason not in emitted_context
|
||||
):
|
||||
print(json.dumps(record, sort_keys=True))
|
||||
emitted_context.add(source_reason)
|
||||
mapping_reason = record.get("mapping", {}).get("sourceReason")
|
||||
if (
|
||||
record["event"] == "tool_call"
|
||||
and not record.get("allowed")
|
||||
and mapping_reason in fault_reasons
|
||||
and mapping_reason not in emitted_tool
|
||||
):
|
||||
print(json.dumps(record, sort_keys=True))
|
||||
emitted_tool.add(mapping_reason)
|
||||
emitted_broker: set[str] = set()
|
||||
for record in generations:
|
||||
reason = record.get("source_reason")
|
||||
if record["event"] == "source_invalidation_revoke" and reason not in emitted_broker:
|
||||
print(json.dumps(record, sort_keys=True))
|
||||
emitted_broker.add(str(reason))
|
||||
|
||||
print("\nP6_PI_CONTEXT_ATOMIC_OBSERVATION:")
|
||||
for record in records:
|
||||
include = (
|
||||
(record["event"] == "context_return" and record.get("injectionDecision") == "ONE_ATOMIC_AGENT_MESSAGE")
|
||||
or (record["event"] == "before_provider_request" and record.get("finalPayloadValid"))
|
||||
or (record["event"] == "message_end" and record.get("exactContextBlockCopied"))
|
||||
)
|
||||
if include and record["seq"] <= 12:
|
||||
print(json.dumps(record, sort_keys=True))
|
||||
|
||||
print("\nRPC_EVENT_COUNTS:")
|
||||
counts: dict[str, int] = {}
|
||||
for line in rpc_lines:
|
||||
try:
|
||||
event = json.loads(line)
|
||||
except json.JSONDecodeError:
|
||||
continue
|
||||
key = str(event.get("type"))
|
||||
counts[key] = counts.get(key, 0) + 1
|
||||
print(json.dumps(counts, sort_keys=True))
|
||||
print("stderr_nonempty=" + str(bool(stderr_lines)))
|
||||
for line in stderr_lines[:10]:
|
||||
print("stderr: " + line[:500])
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user