File size: 7,411 Bytes
f76c374 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 | """Cross-process falsification of the per-session fence.
The unit tests share one interpreter, so they cannot see the failure this whole
change exists to prevent: two SEPARATE gateway processes, each with its own
snapshot of a conversation, both writing to it. That is how the defect was found
and it is the only way to prove it is closed.
Run against the fork's own HERMES_HOME so nothing here touches a real profile:
python scripts/probe_active_session_exclusivity.py
It drives two real ``python -m tui_gateway.entry`` processes over stdio JSON-RPC
and asserts the sequence the reviewer specified:
A resume S, submit -> claims the session
B resume S, submit -> typed SESSION_NOT_OWNED, no row, no turn
A exits -> its lease is pruned as a dead owner
B submit again -> succeeds
No provider is required. The fence is checked BEFORE the agent is built, so a
submit that later fails for want of a model still proves who owns the session --
which is the property under test, and keeps the probe free of credentials and of
inference cost.
"""
from __future__ import annotations
import json
import os
import subprocess
import sys
import time
from pathlib import Path
REPO = Path(__file__).resolve().parent.parent
PYTHON = REPO / "venv" / "Scripts" / "python.exe"
if not PYTHON.exists(): # posix layout
PYTHON = REPO / "venv" / "bin" / "python"
class Gateway:
"""One gateway process, spoken to the way the TUI speaks to it."""
def __init__(self, name: str, home: Path):
env = dict(os.environ)
env["HERMES_HOME"] = str(home)
env["PYTHONUNBUFFERED"] = "1"
self.name = name
self.proc = subprocess.Popen(
[str(PYTHON), "-u", "-m", "tui_gateway.entry"],
cwd=str(REPO),
env=env,
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
encoding="utf-8",
errors="replace",
)
self._next = 1
self.ready()
def _read(self):
line = self.proc.stdout.readline()
if not line:
raise RuntimeError(f"[{self.name}] gateway closed its pipe")
line = line.strip()
if not line:
return None
try:
return json.loads(line)
except json.JSONDecodeError:
return None
def ready(self, timeout: float = 180.0) -> None:
deadline = time.time() + timeout
while time.time() < deadline:
msg = self._read()
if msg and msg.get("method") == "event":
if msg.get("params", {}).get("type") == "gateway.ready":
return
raise RuntimeError(f"[{self.name}] never announced gateway.ready")
def call(self, method: str, params: dict, timeout: float = 180.0) -> dict:
rid = str(self._next)
self._next += 1
self.proc.stdin.write(json.dumps({"jsonrpc": "2.0", "id": rid, "method": method, "params": params}) + "\n")
self.proc.stdin.flush()
deadline = time.time() + timeout
while time.time() < deadline:
msg = self._read()
if msg and msg.get("id") == rid:
return msg
raise RuntimeError(f"[{self.name}] timed out calling {method}")
def close(self):
try:
self.proc.stdin.close()
except Exception:
pass
try:
self.proc.terminate()
self.proc.wait(timeout=15)
except Exception:
try:
self.proc.kill()
except Exception:
pass
def reason_of(response: dict):
return (response.get("error") or {}).get("data", {}).get("reason")
def registry(home: Path):
path = home / "runtime" / "active_sessions.json"
try:
return json.loads(path.read_text(encoding="utf-8")).get("entries", [])
except Exception:
return []
def main() -> int:
home = REPO / ".probe-home"
# A fresh profile each run: a lease left by a previous run would make the
# first check pass or fail for the wrong reason.
import shutil
shutil.rmtree(home, ignore_errors=True)
failures = []
def check(label: str, ok: bool, detail: str = ""):
print(f" {'PASS' if ok else 'FAIL'} {label}{(' -- ' + detail) if detail else ''}")
if not ok:
failures.append(label)
a = Gateway("A", home)
b = None
try:
created = a.call("session.create", {"cols": 80})
sid_a = created["result"]["session_id"]
# Opening a chat must not claim anything -- an idle composer is invisible
# and a slot held by one would fence a real turn for no reason.
check("session.create claims nothing", registry(home) == [], f"{len(registry(home))} entries")
a.call("prompt.submit", {"session_id": sid_a, "text": "probe: A takes the session"})
held = registry(home)
check("A's first turn claims a session", len(held) == 1, json.dumps(held)[:200])
if not held:
raise RuntimeError("A never claimed anything; nothing further can be tested")
# The STORED key, which only materialises when a turn is first submitted --
# and which is what the lease must be keyed on. A lease keyed on the live
# runtime id would fence nothing: two processes resuming one conversation
# have different runtime ids by construction.
key = held[0].get("session_id")
print(f"A live session {sid_a}, stored key {key}")
check("the lease is keyed on the STORED session, not the runtime handle",
bool(key) and key != sid_a, f"key={key} runtime={sid_a}")
b = Gateway("B", home)
resumed = b.call("session.resume", {"session_id": key})
check("B may still RESUME (reading is never fenced)", "result" in resumed,
json.dumps(resumed.get("error", ""))[:160])
sid_b = resumed.get("result", {}).get("session_id")
before = len(registry(home))
refused = b.call("prompt.submit", {"session_id": sid_b, "text": "probe: B must not write"})
check("B's submit is refused", refused.get("error") is not None,
json.dumps(refused.get("result", ""))[:120])
check("refusal is typed SESSION_NOT_OWNED", reason_of(refused) == "SESSION_NOT_OWNED",
str(reason_of(refused)))
check("refusal left the registry untouched", len(registry(home)) == before)
# A dies without releasing -- the crash case, not a clean handoff.
a.proc.kill()
a.proc.wait(timeout=30)
time.sleep(1.0)
retried = b.call("prompt.submit", {"session_id": sid_b, "text": "probe: B may write now"})
check("after A dies, B's retry is accepted", retried.get("error") is None,
json.dumps(retried.get("error", ""))[:200])
held = registry(home)
check("and B now owns the session", len(held) == 1 and held[0].get("session_id") == key,
json.dumps(held)[:160])
finally:
if b is not None:
b.close()
a.close()
print()
if failures:
print(f"FAILED: {len(failures)} check(s): {', '.join(failures)}")
return 1
print("All cross-process checks passed.")
return 0
if __name__ == "__main__":
sys.exit(main())
|