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
203 lines
9.8 KiB
Python
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() |