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
166 lines
10 KiB
Python
166 lines
10 KiB
Python
"""B04 durable authorization reaches real Web Graph execution, offline only."""
|
|
import asyncio
|
|
import hashlib
|
|
import importlib
|
|
import socket
|
|
import sqlite3
|
|
import uuid
|
|
|
|
|
|
import pytest
|
|
|
|
|
|
def _authorization_scenario(tmp_path, monkeypatch, *, construction_failure=False):
|
|
import faulthandler
|
|
faulthandler.dump_traceback_later(15)
|
|
from EvoScientist.config.settings import EvoScientistConfig
|
|
from EvoScientist.llm.contracts import AgentInputV3, WebHostContext, EvoRuntimeError, canonical_json_v1
|
|
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 deepagents.backends.protocol import SandboxBackendProtocol, ExecuteResponse
|
|
from langchain_core.language_models.fake_chat_models import FakeMessagesListChatModel
|
|
from langchain_core.messages import AIMessage, AIMessageChunk
|
|
from langchain_core.outputs import ChatGenerationChunk
|
|
from langgraph.checkpoint.sqlite.aio import AsyncSqliteSaver
|
|
from langgraph.types import Command
|
|
from tests.test_web_model_runtime import _runtime, _preparation, _admission, _Sink
|
|
|
|
def forbidden(*args, **kwargs):
|
|
raise AssertionError("network forbidden")
|
|
monkeypatch.setattr(socket.socket, "connect", forbidden)
|
|
monkeypatch.setattr(socket, "getaddrinfo", forbidden)
|
|
module = importlib.import_module("EvoScientist.EvoScientist")
|
|
monkeypatch.setattr(module, "_load_mcp_tools_cached", lambda **kw: {})
|
|
monkeypatch.setattr(module, "_load_mcp_config_once", lambda: ("b04-offline", {}))
|
|
calls = []
|
|
|
|
class Backend(ScopedFilesystemBackend, SandboxBackendProtocol):
|
|
@property
|
|
def id(self):
|
|
return "b04-no-shell"
|
|
|
|
def execute(self, command, *, timeout=None):
|
|
assert command == "authorized-probe"
|
|
calls.append(command)
|
|
return ExecuteResponse(output="authorized-result", exit_code=0)
|
|
|
|
class Model(FakeMessagesListChatModel):
|
|
def bind_tools(self, tools, **kwargs):
|
|
return self
|
|
|
|
async def _astream(self, messages, stop=None, run_manager=None, **kwargs):
|
|
if any(m.type == "tool" for m in messages):
|
|
yield ChatGenerationChunk(message=AIMessageChunk(content="authorized done"))
|
|
else:
|
|
yield ChatGenerationChunk(message=AIMessageChunk(content="", tool_call_chunks=[{
|
|
"name": "execute", "args": '{"command":"authorized-probe"}',
|
|
"id": "pending-tool", "index": 0}]))
|
|
|
|
async def scenario():
|
|
import tests.test_web_model_runtime as helpers
|
|
payload = helpers.v3_payload()
|
|
payload["providers"]["custom-openai"]["models"][0]["context_window"] = 131072
|
|
monkeypatch.setattr(helpers, "v3_payload", lambda: payload)
|
|
runtime, authority = _runtime(tmp_path, monkeypatch)
|
|
registry = SQLiteHostRegistry(tmp_path / "private" / "host.sqlite", host_id="host", boot_id="one")
|
|
runtime.host_registry = registry
|
|
runtime.model_factory = lambda **kw: Model(responses=[AIMessage(content="unused")])
|
|
cfg = EvoScientistConfig(auto_approve=False, enable_async_subagents=False,
|
|
enable_scheduler=False, memory_workers_enabled=False)
|
|
runtime.agent_factory = lambda snapshot, host, models: create_web_agent(
|
|
snapshot=snapshot, host=host, model_set=models, config=cfg)
|
|
root = tmp_path / "workspace"
|
|
root.mkdir()
|
|
tools, revision = web_tool_registry_manifest()
|
|
async with AsyncSqliteSaver.from_conn_string(str(tmp_path / "graph.sqlite")) as saver:
|
|
host = WebHostContext(str(root), str(root), Backend(root), saver,
|
|
tool_selector_threshold=10000, tool_registry=tools,
|
|
tool_registry_revision=revision, runtime_event_sink=_Sink())
|
|
value = AgentInputV3("request approval", "b04-thread")
|
|
quote = await runtime.prepare_model_run(_preparation(authority, value, title_policy="disabled"), value, host)
|
|
parent = await runtime.start_web_run(_admission(authority, quote))
|
|
assert await parent.wait_stopped(timeout=8) == "awaiting_input"
|
|
assert calls == []
|
|
state = await parent._agent.aget_state({"configurable": {"thread_id": value.checkpoint_thread_id}})
|
|
assert state.values["_verified_review_mode"]["mode"] == "manual"
|
|
details = parent._terminal_event.payload
|
|
identity = {k: details[k] for k in ("checkpoint_thread_id", "checkpoint_id", "checkpoint_ns", "pending_interrupts")}
|
|
pending_hash = hashlib.sha256(canonical_json_v1(identity)).hexdigest()
|
|
approved = AgentInputV3(Command(resume={"decisions": [{"type": "approve"}]}), value.checkpoint_thread_id)
|
|
decision_hash = hashlib.sha256(canonical_json_v1(approved.message.resume)).hexdigest()
|
|
|
|
def grant_for(value=approved, **changes):
|
|
template = _preparation(authority, value, title_policy="disabled")
|
|
fields = ("turn_id", "thread_id", "subject_id", "requested_model_ref", "plan", "roles",
|
|
"requires_vision", "reasoning_effort", "title_policy", "gateway_input_digest",
|
|
"checkpoint_thread_id", "turn_fencing_token")
|
|
return authority.sign_preparation(**{k: getattr(template, k) for k in fields},
|
|
**dict(dict(request_id=str(uuid.uuid4()), checkpoint_snapshot_id=identity["checkpoint_id"],
|
|
predecessor_execution_id=parent.run_id, predecessor_checkpoint_id=identity["checkpoint_id"],
|
|
predecessor_owner_epoch=1, continuation_pending_hash=pending_hash,
|
|
continuation_decision_hash=decision_hash), **changes), ttl_ms=60000)
|
|
|
|
for changes in ({"continuation_decision_hash": "bad"}, {"continuation_pending_hash": "bad"},
|
|
{"predecessor_owner_epoch": 2}):
|
|
with pytest.raises(EvoRuntimeError):
|
|
await runtime.prepare_model_run(grant_for(**changes), approved, host)
|
|
with pytest.raises(EvoRuntimeError):
|
|
await runtime.prepare_model_run(grant_for(value=value), value, host)
|
|
good = grant_for()
|
|
prepared = await runtime.prepare_model_run(good, approved, host)
|
|
approved.message.resume["decisions"][0]["type"] = "reject"
|
|
with pytest.raises(EvoRuntimeError, match="CONTINUATION_DECISION_INVALID"):
|
|
await runtime.start_web_run(_admission(authority, prepared))
|
|
approved.message.resume["decisions"][0]["type"] = "approve"
|
|
if construction_failure:
|
|
def fail_factory(*args):
|
|
raise RuntimeError("controlled construction failure")
|
|
runtime.agent_factory = fail_factory
|
|
with pytest.raises(RuntimeError, match="controlled construction failure"):
|
|
await runtime.start_web_run(_admission(authority, prepared))
|
|
assert registry.inspect(prepared.execution_id)["status"] == "failed"
|
|
with sqlite3.connect(registry.path) as db:
|
|
row = db.execute("SELECT consumed_by, decision_hash FROM pending_continuations WHERE execution_id=?", (parent.run_id,)).fetchone()
|
|
assert row == (prepared.execution_id, decision_hash)
|
|
runtime.host_registry = SQLiteHostRegistry(registry.path, host_id="host", boot_id="two")
|
|
again = await runtime.prepare_model_run(grant_for(), approved, host)
|
|
with pytest.raises(EvoRuntimeError, match="CONTINUATION_CONSUMED_FAILURE_REQUIRES_REAUTHORIZATION"):
|
|
await runtime.start_web_run(_admission(authority, again))
|
|
assert calls == []
|
|
print("B04 consumed construction failure stays consumed across reopen; new admission cannot replay")
|
|
return
|
|
# Mutating the checkpoint after prepare cannot execute its new pending.
|
|
await parent._agent.aupdate_state(state.config, {"_verified_review_mode": {"mode": "auto"}})
|
|
with pytest.raises(EvoRuntimeError):
|
|
await runtime.start_web_run(_admission(authority, prepared))
|
|
assert calls == []
|
|
# Restore the original latest checkpoint pointer in this isolated fixture.
|
|
with sqlite3.connect(tmp_path / "graph.sqlite") as db:
|
|
db.execute("DELETE FROM checkpoints WHERE checkpoint_id > ?", (identity["checkpoint_id"],))
|
|
# Reopen registry before bind to prove pending authority is durable.
|
|
runtime.host_registry = SQLiteHostRegistry(registry.path, host_id="host", boot_id="two")
|
|
child = await runtime.start_web_run(_admission(authority, prepared))
|
|
assert child.run_id != parent.run_id
|
|
assert await child.wait_stopped(timeout=8) == "completed"
|
|
assert calls == ["authorized-probe"]
|
|
with sqlite3.connect(registry.path) as db:
|
|
row = db.execute("SELECT consumed_by, decision_hash FROM pending_continuations WHERE execution_id=?", (parent.run_id,)).fetchone()
|
|
assert row == (child.run_id, decision_hash)
|
|
with pytest.raises(EvoRuntimeError):
|
|
again = await runtime.prepare_model_run(grant_for(), approved, host)
|
|
await runtime.start_web_run(_admission(authority, again))
|
|
assert calls == ["authorized-probe"]
|
|
assert (tmp_path / "private").stat().st_mode & 0o777 == 0o700
|
|
assert (tmp_path / "private" / "host.sqlite").stat().st_mode & 0o777 == 0o600
|
|
print("B04 durable pending -> signed decision -> atomic bind -> Graph tool: exactly one; replay rejected")
|
|
asyncio.run(scenario())
|
|
faulthandler.cancel_dump_traceback_later()
|
|
|
|
|
|
def test_durable_decision_bound_to_graph_execution(tmp_path, monkeypatch):
|
|
_authorization_scenario(tmp_path, monkeypatch)
|
|
|
|
|
|
def test_consumed_graph_child_failure_requires_reauthorization(tmp_path, monkeypatch):
|
|
_authorization_scenario(tmp_path, monkeypatch, construction_failure=True) |