Files
EvoScientist-Multi/tests/test_execution_adapter_stop_contract.py
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

547 lines
27 KiB
Python

"""B02 real resource contracts; run via this file's isolated budget runner."""
from __future__ import annotations
import asyncio
import faulthandler
import json
import os
from pathlib import Path
import signal
import socket
import subprocess
import sys
import threading
from contextlib import asynccontextmanager
import pytest
PLANNED_NODES = (
"test_http_start_without_observer_and_cancel_confirms_eof",
"test_cancel_before_first_observer_never_restarts_graph",
"test_concurrent_cancel_and_cancelled_waiter_preserve_cleanup",
"test_run_timeout_reports_timeout_after_http_exit",
"test_real_execute_cancel_confirms_parent_child_and_worker_exit",
"test_slow_native_worker_retains_run_ownership_until_exit",
"test_owned_clients_close_once_without_closing_borrowed_resources",
"test_parallel_runs_do_not_share_owned_transports",
)
async def build_run(tmp_path, monkeypatch, server):
import importlib
from EvoScientist.config.settings import EvoScientistConfig
from EvoScientist.llm.contracts import AgentInputV3, HmacGrantAuthority, WebHostContext
from EvoScientist.llm.model_config import FileEvoModelConfigStore, EvoModelConfig, endpoint_fingerprint, route_semantics_hash
from EvoScientist.llm.runtime import EvoModelRuntime
from EvoScientist.web_runtime import create_web_agent, web_tool_registry_manifest
from EvoScientist.workspace_files import ScopedFilesystemBackend
from langgraph.checkpoint.memory import InMemorySaver
from tests.v3_fixtures import v3_payload, identity_ring, RUNTIME_SECRET, RUNTIME_KEY_ID
from tests.test_web_model_runtime import _preparation, _admission, _Sink
module = importlib.import_module("EvoScientist.EvoScientist")
monkeypatch.setattr(module, "_load_mcp_tools_cached", lambda **kw: {})
monkeypatch.setattr(module, "_load_mcp_config_once", lambda: ("b02-empty", {}))
monkeypatch.setenv("WEB_RUNTIME_TEST_KEY", "b02-local-dummy")
payload = v3_payload()
payload["providers"]["custom-openai"]["endpoints"][0]["base_url"] = server.url
payload["providers"]["custom-openai"]["models"][0]["context_window"] = 131072
candidate = EvoModelConfig.parse(payload, require_evidence=False)
ring = identity_ring()
for evidence in payload["capability_evidence"]:
route = candidate.concrete_routes("visible-main")[0]
evidence["probe"]["route_semantics_hash"] = route_semantics_hash(candidate, route, ring.derive_current("ai4sci/route-semantics-hash/v3")[1])
evidence["probe"]["endpoint_fingerprint"] = endpoint_fingerprint(candidate, route, ring.derive_current("ai4sci/endpoint-fingerprint/v3")[1])
authority = HmacGrantAuthority(RUNTIME_SECRET, RUNTIME_KEY_ID)
store = FileEvoModelConfigStore(tmp_path / "routes.yaml", admin_verifier=authority)
store.bootstrap_for_development(payload)
config = EvoScientistConfig(enable_async_subagents=False, enable_scheduler=False, memory_workers_enabled=False, auto_approve=False)
def model_factory(**kwargs):
model = EvoModelRuntime._default_model_factory(**kwargs)
server.models.append(model)
assert model.streaming is True and model.disable_streaming is False
return model
runtime = EvoModelRuntime(store, admission_verifier=authority, quote_authority=authority,
model_factory=model_factory,
identity_key_ring=ring, agent_factory=lambda snapshot, host, models:
create_web_agent(snapshot=snapshot, host=host, model_set=models, config=config))
root = tmp_path / "files"
memory = tmp_path / "memory"
root.mkdir()
memory.mkdir()
tools, revision = web_tool_registry_manifest()
sink = _Sink()
host = WebHostContext(str(root), str(memory), ScopedFilesystemBackend(root), InMemorySaver(),
tool_selector_threshold=10000, tool_registry=tools,
tool_registry_revision=revision, runtime_event_sink=sink)
agent_input = AgentInputV3("Wait for a cancellation probe.", "b02:isolated:thread")
quote = await runtime.prepare_model_run(_preparation(authority, agent_input, title_policy="disabled"), agent_input, host)
run = await runtime.start_web_run(_admission(authority, quote))
return run, sink
def test_http_start_without_observer_and_cancel_confirms_eof(tmp_path, monkeypatch, isolated_network):
faulthandler.dump_traceback_later(45, repeat=False)
async def scenario():
async with LoopbackSSE().serve() as server:
run, sink = await build_run(tmp_path, monkeypatch, server)
try:
try:
await asyncio.wait_for(server.requested.wait(), 5)
except TimeoutError:
pytest.fail("accepted start_web_run did not start HTTP without observer")
assert not server.errors
observer = run.stream()
chunks = []
async with asyncio.timeout(5):
async for event in observer:
if event.payload.get("type") == "text":
chunks.append(event.payload.get("content", ""))
if "B02 incremental" in "".join(chunks):
break
assert len(chunks) >= 2, chunks
await observer.aclose()
assert not server.peer_eof.is_set(), "observer owns execution"
assert await run.cancel("b02-http") == "cancelled"
await asyncio.wait_for(run.wait_stopped(), 4)
await asyncio.wait_for(server.peer_eof.wait(), 2)
assert run._agent_task.done()
terminals = [e for e in sink.events if e.payload.get("kind") == "run_terminal"]
assert len(terminals) == 1 and terminals[0].payload["outcome"] == "cancelled"
assert len(server.requests) == 1
print(json.dumps({"evidence": "stream_cancel", "stream": server.requests[0]["stream"],
"chunks": chunks, "peer_eof": server.peer_eof.is_set(),
"agent_done": run._agent_task.done(), "terminals": len(terminals)}), flush=True)
finally:
if run._agent_task is not None and not run._agent_task.done():
run._agent_task.cancel()
await asyncio.gather(run._agent_task, return_exceptions=True)
asyncio.run(scenario())
print("B02 asyncio.run exited; awaiting interpreter teardown", flush=True)
@pytest.fixture
def isolated_network(monkeypatch):
# These must be set BEFORE interpreter startup to prevent import-time leaks.
home = Path(os.environ["HOME"]).resolve()
assert "unified-execution" in str(home)
assert Path(os.environ["EVOSCIENTIST_HOME"]).resolve() == home
assert os.environ.get("PYTHON_DOTENV_DISABLED") == "1"
assert os.environ.get("PYTEST_DISABLE_PLUGIN_AUTOLOAD") == "1"
assert Path(os.environ["EVOSCIENTIST_CONFIG_DIR"]).resolve().is_relative_to(home)
secret_names = [k for k in os.environ if any(
part in k.upper() for part in ("API_KEY", "TOKEN", "SECRET", "PROXY")
)]
assert not secret_names, f"runner inherited credential/proxy keys: {secret_names}"
original_connect = socket.socket.connect
original_connect_ex = socket.socket.connect_ex
original_lookup = socket.getaddrinfo
def require_loopback(address):
assert isinstance(address, tuple) and address[0] == "127.0.0.1", (
f"non-loopback network attempt: {address!r}"
)
def connect(sock, address):
require_loopback(address)
return original_connect(sock, address)
def connect_ex(sock, address):
require_loopback(address)
return original_connect_ex(sock, address)
def lookup(host, *args, **kwargs):
assert host == "127.0.0.1", f"external DNS prohibited: {host!r}"
return original_lookup(host, *args, **kwargs)
monkeypatch.setattr(socket.socket, "connect", connect)
monkeypatch.setattr(socket.socket, "connect_ex", connect_ex)
monkeypatch.setattr(socket, "getaddrinfo", lookup)
class LoopbackSSE:
"""Actual HTTP/1.1 stream, with peer EOF recorded before server teardown."""
def __init__(self, *, tool_call=False):
self.tool_call = tool_call
self.requested = asyncio.Event()
self.peer_eof = asyncio.Event()
self.requests = []
self.tasks = set()
self.writers = set()
self.errors = []
self.models = []
async def handle(self, reader, writer):
task = asyncio.current_task()
self.tasks.add(task)
self.writers.add(writer)
try:
header = await asyncio.wait_for(reader.readuntil(b"\r\n\r\n"), 5)
headers = dict(line.split(b":", 1) for line in header.split(b"\r\n")[1:] if b":" in line)
length = int(next((v for k, v in headers.items() if k.lower() == b"content-length"), b"0"))
assert 0 < length < 2_000_000
body = json.loads(await reader.readexactly(length))
assert header.startswith(b"POST /v1/chat/completions ")
assert body.get("stream") is True, "upstream model HTTP must stream"
assert body.get("tools"), "streaming must retain tools"
self.requests.append(body)
writer.write(b"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nConnection: close\r\n\r\n")
delta: dict = {"role": "assistant", "content": ""}
if self.tool_call:
assert any(t.get("function", {}).get("name") == "execute" for t in body["tools"])
delta["tool_calls"] = [{"index": 0, "id": "b02-execute", "type": "function", "function": {
"name": "execute", "arguments": json.dumps({"command": "b02-heartbeat"}),
}}]
chunk = {"id": "b02-local", "object": "chat.completion.chunk", "created": 1,
"model": body["model"], "choices": [{"index": 0, "delta": delta, "finish_reason": None}]}
writer.write(b"data: " + json.dumps(chunk).encode() + b"\n\n")
await writer.drain()
if not self.tool_call:
for text in ("B02 ", "incremental"):
chunk["choices"][0]["delta"] = {"content": text}
writer.write(b"data: " + json.dumps(chunk).encode() + b"\n\n")
await writer.drain()
await asyncio.sleep(.02)
if self.tool_call:
chunk["choices"] = [{"index": 0, "delta": {}, "finish_reason": "tool_calls"}]
writer.write(b"data: " + json.dumps(chunk).encode() + b"\n\ndata: [DONE]\n\n")
await writer.drain()
self.requested.set()
if self.tool_call:
return
assert await reader.read() == b""
self.peer_eof.set()
except Exception as exc:
self.errors.append(exc)
self.requested.set()
finally:
writer.close()
await writer.wait_closed()
self.writers.discard(writer)
self.tasks.discard(task)
@asynccontextmanager
async def serve(self):
server = await asyncio.start_server(self.handle, "127.0.0.1", 0)
self.url = f"http://127.0.0.1:{server.sockets[0].getsockname()[1]}/v1"
try:
yield self
finally:
# Emergency teardown is deliberately separate from peer_eof evidence.
closed = set()
emergency_closed = 0
for model in self.models:
for attribute in ("root_async_client", "root_client"):
client = getattr(model, attribute, None)
if client is None or id(client) in closed:
continue
closed.add(id(client))
if client.is_closed():
continue
emergency_closed += 1
result = client.close()
if hasattr(result, "__await__"):
await asyncio.wait_for(result, 3)
assert client.is_closed()
print(json.dumps({"evidence": "fixture_sdk_cleanup", "sdk_wrappers_seen": len(closed),
"emergency_clients_closed": emergency_closed,
"threads": [t.name for t in threading.enumerate()]}), flush=True)
server.close()
await server.wait_closed()
for writer in tuple(self.writers):
writer.close()
tasks = tuple(self.tasks)
for task in tasks:
task.cancel()
await asyncio.gather(*tasks, return_exceptions=True)
class HeartbeatExecutor:
"""Whitelist-only launch fixture; production aexecute/collector stay real."""
def __init__(self, root):
self.root = root
self.process = None
self.started = threading.Event()
self.finished = threading.Event()
self.worker_release = threading.Event()
self.worker_release.set()
self.child_pid = root / "child.pid"
self.heartbeat = root / "heartbeat"
def execute(self, command, *, timeout, cancel_event):
from EvoScientist.native_sandbox import _collect_process
from deepagents.backends.protocol import ExecuteResponse
assert command == "b02-heartbeat"
child = (
"import pathlib,time; p=pathlib.Path('heartbeat'); "
"\nwhile True:\n p.write_text(str(time.monotonic_ns())); time.sleep(.03)\n"
)
parent = (
"import pathlib,signal,subprocess,sys,time\n"
f"p=subprocess.Popen([sys.executable,'-I','-c',{child!r}])\n"
"pathlib.Path('child.pid').write_text(str(p.pid))\n"
"def stop(*args):\n p.wait(timeout=2); sys.exit(0)\n"
"signal.signal(signal.SIGTERM,stop)\n"
"while True: time.sleep(.03)\n"
)
self.process = subprocess.Popen(
[sys.executable, "-I", "-c", parent], cwd=self.root,
env={"HOME": str(self.root), "PATH": "/usr/bin:/bin"},
stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
start_new_session=True, close_fds=True,
)
self.started.set()
try:
raw, code, truncated, _, _ = _collect_process(
self.process, timeout=8, output_limit=4096, cancel_event=cancel_event,
)
# Slow-worker node clears this gate before launch, then releases it
# only after verifying no successful cancellation was committed.
assert self.worker_release.wait(8), "test worker release deadline"
return ExecuteResponse(output=raw.decode(), exit_code=code, truncated=truncated)
finally:
self.finished.set()
async def assert_stopped_before_teardown(self):
assert self.finished.is_set(), "native worker not finished"
assert self.process is not None and self.process.returncode is not None
with pytest.raises(ProcessLookupError):
os.killpg(self.process.pid, 0)
with pytest.raises(ProcessLookupError):
os.kill(int(self.child_pid.read_text()), 0)
before = self.heartbeat.read_bytes()
await asyncio.sleep(.15)
assert self.heartbeat.read_bytes() == before
def emergency_cleanup(self):
self.worker_release.set()
if self.process is not None:
try:
os.killpg(self.process.pid, signal.SIGKILL)
except ProcessLookupError:
pass
self.process.wait(timeout=3)
def test_cancel_before_first_observer_never_restarts_graph(tmp_path, monkeypatch, isolated_network):
async def scenario():
async with LoopbackSSE().serve() as server:
run, sink = await build_run(tmp_path, monkeypatch, server)
assert await run.cancel("before-observer") == "cancelled"
events = [event async for event in run.stream()]
await asyncio.sleep(.1)
assert len([e for e in events if e.payload.get("kind") == "run_terminal"]) == 1
assert not any(e.payload.get("kind") == "run_started" for e in sink.events)
assert not server.requests
assert run._agent_task is None or run._agent_task.done()
asyncio.run(scenario())
def test_slow_native_worker_retains_run_ownership_until_exit(tmp_path, isolated_network):
from EvoScientist.native_sandbox import NativeWorkspaceBackend
async def scenario():
executor = HeartbeatExecutor(tmp_path)
executor.worker_release.clear()
backend = object.__new__(NativeWorkspaceBackend)
backend._executor = executor
task = asyncio.create_task(backend.aexecute("b02-heartbeat"))
try:
async with asyncio.timeout(5):
while not executor.heartbeat.exists() or not executor.child_pid.exists():
await asyncio.sleep(.02)
task.cancel()
await asyncio.sleep(2.2)
assert not task.done(), "cancel abandoned an unconfirmed native worker"
assert not executor.finished.is_set()
executor.worker_release.set()
with pytest.raises(asyncio.CancelledError):
await asyncio.wait_for(task, 3)
await executor.assert_stopped_before_teardown()
finally:
executor.emergency_cleanup()
await asyncio.gather(task, return_exceptions=True)
asyncio.run(scenario())
def test_owned_clients_close_once_without_closing_borrowed_resources(tmp_path, monkeypatch, isolated_network):
async def scenario():
async with LoopbackSSE().serve() as server:
run, sink = await build_run(tmp_path, monkeypatch, server)
await asyncio.wait_for(server.requested.wait(), 5)
assert not server.errors
clients = {}
counts = {}
for model in server.models:
for name in ("root_client", "root_async_client"):
client = getattr(model, name)
transport = client._client
key = id(transport)
if key in clients:
continue
clients[key] = transport
counts[key] = 0
original = transport.aclose if name == "root_async_client" else transport.close
if name == "root_async_client":
async def close(original=original, key=key):
counts[key] += 1
await original()
monkeypatch.setattr(transport, "aclose", close)
else:
def close(original=original, key=key):
counts[key] += 1
original()
monkeypatch.setattr(transport, "close", close)
borrowed_closed = []
for resource in (run._host.checkpointer, run._host.workspace_backend):
monkeypatch.setattr(resource, "close", lambda: borrowed_closed.append(True), raising=False)
# LangChain resource aliases refer to the same SDK transport.
assert await run.cancel("owned-resource-contract") == "cancelled"
assert await run.wait_stopped() == "cancelled"
assert all(client.is_closed for client in clients.values()), "run left SDK transports open"
assert set(counts.values()) == {1}, counts
assert not borrowed_closed
print(json.dumps({"evidence": "production_owned_cleanup", "counts": list(counts.values()),
"borrowed_closed": borrowed_closed}), flush=True)
asyncio.run(scenario())
def test_concurrent_cancel_and_cancelled_waiter_preserve_cleanup(tmp_path, monkeypatch, isolated_network):
async def scenario():
async with LoopbackSSE().serve() as server:
run, sink = await build_run(tmp_path, monkeypatch, server)
await asyncio.wait_for(server.requested.wait(), 5)
transport = server.models[0].root_async_client._client
original = transport.aclose
entered, release = asyncio.Event(), asyncio.Event()
calls = []
async def delayed_close():
calls.append(True)
entered.set()
await release.wait()
await original()
monkeypatch.setattr(transport, "aclose", delayed_close)
cancellations = [asyncio.create_task(run.cancel(str(i))) for i in range(4)]
try:
await asyncio.wait_for(entered.wait(), 4)
waiter = asyncio.create_task(run.wait_stopped(timeout=None))
await asyncio.sleep(0)
waiter.cancel()
cancellations[0].cancel()
await asyncio.gather(waiter, cancellations[0], return_exceptions=True)
assert await run.wait_stopped(timeout=.02) == "unknown"
assert id(transport) in run._owned_clients
assert run._terminal_event is None
assert not run._cleanup_task.done()
assert not transport.is_closed
release.set()
results = await asyncio.gather(*cancellations[1:])
assert results == ["cancelled"] * 3
assert calls == [True]
assert not run._owned_clients
terminals = [e for e in sink.events if e.payload.get("kind") == "run_terminal"]
assert len(terminals) == 1
print(json.dumps({"evidence": "cleanup_race", "bounded_wait": "unknown",
"retained_until_close": True, "close_calls": len(calls),
"terminals": len(terminals), "results": results}), flush=True)
finally:
release.set()
await asyncio.gather(*cancellations, return_exceptions=True)
asyncio.run(scenario())
def test_run_timeout_reports_timeout_after_http_exit(tmp_path, monkeypatch, isolated_network):
from dataclasses import replace
async def scenario():
async with LoopbackSSE().serve() as server:
run, sink = await build_run(tmp_path, monkeypatch, server)
# start_web_run schedules the task; replace before yielding to it.
run._snapshot = replace(run._snapshot, active_run_timeout_seconds=.5)
assert await run.wait_stopped(timeout=5) == "failed"
await asyncio.wait_for(server.peer_eof.wait(), 2)
assert server.requests and not server.errors
assert run._terminal_event.payload["error_code"] == "RUN_TIMEOUT"
assert not run._owned_clients
assert all(m.root_async_client.is_closed() and m.root_client.is_closed() for m in server.models)
print(json.dumps({"evidence": "run_timeout", "error_code": run._terminal_event.payload["error_code"],
"peer_eof": server.peer_eof.is_set(), "resources_remaining": len(run._owned_clients)}), flush=True)
asyncio.run(scenario())
def test_parallel_runs_do_not_share_owned_transports(tmp_path, monkeypatch, isolated_network):
async def scenario():
async with LoopbackSSE().serve() as server:
roots = [tmp_path / str(i) for i in range(2)]
for root in roots:
root.mkdir()
first, _ = await build_run(roots[0], monkeypatch, server)
second, _ = await build_run(roots[1], monkeypatch, server)
try:
async with asyncio.timeout(5):
while len(server.requests) < 2:
await asyncio.sleep(.01)
first_ids = set(first._owned_clients)
second_ids = set(second._owned_clients)
assert first_ids.isdisjoint(second_ids), "independent runs share owned transports"
assert await first.cancel("first-only") == "cancelled"
assert all(not client.is_closed for client in second._owned_clients.values())
assert second._terminal_event is None
print(json.dumps({"evidence": "run_transport_isolation", "first": len(first_ids),
"second": len(second_ids), "disjoint": True}), flush=True)
finally:
await first.cancel("teardown")
await second.cancel("teardown")
asyncio.run(scenario())
if __name__ == "__main__":
repo = Path(__file__).resolve().parents[1]
root = repo.parent / ".hermes/test-runtime/unified-execution/b02"
ledger = root / "attempts.jsonl"
node = "tests/test_execution_adapter_stop_contract.py::" + sys.argv[1]
records = [json.loads(line) for line in ledger.read_text().splitlines()]
attempt = 1 + sum(r.get("phase") == "start" and r.get("node") == node for r in records)
assert sys.argv[1] in PLANNED_NODES, "unregistered node"
limit = 8 if sys.argv[1] == "test_http_start_without_observer_and_cancel_confirms_eof" else 3
assert attempt <= limit, "node budget exhausted"
home = root / (sys.argv[1] + f"-{attempt}-home")
config = home / "config"
config.mkdir(parents=True, exist_ok=True)
env = {
"HOME": str(home), "PATH": str(repo / ".venv/bin") + ":/usr/bin:/bin",
"PYTHONPATH": str(repo), "PYTHON_DOTENV_DISABLED": "1",
"PYTEST_DISABLE_PLUGIN_AUTOLOAD": "1", "PYTHONDONTWRITEBYTECODE": "1",
"EVOSCIENTIST_HOME": str(home), "EVOSCIENTIST_CONFIG_DIR": str(config),
"XDG_CONFIG_HOME": str(config), "EVOSCIENTIST_DATA_DIR": str(home / "data"),
"EVOSCIENTIST_WORKSPACE_DIR": str(home / "workspace"),
"EVOSCIENTIST_SKILLS_DIR": str(home / "skills"),
"EVOSCIENTIST_MEMORIES_DIR": str(home / "memory"),
}
command = [str(repo / ".venv/bin/python"), "-m", "pytest", "--noconftest",
"-p", "no:cacheprovider", "-q", "-s", "--tb=short", "--disable-warnings", node]
with ledger.open("a") as handle:
if attempt == 1:
handle.write(json.dumps({"phase": "registered", "node": node, "used": 0, "limit": limit, "command": command}) + "\n")
handle.write(json.dumps({"phase": "start", "node": node, "attempt": attempt, "command": command}) + "\n")
try:
result = subprocess.run(command, cwd=repo, env=env, capture_output=True, text=True, timeout=60)
except subprocess.TimeoutExpired as exc:
output = exc.stdout or b""
if isinstance(output, bytes):
output = output.decode(errors="replace")
result = subprocess.CompletedProcess(command, 124, output, "runner timeout; child killed and reaped\n")
log = root / (sys.argv[1] + f"-{attempt}.log")
log.write_text(result.stdout + result.stderr)
with ledger.open("a") as handle:
handle.write(json.dumps({"phase": "result", "node": node, "attempt": attempt,
"exit_code": result.returncode, "log": str(log), "remaining": limit-attempt}) + "\n")
print(result.stdout + result.stderr)
sys.exit(result.returncode)