Files
m4 d4b53bfb08
Docker / build (push) Has been cancelled
Test / pytest (ubuntu-latest, 3.12) (push) Has been cancelled
Test / pytest (windows-latest, 3.11) (push) Has been cancelled
Test / pytest (windows-latest, 3.12) (push) Has been cancelled
Lint / ruff (push) Has been cancelled
Test / pytest (ubuntu-latest, 3.11) (push) Has been cancelled
Build / build (push) Has been cancelled
test: cover stop contract, execution adapters, checkpointer race and runtime identity
2026-09-13 15:12:17 +08:00

203 lines
9.8 KiB
Python

"""Isolated B04 process-loss probe. Not a production owner or adapter."""
from __future__ import annotations
import asyncio
import dataclasses
import hashlib
import json
import os
from pathlib import Path
import pickle
import socket
import subprocess
import sys
import time
ROOT = Path(__file__).resolve().parents[3]
REPO = ROOT / "EvoScientist"
BASE = ROOT / ".hermes/test-runtime/unified-execution/b04"
TARGET = "b04_current_host_restart_identity_no_automatic_replay"
def append(path, value):
with path.open("a") as out:
out.write(json.dumps(value, sort_keys=True) + "\n")
out.flush()
os.fsync(out.fileno())
async def child(mode, directory):
started = time.perf_counter()
def blocked(*args, **kwargs):
raise AssertionError("B04 prohibits network")
socket.socket.connect = blocked
socket.socket.connect_ex = blocked
socket.getaddrinfo = blocked
sys.path.insert(0, str(REPO))
import importlib
from langchain_core.language_models.fake_chat_models import FakeMessagesListChatModel
from langchain_core.messages import AIMessage
from langgraph.checkpoint.memory import InMemorySaver
from EvoScientist.config.settings import EvoScientistConfig
from EvoScientist.llm.contracts import AgentInputV3, HmacGrantAuthority, WebHostContext
from EvoScientist.llm.model_config import FileEvoModelConfigStore
from EvoScientist.llm.runtime import EvoModelRuntime
from EvoScientist.llm.host_execution_registry import SQLiteHostRegistry
from EvoScientist.web_runtime import create_web_agent, web_tool_registry_manifest
from EvoScientist.workspace_files import ScopedFilesystemBackend
from tests.v3_fixtures import v3_payload, identity_ring, RUNTIME_SECRET, RUNTIME_KEY_ID
from tests.test_web_model_runtime import _preparation, _admission
import resource
imported = time.perf_counter()
module = importlib.import_module("EvoScientist.EvoScientist")
module._load_mcp_tools_cached = lambda **kw: {}
module._load_mcp_config_once = lambda: ("b04-empty", {})
entered = asyncio.Event()
class LocalBlockingModel(FakeMessagesListChatModel):
def bind_tools(self, tools, **kwargs):
return self
async def _agenerate(self, *args, **kwargs):
append(directory / "model-entries.jsonl", {"pid": os.getpid(), "mode": mode})
entered.set()
await asyncio.Event().wait()
class EvidenceSink:
async def commit(self, event):
append(directory / "events.jsonl", {
"event_id": event.event_id, "run_id": event.run_id,
"kind": event.kind, "payload": dict(event.payload),
})
return "committed"
async def confirm(self, *args):
raise AssertionError("unexpected commit ambiguity")
authority = HmacGrantAuthority(RUNTIME_SECRET, RUNTIME_KEY_ID)
store = FileEvoModelConfigStore(directory / "routes.yaml", admin_verifier=authority)
if mode == "start":
payload = v3_payload()
for provider in payload["providers"].values():
for model in provider["models"]:
model["context_window"] = 131072
store.bootstrap_for_development(payload)
config = EvoScientistConfig(enable_async_subagents=False, enable_scheduler=False,
memory_workers_enabled=False, auto_approve=False)
runtime = EvoModelRuntime(store, admission_verifier=authority, quote_authority=authority,
host_registry=SQLiteHostRegistry(directory / "host.sqlite", host_id="b04-host", boot_id=__import__("uuid").uuid4().hex),
identity_key_ring=identity_ring(),
model_factory=lambda **kw: LocalBlockingModel(responses=[AIMessage(content="unused")]),
agent_factory=lambda snapshot, host, models: create_web_agent(
snapshot=snapshot, host=host, model_set=models, config=config))
constructed = time.perf_counter()
metrics = {"mode": mode, "pid": os.getpid(), "runtime_instance_id": runtime.runtime_instance_id,
"import_seconds": imported-started, "runtime_setup_seconds": constructed-imported}
if mode == "start":
for name in ("files", "memory"):
(directory / name).mkdir(exist_ok=True)
tools, revision = web_tool_registry_manifest()
host = WebHostContext(str(directory / "files"), str(directory / "memory"),
ScopedFilesystemBackend(directory / "files"), InMemorySaver(),
tool_registry=tools, tool_registry_revision=revision,
tool_selector_threshold=10000, runtime_event_sink=EvidenceSink())
agent_input = AgentInputV3("B04 controlled host loss", "b04:isolated:thread")
t = time.perf_counter()
quote = await runtime.prepare_model_run(_preparation(authority, agent_input, title_policy="disabled"), agent_input, host)
metrics["prepare_seconds"] = time.perf_counter()-t
admission = _admission(authority, quote)
t = time.perf_counter()
run = await runtime.start_web_run(admission)
metrics["start_graph_seconds"] = time.perf_counter()-t
metrics["run_id"] = run.run_id
metrics["same_grant_same_object"] = await runtime.start_web_run(admission) is run
await asyncio.wait_for(entered.wait(), 10)
metrics["start_to_model_entry_seconds"] = time.perf_counter()-t
metrics["task_done_before_loss"] = run._agent_task.done()
with (directory / "admission.pickle").open("wb") as out:
pickle.dump(admission, out)
out.flush()
os.fsync(out.fileno())
else:
prior = json.loads((directory / "start.json").read_text())
metrics["previous_run_id"] = prior["run_id"]
metrics["incarnation_changed"] = prior["runtime_instance_id"] != runtime.runtime_instance_id
metrics["public_identity_query_methods"] = [name for name in ("inspect", "get_run", "check_run", "get_execution") if callable(getattr(runtime, name, None))]
metrics["prepared_count"] = len(runtime._prepared)
metrics["started_grants_count"] = len(runtime._started_grants)
# Trusted local test artifact, never an untrusted or production pickle.
with (directory / "admission.pickle").open("rb") as inp:
admission = pickle.load(inp)
try:
await runtime.start_web_run(admission)
except Exception as exc:
metrics["old_admission_result"] = str(exc)
else:
raise AssertionError("old admission unexpectedly started")
await asyncio.sleep(0.1)
metrics["model_entries_total"] = len((directory / "model-entries.jsonl").read_text().splitlines())
assert metrics["incarnation_changed"]
metrics["inspection"] = runtime.inspect(prior["run_id"])
assert metrics["old_admission_result"] == "EXECUTION_UNKNOWN"
assert metrics["inspection"]["execution_id"] == prior["run_id"]
assert metrics["inspection"]["status"] == "unknown"
assert metrics["inspection"]["resources_confirmed_exited"] is False
assert metrics["model_entries_total"] == 1
assert metrics["public_identity_query_methods"] == ["inspect"]
metrics["max_rss_bytes_macos"] = resource.getrusage(resource.RUSAGE_SELF).ru_maxrss
(directory / f"{mode}.json").write_text(json.dumps(metrics, indent=2))
print(json.dumps(metrics), flush=True)
if mode == "start":
# Deliberate abrupt loss while Graph is active: no cancellation/teardown.
os._exit(0)
def main():
if len(sys.argv) > 1:
asyncio.run(child(sys.argv[1], Path(sys.argv[2])))
print("B04 restart asyncio.run and interpreter exit", flush=True)
return
os.umask(0o077)
BASE.mkdir(parents=True, exist_ok=True)
ledger = BASE / "attempts.jsonl"
records = [json.loads(line) for line in ledger.read_text().splitlines()] if ledger.exists() else []
used = sum(r.get("phase") == "start" for r in records)
if used >= 5:
raise SystemExit("B04 cumulative budget exhausted")
number = used + 1
directory = BASE / f"attempt-{number}"
directory.mkdir()
home = directory / "home"
home.mkdir()
env = {"PATH": "/usr/bin:/bin", "HOME": str(home), "EVOSCIENTIST_HOME": str(home),
"EVOSCIENTIST_CONFIG_DIR": str(home / "config"), "PYTHON_DOTENV_DISABLED": "1",
"PYTEST_DISABLE_PLUGIN_AUTOLOAD": "1", "WEB_RUNTIME_TEST_KEY": "b04-local-dummy",
"PYTHONPATH": str(REPO), "LANG": "en_US.UTF-8"}
command = [str(REPO / ".venv/bin/python"), str(Path(__file__).resolve())]
snapshot = {name: hashlib.sha256((REPO / name).read_bytes()).hexdigest() for name in (
"EvoScientist/llm/runtime.py", "EvoScientist/llm/contracts.py", "EvoScientist/web_runtime.py",
"tests/support/b04_host_restart_probe.py")}
append(ledger, {"phase": "start", "target": TARGET, "attempt": number, "limit": 5,
"time": time.time(), "command": command, "snapshot_sha256": snapshot})
code = 0
for mode in ("start", "restart"):
t = time.perf_counter()
with (directory / f"{mode}.log").open("w") as out:
try:
result = subprocess.run(command + [mode, str(directory)], cwd=REPO, env=env,
stdout=out, stderr=subprocess.STDOUT, timeout=40)
code = result.returncode
except subprocess.TimeoutExpired:
code = 124
append(ledger, {"phase": "child_exit", "attempt": number, "mode": mode,
"exit_code": code, "wall_seconds": time.perf_counter()-t})
print((directory / f"{mode}.log").read_text())
if code:
break
append(ledger, {"phase": "result", "attempt": number, "exit_code": code,
"time": time.time(), "remaining": 5-number, "directory": str(directory)})
raise SystemExit(code)
if __name__ == "__main__":
main()