From ddf337fc407bd21140d9511d4ffc713205c279e0 Mon Sep 17 00:00:00 2001 From: yangbo Date: Mon, 14 Sep 2026 16:53:30 +0800 Subject: [PATCH] =?UTF-8?q?=E5=AE=A1=E6=89=B9=E6=9A=82=E5=81=9C=E4=B8=8E?= =?UTF-8?q?=E7=BB=A7=E7=BB=AD=EF=BC=9A=E9=94=9A=E7=82=B9=E6=81=A2=E5=A4=8D?= =?UTF-8?q?=E5=87=86=E5=85=A5=20+=20=E4=B8=8D=E6=89=B9=E5=87=86=E7=BB=88?= =?UTF-8?q?=E6=AD=A2=E6=9C=AC=E8=BD=AE=EF=BC=88=E5=90=AB=20HITL=20?= =?UTF-8?q?=E8=A3=85=E9=85=8D=E6=94=B6=E6=95=9B=E4=B8=8E=E5=B7=A5=E4=BD=9C?= =?UTF-8?q?=E5=8C=BA=E8=8C=83=E5=9B=B4=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit langgraph_dev/http.py - _compatible_checkpoint/_history_admission 支持带锚点准入:以承载 __interrupt__ 的检查点 为父状态读取并校验中断标识;payload 保留锚点作为父状态 - 不带锚点时行为与原来完全一致(兼容性只靠这一条) middleware/dynamic_review.py - _abort_requested/_finish_turn + hook_config(can_jump_to=["end"]):在库消费中断后结束本轮 (工具不执行、不再发起模型调用);无标记时完全休眠 随本批一并落定(早前改动):EvoScientist.py / workspace_scope.py 的 HITL 装配收敛与工作区范围, 以及对应测试调整。 --- EvoScientist/EvoScientist.py | 60 +++++++-------- EvoScientist/langgraph_dev/http.py | 78 +++++++++++++++++-- EvoScientist/middleware/dynamic_review.py | 45 ++++++++--- EvoScientist/workspace_scope.py | 78 +++++++++++++++++-- tests/test_agent_factory_extensions.py | 3 +- tests/test_hitl.py | 94 ++++++++++++++++++----- 6 files changed, 282 insertions(+), 76 deletions(-) diff --git a/EvoScientist/EvoScientist.py b/EvoScientist/EvoScientist.py index 176b5f3..75198f2 100644 --- a/EvoScientist/EvoScientist.py +++ b/EvoScientist/EvoScientist.py @@ -1043,10 +1043,12 @@ def _get_default_backend( return _get_legacy_backend( guard_dangerous=guard_dangerous, refuse_delete=refuse_delete ) - from .workspace_scope import create_workspace_backend_factory + from .workspace_scope import create_deferred_scoped_backend cfg = _ensure_config() - return create_workspace_backend_factory( + # deepagents 0.7 removed backend factories, so hand the middleware an + # instance that resolves this run's scope from the runnable config. + return create_deferred_scoped_backend( _get_legacy_backend, dangerous=cfg.dangerous_mode, allow_unscoped_legacy=False, @@ -1365,20 +1367,18 @@ def _get_default_agent(): else _get_default_middleware() ) - # HITL on main agent only (mirrors create_cli_agent). Use middleware, - # not interrupt_on= kwarg — the kwarg propagates to every subagent and - # breaks parallel execute calls (multi-pending-interrupt LangGraph - # error). See PR #202. + # HITL on main agent only (mirrors create_cli_agent). Arm it through the + # review middleware and NEVER through `create_deep_agent(interrupt_on=)`: + # that kwarg makes deepagents append its own plain + # HumanInTheLoopMiddleware, and the second, auto-blind layer interrupts + # even on a run the gateway verified as auto — silently disabling + # automatic approval. It also propagates to every subagent and breaks + # parallel execute calls (multi-pending-interrupt LangGraph error). + # See PR #202. `HITL_INTERRUPT_ON` is the single source of the tool set. from .middleware import DynamicReviewMiddleware mw.append( - DynamicReviewMiddleware( - interrupt_on={ - "execute": True, - "run_in_background": True, - "schedule_task": True, - } - ) + DynamicReviewMiddleware(interrupt_on=dict(HITL_INTERRUPT_ON)) ) if web_full: @@ -1412,9 +1412,10 @@ def _get_default_agent(): ) kwargs = _apply_budgeted_skill_context(kwargs, be) + # No `interrupt_on=` here: HITL is armed above, as a single layer, by the + # auto-aware review middleware (see the comment there). _EvoScientist_agent = create_deep_agent( **kwargs, - interrupt_on=_build_hitl_interrupt_on(auto_approve=cfg.auto_approve), ).with_config({"recursion_limit": cfg.recursion_limit}) return _EvoScientist_agent @@ -1694,25 +1695,20 @@ def create_cli_agent( if main_agent_outer_middlewares: mw = [*main_agent_outer_middlewares, *mw] - # HITL on main agent only — passing `interrupt_on=` to create_deep_agent - # would propagate it to every subagent, breaking parallel execute calls - # (multi-pending-interrupt LangGraph error). + # HITL on main agent only. Arm it as a middleware and NEVER through + # `create_deep_agent(interrupt_on=)`: that kwarg makes deepagents append its + # own plain HumanInTheLoopMiddleware, and that second, auto-blind layer + # interrupts even on a run the gateway verified as auto — silently disabling + # automatic approval. It also propagates to every subagent, breaking parallel + # execute calls (multi-pending-interrupt LangGraph error). Web runs defer to + # the gateway's per-thread review mode, so they always arm the auto-aware + # middleware; `HITL_INTERRUPT_ON` is the single source of the tool set. if is_web: from .middleware.dynamic_review import DynamicReviewMiddleware - mw.append(DynamicReviewMiddleware(interrupt_on={ - "execute": True, "run_in_background": True, "schedule_task": True, - })) - elif not cfg.auto_approve: - mw.append( - HumanInTheLoopMiddleware( - interrupt_on={ - "execute": True, - "run_in_background": True, - "schedule_task": True, - } - ) - ) + mw.append(DynamicReviewMiddleware(interrupt_on=dict(HITL_INTERRUPT_ON))) + elif _build_hitl_interrupt_on(auto_approve=cfg.auto_approve) is not None: + mw.append(HumanInTheLoopMiddleware(interrupt_on=dict(HITL_INTERRUPT_ON))) # Re-load MCP tools from current config (picks up /mcp add changes) kwargs = load_mcp_and_build_kwargs( @@ -1757,8 +1753,10 @@ def create_cli_agent( from deepagents import create_deep_agent as original_factory create_deep_agent.__kwdefaults__ = original_factory.__kwdefaults__ + # No `interrupt_on=` here: HITL is armed above, as a single middleware layer + # (see the comment there). Passing the kwarg would add a second, auto-blind + # layer that defeats verified auto approval. return create_deep_agent( **kwargs, checkpointer=checkpointer, - interrupt_on=_build_hitl_interrupt_on(auto_approve=cfg.auto_approve), ).with_config({"recursion_limit": cfg.recursion_limit}) diff --git a/EvoScientist/langgraph_dev/http.py b/EvoScientist/langgraph_dev/http.py index bec48be..fff5710 100644 --- a/EvoScientist/langgraph_dev/http.py +++ b/EvoScientist/langgraph_dev/http.py @@ -275,8 +275,46 @@ async def bind_workspace_run(request: Request) -> JSONResponse: return JSONResponse(_run_payload(run)) +def _interrupt_ids(value: Any) -> set[str]: + """从 __interrupt__ 写入值里取出中断标识(Interrupt 对象或字典都兼容)。""" + + items = value if isinstance(value, (list, tuple)) else [value] + found: set[str] = set() + for item in items: + ident = getattr(item, "id", None) + if ident is None and isinstance(item, dict): + ident = item.get("id") or item.get("interrupt_id") + if ident: + found.add(str(ident)) + return found + + +def _anchor_has_interrupt(checkpoint: Any, expected: str) -> bool: + """锚点处是否仍承载待审批中断(并核对中断标识)。 + + 这是"按锚点恢复"的准入判据:只要承载该中断的检查点写入还在 PG,暂停就仍然 + 有效 —— 与运行时进程是否重启过、距暂停多久都无关。 + """ + + found: set[str] = set() + for write in getattr(checkpoint, "pending_writes", None) or (): + try: + if len(write) < 3 or str(write[1]) != "__interrupt__": + continue + found |= _interrupt_ids(write[2]) + except TypeError: + continue + if not found: + return False + # 网关未提供标识(或回退值 default)时只做存在性判定。 + if not expected or expected == "default": + return True + return expected in found + + async def _compatible_checkpoint(conn, thread_id: str, assistant_id: str, - config: dict) -> tuple[bool, bool]: + config: dict, anchor: dict | None = None + ) -> tuple[bool, bool]: """Read through the API-owned saver and graph factory, never execute here.""" from langgraph_api._checkpointer import get_checkpointer from langgraph_api.graph import get_graph, graph_exists @@ -285,10 +323,15 @@ async def _compatible_checkpoint(conn, thread_id: str, assistant_id: str, saver = await get_checkpointer(conn=conn) read_config = {**config, "configurable": { **config.get("configurable", {}), "thread_id": thread_id, - "checkpoint_ns": "", + "checkpoint_ns": str((anchor or {}).get("checkpoint_ns") or ""), }} - # Admission always checks the current head, never a caller-selected ancestor. - read_config["configurable"].pop("checkpoint_id", None) + if anchor and str(anchor.get("checkpoint_id") or ""): + # 按锚点恢复:调用方(网关)指定了承载该中断的祖先检查点,暂停时就已落库。 + # 只有该锚点处确实还有待审批写入才放行 —— 不依赖运行时当前头部。 + read_config["configurable"]["checkpoint_id"] = str(anchor["checkpoint_id"]) + else: + # Admission always checks the current head, never a caller-selected ancestor. + read_config["configurable"].pop("checkpoint_id", None) checkpoint = await saver.aget_tuple(read_config) if checkpoint is None: return False, False @@ -300,6 +343,10 @@ async def _compatible_checkpoint(conn, thread_id: str, assistant_id: str, graph_id = assistant["graph_id"] if checkpoint.metadata.get("graph_id", graph_id) != graph_id: raise ValueError("checkpoint graph mismatch") + if anchor and str(anchor.get("checkpoint_id") or ""): + return True, _anchor_has_interrupt( + checkpoint, str(anchor.get("interrupt_id") or "") + ) # get_graph enters coroutine/async-context-manager factories and binds the # same API saver used by the worker. aget_state also validates delta seeds. async with get_graph(graph_id, read_config, checkpointer=saver, @@ -337,8 +384,11 @@ async def _has_legacy_history(conn, thread_id: str) -> bool: return True -async def _history_admission(conn, thread_id, assistant_id, config, operation, history): - exists, pending = await _compatible_checkpoint(conn, thread_id, assistant_id, config) +async def _history_admission(conn, thread_id, assistant_id, config, operation, history, + anchor: dict | None = None): + exists, pending = await _compatible_checkpoint( + conn, thread_id, assistant_id, config, anchor + ) if operation == "resume": return "resume" if exists and pending else "CHECKPOINT_RESUME_UNAVAILABLE" if pending: @@ -398,6 +448,11 @@ async def create_recoverable_run(request: Request) -> JSONResponse: return JSONResponse({"code": "INVALID_RESUME_REQUEST"}, status_code=400) elif command is not None: return JSONResponse({"code": "INVALID_START_REQUEST"}, status_code=400) + # 按锚点恢复(可选,仅 resume):网关在暂停时把"承载该中断的检查点"落库, + # 继续时回传。这样续接只依赖该检查点仍在 PG,而不依赖运行时当前头部。 + anchor = value.get("anchor") if operation == "resume" else None + if not isinstance(anchor, dict) or not str(anchor.get("checkpoint_id") or ""): + anchor = None history = value.get("history") # History validation is performed under the create lock, after idempotent # lookup. Only that branch can attest this request did not create a Run. @@ -427,12 +482,18 @@ async def create_recoverable_run(request: Request) -> JSONResponse: "run_request_id": run_request_id, "request_hash": request_hash, } - # The admission read and the worker must address the same current head. - # Forking from an ancestor is not part of the recoverable-run contract. + # The admission read and the worker must address the same checkpoint. + # Without an anchor that is the current head (forking from an ancestor is + # not part of the recoverable-run contract); with an anchor it is the + # ancestor recorded at pause time — which is exactly what makes the + # continuation independent of restarts and elapsed time. configurable = dict(payload["config"].get("configurable", {})) for key in ("checkpoint_id", "checkpoint_map"): configurable.pop(key, None) configurable.update(thread_id=thread_id, checkpoint_ns="") + if anchor is not None: + configurable["checkpoint_id"] = str(anchor["checkpoint_id"]) + configurable["checkpoint_ns"] = str(anchor.get("checkpoint_ns") or "") payload["config"] = {**payload["config"], "configurable": configurable} async with _recoverable_run_lock: async with connect() as conn: @@ -459,6 +520,7 @@ async def create_recoverable_run(request: Request) -> JSONResponse: try: admission = await _history_admission( conn, thread_id, assistant_id, payload["config"], operation, history, + anchor, ) except Exception: return JSONResponse({ diff --git a/EvoScientist/middleware/dynamic_review.py b/EvoScientist/middleware/dynamic_review.py index 7ff83d4..cf2bb8f 100644 --- a/EvoScientist/middleware/dynamic_review.py +++ b/EvoScientist/middleware/dynamic_review.py @@ -9,7 +9,7 @@ from typing import Annotated, Any, NotRequired import httpx from EvoScientist.internal_service import internal_service_headers -from langchain.agents.middleware import HumanInTheLoopMiddleware +from langchain.agents.middleware import HumanInTheLoopMiddleware, hook_config from langchain.agents.middleware.types import AgentState, OmitFromSchema from langgraph.config import get_config @@ -137,6 +137,31 @@ class DynamicReviewMiddleware(HumanInTheLoopMiddleware): state_schema = DynamicReviewState + @staticmethod + def _abort_requested() -> bool: + """本轮是否被要求"终止"(网关放弃某个待审批时注入的标记)。 + + 方案 A:中断仍由库正常消费(该工具不执行),但随后不再发起模型调用, + 直接把本轮跳到 end —— 即用户口径里的"不批准 → 停止这条对话"。 + """ + + try: + config = get_config() + except RuntimeError: + return False + configurable = config.get("configurable") if isinstance(config, Mapping) else None + if not isinstance(configurable, Mapping): + return False + return bool(configurable.get("ai4sci_abort_turn")) + + @classmethod + def _finish_turn(cls, update: dict[str, Any] | None) -> dict[str, Any] | None: + """终止本轮:保留库的决议结果(中断被消费掉),再结束本轮。""" + + if not cls._abort_requested(): + return update + return {**(update or {}), "jump_to": "end"} + def before_agent(self, state: DynamicReviewState, runtime: Any) -> dict[str, Any]: del state, runtime run_id, review = _review_context() @@ -153,6 +178,7 @@ class DynamicReviewMiddleware(HumanInTheLoopMiddleware): return {"_verified_review_mode": _manual_state(run_id, review)} return {"_verified_review_mode": await _resolve_async(run_id, review)} + @hook_config(can_jump_to=["end"]) def after_model( self, state: DynamicReviewState, runtime: Any ) -> dict[str, Any] | None: @@ -177,12 +203,12 @@ class DynamicReviewMiddleware(HumanInTheLoopMiddleware): return None if review.get("requested_mode") != "auto": # The gateway now requires manual approval for this turn. - return super().after_model(state, runtime) + return self._finish_turn(super().after_model(state, runtime)) try: _resolve_sync(current_run_id, review) except AutoReviewVerificationError: # Auto approval could not be re-verified; fall back to review. - return super().after_model(state, runtime) + return self._finish_turn(super().after_model(state, runtime)) return None if mode == "manual": # A LangGraph resume continues at this interrupted node and does not @@ -197,11 +223,12 @@ class DynamicReviewMiddleware(HumanInTheLoopMiddleware): try: _resolve_sync(current_run_id, review) except AutoReviewVerificationError: - return super().after_model(state, runtime) + return self._finish_turn(super().after_model(state, runtime)) return None - return super().after_model(state, runtime) + return self._finish_turn(super().after_model(state, runtime)) raise AutoReviewVerificationError("REVIEW_MODE_STATE_INVALID") + @hook_config(can_jump_to=["end"]) async def aafter_model( self, state: DynamicReviewState, runtime: Any ) -> dict[str, Any] | None: @@ -218,11 +245,11 @@ class DynamicReviewMiddleware(HumanInTheLoopMiddleware): if review is None: return None if review.get("requested_mode") != "auto": - return super().after_model(state, runtime) + return self._finish_turn(super().after_model(state, runtime)) try: await _resolve_async(current_run_id, review) except AutoReviewVerificationError: - return super().after_model(state, runtime) + return self._finish_turn(super().after_model(state, runtime)) return None if mode == "manual": # Mirror the sync path: an injected auto context on a resume child @@ -232,7 +259,7 @@ class DynamicReviewMiddleware(HumanInTheLoopMiddleware): try: await _resolve_async(current_run_id, review) except AutoReviewVerificationError: - return super().after_model(state, runtime) + return self._finish_turn(super().after_model(state, runtime)) return None - return super().after_model(state, runtime) + return self._finish_turn(super().after_model(state, runtime)) raise AutoReviewVerificationError("REVIEW_MODE_STATE_INVALID") diff --git a/EvoScientist/workspace_scope.py b/EvoScientist/workspace_scope.py index 37bae8e..e423395 100644 --- a/EvoScientist/workspace_scope.py +++ b/EvoScientist/workspace_scope.py @@ -354,21 +354,48 @@ class DeferredScopedBackend(SandboxBackendProtocol): def __init__( self, - config: _RuntimeScopeConfig, + config: _RuntimeScopeConfig | None, *, dangerous: bool, + legacy_backend: Callable[[], Any] | None = None, + allow_unscoped_legacy: bool = False, ) -> None: self._config = config self._dangerous = dangerous + self._legacy_backend = legacy_backend + self._allow_unscoped_legacy = allow_unscoped_legacy self._backend: Any | None = None self._backend_key: tuple[str, str, str, int] | None = None self._lock = threading.RLock() + def _scope_config(self) -> _RuntimeScopeConfig | None: + """Return this proxy's scope identity, resolving it lazily if needed. + + deepagents 0.7 removed backend factories: the middleware only accepts + an initialized ``BackendProtocol``, so a deploy-scoped proxy is built + before any run exists and must read its scope from the active runnable + config on each operation. ``asyncio.to_thread`` copies the contextvars + context, so the runnable config remains readable from the worker + threads the inherited async methods dispatch into. + """ + + if self._config is not None: + return self._config + return _runtime_scope_config(None, kind="filesystem backend") + @property def id(self) -> str: # This is queried while composing the model request; do not initialize # the real backend or touch the Registry here. - return f"scope-{self._config.scope_id[:8]}-{self._config.owner_id[:8]}" + config = self._config + if config is None: + try: + config = _runtime_scope_config(None, kind="filesystem backend") + except ScopeAccessError: + config = None + if config is None: + return "scope-deferred" + return f"scope-{config.scope_id[:8]}-{config.owner_id[:8]}" def _delegate(self) -> Any: """Validate the current scope and return a concrete backend. @@ -377,10 +404,18 @@ class DeferredScopedBackend(SandboxBackendProtocol): keep using a backend constructed before the lifecycle transition. """ - # Async backend methods run this code in a worker thread. LangGraph's - # RunnableConfig context variable is not available there, so validate - # the immutable scope parsed by the factory on the graph thread. - context = _resolve_scope_context(self._config) + # ``asyncio.to_thread`` copies the contextvars context, so the active + # runnable config is still readable here even though the inherited + # async methods dispatch this work to a worker thread. Registry + # validation therefore stays off the Agent event loop. + config = self._scope_config() + if config is None: + if not self._allow_unscoped_legacy or self._legacy_backend is None: + raise ScopeAccessError( + "deployed graph runs require a workspace scope" + ) + return self._legacy_backend() + context = _resolve_scope_context(config) if context is None: raise ScopeAccessError("scoped backend lost its workspace scope") if (is_required() or os.getenv("EVOSCIENTIST_DEPLOY_MODE", "").lower() == "full") and self._dangerous: @@ -478,6 +513,37 @@ def create_workspace_backend_factory( return factory +def create_deferred_scoped_backend( + legacy_backend: Callable[[], Any], + *, + dangerous: bool = False, + allow_unscoped_legacy: bool = False, +) -> Any: + """Return a deploy-scoped backend INSTANCE for deepagents >= 0.7. + + deepagents 0.7 removed backend factories: ``FilesystemMiddleware`` now + rejects any ``backend`` that is callable without being a + ``BackendProtocol`` instance. Return an instance instead and let it read + the run's workspace scope from the active runnable config on every + operation. Unscoped runs stay fail-closed unless the caller explicitly + allows the legacy backend. + """ + + if ( + is_required() + or os.getenv("EVOSCIENTIST_DEPLOY_MODE", "").lower() == "full" + ) and dangerous: + raise ScopeAccessError( + "dangerous_mode is incompatible with required isolation" + ) + return DeferredScopedBackend( + None, + dangerous=dangerous, + legacy_backend=legacy_backend, + allow_unscoped_legacy=allow_unscoped_legacy, + ) + + def workspace_metadata(record: ScopeRecord) -> dict[str, str | int]: """Metadata mirrored onto the LangGraph primary thread by trusted callers.""" diff --git a/tests/test_agent_factory_extensions.py b/tests/test_agent_factory_extensions.py index eebd5d2..374355f 100644 --- a/tests/test_agent_factory_extensions.py +++ b/tests/test_agent_factory_extensions.py @@ -9,9 +9,10 @@ def test_create_cli_agent_accepts_host_backend_and_memory_options( chat_model = object() class _CompositeBackend: - def __init__(self, *, default, routes): + def __init__(self, *, default, routes, artifacts_root=None): calls["default_backend"] = default calls["routes"] = routes + calls["artifacts_root"] = artifacts_root class _MemoryBackend: def __init__(self, **kwargs): diff --git a/tests/test_hitl.py b/tests/test_hitl.py index 3a38745..7fb9307 100644 --- a/tests/test_hitl.py +++ b/tests/test_hitl.py @@ -742,27 +742,57 @@ class TestInterruptOnWiring: assert cfg.auto_approve is True assert _build_hitl_interrupt_on(auto_approve=cfg.auto_approve) is None - def test_hitl_interrupt_on_reaches_create_deep_agent(self): - """The kwarg must actually reach ``create_deep_agent`` — not just the - pure helper — so a future edit that drops it or re-adds a bare - ``HumanInTheLoopMiddleware`` append gets caught.""" - import EvoScientist.EvoScientist as es_mod - from EvoScientist.EvoScientist import _build_hitl_interrupt_on + def test_hitl_is_armed_by_the_review_middleware_only(self): + """HITL must be exactly one layer, and it must be the auto-aware one. - captured = [] + ``create_deep_agent(interrupt_on=...)`` makes deepagents append its own + plain ``HumanInTheLoopMiddleware``. That second, auto-blind layer + interrupts even on a run the gateway verified as auto — which silently + disables automatic approval, because ``DynamicReviewMiddleware`` can only + decline to interrupt for itself. So the tool set is armed through the + middleware and the kwarg must stay out. + """ + import EvoScientist.EvoScientist as es_mod + from EvoScientist.EvoScientist import HITL_INTERRUPT_ON + from EvoScientist.middleware import DynamicReviewMiddleware + from langchain.agents.middleware import HumanInTheLoopMiddleware + + captured_kwargs = [] + captured_middleware = [] + + class _WebProfile(MagicMock): + """A Web execution profile: fixed name, every other field a mock.""" + + name = "web_v3" + + web_profile = _WebProfile() def fake_create_deep_agent(**kwargs): - captured.append(kwargs.get("interrupt_on", "MISSING")) + captured_kwargs.append(kwargs) agent = MagicMock() agent.with_config.return_value = agent return agent - for auto_approve in (False, True): + def fake_build_kwargs(_backend, middleware, **_kwargs): + captured_middleware.append(list(middleware)) + return {"name": "x"} + + def build(auto_approve, execution_profile=None): cfg = MagicMock() cfg.auto_approve = auto_approve cfg.dangerous_mode = False cfg.sandbox_execute_timeout = 300 cfg.recursion_limit = 100 + # A Web run must supply the host-provided backend and checkpointer. + web_only = ( + { + "memory_dir": "/tmp/test-interrupt-on-wiring-memory", + "workspace_backend": MagicMock(), + "checkpointer": MagicMock(), + } + if execution_profile is not None + else {} + ) with patch( "deepagents.create_deep_agent", side_effect=fake_create_deep_agent @@ -774,25 +804,47 @@ class TestInterruptOnWiring: with patch.object( es_mod, "load_mcp_and_build_kwargs", - return_value={"name": "x"}, + side_effect=fake_build_kwargs, ): es_mod.create_cli_agent( workspace_dir="/tmp/test-interrupt-on-wiring", config=cfg, chat_model=MagicMock(), + execution_profile=execution_profile, + **web_only, ) - assert captured == [ - _build_hitl_interrupt_on(auto_approve=False), - _build_hitl_interrupt_on(auto_approve=True), - ] - assert captured[0] == { - "execute": True, - "run_in_background": True, - "schedule_task": True, - "delete": True, - } - assert captured[1] is None + build(auto_approve=False) # attended CLI + build(auto_approve=False, execution_profile=web_profile) # Web run + build(auto_approve=True) # auto: nothing armed + + # The kwarg would add a second, auto-blind HITL layer: never pass it. + assert all("interrupt_on" not in kwargs for kwargs in captured_kwargs) + + def hitl_layers(index): + return [ + item + for item in captured_middleware[index] + if isinstance( + item, (DynamicReviewMiddleware, HumanInTheLoopMiddleware) + ) + ] + + # Attended CLI: one plain HITL layer covering the reviewed tool set. + cli_layers = hitl_layers(0) + assert len(cli_layers) == 1 + assert isinstance(cli_layers[0], HumanInTheLoopMiddleware) + assert set(cli_layers[0].interrupt_on) == set(HITL_INTERRUPT_ON) + + # Web: one layer, and it must be the auto-aware one — it is what lets a + # gateway-verified auto run proceed without interrupting. + web_layers = hitl_layers(1) + assert len(web_layers) == 1 + assert isinstance(web_layers[0], DynamicReviewMiddleware) + assert set(web_layers[0].interrupt_on) == set(HITL_INTERRUPT_ON) + + # auto_approve: nothing is armed at all. + assert hitl_layers(2) == [] # =============================================================================