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
222 lines
12 KiB
Python
222 lines
12 KiB
Python
"""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) |