"""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()