"""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)