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
547 lines
27 KiB
Python
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) |