"""B03 host boundaries. Only run through the isolated budget runner.""" import asyncio import base64 import importlib import json import os from pathlib import Path import subprocess import sys import faulthandler import re def stage(name): print('PENDING_STAGE ' + name, flush=True) if __name__ != '__main__': faulthandler.enable() faulthandler.dump_traceback_later(30, repeat=False) stage('module-import') if __name__ != "__main__": from tests.test_execution_adapter_stop_contract import isolated_network NODES = ("test_main_mcp_registry_and_tenants", "test_media_provider_workspace_boundary", "test_manual_pending_public_projection") def factory(tmp_path, responses, tenant="alice", model_class=None, sandbox=False): stage('factory-import-start') from EvoScientist.config.settings import EvoScientistConfig from EvoScientist.llm.contracts import AgentModelSet, WebHostContext 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 langchain_core.language_models.fake_chat_models import FakeMessagesListChatModel class Model(FakeMessagesListChatModel): def bind_tools(self, tools, **kwargs): return self root = tmp_path / tenant root.mkdir(exist_ok=True) backend = ScopedFilesystemBackend(root) if sandbox: from deepagents.backends.protocol import SandboxBackendProtocol class NoExecute(ScopedFilesystemBackend, SandboxBackendProtocol): @property def id(self): return "b03-no-execute" def execute(self, command, *, timeout=None): raise AssertionError("manual pending must not execute") backend = NoExecute(root) model = (model_class or Model)(responses=responses) stage('registry-start') _, revision = web_tool_registry_manifest() stage('graph-build-start') graph = create_web_agent(snapshot=None, host=WebHostContext( workspace_dir=str(root), memory_dir=str(tmp_path / (tenant + "-memory")), workspace_backend=backend, checkpointer=InMemorySaver(), tool_registry_revision=revision, tool_selector_threshold=10000), model_set=AgentModelSet(model, model, model), config=EvoScientistConfig( auto_approve=False, enable_async_subagents=False, enable_scheduler=False, memory_workers_enabled=False)) stage('graph-built') return graph, backend def call(name, args, ident): from langchain_core.messages import AIMessage return AIMessage(content="", tool_calls=[dict(name=name, args=args, id=ident, type="tool_call")]) def config(tenant): return {"configurable": {"thread_id": tenant + ":same-thread", "ai4sci_run_id": tenant}} def test_main_mcp_registry_and_tenants(tmp_path, isolated_network): from EvoScientist.mcp.client import USER_MCP_CONFIG from langchain_core.messages import AIMessage, ToolMessage from EvoScientist.llm.contracts import EvoRuntimeError import pytest server = tmp_path / "server.py" server.write_text('from mcp.server.fastmcp import FastMCP\nm = FastMCP("b03")\n@m.tool()\ndef host_echo(value: str) -> str:\n """Return the supplied value without side effects."""\n return "sdk:" + value\nm.run(transport="stdio")\n') USER_MCP_CONFIG.parent.mkdir(parents=True, exist_ok=True) settings = {"b03": {"transport": "stdio", "command": sys.executable, "args": [str(server)], "expose_to": ["main"]}} USER_MCP_CONFIG.write_text(json.dumps(settings)) async def scenario(): graphs = {t: factory(tmp_path, [call("host_echo", {"value": t}, t), AIMessage(content="done"), call("host_echo", {"value": "blocked"}, t + "-stale")], t)[0] for t in ("alice", "bob")} for tenant, graph in graphs.items(): result = await graph.ainvoke({"messages": [("user", "echo")]}, config(tenant)) tools = [m for m in result["messages"] if isinstance(m, ToolMessage)] assert len(tools) == 1 and "sdk:" + tenant in str(tools[0].content) assert (await graph.aget_state(config("other"))).values == {} settings["b03"]["expose_to"] = ["research-agent"] USER_MCP_CONFIG.write_text(json.dumps(settings)) with pytest.raises(EvoRuntimeError, match="TOOL_REGISTRY_STALE"): await graphs["alice"].ainvoke({"messages": [("user", "echo again")]}, config("alice")) asyncio.run(scenario()) def test_media_provider_workspace_boundary(tmp_path, isolated_network): from langchain_core.messages import AIMessage from langchain_core.language_models.fake_chat_models import FakeMessagesListChatModel seen = [] class Model(FakeMessagesListChatModel): def bind_tools(self, tools, **kwargs): return self def _generate(self, messages, stop=None, run_manager=None, **kwargs): seen.append(messages) return super()._generate(messages, stop=stop, run_manager=run_manager, **kwargs) raw = b"b03-local-media" media = AIMessage(content=[{"type": "image", "base64": base64.b64encode(raw).decode(), "mime_type": "image/png"}]) async def scenario(): graph, alice = factory(tmp_path, [AIMessage(content="done")], model_class=Model) _, bob = factory(tmp_path, [AIMessage(content="done")], "bob") result = await graph.ainvoke({"messages": [media, ("user", "describe prior media")]}, config("alice")) assert seen and "base64" not in str(seen[-1]) stored = next(m for m in result["messages"] if isinstance(m, AIMessage) and isinstance(m.content, list)) path = stored.content[0]["url"] assert path in str(seen[-1]) and "generated_image" in str(seen[-1]) assert alice.download_files([path])[0].content == raw assert bob.download_files([path])[0].error assert bob.download_files(["/../alice/" + path.lstrip("/")])[0].error asyncio.run(scenario()) def test_manual_pending_public_projection(tmp_path, isolated_network): stage('real-import-start') from langchain_core.messages import AIMessage from gateway.services.recoverable_runs import normalize_interrupt, content_hash, _safe_json from gateway.services.projection_event_adapter import pending_input_projection_event stage('real-import-done') async def scenario(): sentinel = "b03-private-" + "canary-7f9a" graph, _ = factory(tmp_path, [call("execute", {"command": "printf " + sentinel}, "private-call"), AIMessage(content="done")], sandbox=True) stage('invoke-start') await graph.ainvoke({"messages": [("user", "manual approval")]}, config("alice")) stage('invoke-done') before = await graph.aget_state(config("alice")) interrupt = before.tasks[0].interrupts[0] internal = {"id": interrupt.id, "value": {**interrupt.value, "checkpoint_details": "internal-only"}} internal["value"]["action_requests"] = json.loads(json.dumps(interrupt.value["action_requests"])) internal["value"]["review_configs"] = json.loads(json.dumps(interrupt.value["review_configs"])) internal["value"]["action_requests"][0]["checkpoint_details"] = {"nested": [sentinel]} internal["value"]["review_configs"][0]["checkpoint_details"] = {"nested": [sentinel]} frozen = json.loads(json.dumps(internal)) expected_payload = _safe_json(frozen["value"]) expected_hash = content_hash(expected_payload) raw_hash = content_hash(frozen) stage('projection-start') pending = normalize_interrupt(internal) event = pending_input_projection_event(pending, message_id="message", run_id="run", source_sequence=1) assert set(event.payload) == {"parent_run_id", "interrupt_id", "payload_hash", "display_payload", "questions", "call_id"} assert "internal-only" not in str(event.payload) assert event.payload["interrupt_id"] == interrupt.id assert event.payload["payload_hash"] == pending["payload_hash"] assert internal == frozen assert pending["payload"]["checkpoint_details"] == "internal-only" assert before == await graph.aget_state(config("alice")) public = json.dumps({"safe_payload": pending["safe_payload"], "event": pending["event"], "projection": event.model_dump(mode="json")}) display = event.payload["display_payload"] checks = { "public_has_no_sentinel": sentinel not in public, "no_unknown_nested_fields": "checkpoint_details" not in public, "description_regenerated": sentinel in frozen["value"]["action_requests"][0]["description"] and "printf" in display["action_requests"][0]["description"], "command_visible": "printf" in json.dumps(display["action_requests"][0]["args"]), "internal_payload_unchanged": pending["payload"] == expected_payload, "decision_hash_unchanged": pending["payload_hash"] == expected_hash, "raw_hash_unchanged": content_hash(internal) == raw_hash, "reviews_compatible": display["review_configs"][0]["allowed_decisions"] == ["reject"], } print(json.dumps(checks), flush=True) assert all(checks.values()), "public pending confidentiality/identity contract failed" assert before.values["_verified_review_mode"]["mode"] == "manual" stage('assertions-done') asyncio.run(scenario()) stage('asyncio-exit') if __name__ == "__main__": repo = Path(__file__).resolve().parents[1] root = repo.parent / ".hermes/test-runtime/unified-execution/b03-host" root.mkdir(parents=True, exist_ok=True) ledger = root / "attempts.jsonl" name = sys.argv[1] assert name in NODES records = [json.loads(x) for x in ledger.read_text().splitlines()] if ledger.exists() else [] attempt = 1 + sum(r.get("phase") == "start" and r["node"] == name for r in records) assert attempt <= 5 home = root / f"{name}-{attempt}-home" home.mkdir(parents=True, exist_ok=True) env = {"HOME": str(home), "PATH": str(repo / ".venv/bin") + ":/usr/bin:/bin", "PYTHONPATH": str(repo) + os.pathsep + str(repo.parent / "Ai4Sci-Web"), "PYTHON_DOTENV_DISABLED": "1", "PYTEST_DISABLE_PLUGIN_AUTOLOAD": "1", "EVOSCIENTIST_HOME": str(home), "EVOSCIENTIST_CONFIG_DIR": str(home / "config"), "XDG_CONFIG_HOME": str(home / "config"), "PYTHONDONTWRITEBYTECODE": "1"} for key in ("DATA", "WORKSPACE", "SKILLS", "MEMORIES"): env[f"EVOSCIENTIST_{key}_DIR"] = str(home / key.lower()) cmd = [str(repo / ".venv/bin/python"), "-m", "pytest", "--noconftest", "-p", "no:cacheprovider", "-q", "-s", "--tb=short", "--disable-warnings", __file__ + "::" + name] with ledger.open("a") as f: f.write(json.dumps({"phase": "start", "node": name, "attempt": attempt, "limit": 5, "command": cmd, "overlap": "B01/B02 fixtures only; new host contract, historical gates not run", "snapshot": __import__("hashlib").sha256(Path(__file__).read_bytes()).hexdigest()}) + "\n") try: result = subprocess.run(cmd, cwd=repo, env=env, capture_output=True, text=True, timeout=180) output, code = result.stdout + result.stderr, result.returncode except subprocess.TimeoutExpired as exc: def text(value): return value.decode(errors='replace') if isinstance(value, bytes) else (value or '') output, code = text(exc.stdout) + text(exc.stderr) + '\nstartup/execution timeout\n', 124 output = output.replace('b03-private-canary-7f9a', '[REDACTED]') log = root / f"{name}-{attempt}.log" log.write_text(output) with ledger.open("a") as f: f.write(json.dumps({"phase": "result", "node": name, "attempt": attempt, "exit": code, "log": str(log)}) + "\n") print(output) sys.exit(code)